@solncebro/websocket-engine
v0.6.1
Published
Reliable WebSocket client with reconnect, heartbeat and typed messages, and a reusable signal relay (server + client) under ./relay
Downloads
325
Maintainers
Readme
@solncebro/websocket-engine
Reliable, environment-agnostic WebSocket client with automatic reconnection, application-level heartbeat plus idle (stale) detection, typed messages, optional authentication phase, and request/response pattern. Built on the global WebSocket, so it runs unchanged in Node.js 22+, browsers, and React Native — no ws dependency.
Installation
yarn add @solncebro/websocket-enginenpm install @solncebro/websocket-engineFeatures
- Automatic reconnection — exponential backoff with jitter; fast reconnect for specific close codes (1001, 1006, 1011–1014)
- Heartbeat — application-level ping, either JSON (e.g. Bybit
{ op: 'ping' }, Binance{ method: 'LIST_SUBSCRIPTIONS' }) or bare text (e.g. a plainpingstring), paired with idle detection: the socket reconnects if no message arrives withinstaleThreshold. No TCP control frames are used, so it works on any WHATWGWebSocket. - Connection timeout — handshake timeout with retry
- Auth phase — optional
onOpenasync callback for authentication before the connection is considered ready - Send —
ws.sendToConnectedSocket(data)for outbound messages - Request/response —
ws.waitForMessage(predicate, timeout)to await a specific incoming message by any criteria (e.g.reqId) - Typed messages — optional
parseMessagefor type-safe payloads - Notifications —
onNotifycallback for alerts on connection issues and max retries exceeded
Requirements
- Node.js 22+ (built-in global
WebSocket), or any environment with a WHATWGWebSocket(browsers, React Native) - TypeScript 5.x (optional, for types)
Usage
Basic (no auth)
import { ReliableWebSocket } from "@solncebro/websocket-engine";
interface StreamMessage {
type: string;
data: unknown;
}
const ws = new ReliableWebSocket<StreamMessage>({
url: "wss://stream.example.com",
label: "market-stream",
logger: pinoLogger,
parseMessage: (rawData) => JSON.parse(rawData.toString()) as StreamMessage,
onMessage: (message) => {
console.log(message.type, message.data);
},
onNotify: async (message) => {
await sendTelegramAlert(message);
},
});
ws.close();With authentication (e.g. Bybit trading WebSocket)
import crypto from "crypto";
import {
ReliableWebSocket,
WebSocketOpenContext,
} from "@solncebro/websocket-engine";
interface BybitMessage {
op?: string;
retCode?: number;
retMsg?: string;
reqId?: string;
data?: unknown;
}
const ws = new ReliableWebSocket<BybitMessage>({
url: "wss://stream.bybit.com/v5/trade",
label: "bybit-trade",
logger: pinoLogger,
parseMessage: (rawData) => JSON.parse(rawData.toString()) as BybitMessage,
onOpen: async ({ send, waitForMessage }: WebSocketOpenContext<BybitMessage>) => {
const expires = Date.now() + 10000;
const signature = crypto
.createHmac("sha256", SECRET)
.update(`GET/realtime${expires}`)
.digest("hex");
send({ op: "auth", args: [API_KEY, expires, signature] });
const response = await waitForMessage((message) => message.op === "auth", 10000);
if (response.retMsg !== "OK") {
throw new Error(`Auth failed: ${response.retMsg}`);
}
},
heartbeat: {
buildPayload: () => ({ op: "ping" }),
isResponse: (msg) => msg.op === "pong",
},
onMessage: (message) => {
if (message.op === "order.create") {
// handle order response
}
},
onReconnectSuccess: () => {
console.log("Reconnected and re-authenticated");
},
onNotify: async (message) => {
await sendTelegramAlert(message);
},
});
// Send an order and await its specific response by reqId
const sendOrder = async (orderParams: Record<string, unknown>) => {
const reqId = `req_${Date.now()}`;
ws.sendToConnectedSocket({
reqId,
op: "order.create",
args: [orderParams],
});
return ws.waitForMessage((message) => message.reqId === reqId, 30000);
};API
new ReliableWebSocket<TMessage>(args)
Creates a ReliableWebSocket<TMessage> instance. Connection starts immediately on construction.
Arguments
| Property | Type | Required | Description |
|----------|------|----------|-------------|
| url | string | Yes | WebSocket URL |
| label | string | Yes | Identifier for logs and notifications |
| logger | WebSocketLogger | Yes | Logger with debug, info, warn, error, fatal |
| onMessage | (message: TMessage) => void | Yes | Called for each incoming message (not intercepted by waitForMessage or heartbeat) |
| parseMessage | (rawData: RawData) => TMessage | No | Parse raw data to TMessage; default: pass-through |
| onOpen | (context: WebSocketOpenContext<TMessage>) => Promise<void> | No | Async setup phase after connect (e.g. auth). Connection is not considered ready until this resolves. |
| onReconnectSuccess | () => void | No | Called after a successful reconnection (not on first connect) |
| onClose | (context: WebSocketCloseContext) => void | No | Called on every disruption, before the reconnect is scheduled, with the close/error details |
| onNotify | (message: string) => void \| Promise<void> | No | Called on connection issues and when max retries exceeded |
| heartbeat | WebSocketHeartbeatOptions<TMessage> | No | Application-level heartbeat (JSON or bare-text ping/pong). When provided, an active ping is sent every pingInterval and idle detection is enabled. Without it, the socket relies on close/error events. |
| configuration | Partial<WebSocketConfiguration> | No | Override default timeouts and retry behaviour |
Instance Methods
| Method | Description |
|--------|-------------|
| close() | Stops reconnection, clears timers, rejects pending waiters, closes the socket |
| getStatus() | Returns current WebSocketStatus |
| getUrl() | Returns the WebSocket URL |
| sendToConnectedSocket(data) | Send data; string is sent as-is, anything else is JSON.stringify-ed. Throws if not connected. |
| waitForMessage(predicate, timeoutMilliseconds) | Returns a Promise<TMessage> that resolves with the first incoming message matching predicate. The message is not passed to onMessage. Rejects on timeout, connection close, or if predicate throws. |
WebSocketStatus
CONNECTING— initial connection attemptCONNECTED— connected (and auth passed ifonOpenwas provided)DISCONNECTED— disrupted, reconnect scheduledRECONNECTING— reconnect in progress (includesonOpenphase)FAILED— closed by user or max retries exceeded
WebSocketOpenContext
Passed to onOpen:
| Property | Description |
|----------|-------------|
| send | Send to the open socket (for use during onOpen) |
| waitForMessage | Same as instance waitForMessage |
WebSocketHeartbeatOptions
| Property | Description |
|----------|-------------|
| buildPayload | Returns the ping payload. An object is serialized to JSON; a string is sent verbatim, for exchanges whose keepalive is bare text such as ping. |
| isResponse | Returns true if the message is a pong. Matching messages are not passed to onMessage. |
WebSocketCloseContext
Passed to onClose:
| Property | Type | Description |
|----------|------|-------------|
| closeCode | number \| undefined | Close code, when the disruption came from a close event |
| closeReason | string \| undefined | CloseEvent.reason, when the runtime supplies one |
| isCleanClose | boolean \| undefined | CloseEvent.wasClean, when the disruption came from a close event |
| errorMessage | string \| undefined | Error text, when the disruption came from an error event |
| errorName | string \| undefined | ErrorEvent.error.name, when the runtime supplies one |
| errorCause | string \| undefined | The error's cause chain, walked recursively (code, message, nested cause); this is what a bare TypeError from Node's global WebSocket otherwise hides |
| errorCode | string \| undefined | The error's own code (e.g. ECONNRESET, UND_ERR_SOCKET), read off the error itself or its event |
| errorOrigin | string \| undefined | The first stack frame of the error, so the code that produced the failure is visible without the full stack |
| errorStack | string \| undefined | The error's full stack trace |
| errorPropertyList | string \| undefined | Comma-separated own property names of the error object; "none" when the error has no code, cause or stack to go on, so an empty error is distinguishable from one the runtime simply did not fill in |
| wasEverOpen | boolean | Whether the socket had opened at least once before this disruption |
| consecutiveFailures | number | Consecutive failed attempts so far |
| missedPongCount | number | Consecutive stale checks with no incoming message |
| isPongTimeout | boolean | true when missedPongCount has reached missedPongThreshold |
Configuration (configuration)
| Option | Default | Description |
|--------|---------|-------------|
| maxRetryAttempts | 15 | Max reconnection attempts before entering FAILED status |
| initialRetryDelay | 1000 | Initial delay (ms) for exponential backoff |
| maxRetryDelay | 30000 | Cap (ms) for backoff delay |
| retryDelayMultiplier | 1.8 | Backoff multiplier |
| connectionTimeout | 30000 | Handshake timeout (ms) |
| pingInterval | 15000 | Active heartbeat ping interval (ms); only sent when heartbeat is set |
| pongTimeout | 10000 | Pong wait timeout (ms) |
| heartbeatGracePeriod | 3000 | Delay before first ping (ms) |
| staleThreshold | 60000 | Reconnect if no incoming message (including ping responses) for this long (ms); requires heartbeat |
| staleCheckInterval | 5000 | How often idle is checked (ms) |
| fastReconnectCodes | [1001, 1006, 1011, 1012, 1013, 1014] | Close codes that use short reconnect delay |
| missedPongThreshold | 3 | Consecutive stale checks reflected in isPongTimeout of the close context |
Behaviour
- On close or error, reconnection is scheduled with exponential backoff.
- If
onOpenthrows, the connection is terminated and reconnect is triggered (with the same backoff logic). Retry counters are only reset afteronOpenresolves successfully. - After maxRetryAttempts failed attempts,
onNotifyis awaited with a critical message and status becomesFAILED. The process is not terminated — the caller can checkgetStatus()and decide how to handle the failure. waitForMessagepending promises are rejected when the connection is disrupted orclose()is called.- Heartbeat response messages and
waitForMessage-intercepted messages are never passed toonMessage. - With a
heartbeat, the connection is also reconnected if no incoming message (including ping responses) arrives withinstaleThreshold. This catches a silently dead connection without relying on TCP-level control frames, which the globalWebSocketcannot send from code.
Signal relay: source and consumer
@solncebro/websocket-engine/relay is a ready-made scheme for "one process decides, other processes execute". The deciding process (source) publishes a signal once; every executing process (consumer) receives it, decides whether it can act on it right now, and acts. A consumer that starts late, restarts, or loses the network is caught up automatically: on every successful handshake the source replays the signals it still considers alive.
The relay entry runs on Node.js only — it uses ws for the server side. The main entry (@solncebro/websocket-engine) is unchanged and still has no ws in its import graph, so browsers and React Native are unaffected.
1. Read the role from the environment
import { readSignalRelayConfig } from "@solncebro/websocket-engine/relay";
const relayConfig = readSignalRelayConfig({ env: process.env });
if (relayConfig.role === "source") {
await startAsSource(relayConfig.host, relayConfig.port, relayConfig.secret);
} else if (relayConfig.role === "consumer") {
startAsConsumer(relayConfig.url, relayConfig.secret);
} else {
await startStandalone();
}SIGNAL_ROLE selects the role, SIGNAL_RELAY_HOST / SIGNAL_RELAY_PORT / SIGNAL_RELAY_SECRET / SIGNAL_RELAY_URL configure it (roleVariable and prefix override both names). Relay variables set without a role are an error rather than a silent standalone start — an operator who filled them in believes the relay is on.
| SIGNAL_ROLE | Required variables | Result |
|---------------|--------------------|--------|
| unset or standalone | none (any relay variable set is an error) | { role: "standalone" } |
| source | SIGNAL_RELAY_PORT (1–65535), SIGNAL_RELAY_SECRET; SIGNAL_RELAY_HOST defaults to 0.0.0.0; SIGNAL_RELAY_URL is an error | { role: "source", host, port, secret } |
| consumer | SIGNAL_RELAY_URL (ws:// or wss://), SIGNAL_RELAY_SECRET; SIGNAL_RELAY_HOST / SIGNAL_RELAY_PORT are an error | { role: "consumer", url, secret } |
Any other value (including Source — the comparison is exact) throws.
2. Describe your signal
The payload carries the facts of the decision and nothing else. It is your type; the relay only wraps it into RelayedSignalEnvelope<TPayload> (signalId, protocolVersion, sourceLabel, emittedAtMs, payload).
interface EntrySignalPayload {
symbol: string;
barTime: number;
closePrice: number;
klineHeightPercent: number;
}3. Source: publish at the moment of the decision
import { SignalSourceHub } from "@solncebro/websocket-engine/relay";
const SIGNAL_LIFETIME_MS = 30 * 60 * 1000;
const startAsSource = async (
host: string,
port: number,
secret: string
): Promise<SignalSourceHub<EntrySignalPayload>> => {
const sourceHub = new SignalSourceHub<EntrySignalPayload>({
server: { host, port, secret, serverLabel: "signal-source", logger },
sourceLabel: "signal-source",
isSignalAlive: (envelope) =>
Date.now() - envelope.emittedAtMs < SIGNAL_LIFETIME_MS,
logger,
});
await sourceHub.start();
return sourceHub;
};
const sourceHub = await startAsSource(host, port, secret);
// at the moment the strategy decides
sourceHub.publish({
symbol: "BTCUSDT",
barTime: 1755000000000,
closePrice: 61234.5,
klineHeightPercent: 6.2,
});isSignalAlive is the only thing the source knows about your domain: it decides what a late consumer still deserves to receive. A predicate that throws counts as dead and is logged.
4. Consumer: accept when your own state allows it
import { SignalConsumerHub } from "@solncebro/websocket-engine/relay";
const startAsConsumer = (
url: string,
secret: string
): SignalConsumerHub<EntrySignalPayload> => {
const consumerHub = new SignalConsumerHub<EntrySignalPayload>({
client: { url, secret, label: "signal-consumer", logger },
isReady: (envelope) => {
if (Date.now() - envelope.emittedAtMs > SIGNAL_LIFETIME_MS) {
return "expired";
}
return hasOpenPosition(envelope.payload.symbol) ? "wait" : "ready";
},
onAccepted: async (envelope) => {
await placeLadder(envelope.payload);
},
keyOf: (envelope) => envelope.payload.symbol,
onExpired: (envelope, reason) => {
logger.info(`Dropped ${envelope.payload.symbol}: ${reason}`);
},
logger,
});
consumerHub.start();
return consumerHub;
};A signal answered with wait is parked in the pending book under keyOf (default: the signal id). When the state that blocked it changes, the host says so:
const handlePositionClosed = async (symbol: string) => {
await consumerHub.retryPending(symbol);
};onAccepted runs exactly once per signal: the id is deduplicated across message and replay frames, and a pending entry is removed from the book before the handler is awaited, so concurrent retryPending calls cannot double-execute it.
Decision block / execution block pattern
The application splits in two. The decision block answers "is there a signal?" — in the source it is your own detection code, in the consumer it is the hub, which turns an arriving envelope into a decision to act. The execution block answers "can I act, and how?" and is identical in both roles: the same sizing, the same orders, the same journal.
That split is what keeps the standalone role honest — with SIGNAL_ROLE unset the same process runs both blocks in memory, with no relay at all, and the execution block cannot tell the difference.
Hence the rule for the payload: it carries only the facts of the decision (symbol, bar time, prices, the measurements behind the signal). Never put execution settings in it — position size, leverage, take mode, fee rates. Each consumer owns its own account and its own settings; a signal that dictates them makes every consumer a copy of the source and turns a settings change into a protocol change.
Consumers are not obliged to act. isReady is the consumer's own veto: a symbol it already holds returns wait, a signal older than its own tolerance returns expired, and the source never learns about either.
Protocol
One JSON frame per message, both directions.
| Frame | Direction | Meaning |
|-------|-----------|---------|
| hello | client → server | { secret, protocolVersion }, sent immediately after the socket opens |
| welcome | server → client | Handshake accepted; carries serverLabel |
| replay | server → client | Sent once right after welcome when the replay list is not empty |
| message | server → client | One broadcast payload |
| ping / pong | both | Application-level heartbeat; either side may ping |
| Close code | Reason text | Cause |
|------------|-------------|-------|
| 1008 | unauthorized | Wrong secret. Listed in slowReconnectCodes, so the client backs off instead of hammering |
| 1008 | hello timeout | No hello within helloTimeoutMs |
| 1002 | unsupported protocol version | protocolVersion differs from RELAY_PROTOCOL_VERSION |
| 1001 | heartbeat missed | Client missed more than maxMissedHeartbeatCount pings |
| 1001 | server stopped | stop() was called, or a frame could not be sent to this one client (its socket is treated as dead) |
new SignalSourceHub<TPayload>(args)
| Property | Type | Required | Description |
|----------|------|----------|-------------|
| server | Omit<RelayServerArgs<RelayedSignalEnvelope<TPayload>>, "buildReplayList"> | Yes | Transport options: host, port, secret, serverLabel, logger, plus the optional heartbeatIntervalMs, maxMissedHeartbeatCount, helloTimeoutMs, maxPayloadBytes, onClientConnected, onClientDisconnected. The replay list is owned by the hub |
| sourceLabel | string | Yes | Written into every envelope and into the hub's own log lines |
| isSignalAlive | (envelope: RelayedSignalEnvelope<TPayload>) => boolean | Yes | Is this signal still worth replaying? Called on publish, on every handshake and on describeStatus() |
| logger | WebSocketLogger | Yes | Logger with debug, info, warn, error, fatal |
| onClientCountChanged | (count: number, change: "connected" \| "disconnected", info: RelayClientInfo) => void | No | Called after each authenticated connect and each disconnect, with the count that remains |
| maxBookSize | number | No | Cap on the book of live signals (default 500); the oldest are dropped first |
| nowMs | () => number | No | Clock for emittedAtMs (default Date.now) |
| Method | Description |
|--------|-------------|
| start() | Starts the relay server (Promise<void>) |
| stop() | Closes every client and frees the port (Promise<void>) |
| publish(payload) | Prunes the book, wraps the payload, books it and broadcasts it; returns the envelope. Publishing before start() is allowed — the signal is booked and replayed to the first consumer |
| getClientCount() | Number of authenticated consumers |
| describeStatus() | { role: "source", clientCount, publishedCount, lastPublishedAtMs, aliveCount } |
new SignalConsumerHub<TPayload>(args)
| Property | Type | Required | Description |
|----------|------|----------|-------------|
| client | Omit<RelayClientArgs<RelayedSignalEnvelope<TPayload>>, "onMessage" \| "onReplay" \| "onConnected" \| "onDisconnected"> | Yes | Transport options: url, secret, label, logger, plus the optional configuration, welcomeTimeoutMs, onNotify. The four frame handlers are owned by the hub |
| isReady | (envelope: RelayedSignalEnvelope<TPayload>) => "ready" \| "wait" \| "expired" | Yes | The consumer's own veto. Called on arrival and on every retry, so it must be cheap and synchronous |
| onAccepted | (envelope: RelayedSignalEnvelope<TPayload>) => Promise<void> | Yes | Runs once per signal, only after ready. A rejection is logged, not thrown |
| logger | WebSocketLogger | Yes | Logger with debug, info, warn, error, fatal |
| keyOf | (envelope: RelayedSignalEnvelope<TPayload>) => string | No | Pending book key (default: signalId). A newer signal with the same key replaces the parked one |
| onExpired | (envelope, reason: "expired_on_arrival" \| "expired_while_waiting") => void | No | A signal was dropped without being executed |
| onDuplicate | (envelope) => void | No | The same signalId arrived again (a replay after a reconnect, normally) |
| onConnectionLost | (reason: string) => void | No | Once per episode, on the transition from connected to lost |
| onConnectionRestored | (serverLabel: string) => void | No | Once per episode, on the first welcome after a loss. The very first connection does not call it |
| maxSeenSignalCount | number | No | Size of the deduplication memory (default 2000); the oldest ids are forgotten first |
| nowMs | () => number | No | Clock for lastSignalAtMs (default Date.now) |
| Method | Description |
|--------|-------------|
| start() | Connects and keeps reconnecting (the underlying ReliableWebSocket owns the backoff) |
| close() | Stops reconnection and closes the socket |
| retryPending(key?) | Re-runs isReady for one parked signal or for all of them; accepts what became ready, drops what expired (Promise<void>) |
| expirePending() | Re-runs isReady for every parked signal and drops the expired ones |
| getPendingCount() | Size of the pending book |
| describeStatus() | { role: "consumer", connectionStatus, isConnected, lastSignalAtMs, acceptedCount, pendingCount, duplicateCount, expiredCount } |
License
ISC
