npm package discovery and stats viewer.

Discover Tips

  • General search

    [free text search, go nuts!]

  • Package details

    pkg:[package-name]

  • User packages

    @[username]

Sponsor

Optimize Toolset

I’ve always been into building performant and accessible sites, but lately I’ve been taking it extremely seriously. So much so that I’ve been building a tool to help me optimize and monitor the sites that I build to make sure that I’m making an attempt to offer the best experience to those who visit them. If you’re into performant, accessible and SEO friendly sites, you might like it too! You can check it out at Optimize Toolset.

About

Hi, 👋, I’m Ryan Hefner  and I built this site for me, and you! The goal of this site was to provide an easy way for me to check the stats on my npm packages, both for prioritizing issues and updates, and to give me a little kick in the pants to keep up on stuff.

As I was building it, I realized that I was actually using the tool to build the tool, and figured I might as well put this out there and hopefully others will find it to be a fast and useful way to search and browse npm packages as I have.

If you’re interested in other things I’m working on, follow me on Twitter or check out the open source projects I’ve been publishing on GitHub.

I am also working on a Twitter bot for this site to tweet the most popular, newest, random packages from npm. Please follow that account now and it will start sending out packages soon–ish.

Open Software & Tools

This site wouldn’t be possible without the immense generosity and tireless efforts from the people who make contributions to the world and share their work via open source initiatives. Thank you 🙏

© 2026 – Pkg Stats / Ryan Hefner

@push.rocks/smartipc

v3.0.0

Published

A library for node inter process communication, providing an easy-to-use API for IPC.

Readme

@push.rocks/smartipc 🚀

Rock-solid IPC and local named mutexes for Node.js

npm version TypeScript License: MIT

SmartIPC delivers Inter-Process Communication for Node.js applications with automatic reconnection, heartbeat monitoring, clean shutdowns, streaming, and native local locking primitives. A packaged, statically linked Rust helper runs only when NamedMutex or probeFileOpenExclusivity() is used; transport-only consumers never spawn it.

Issue Reporting and Security

For reporting bugs, issues, or security vulnerabilities, please visit community.foss.global/. This is the central community hub for all issue reporting. Developers who sign and comply with our contribution agreement and go through identification can also get a code.foss.global/ account to submit Pull Requests directly.

🎯 Why SmartIPC?

  • ABI-Independent Native Locking - NamedMutex uses the packaged Rust helper instead of a Node native addon and reports NATIVE_BACKEND_UNAVAILABLE if the helper cannot start
  • Battle-tested Reliability - Automatic reconnection, graceful degradation, and timeout handling
  • Type-Safe - Full TypeScript support with generics for compile-time safety
  • CI/Test Ready - Built-in helpers and race condition prevention for testing
  • Observable - Real-time metrics, connection tracking, and health monitoring
  • Multiple Patterns - Request/Response, Pub/Sub, and Fire-and-Forget messaging
  • Streaming Support - Efficient, backpressure‑aware streaming for large data and files

📦 Installation

pnpm add @push.rocks/smartipc

🚀 Quick Start

import { SmartIpc } from '@push.rocks/smartipc';

// Create a server
const server = SmartIpc.createServer({
  id: 'my-service',
  socketPath: '/tmp/my-service.sock', // a stale socket file is replaced safely on start
});

// Handle incoming messages
server.onMessage('greet', async (data, clientId) => {
  console.log(`Client ${clientId} says:`, data.message);
  return { response: `Hello ${data.name}!` };
});

// Start the server
await server.start({ readyWhen: 'accepting' });  // Wait until fully ready
console.log('Server is ready to accept connections! ✨');

// Create a client
const client = SmartIpc.createClient({
  id: 'my-service',
  socketPath: '/tmp/my-service.sock',
  connectRetry: {
    enabled: true,
    maxAttempts: 10
  }
});

// Connect with automatic retry
await client.connect();

// Send a request and get a response
const response = await client.request('greet', { 
  name: 'World',
  message: 'Hi there!' 
});
console.log('Server said:', response.response);  // "Hello World!"

Local named mutex

NamedMutex coordinates cooperating Node.js processes running as the same operating-system identity on one machine and one local filesystem. It uses a permanent, hashed lock-anchor filename in a private per-user directory. The anchor is never deleted by SmartIPC, so a process cannot replace another process's active lock by recreating a path.

import { NamedMutex, NamedMutexError } from '@push.rocks/smartipc';

const migrationMutex = new NamedMutex('my-service:database-migrations');

try {
  const lease = await migrationMutex.acquire({
    timeoutMs: 5_000,
    signal: shutdownController.signal,
  });
  try {
    await runMigrations();
  } finally {
    await lease.release();
  }
} catch (error) {
  if (error instanceof NamedMutexError && error.code === 'TIMEOUT') {
    // Another same-user process retained the lease for the whole wait bound.
  } else {
    throw error;
  }
}

tryAcquire() makes one immediate nonblocking attempt and returns undefined on contention. acquire() keeps same-isolate callers in FIFO order and uses nonblocking native attempts until its timeout or AbortSignal fires. timeoutMs defaults to 10 seconds and pollIntervalMs defaults to 25 milliseconds; both use monotonic time, the native helper enforces a 1 millisecond polling floor, and polling delays are capped at the largest delay Node can represent safely.

Each NamedMutexLease exposes an active, releasing, released, or releaseFailed state. Concurrent and repeated release() calls share one result, so they cannot release a later lease. If helper termination or descriptor close has an ambiguous outcome, the anchor remains poisoned within the isolate; an ambiguous active-lease close rejects with RELEASE_FAILED, and subsequent acquisitions fail closed instead of assuming ownership is safe. Other public error codes distinguish invalid arguments, unsupported runtimes, unsafe anchors, an unavailable Rust backend, native lock failures, timeouts, and cancellation.

The guarantee is intentionally local: it does not cover different users, different machines, containers without a shared local anchor, NFS or other network filesystems, or uncooperative processes that ignore advisory locks. The supported runtime contract is Node.js 24 through 26 on Linux x64. Anchor directories and files must remain owned by the current user with modes 0700 and 0600 respectively. A custom directoryPath must be an absolute path on a trusted local filesystem.

Point-in-time file open probe

probeFileOpenExclusivity() uses the packaged Rust helper to ask the Linux kernel whether another read/write open description exists for an already-open regular file. It returns exclusive when none exists and contended when F_SETLEASE reports a conflict. This is a point-in-time observation, not a retained lock.

import { open } from 'node:fs/promises';
import { probeFileOpenExclusivity } from '@push.rocks/smartipc';

const handle = await open('/run/my-app/legacy.lock', 'r+');
const result = await probeFileOpenExclusivity(handle);

The input must be an open read/write regular-file FileHandle. The function consumes and closes the handle on every conclusive result and error path. The caller must surrender the handle when calling and must reopen and revalidate the path if later work needs another descriptor. CLEANUP_UNCERTAIN means helper termination or handle closure could not be proven, so callers must fail closed. Linux O_PATH references do not contend because they cannot represent a read/write ownership handle. The API has the same Node.js 24 through 26, Linux x64, local-filesystem runtime contract as NamedMutex.

FileOpenProbeError.code distinguishes INVALID_ARGUMENT, UNSUPPORTED_RUNTIME, NATIVE_BACKEND_UNAVAILABLE, HANDLE_CLOSE_FAILED, PROBE_FAILED, and CLEANUP_UNCERTAIN.

Peer credentials (Linux)

A server can require the kernel-verified identity of every connecting process. With requirePeerCredentials: true, each accepted Unix-domain connection is held paused, its SO_PEERCRED triple (pid, uid, gid) is read from the kernel, an optional authorizer decides, and only then is any byte read from the connection. A rejected connection is closed without a single message being parsed or dispatched.

import { SmartIpc, IpcPeerCredentialsError } from '@push.rocks/smartipc';

const server = SmartIpc.createServer({
  id: 'renderer-control',
  socketPath: '\0my-app-control', // Linux abstract address (leading NUL)
  requirePeerCredentials: true,
  authorizeConnection: async ({ connection, peerCredentials, signal }) => {
    // Only `true` admits. `signal` aborts on timeout, server stop or peer close.
    return peerCredentials.uid === process.getuid!() && (await isExpectedChild(peerCredentials.pid, signal));
  },
  connectionAdmissionTimeoutMs: 5_000, // lookup + authorization bound (default 5000)
  maxPendingConnectionAdmissions: 16,  // excess connections are rejected (default 16)
});

server.on('connectionAdmitted', (connection) => {
  console.log(connection.id, connection.peerCredentials); // { pid, uid, gid }, frozen
});
server.on('connectionRejected', ({ connectionId, reason, peerCredentials, error }) => {
  // reason: 'unauthorized' | 'authorizer-failed' | 'credentials-unavailable'
  //       | 'admission-timeout' | 'admission-limit' | 'peer-closed' | 'server-stopping'
});

server.onMessage('startupReport', (payload, clientId, connection) => {
  // Authenticate with connection.peerCredentials, never with clientId:
  // clientId and every other envelope field are client-supplied.
  return { accepted: verifyNonce(payload.nonce, connection!.peerCredentials!.pid) };
});

await server.start(); // rejects with IpcPeerCredentialsError on unsupported setups

API:

  • IIpcServerOptions.requirePeerCredentials?: boolean, authorizeConnection?: TIpcConnectionAuthorizer, connectionAdmissionTimeoutMs?: number, maxPendingConnectionAdmissions?: number. The last three require requirePeerCredentials; otherwise start() rejects with code INVALID_OPTIONS.
  • TIpcConnectionAuthorizer = (context: IIpcConnectionAdmissionContext) => boolean | Promise<boolean> with context.connection, context.peerCredentials and context.signal.
  • IpcServerConnection: id (server-generated), peerCredentials (IIpcPeerCredentials | undefined), acceptedAt, getIsClosed(), and close(), which resolves once the socket has closed. The instance and its credentials are frozen. Servers (Unix-domain and TCP) pass the connection as the third argument of onMessage handlers, of 'message' (message, clientId, connection) and of 'clientConnect' events, and report it through 'connectionOpened' and 'connectionClosed' (IIpcConnectionClose). getClientInfo(clientId) additionally reports connectionId and peerCredentials of the connection that registered that client ID.
  • Server events 'connectionAdmitted' (IpcServerConnection) and 'connectionRejected' (IIpcConnectionRejection).
  • getPeerCredentialsSupport({ socketPath?, host?, port? }) returns { supported: true } or { supported: false, code, reason } without starting anything.
  • IpcPeerCredentialsError.code: INVALID_OPTIONS, UNSUPPORTED_PLATFORM, UNSUPPORTED_TRANSPORT, NATIVE_HELPER_UNAVAILABLE, LOOKUP_FAILED, LOOKUP_ABORTED.

Behaviour:

  • Explicit bind. Like every IpcServer, a server that requires peer credentials only ever listens. An address held by a live server makes start() reject with IpcServerBindError code ADDRESS_IN_USE; it never falls back to connecting as a client (see "Explicit Server Binding and Stale Sockets").
  • Platform support. Linux Unix-domain stream sockets on x64, where the native helper ships. TCP (host/port) reports UNSUPPORTED_TRANSPORT; other platforms, including Windows named pipes and non-x64 Linux, report UNSUPPORTED_PLATFORM. A missing or non-executable helper reports NATIVE_HELPER_UNAVAILABLE. Servers that do not require credentials expose connections whose peerCredentials is absent; SmartIPC never fabricates values.
  • Abstract addresses. A socketPath beginning with \0 binds in the Linux abstract namespace (isAbstractSocketPath() tells them apart). Nothing is created or unlinked on disk. socketMode cannot apply to an abstract address and makes start() reject. Abstract addresses have no file permissions, so any process in the same network namespace can connect; peer credentials are how such a server tells peers apart.
  • Cost. Node.js has no API for SO_PEERCRED. Each accepted connection is shared, as a documented child_process stdio stream, with one short-lived run of the packaged static Rust helper, which calls getsockopt(SO_PEERCRED) on the duplicated descriptor, prints the triple and exits. Expect a process spawn (typically a few milliseconds) per connection; the feature is opt-in per server and suited to control channels, not high connection rates.
  • Cancellation and closure. Server stop() stops accepting, aborts every pending admission (killing its helper with SIGKILL), waits until each helper has exited and its pipes are closed, then closes admitted connections. Admission timeouts and the pending-admission bound make the work per server finite.
  • Bounded framing. Admitted connections use the framing below. A frame whose declared length exceeds maxMessageSize closes that connection before the body is buffered.

Security notes:

  • The kernel records the credentials when the peer calls connect(). They describe that process even if the descriptor is later passed to another process, and a pid can be reused after the process exits. Validate whatever else identifies the peer (parent, executable, start time, a per-launch nonce) in the authorizer or in the first message.
  • A peer in another PID namespace makes the lookup fail (credentials-unavailable) instead of yielding pid 0. uid/gid values come from the kernel's mapping for the server's user namespace; unmapped IDs appear as the overflow IDs (normally 65534).
  • clientId, headers and payloads remain client-supplied. Replies go to the requesting connection and a registered client ID cannot be taken over by another connection (see the wire protocol), but the ID itself proves nothing about the peer. Authenticate with the connection argument.

Wire protocol

All SmartIPC transports use the same byte stream framing, so a non-Node client can talk to an IpcServer without the TypeScript client.

Connection. Connect to the server's Unix-domain stream socket (AF_UNIX, SOCK_STREAM). For an abstract address, the socket address is a NUL byte followed by the name bytes, and the address length covers exactly offsetof(sun_path) + 1 + name length (no trailing NUL). There are no hello bytes, magic number or protocol version marker: the first bytes on the stream are the first frame. A server that requires peer credentials does not read anything until the connection is admitted; data written earlier stays in the kernel socket buffer and is delivered after admission, or discarded when the connection is rejected and closed.

Frame. Every message in either direction is one frame:

| Offset | Size | Content | | --- | --- | --- | | 0 | 4 bytes | Body length N in bytes, unsigned 32-bit, big-endian (network byte order). The header itself is not counted. | | 4 | N bytes | UTF-8 encoded JSON object: the message envelope. |

N must not exceed the receiver's maxMessageSize (default 8388608, 8 MiB). Frames are not padded and there is no trailer. Each server connection keeps its own framing state, on Unix-domain sockets and TCP alike. Invalid UTF-8 is decoded with replacement characters.

Protocol errors. Both sides apply the same envelope validation. A client that receives an oversized frame, invalid JSON or an invalid envelope from the server closes its connection: 'disconnect' carries the reason protocol-error and the error, and the non-crashing 'connectionError' event (IIpcClientConnectionError) reports it. The client never skips bytes and continues, so frame alignment cannot be lost. Its configured autoReconnect policy then applies. A server closes the one connection that announces a frame larger than maxMessageSize, sends a body that is not valid JSON, or sends an invalid envelope (not a JSON object, id not a string, type not a non-empty string, or a present timestamp, correlationId or headers of the wrong type). The server reports it with a 'connectionClosed' event (reason protocol-error, the error attached); it never raises a server 'error' for peer input, and other connections are unaffected.

Envelope.

| Field | Type | Meaning | | --- | --- | --- | | id | string | Message ID chosen by the sender; unique per request. Required. | | type | string | Non-empty message type; the server dispatches onMessage(type, …) handlers on it. Required. | | timestamp | number, optional | Sender time in milliseconds since the Unix epoch (informational). SmartIPC always sends it. | | payload | any JSON value | Application data. | | correlationId | string, optional | Set on responses to the id of the request. | | headers | object, optional | clientId (string): the client ID the sender claims (informational, see below); requiresResponse (boolean): the sender awaits a response; error (string): set on error responses. |

Requests. Send headers.requiresResponse: true with an id. The server answers on the same connection with a frame whose type is <request type>_response, correlationId is the request id and payload is the handler result; if the handler throws, payload is null and headers.error carries the message. Replies always go to the connection the request arrived on, whatever headers.clientId says, and are never broadcast; a clientId is not needed to receive them.

Registration is optional for sending messages. It binds a client ID to the connection and makes the client visible to getClientIds(), sendToClient(), requestFromClient(), broadcast(), topic subscriptions and getClientInfo(): send a request of type __register__ with payload { "clientId": "<id>", "metadata": { … } }; the response type is __register___response with payload { "success": true, "clientId": "<id>" }. Registration fails (success: false) while another open connection holds that client ID, and a connection can hold only one client ID. Handlers receive the registered client ID of the connection; an unregistered connection's headers.clientId claim is used only when no other connection owns that ID, otherwise the message is attributed to unknown. Server pushes (sendToClient(), broadcast(), requestFromClient(), streams) are delivered only to the connection that registered the client ID.

Reserved types. __register__, __heartbeat__, __heartbeat_response__, __subscribe__, __unsubscribe__, __publish__ and __stream_init__, __stream_chunk__, __stream_end__, __stream_error__, __stream_cancel__ (stream chunks carry base64 data in payload.chunk). Ignore reserved types you do not implement.

Heartbeats. With heartbeat enabled (the default), server heartbeats are per connection. The server sends each connection a __heartbeat__ frame every heartbeatInterval and answers every __heartbeat__ it receives with __heartbeat_response__ (correlationId = heartbeat id). Any frame a connection sends counts as liveness. A connection that sends nothing for longer than heartbeatTimeout (after heartbeatInitialGracePeriodMs from opening) is closed with reason heartbeat-timeout; the listener and other connections are unaffected. A long-lived client therefore answers __heartbeat__ with __heartbeat_response__ (or sends its own __heartbeat__) at least once per heartbeatTimeout.

Example startup report (fire-and-forget) as bytes on the wire:

00 00 00 7a  {"id":"r1","type":"startupReport","timestamp":1767225600000,"payload":{"nonce":"abc"},"headers":{"clientId":"renderer-1"}}

Migrating from 2.x

SmartIPC 3.0.0 makes every server connection a first-class, isolated object. The changes below are breaking:

  • Targeted sends. In server mode, IpcChannel.sendMessage(), request() and sendStream() and the transports' send() take a target IpcServerConnection: an optional fourth argument, { connection }, IStreamSendOptions.connection, and a third argument for cancelIncomingStream()/cancelOutgoingStream(). A server never broadcasts an untargeted message; such a send fails. IpcServer.sendToClient(), broadcast(), broadcastTo(), requestFromClient() and sendStreamToClient() remain the explicit push APIs and target the connection that registered the client ID. cancelIncomingStreamFromClient() and cancelOutgoingStreamToClient() throw for an unknown client ID.
  • Replies follow the request's connection. Responses, heartbeat answers and stream errors go to the connection the request arrived on. headers.clientId no longer routes anything. A response to a server-initiated request is accepted only from the connection it was sent to. Stream state is scoped to its connection.
  • Registration rules. __register__ fails while another open connection holds the client ID. A connection can register only one client ID. Handlers receive the client ID registered on the connection; an unregistered header claim of another connection's ID is attributed to 'unknown'. Topic subscriptions need a registered client on that connection.
  • Removed IpcChannel event 'clientDisconnected'. Use 'connectionClosed' (IIpcConnectionClose { connection, reason, error? }), which IpcServer also emits. IpcServer's 'clientDisconnect' is unchanged.
  • Server heartbeats are per connection. A silent connection is closed with reason heartbeat-timeout; the listener is never closed, and a server no longer emits 'error' for heartbeat timeouts. The server 'heartbeatTimeout' event is now (error, clientId | undefined, connection) instead of (error, 'server').
  • Peer input never raises a server 'error'. Invalid JSON, invalid envelopes and oversized frames close only that connection ('connectionClosed', reason protocol-error). Envelopes need a string id and a non-empty string type. Server-side frames are limited to maxMessageSize.
  • Handler failures are contained. Synchronous throws and rejections of fire-and-forget handlers emit 'handlerError' (IIpcHandlerError) on IpcServer, IpcClient and IpcChannel. They no longer propagate from socket callbacks or surface as 'error' events.
  • No client fallback for servers. IpcServer.start() always binds as a server (Unix path, abstract address, named pipe or TCP). A live address rejects with IpcServerBindError (ADDRESS_IN_USE). A stale socket file is replaced, and a non-socket path is refused with PATH_NOT_SOCKET. The autoCleanupSocketFile option is removed because stale socket handling is now always on and safe. createTransport() and the IpcChannel constructor take an IIpcServerBinding as their second argument.
  • Clients never become servers. The IpcClient (and client transport) fallback that started a server when nothing was listening is removed, together with the clientOnly option and the SMARTIPC_CLIENT_ONLY environment variable; every client is client-only. With no server, connect() rejects with IpcClientConnectError (SERVER_UNAVAILABLE) after its connectRetry attempts. IpcChannel.connect() now rejects on failure instead of emitting 'error' and scheduling a background reconnect. A failed reconnect attempt of a dropped connection now schedules the next attempt within maxReconnectAttempts; in 2.x it stopped after one failed attempt.
  • No 'error' events. IpcServer, IpcClient, IpcChannel and the transports never emit 'error'. Remove 'error' listeners that only kept processes alive. Failures arrive as promise rejections or as these typed events:
    • Client connection failures. Socket errors (ECONNRESET, EPIPE, ETIMEDOUT and the like), the TCP timeout option, protocol violations by the server (invalid JSON, invalid envelopes, oversized frames) and a failed re-registration after a reconnect all close the client connection. 'disconnect' is emitted as (reason, error) with a TIpcClientDisconnectReason (peer-closed, socket-error, socket-timeout, protocol-error, heartbeat-timeout, registration-failed, closed-by-client), and 'connectionError' (IIpcClientConnectionError { reason, error }) is emitted on IpcClient and IpcChannel. The autoReconnect policy applies. An oversized frame no longer resets the buffer and continues mid-stream, and incoming envelopes are validated like on the server.
    • Client heartbeat timeouts. These emit 'heartbeatTimeout' instead of 'error'. heartbeatThrowOnTimeout is renamed to heartbeatCloseOnTimeout, with no alias, and only affects clients. With true (the default), the connection is dropped with reason heartbeat-timeout and the reconnect policy applies. With false, only the event is emitted and the connection is kept.
    • Reconnect exhaustion. Reaching maxReconnectAttempts emits 'reconnectFailed' instead of 'error'.
    • Server broadcasts. broadcast() and broadcastTo() resolve to the client IDs they could not reach instead of emitting 'error' per client. IpcServer no longer forwards channel errors as 'error' (clientId 'server').
    • Disconnect. IpcClient.disconnect() also cancels a pending automatic reconnect, and emits 'disconnect' once, with reason closed-by-client.
  • Socket file ownership. Filesystem sockets are bound under a 9-byte temporary name, with socketMode applied, then hard-linked to the requested path. stop() removes the path only while it still holds the socket this server bound. The requested path must fit the socket address limit (107 bytes on Linux), and its directory part must leave room for the temporary name (at most 97 bytes on Linux). Otherwise start() rejects with IpcServerBindError PATH_TOO_LONG. A socketMode that cannot be applied now fails start(); in 2.x the failure was ignored.
  • disconnectClient() scope. It, and clientIdleTimeout, close only that client's connection (reason closed-by-server). In 2.x they disconnected the shared server channel.
  • TCP connection objects. TCP servers now create an IpcServerConnection per accepted socket (peerCredentials absent), with per-connection framing. One client's socket error or timeout closes only that connection (socket-error, socket-timeout). TCP pushes are no longer broadcast to every client.
  • New constructor argument. IpcServerConnection takes a fifth constructor argument (its close callback). Connections are created by the transports; application code receives them from events and handlers.
  • Transport API. Custom code that drives the transports directly sees these signatures:
    • connect(signal?: AbortSignal): in client mode, aborting the signal stops the attempt; a socket that connects afterwards is destroyed and never attached.
    • send(message, target?): server mode requires the target connection.
    • getRole(): returns 'client' | 'server' | 'none'.
    • Server mode: getServerConnections() and closeConnection(connection, reason, error?).
    • Client mode: dropClientConnection(reason, error?) closes the current connection with a TIpcClientDisconnectReason, and the owner's reconnect policy applies.
  • New events. IpcServer emits 'connectionOpened', 'connectionClosed', 'handlerError', and for peer-credential servers 'connectionAdmitted' and 'connectionRejected'.

🎮 Core Concepts

Transport Types

SmartIPC supports multiple transport mechanisms, automatically selecting the best one for your platform:

// TCP Socket (cross-platform, network-capable)
const tcpServer = SmartIpc.createServer({
  id: 'tcp-service',
  host: 'localhost',
  port: 9876
});

// Unix Domain Socket (Linux/macOS, fastest local IPC)
const unixServer = SmartIpc.createServer({
  id: 'unix-service',
  socketPath: '/tmp/my-app.sock'
});

// Windows Named Pipe (Windows optimal)
// Automatically used on Windows when socketPath is provided
const windowsServer = SmartIpc.createServer({
  id: 'pipe-service',
  socketPath: '\\\\.\\pipe\\my-app-pipe'
});

Message Patterns

🔥 Fire and Forget

Send messages without waiting for a response:

// Server
server.onMessage('log', (data, clientId) => {
  console.log(`[${clientId}] ${data.level}:`, data.message);
  // No return needed
});

// Client
await client.sendMessage('log', { 
  level: 'info',
  message: 'User logged in',
  timestamp: Date.now()
});

📞 Request/Response

RPC-style communication with type safety:

interface UserRequest {
  userId: string;
  fields?: string[];
}

interface UserResponse {
  id: string;
  name: string;
  email?: string;
  createdAt: number;
}

// Server
server.onMessage<UserRequest, UserResponse>('getUser', async (data) => {
  const user = await db.getUser(data.userId);
  return {
    id: user.id,
    name: user.name,
    email: data.fields?.includes('email') ? user.email : undefined,
    createdAt: user.createdAt
  };
});

// Client - with timeout
const user = await client.request<UserRequest, UserResponse>(
  'getUser',
  { userId: '123', fields: ['email'] },
  { timeout: 5000 }
);

📢 Pub/Sub Pattern

Topic-based message broadcasting:

// Subscribers
const subscriber1 = SmartIpc.createClient({
  id: 'events-service',
  socketPath: '/tmp/events.sock'
});

await subscriber1.connect();
await subscriber1.subscribe('user.login', (data) => {
  console.log('User logged in:', data);
});

// Publisher
const publisher = SmartIpc.createClient({
  id: 'events-service',
  socketPath: '/tmp/events.sock'
});

await publisher.connect();
await publisher.publish('user.login', { 
  userId: '123',
  ip: '192.168.1.1',
  timestamp: Date.now()
});

💪 Advanced Features

📦 Streaming Large Data & Files

SmartIPC supports efficient, backpressure-aware streaming of large payloads using chunked messages. Streams work both directions and emit a high-level stream event for consumption.

Client → Server streaming:

// Server side: receive stream
server.on('stream', async (info, readable) => {
  if (info.meta?.type === 'file') {
    console.log('Receiving file', info.meta.basename, 'from', info.clientId);
  }
  // Pipe to disk or process chunks
  await SmartIpc.pipeStreamToFile(readable, '/tmp/incoming.bin');
});

// Client side: send a stream
const readable = fs.createReadStream('/path/to/local.bin');
await client.sendStream(readable, {
  meta: { type: 'file', basename: 'local.bin' },
  chunkSize: 64 * 1024 // optional, defaults to 64k
});

Server → Client streaming:

client.on('stream', async (info, readable) => {
  console.log('Got stream from server', info.meta);
  await SmartIpc.pipeStreamToFile(readable, '/tmp/from-server.bin');
});

await server.sendStreamToClient(client.getClientId(), fs.createReadStream('/path/server.bin'), {
  meta: { type: 'file', basename: 'server.bin' }
});

High-level helpers for files:

// Client → Server
await client.sendFile('/path/to/bigfile.iso');

// Server → Client
await server.sendFileToClient(clientId, '/path/to/backup.tar');

// Save an incoming stream to a file (both sides)
server.on('stream', async (info, readable) => {
  await SmartIpc.pipeStreamToFile(readable, '/data/uploaded/' + info.meta?.basename);
});

Events & metadata:

  • channel/server/client emit stream with (info, readable)
  • info contains: streamId, meta (your metadata, e.g., filename/size), headers, and clientId (if available)

API summary:

  • Client: sendStream(readable, opts), sendFile(filePath, opts), cancelOutgoingStream(id), cancelIncomingStream(id)
  • Server: sendStreamToClient(clientId, readable, opts), sendFileToClient(clientId, filePath, opts), cancelIncomingStreamFromClient(clientId, id), cancelOutgoingStreamToClient(clientId, id)
  • Utility: SmartIpc.pipeStreamToFile(readable, filePath)

Concurrency and cancelation:

// Limit concurrent streams per connection
const server = SmartIpc.createServer({
  id: 'svc', socketPath: '/tmp/svc.sock', maxConcurrentStreams: 2
});

// Cancel a stream from the receiver side
server.on('stream', (info, readable) => {
  if (info.meta?.shouldCancel) {
    (server as any).primaryChannel.cancelIncomingStream(info.streamId, { clientId: info.clientId });
  }
});

Notes:

  • Streaming uses chunked messages under the hood and respects socket backpressure.
  • Include meta to share context like filename/size; it’s delivered with the stream event.
  • Configure maxConcurrentStreams (default: 32) to guard resources.

🏁 Server Readiness Detection

Eliminate race conditions in tests and production:

const server = SmartIpc.createServer({
  id: 'my-service',
  socketPath: '/tmp/my-service.sock',
});

// Option 1: Wait for full readiness
await server.start({ readyWhen: 'accepting' });
// Server is now FULLY ready to accept connections

// Option 2: Use ready event
server.on('ready', () => {
  console.log('Server is ready!');
  startClients();
});

await server.start();

// Option 3: Check readiness state
if (server.getIsReady()) {
  console.log('Ready to rock! 🎸');
}

🔄 Smart Connection Retry

Never lose messages due to temporary connection issues:

const client = SmartIpc.createClient({
  id: 'resilient-client',
  socketPath: '/tmp/service.sock',
  connectRetry: {
    enabled: true,
    initialDelay: 100,      // Start with 100ms
    maxDelay: 1500,         // Cap at 1.5 seconds
    maxAttempts: 20,        // Try 20 times
    totalTimeout: 15000     // Give up after 15 seconds total
  },
  registerTimeoutMs: 8000   // Registration handshake timeout
});

// Will retry automatically if server isn't ready yet
await client.connect({ 
  waitForReady: true,       // Wait for server to exist
  waitTimeout: 10000        // Wait up to 10 seconds
});

🛑 Clients Never Become Servers

An IpcClient only ever connects; it never binds an address or starts a server. When nothing is listening, connect() rejects with IpcClientConnectError (code SERVER_UNAVAILABLE), after the attempts allowed by connectRetry.

import { SmartIpc, IpcClientConnectError } from '@push.rocks/smartipc';

const client = SmartIpc.createClient({
  id: 'my-service',
  socketPath: '/tmp/my-service.sock',
  clientId: 'my-cli',
  connectRetry: { enabled: false }  // fail fast; enable to wait for a starting server
});

try {
  await client.connect();
} catch (err) {
  if (err instanceof IpcClientConnectError && err.code === 'SERVER_UNAVAILABLE') {
    console.error(err.message); // "No server is listening at /tmp/my-service.sock (ENOENT)"
  }
}
  • connectRetry governs the initial connect: with enabled: true, the client retries with backoff up to maxAttempts and totalTimeout, then rejects with the last error.
    • Concurrent connect() calls share one attempt.
    • disconnect() aborts a connect() in progress, including its retry delay. That connect() rejects with IpcClientConnectError code CONNECT_ABORTED, never registers, and never emits 'connect'.
    • Registration is single-flight per connection, so a manual connect() during a reconnect's registration shares it and 'connect' is emitted once.
  • autoReconnect governs established connections that drop: the client retries with the configured backoff until maxReconnectAttempts, and re-registers once reconnected.
    • Attempt counting. The attempt counter and its backoff reset only once a reconnect is fully established, meaning the socket is connected and registration succeeded. A refused registration, or a protocol error right after connecting, counts as a failed attempt. When the limit is reached, 'reconnectFailed' is emitted.
    • One attempt at a time. A manual connect() cancels a pending reconnect, and disconnect() cancels an in-flight one. A socket that finishes connecting after that is closed immediately, never attached, so a client never holds two sockets or registers twice.

💓 Graceful Heartbeat Monitoring

Keep connections alive without crashing on timeouts:

const server = SmartIpc.createServer({
  id: 'monitored-service',
  socketPath: '/tmp/monitored.sock',
  heartbeat: true,
  heartbeatInterval: 3000,
  heartbeatTimeout: 10000,
  heartbeatInitialGracePeriodMs: 5000,    // Grace period for startup
});

// Server heartbeats are per connection: a connection that stays silent past
// heartbeatTimeout is closed; the server keeps listening.
server.on('heartbeatTimeout', (error, clientId, connection) => {
  console.log(`Connection ${connection.id} (${clientId ?? 'unregistered'}) timed out`);
});

// Client configuration
const client = SmartIpc.createClient({
  id: 'monitored-service',
  socketPath: '/tmp/monitored.sock',
  heartbeat: true,
  heartbeatInterval: 3000,
  heartbeatTimeout: 10000,
  heartbeatInitialGracePeriodMs: 5000,
  // true (default): drop the connection on timeout and reconnect per policy;
  // false: only emit 'heartbeatTimeout' and keep the connection
  heartbeatCloseOnTimeout: true
});

client.on('heartbeatTimeout', () => {
  console.log('Heartbeat timeout detected; the client reconnects');
});

🧹 Explicit Server Binding and Stale Sockets

An IpcServer always binds as a server; it never connects to an existing server instead. When a Unix socket path already exists, start() probes it:

  • a live server accepts the probe: start() rejects with IpcServerBindError code ADDRESS_IN_USE;
  • the connection is refused (a stale socket file left by a crashed process): that socket file is removed and the bind is retried once;
  • the path is not a socket (for example a regular file): start() rejects with code PATH_NOT_SOCKET and the file is left untouched.

Abstract addresses (\0name) and TCP ports have no stale state: if they are taken, start() rejects with ADDRESS_IN_USE.

A filesystem socket is published atomically:

  • Bind. The server listens on a private temporary name in the same directory (. plus 8 random base64url characters, 9 bytes) and applies socketMode to it. It then hard-links the socket to the requested path, which fails if anything already exists there, and removes the temporary name. The public path therefore never exists with default permissions, and libuv's unconditional unlink-on-close never touches it.
  • Path length. Both the requested path and the temporary path must fit the Unix socket address limit L (107 bytes on Linux, 103 on macOS). The requested path must be at most L bytes, and the directory part must be at most L - 1 - 9 bytes (97 on Linux, 93 on macOS). Otherwise start() rejects with IpcServerBindError code PATH_TOO_LONG before anything is created, and the message states the limit. Relative paths count as written.
  • Failure cleanup. If socketMode cannot be applied, or any step after listening fails, start() rejects. The listener is closed, and neither the temporary name nor a public path this server created is left behind.
  • Stop. stop() removes the socket file only while it is still the socket this server bound (same device, inode and change time). A path that another server has since taken over is left alone.
import { SmartIpc, IpcServerBindError } from '@push.rocks/smartipc';

const server = SmartIpc.createServer({
  id: 'clean-service',
  socketPath: '/tmp/service.sock',
  socketMode: 0o600 // Set socket permissions (Unix only)
});

try {
  await server.start();
} catch (error) {
  if (error instanceof IpcServerBindError && error.code === 'ADDRESS_IN_USE') {
    // Another instance is already running.
  }
  throw error;
}

📊 Real-time Metrics

Monitor your IPC performance:

// Server stats
const serverStats = server.getStats();
console.log({
  isRunning: serverStats.isRunning,
  connectedClients: serverStats.connectedClients,
  totalConnections: serverStats.totalConnections,
  metrics: {
    messagesSent: serverStats.metrics.messagesSent,
    messagesReceived: serverStats.metrics.messagesReceived,
    errors: serverStats.metrics.errors
  }
});

// Client stats
const clientStats = client.getStats();
console.log({
  connected: clientStats.connected,
  reconnectAttempts: clientStats.reconnectAttempts,
  metrics: clientStats.metrics
});

// Get specific client info
const clientInfo = server.getClientInfo('client-123');
console.log({
  connectedAt: new Date(clientInfo.connectedAt),
  lastActivity: new Date(clientInfo.lastActivity),
  metadata: clientInfo.metadata
});

🎯 Broadcasting

Send messages to multiple clients:

// Broadcast to all connected clients
await server.broadcast('announcement', { 
  message: 'Server will restart in 5 minutes',
  severity: 'warning'
});

// Send to specific clients
await server.broadcastTo(
  ['client-1', 'client-2'],
  'private-message',
  { content: 'This is just for you two' }
);

// Send to one client
await server.sendToClient('client-1', 'direct', {
  data: 'Personal message'
});

🧪 Testing Utilities

SmartIPC includes powerful helpers for testing:

Wait for Server

import { SmartIpc } from '@push.rocks/smartipc';

// Start your server in another process
const serverProcess = spawn('node', ['server.js']);

// Wait for it to be ready
const abortController = new AbortController();
await SmartIpc.waitForServer({
  socketPath: '/tmp/test.sock',
  timeoutMs: 10000,
  signal: abortController.signal // optional
});

// Now safe to connect clients
const client = SmartIpc.createClient({
  id: 'test-client',
  socketPath: '/tmp/test.sock'
});
await client.connect();
  • waitForServer() probes with short-lived clients that connect and register, one every 200 ms, until one succeeds or timeoutMs elapses.
  • The optional signal is checked on every pass. Aborting it disconnects the running probe immediately, starts no further probe, and rejects with IpcClientConnectError code CONNECT_ABORTED.
  • client.connect({ waitForReady: true }) passes its own attempt's signal, so client.disconnect() ends the wait the same way.

Spawn and Connect

// Helper that spawns a server and connects a client
const { client, serverProcess } = await SmartIpc.spawnAndConnect({
  serverScript: './server.js',
  socketPath: '/tmp/test.sock',
  clientId: 'test-client',
  connectRetry: {
    enabled: true,
    maxAttempts: 10
  }
});

// Use the client
const response = await client.request('ping', {});

// Cleanup
await client.disconnect();
serverProcess.kill();

🎭 Event Handling

SmartIPC provides comprehensive event emitters:

// Server events
server.on('start', () => console.log('Server started'));
server.on('ready', () => console.log('Server ready for connections'));
server.on('clientConnect', (clientId, metadata) => {
  console.log(`Client ${clientId} connected with metadata:`, metadata);
});
server.on('clientDisconnect', (clientId) => {
  console.log(`Client ${clientId} disconnected`);
});
// Every server connection, registered or not (IpcServerConnection)
server.on('connectionOpened', (connection) => console.log('opened', connection.id));
server.on('connectionClosed', ({ connection, reason, error }) => {
  // reason: 'peer-closed' | 'socket-error' | 'protocol-error' | 'heartbeat-timeout'
  //       | 'socket-timeout' | 'closed-by-server' | 'server-stopping'
  console.log('closed', connection.id, reason, error?.message);
});
server.on('handlerError', ({ error, messageType, clientId }) => {
  console.error(`Handler ${messageType} failed for ${clientId}:`, error);
});

// Client events  
client.on('connect', () => console.log('Connected to server'));
client.on('disconnect', (reason, error) => {
  // reason: 'peer-closed' | 'socket-error' | 'socket-timeout' | 'protocol-error'
  //       | 'heartbeat-timeout' | 'registration-failed' | 'closed-by-client'
  console.log('Disconnected from server:', reason, error?.message);
});
client.on('connectionError', ({ reason, error }) => {
  console.warn(`Connection failed (${reason}):`, error.message);
});
client.on('reconnecting', ({ attempt, delay }) => {
  console.log(`Reconnection attempt ${attempt} in ${delay} ms`);
});
client.on('reconnectFailed', (error) => {
  console.error('Gave up reconnecting:', error.message);
});
client.on('heartbeatTimeout', (error) => {
  console.warn('Heartbeat timeout:', error);
});

What can still emit 'error'. None of SmartIPC's own objects: IpcServer, IpcClient, IpcChannel and the transports never emit 'error', so a process needs no 'error' listener to survive peer or network behaviour. Failures of calls you make (start(), connect(), request(), sendMessage(), the send and stream APIs) reject their promises. Everything caused by a peer or the network is reported through the typed events above. The only 'error' events left are standard Node.js stream semantics: an incoming stream ('stream' event readable) is destroyed with an error when the sender reports a stream error, the stream is cancelled, or its connection closes. Attach an 'error' listener or use stream.pipeline() on readables you consume, as with any Node.js stream.

🛡️ Error Handling

Robust error handling with detailed error information:

// Client-side error handling
try {
  const response = await client.request('riskyOperation', data, {
    timeout: 5000
  });
} catch (error) {
  if (error.message.includes('timeout')) {
    console.error('Request timed out');
  } else if (error.message.includes('Failed to register')) {
    console.error('Could not register with server');
  } else {
    console.error('Unknown error:', error);
  }
}

// Server-side error boundaries
server.onMessage('process', async (data, clientId) => {
  try {
    return await riskyProcessing(data);
  } catch (error) {
    console.error(`Processing failed for ${clientId}:`, error);
    throw error;  // Will be sent back to client as error
  }
});

// Handler failures never crash the process and never surface as 'error'
server.on('handlerError', ({ error, messageType, requiresResponse, clientId, connection }) => {
  console.error(`Handler for ${messageType} failed for ${clientId}:`, error);
});

A handler that throws synchronously or returns a rejected promise is contained:

  • Requests: the peer receives an error response on its own connection (payload: null, headers.error set). client.request() rejects with that message.
  • Fire-and-forget messages: the failure is reported through 'handlerError' (IIpcHandlerError). Unlike 'error', this event has no crash semantics when nobody listens.
  • The connection stays open in both cases. A handler failure is a fault in the receiving application, not a protocol violation by the peer. Protocol violations close the connection ('connectionClosed', reason protocol-error). IpcClient handlers behave the same way and emit 'handlerError' on the client.

🏗️ Architecture

SmartIPC uses a clean, layered architecture:

┌─────────────────────────────────────────┐
│          Your Application               │
│        (Business logic)                 │
└─────────────────────────────────────────┘
                   ↕
┌─────────────────────────────────────────┐
│         IpcServer / IpcClient           │
│   (High-level API, Message routing)     │
└─────────────────────────────────────────┘
                   ↕
┌─────────────────────────────────────────┐
│            IpcChannel                   │
│  (Connection management, Heartbeat,     │
│   Reconnection, Request/Response)       │
└─────────────────────────────────────────┘
                   ↕
┌─────────────────────────────────────────┐
│          Transport Layer                │
│  (TCP, Unix Socket, Named Pipe)         │
│     (Framing, buffering, I/O)           │
└─────────────────────────────────────────┘

🎯 Common Use Cases

Microservices Communication

// API Gateway
const gateway = SmartIpc.createServer({
  id: 'api-gateway',
  socketPath: '/tmp/gateway.sock'
});

// User Service
const userService = SmartIpc.createClient({
  id: 'api-gateway',
  socketPath: '/tmp/gateway.sock',
  clientId: 'user-service'
});

// Order Service  
const orderService = SmartIpc.createClient({
  id: 'api-gateway',
  socketPath: '/tmp/gateway.sock',
  clientId: 'order-service'
});

Worker Process Management

// Main process
const server = SmartIpc.createServer({
  id: 'main',
  socketPath: '/tmp/workers.sock'
});

server.onMessage('job-complete', (result, workerId) => {
  console.log(`Worker ${workerId} completed job:`, result);
});

// Worker process
const worker = SmartIpc.createClient({
  id: 'main',
  socketPath: '/tmp/workers.sock',
  clientId: `worker-${process.pid}`
});

await worker.sendMessage('job-complete', {
  jobId: '123',
  result: processedData
});

Real-time Event Distribution

// Event bus
const eventBus = SmartIpc.createServer({
  id: 'event-bus',
  socketPath: '/tmp/events.sock'
});

// Services subscribe to events
const analyticsService = SmartIpc.createClient({
  id: 'event-bus',
  socketPath: '/tmp/events.sock'
});

await analyticsService.subscribe('user.*', (event) => {
  trackEvent(event);
});

📈 Performance

SmartIPC is optimized for high throughput and low latency:

| Transport | Messages/sec | Avg Latency | Use Case | |-----------|-------------|-------------|----------| | Unix Socket | 150,000+ | < 0.1ms | Local high-performance IPC (Linux/macOS) | | Named Pipe | 120,000+ | < 0.15ms | Windows local IPC | | TCP (localhost) | 100,000+ | < 0.2ms | Local network-capable IPC | | TCP (network) | 50,000+ | < 1ms | Distributed systems |

  • Memory efficient: Streaming support for large payloads
  • CPU efficient: Event-driven, non-blocking I/O

🔧 Requirements

  • Node.js >= 14.x for the IPC transports; Node.js 24 through 26 on Linux x64 for NamedMutex and probeFileOpenExclusivity()
  • Linux x64 for requirePeerCredentials (Unix-domain sockets only); abstract addresses need Node.js >= 20.8
  • TypeScript 6.x via the declared GitZone toolchain (for development)
  • Unix-like OS (Linux, macOS) or Windows

License and Legal Information

This repository contains open-source code licensed under the MIT License. A copy of the license can be found in the license.md file. Notices for dependencies incorporated into the packaged Rust helper are provided in third-party-notices.md.

Please note: The MIT License does not grant permission to use the trade names, trademarks, service marks, or product names of the project, except as required for reasonable and customary use in describing the origin of the work and reproducing the content of the NOTICE file.

Trademarks

This project is owned and maintained by Task Venture Capital GmbH. The names and logos associated with Task Venture Capital GmbH and any related products or services are trademarks of Task Venture Capital GmbH or third parties, and are not included within the scope of the MIT license granted herein.

Use of these trademarks must comply with Task Venture Capital GmbH's Trademark Guidelines or the guidelines of the respective third-party owners, and any usage must be approved in writing. Third-party trademarks used herein are the property of their respective owners and used only in a descriptive manner, e.g. for an implementation of an API or similar.

Company Information

Task Venture Capital GmbH Registered at District Court Bremen HRB 35230 HB, Germany

For any legal inquiries or further information, please contact us via email at [email protected].

By using this repository, you acknowledge that you have read this section, agree to comply with its terms, and understand that the licensing of the code does not imply endorsement by Task Venture Capital GmbH of any derivative works.