@push.rocks/smartipc
v3.0.0
Published
A library for node inter process communication, providing an easy-to-use API for IPC.
Maintainers
Readme
@push.rocks/smartipc 🚀
Rock-solid IPC and local named mutexes for Node.js
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 -
NamedMutexuses the packaged Rust helper instead of a Node native addon and reportsNATIVE_BACKEND_UNAVAILABLEif 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 setupsAPI:
IIpcServerOptions.requirePeerCredentials?: boolean,authorizeConnection?: TIpcConnectionAuthorizer,connectionAdmissionTimeoutMs?: number,maxPendingConnectionAdmissions?: number. The last three requirerequirePeerCredentials; otherwisestart()rejects with codeINVALID_OPTIONS.TIpcConnectionAuthorizer = (context: IIpcConnectionAdmissionContext) => boolean | Promise<boolean>withcontext.connection,context.peerCredentialsandcontext.signal.IpcServerConnection:id(server-generated),peerCredentials(IIpcPeerCredentials | undefined),acceptedAt,getIsClosed(), andclose(), 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 ofonMessagehandlers, of'message'(message, clientId, connection) and of'clientConnect'events, and report it through'connectionOpened'and'connectionClosed'(IIpcConnectionClose).getClientInfo(clientId)additionally reportsconnectionIdandpeerCredentialsof 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 makesstart()reject withIpcServerBindErrorcodeADDRESS_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) reportsUNSUPPORTED_TRANSPORT; other platforms, including Windows named pipes and non-x64 Linux, reportUNSUPPORTED_PLATFORM. A missing or non-executable helper reportsNATIVE_HELPER_UNAVAILABLE. Servers that do not require credentials expose connections whosepeerCredentialsis absent; SmartIPC never fabricates values. - Abstract addresses. A
socketPathbeginning with\0binds in the Linux abstract namespace (isAbstractSocketPath()tells them apart). Nothing is created or unlinked on disk.socketModecannot apply to an abstract address and makesstart()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 documentedchild_processstdio stream, with one short-lived run of the packaged static Rust helper, which callsgetsockopt(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 withSIGKILL), 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
maxMessageSizecloses 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 pid0.uid/gidvalues come from the kernel's mapping for the server's user namespace; unmapped IDs appear as the overflow IDs (normally65534). clientId,headersand 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 theconnectionargument.
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()andsendStream()and the transports'send()take a targetIpcServerConnection: an optional fourth argument,{ connection },IStreamSendOptions.connection, and a third argument forcancelIncomingStream()/cancelOutgoingStream(). A server never broadcasts an untargeted message; such a send fails.IpcServer.sendToClient(),broadcast(),broadcastTo(),requestFromClient()andsendStreamToClient()remain the explicit push APIs and target the connection that registered the client ID.cancelIncomingStreamFromClient()andcancelOutgoingStreamToClient()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.clientIdno 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
IpcChannelevent'clientDisconnected'. Use'connectionClosed'(IIpcConnectionClose { connection, reason, error? }), whichIpcServeralso 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', reasonprotocol-error). Envelopes need a stringidand a non-empty stringtype. Server-side frames are limited tomaxMessageSize. - Handler failures are contained. Synchronous throws and rejections of fire-and-forget handlers emit
'handlerError'(IIpcHandlerError) onIpcServer,IpcClientandIpcChannel. 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 withIpcServerBindError(ADDRESS_IN_USE). A stale socket file is replaced, and a non-socket path is refused withPATH_NOT_SOCKET. TheautoCleanupSocketFileoption is removed because stale socket handling is now always on and safe.createTransport()and theIpcChannelconstructor take anIIpcServerBindingas 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 theclientOnlyoption and theSMARTIPC_CLIENT_ONLYenvironment variable; every client is client-only. With no server,connect()rejects withIpcClientConnectError(SERVER_UNAVAILABLE) after itsconnectRetryattempts.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 withinmaxReconnectAttempts; in 2.x it stopped after one failed attempt. - No
'error'events.IpcServer,IpcClient,IpcChanneland 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
timeoutoption, 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 aTIpcClientDisconnectReason(peer-closed,socket-error,socket-timeout,protocol-error,heartbeat-timeout,registration-failed,closed-by-client), and'connectionError'(IIpcClientConnectionError { reason, error }) is emitted onIpcClientandIpcChannel. TheautoReconnectpolicy 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'.heartbeatThrowOnTimeoutis renamed toheartbeatCloseOnTimeout, with no alias, and only affects clients. Withtrue(the default), the connection is dropped with reasonheartbeat-timeoutand the reconnect policy applies. Withfalse, only the event is emitted and the connection is kept. - Reconnect exhaustion. Reaching
maxReconnectAttemptsemits'reconnectFailed'instead of'error'. - Server broadcasts.
broadcast()andbroadcastTo()resolve to the client IDs they could not reach instead of emitting'error'per client.IpcServerno longer forwards channel errors as'error'(clientId 'server'). - Disconnect.
IpcClient.disconnect()also cancels a pending automatic reconnect, and emits'disconnect'once, with reasonclosed-by-client.
- Client connection failures. Socket errors (ECONNRESET, EPIPE, ETIMEDOUT and the like), the TCP
- Socket file ownership. Filesystem sockets are bound under a 9-byte temporary name, with
socketModeapplied, 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). Otherwisestart()rejects withIpcServerBindErrorPATH_TOO_LONG. AsocketModethat cannot be applied now failsstart(); in 2.x the failure was ignored. disconnectClient()scope. It, andclientIdleTimeout, close only that client's connection (reasonclosed-by-server). In 2.x they disconnected the shared server channel.- TCP connection objects. TCP servers now create an
IpcServerConnectionper accepted socket (peerCredentialsabsent), with per-connection framing. One client's socket error ortimeoutcloses only that connection (socket-error,socket-timeout). TCP pushes are no longer broadcast to every client. - New constructor argument.
IpcServerConnectiontakes 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()andcloseConnection(connection, reason, error?). - Client mode:
dropClientConnection(reason, error?)closes the current connection with aTIpcClientDisconnectReason, and the owner's reconnect policy applies.
- New events.
IpcServeremits'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/clientemitstreamwith(info, readable)infocontains:streamId,meta(your metadata, e.g., filename/size),headers, andclientId(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
metato share context like filename/size; it’s delivered with thestreamevent. - 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)"
}
}connectRetrygoverns the initial connect: withenabled: true, the client retries with backoff up tomaxAttemptsandtotalTimeout, then rejects with the last error.- Concurrent
connect()calls share one attempt. disconnect()aborts aconnect()in progress, including its retry delay. Thatconnect()rejects withIpcClientConnectErrorcodeCONNECT_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.
- Concurrent
autoReconnectgoverns established connections that drop: the client retries with the configured backoff untilmaxReconnectAttempts, 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, anddisconnect()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.
- 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,
💓 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 withIpcServerBindErrorcodeADDRESS_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 codePATH_NOT_SOCKETand 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 appliessocketModeto 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 mostLbytes, and the directory part must be at mostL - 1 - 9bytes (97 on Linux, 93 on macOS). Otherwisestart()rejects withIpcServerBindErrorcodePATH_TOO_LONGbefore anything is created, and the message states the limit. Relative paths count as written. - Failure cleanup. If
socketModecannot 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 ortimeoutMselapses.- The optional
signalis checked on every pass. Aborting it disconnects the running probe immediately, starts no further probe, and rejects withIpcClientConnectErrorcodeCONNECT_ABORTED. client.connect({ waitForReady: true })passes its own attempt's signal, soclient.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.errorset).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', reasonprotocol-error).IpcClienthandlers 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
NamedMutexandprobeFileOpenExclusivity() - 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.
