src/ws/client.ts
v0.3.0 · 35.3 KB
// Long-lived Hyperliquid WebSocket client: subscription registry with refcount, ping/silence
// watchdog, ack tracking, infinite reconnect with backoff, per-IP budget checks before sending.
import { WS_URL, type Network } from '../transport/types.js';
import { DEFAULT_RECONNECT_POLICY, reconnectDelay, WS_CLOSE_POLICY_VIOLATION, type ReconnectPolicy } from './backoff.js';
import { classifyWsServerError, WsLimitError, type WsClientError } from './errors.js';
import { createWsIpBudget, WS_LIMITS, type WsIpBudget, type WsIpBudgetLimits, type WsIpBudgetStats } from './limits.js';
import {
canonicalizeSubscription,
channelsOf,
frameMatchesSubscription,
matchEcho,
trackedUserOf,
} from './subscriptionKey.js';
import type { WsEventFor, WsFrame, WsSubscription } from './types.js';
/**
* Structural type satisfied by the `ws` package class (Node 20), `globalThis.WebSocket` (browsers,
* Node 22+) and test fakes. Only the `on*` handler properties are used.
*/
export interface WebSocketLike {
readonly readyState: number;
send(data: string): void;
close(code?: number, reason?: string): void;
/** Present on the `ws` package: drops the TCP connection without the closing handshake. */
terminate?: () => void;
onopen: ((event: any) => void) | null;
onmessage: ((event: any) => void) | null;
onclose: ((event: any) => void) | null;
onerror: ((event: any) => void) | null;
}
/** Constructor of a {@link WebSocketLike}. */
// `any`: the options argument differs between `ws` (ClientOptions) and the WHATWG API (protocols).
export type WebSocketConstructorLike = new (url: string, options?: any) => WebSocketLike;
/**
* Recommended second constructor argument for the `ws` package (NOT for browser/undici WebSocket,
* whose second argument means protocols):
* - `perMessageDeflate: false` - do not spend CPU/GC inflating every frame in the hot path;
* - `handshakeTimeout: 15_000` - a hung handshake must not hang the client.
* Add `localAddress` to pin an egress IP (the `ws` package forwards it down to `tls.connect`);
* clients on different egress IPs need different budgets.
*/
export const NODE_WS_SOCKET_OPTIONS = Object.freeze({ perMessageDeflate: false, handshakeTimeout: WS_LIMITS.connectTimeoutMs });
export type WsClientState =
/** Opening a socket. */
| 'connecting'
/** Socket open, some subscriptions not acknowledged yet. */
| 'subscribing'
/** Socket open and every subscription acknowledged. */
| 'connected'
/** Waiting before the next connect attempt. */
| 'reconnecting'
/** Closed by `close()` or reconnect disabled/exhausted. */
| 'closed';
export type WsHandler<S> = (data: WsEventFor<S>, frame: WsFrame) => void;
/** Removes the handler; the server subscription is dropped (on the same socket) with its last handler. Idempotent. */
export type WsUnsubscribe = () => void;
export interface WsClientOptions {
/** Full URL; overrides `network`. */
url?: string;
/** Default `mainnet`. */
network?: Network;
/**
* WebSocket constructor. Node 20 has no global WebSocket: pass `WebSocket` from the `ws` package.
* Defaults to `globalThis.WebSocket` (browsers, Node 22+).
*/
WebSocket?: WebSocketConstructorLike;
/** Second constructor argument, passed only when set. See {@link NODE_WS_SOCKET_OPTIONS}. */
socketOptions?: unknown;
/** Ping period. Default 30 000 ms (server idle timeout is 60 s). */
pingIntervalMs?: number;
/** No inbound frame (pong included) for this long => terminate + reconnect. Default ping x 2 + 5 000 ms. */
silenceTimeoutMs?: number;
/** Socket must open within this time. Default 15 000 ms; 0 disables. */
connectTimeoutMs?: number;
/** Every sent subscription must be acknowledged within this time, else terminate. Default 30 000 ms; 0 disables. */
confirmTimeoutMs?: number;
/** `false` disables reconnect; an object overrides parts of {@link DEFAULT_RECONNECT_POLICY}. */
reconnect?: boolean | Partial<ReconnectPolicy>;
/** Per-egress-IP budget shared with other clients on the same IP. Default: a private budget. */
budget?: WsIpBudget;
/** Limits of the private budget (ignored when `budget` is given). */
budgetLimits?: WsIpBudgetLimits;
/**
* Where problems go when no `onError` listener is registered. Default `console`. The error channel
* is never dropped silently: an unlogged capacity refusal is completely invisible.
*/
logger?: Pick<Console, 'warn' | 'error'> | false;
/** Label used in the client id and log lines. */
name?: string;
}
export interface WsClientStats {
state: WsClientState;
/** URL without credentials and query. */
url: string;
subscriptions: number;
handlers: number;
/** Distinct users subscribed through this client. */
uniqueUsers: number;
pendingAcks: number;
/**
* Capacity refusals (`Cannot track more than N total users`) received by this client. Compare per
* client/connection: refusals on busy connections only suggest a per-connection cap, refusals on
* lightly used ones while the IP total is high point to the per-IP pool.
*/
trackingRejects: number;
/** Subscriptions suspected of making the server drop the socket, probed one at a time. */
suspects: number;
/** Subscriptions removed after the server dropped the socket right after them repeatedly. */
quarantined: number;
/** Current backoff attempt counter. */
reconnectAttempt: number;
/** Total reconnects scheduled. */
reconnects: number;
/** Delay of the last scheduled reconnect. */
nextReconnectDelayMs: number;
/** ms timestamp of the last inbound frame, 0 before the first one. */
lastMessageAt: number;
/** Budget of the IP. Counters only: never put raw server errors or refused addresses on a public health page. */
budget: WsIpBudgetStats;
}
export interface WsClient {
/**
* Subscribes `handler` to `subscription` and returns its unsubscribe function.
*
* - Validates and lowercases `user` (invalid address throws, nothing is sent).
* - Checks the per-IP budget BEFORE sending: subscriptions (~1000/IP) and unique tracked users
* (default cap 10 documented; ghost slots included). Throws {@link WsLimitError} on refusal.
* - Handlers of the same subscription are refcounted: one server subscription, unsubscribed with the last handler.
* - The subscription is remembered and re-sent after every reconnect.
* - Coins are case-sensitive and not validated here: an unknown or wrong-case coin in `l2Book`/`trades`/
* `bbo` makes the server drop the socket (live, 2026-09-15). The client isolates and removes such a
* subscription after two drops (reported as `invalidSubscription`), at the cost of a few reconnects.
*
* @throws {WsLimitError} limit exceeded or client closed.
* @throws {import('./subscriptionKey.js').WsSubscriptionError} malformed subscription.
*/
subscribe<S extends WsSubscription>(subscription: S, handler: WsHandler<S>): WsUnsubscribe;
/**
* Re-subscribes an existing subscription as `unsubscribe` + `subscribe` on the SAME socket (the only
* safe way to heal a silent address: moving it to another connection holds a ghost slot ~60 s and
* can start refusal storms). Returns `false` when the subscription is unknown or the socket is not open
* (it will be re-sent on the next open anyway). Call it from a timer, never from the event path.
*/
resubscribe(subscription: WsSubscription): boolean;
/** Registers an error/warning listener. Returns a function that removes it. */
onError(listener: (error: WsClientError) => void): () => void;
/** Registers a state listener. Returns a function that removes it. */
onStateChange(listener: (state: WsClientState, previous: WsClientState) => void): () => void;
/**
* Closes the socket for good and releases the budget. Close clients on process exit: with frequent
* restarts lingering connections pile up against the per-IP connection limit.
*/
close(): void;
readonly state: WsClientState;
/** Socket open and all subscriptions acknowledged. */
readonly isHealthy: boolean;
stats(): WsClientStats;
}
const OPEN = 1;
let clientSeq = 0;
/** Removes credentials, query and fragment from a URL for logs. */
export function redactWsUrl(url: string): string {
try {
const u = new URL(url);
u.username = '';
u.password = '';
u.search = '';
u.hash = '';
return u.toString();
} catch {
return url.replace(/\/\/[^@/]*@/, '//').replace(/[?#].*$/, '');
}
}
function unref(timer: unknown): void {
(timer as { unref?: () => void } | undefined)?.unref?.();
}
const textDecoder = new TextDecoder();
function decodeFrame(data: unknown): string | undefined {
if (typeof data === 'string') return data;
if (data instanceof ArrayBuffer) return textDecoder.decode(data);
if (ArrayBuffer.isView(data)) return textDecoder.decode(data);
if (Array.isArray(data) && data.every((d) => ArrayBuffer.isView(d))) {
return data.map((d) => textDecoder.decode(d as ArrayBufferView, { stream: true })).join('') + textDecoder.decode();
}
return undefined;
}
interface Entry {
key: string;
subscription: WsSubscription;
user: string | undefined;
channels: string[];
handlers: Set<(data: unknown, frame: WsFrame) => void>;
/** Server-side drops that happened while this subscription was the ONLY unacknowledged one. */
strikes: number;
}
/**
* Consecutive isolated drops after which a subscription is quarantined (removed from the registry).
* Two, so a single unrelated network drop right after a subscribe cannot evict a valid subscription.
*/
const POISON_STRIKES = 2;
/**
* Creates a WebSocket client and starts connecting immediately.
*
* Behavior (each point has a reason):
* - `{"method":"ping"}` every 30 s; the server closes idle sockets after 60 s. Liveness is judged by
* the silence of ALL frames (pong included) for 65 s, never by one channel: `l2Book` legitimately
* stays silent on a quiet market.
* - Reconnect forever (1 s -> x2 -> 30 s; >= 60 s after close 1008). A transport that gives up after
* 3 retries leaves a bot silently without data.
* - The backoff resets only after {@link ReconnectPolicy.stableMs} of a healthy socket (all acks, no
* server error), never on `open` or on the first ack.
* - `connected` requires an exact `subscriptionResponse` (method `subscribe`) for every subscription.
* - Every socket handler starts with a stale-socket guard, so a late `close` of a replaced socket
* cannot reset the new one and start a double reconnect.
* - Invalid JSON is reported and skipped, the socket stays open; a throwing handler does not affect others.
* - `pong` and acks are filtered out and never reach handlers.
* - Poison subscriptions: live, an `l2Book`/`trades`/`bbo` subscription to an unknown or wrong-case coin
* makes the server drop the socket immediately (close 1006, no error frame). Re-sending it on every
* reconnect would kill every other subscription forever. When the server drops a socket with
* unacknowledged subscriptions, those become suspects; after the next open the suspects are sent one
* at a time (only when nothing else awaits an ack), so a drop is attributed to exactly one of them.
* A suspect dropped in isolation twice in a row is removed and reported as `invalidSubscription`.
*/
export function createWsClient(options: WsClientOptions = {}): WsClient {
const MaybeCtor = options.WebSocket ?? (globalThis as { WebSocket?: WebSocketConstructorLike }).WebSocket;
if (typeof MaybeCtor !== 'function') {
throw new Error(
'hl-kit ws: no WebSocket implementation. On Node 20 pass `WebSocket` from the `ws` package: createWsClient({ WebSocket })',
);
}
const Ctor: WebSocketConstructorLike = MaybeCtor;
const url = options.url ?? WS_URL[options.network ?? 'mainnet'];
const pingIntervalMs = positive(options.pingIntervalMs, WS_LIMITS.pingIntervalMs, 'pingIntervalMs');
const silenceTimeoutMs = positive(options.silenceTimeoutMs, pingIntervalMs * 2 + 5_000, 'silenceTimeoutMs');
const connectTimeoutMs = nonNegative(options.connectTimeoutMs, WS_LIMITS.connectTimeoutMs, 'connectTimeoutMs');
const confirmTimeoutMs = nonNegative(options.confirmTimeoutMs, 30_000, 'confirmTimeoutMs');
const reconnectEnabled = options.reconnect !== false;
const policy: ReconnectPolicy = {
...DEFAULT_RECONNECT_POLICY,
...(typeof options.reconnect === 'object' ? options.reconnect : {}),
};
const budget = options.budget ?? createWsIpBudget(options.budgetLimits);
const name = options.name ?? 'ws';
const clientId = `${name}#${++clientSeq}`;
const logger = options.logger === false ? undefined : (options.logger ?? console);
const tickMs = Math.min(pingIntervalMs, 5_000);
const registry = new Map<string, Entry>();
const byChannel = new Map<string, Set<Entry>>();
/** key -> time the subscribe frame actually left (undefined while queued). */
const pending = new Map<string, number | undefined>();
/** Keys suspected of making the server drop the socket; probed one at a time after the next open. */
const suspects = new Set<string>();
/** Suspect keys of the current socket that are pending but not sent yet (waiting for their turn). */
const held = new Set<string>();
let socketOpened = false;
const errorListeners = new Set<(error: WsClientError) => void>();
const stateListeners = new Set<(state: WsClientState, previous: WsClientState) => void>();
let outQueue: Array<{ payload: string; onSent?: () => void }> = [];
let socket: WebSocketLike | null = null;
let state: WsClientState = 'connecting';
let closedByUser = false;
let attempt = 0;
let reconnects = 0;
let nextReconnectDelayMs = 0;
let lastMessageAt = 0;
let lastPingAt = 0;
let serverErrorSinceOpen = false;
let capacityRejectsSinceOpen = 0;
let trackingRejects = 0;
let quarantined = 0;
let tickTimer: ReturnType<typeof setInterval> | undefined;
let connectTimer: ReturnType<typeof setTimeout> | undefined;
let confirmTimer: ReturnType<typeof setTimeout> | undefined;
let stableTimer: ReturnType<typeof setTimeout> | undefined;
let reconnectTimer: ReturnType<typeof setTimeout> | undefined;
let flushTimer: ReturnType<typeof setTimeout> | undefined;
const now = () => Date.now();
const isOpen = () => socket !== null && socket.readyState === OPEN;
function emit(error: WsClientError): void {
if (errorListeners.size === 0) {
const line = `[hl-kit ${name}] ${error.kind}: ${error.message}`;
if (error.severity === 'warning') logger?.warn(line);
else logger?.error(line);
return;
}
for (const listener of [...errorListeners]) {
try { listener(error); } catch { /* a throwing listener must not break the client */ }
}
}
function setState(next: WsClientState): void {
if (state === next) return;
const previous = state;
state = next;
for (const listener of [...stateListeners]) {
try { listener(next, previous); } catch { /* ignore */ }
}
}
// ---- outgoing ------------------------------------------------------------
function enqueue(payload: string, onSent?: () => void): void {
outQueue.push(onSent ? { payload, onSent } : { payload });
flush();
}
function flush(): void {
const ws = socket;
if (!ws || ws.readyState !== OPEN) return;
while (outQueue.length > 0) {
const wait = budget.reserveMessage(now());
if (wait > 0) {
if (!flushTimer) {
flushTimer = setTimeout(() => { flushTimer = undefined; flush(); }, wait);
unref(flushTimer);
}
return;
}
const item = outQueue.shift() as { payload: string; onSent?: () => void };
try {
ws.send(item.payload);
} catch (cause) {
emit({ kind: 'socket', severity: 'error', message: 'send failed; terminating socket', cause });
terminate(ws);
return;
}
item.onSent?.();
}
}
function sendSubscribe(entry: Entry): void {
pending.set(entry.key, undefined);
enqueue(JSON.stringify({ method: 'subscribe', subscription: entry.subscription }), () => {
if (pending.has(entry.key) && registry.get(entry.key) === entry) {
pending.set(entry.key, now());
scheduleConfirmCheck();
}
});
}
// ---- health --------------------------------------------------------------
function clearStable(): void {
if (stableTimer) { clearTimeout(stableTimer); stableTimer = undefined; }
}
/** Sends the next held suspect when no other subscription awaits an ack. */
function releaseProbe(): void {
if (!isOpen() || held.size === 0) return;
for (const key of pending.keys()) if (!held.has(key)) return;
for (const key of held) {
held.delete(key);
const entry = registry.get(key);
if (!entry) { pending.delete(key); continue; }
sendSubscribe(entry);
return;
}
}
function acknowledged(key: string): void {
pending.delete(key);
suspects.delete(key);
const entry = registry.get(key);
if (entry) entry.strikes = 0;
}
function updateHealth(): void {
if (!isOpen()) return;
releaseProbe();
if (pending.size > 0) {
clearStable();
setState('subscribing');
return;
}
setState('connected');
if (attempt > 0 && !serverErrorSinceOpen && !stableTimer) {
const ws = socket;
stableTimer = setTimeout(() => {
stableTimer = undefined;
if (socket === ws && isOpen() && pending.size === 0 && !serverErrorSinceOpen) attempt = 0;
}, policy.stableMs);
unref(stableTimer);
}
}
function scheduleConfirmCheck(): void {
if (confirmTimeoutMs <= 0 || confirmTimer || !socket) return;
let earliest = Number.POSITIVE_INFINITY;
for (const sentAt of pending.values()) if (sentAt !== undefined && sentAt < earliest) earliest = sentAt;
if (!Number.isFinite(earliest)) return;
const ws = socket;
confirmTimer = setTimeout(() => {
confirmTimer = undefined;
if (socket !== ws) return;
checkConfirmations(ws);
}, Math.max(0, earliest + confirmTimeoutMs - now()));
unref(confirmTimer);
}
function checkConfirmations(ws: WebSocketLike): void {
const t = now();
const overdue = [...pending].filter(([, sentAt]) => sentAt !== undefined && t - sentAt >= confirmTimeoutMs);
if (overdue.length === 0) { scheduleConfirmCheck(); return; }
if (overdue.length <= capacityRejectsSinceOpen) {
// Unacked subscriptions explained by capacity refusals on this connection: a reconnect cannot
// add capacity (the pool is per IP) and would leave ghost slots, so keep the socket.
capacityRejectsSinceOpen -= overdue.length;
for (const [key] of overdue) pending.delete(key);
emit({
kind: 'confirmTimeout',
severity: 'warning',
message: `${overdue.length} subscription(s) never acknowledged after capacity refusals; socket kept`,
});
updateHealth();
scheduleConfirmCheck();
return;
}
emit({
kind: 'confirmTimeout',
severity: 'error',
message: `subscription confirmation timeout: ${overdue.length} of ${pending.size} pending after ${confirmTimeoutMs} ms`,
});
terminate(ws);
}
// ---- connection lifecycle -----------------------------------------------
function connect(): void {
reconnectTimer = undefined;
if (closedByUser || socket) return;
socketOpened = false;
const check = budget.tryOpenConnection(clientId, now());
if (!check.ok) {
const delay = Math.max(check.retryInMs, reconnectDelay(attempt, undefined, policy));
emit({
kind: 'connectionLimit',
severity: 'warning',
message: `per-IP budget refused a new connection (${check.reason}); retry in ${delay} ms`,
});
scheduleConnect(delay);
return;
}
setState('connecting');
let ws: WebSocketLike;
try {
ws = options.socketOptions !== undefined ? new Ctor(url, options.socketOptions) : new Ctor(url);
} catch (cause) {
budget.connectionClosed(clientId, now(), false);
emit({ kind: 'socket', severity: 'error', message: `cannot create WebSocket for ${redactWsUrl(url)}`, cause });
scheduleReconnectAfterFailure(undefined);
return;
}
socket = ws;
ws.onopen = () => handleOpen(ws);
ws.onmessage = (event: { data?: unknown }) => handleMessage(ws, event);
ws.onerror = (event: unknown) => handleSocketError(ws, event);
ws.onclose = (event: { code?: number; reason?: unknown } | undefined) => handleClose(ws, event?.code, true);
if (connectTimeoutMs > 0) {
connectTimer = setTimeout(() => {
connectTimer = undefined;
if (socket !== ws || ws.readyState === OPEN) return;
emit({ kind: 'connectTimeout', severity: 'error', message: `socket did not open within ${connectTimeoutMs} ms` });
terminate(ws);
}, connectTimeoutMs);
unref(connectTimer);
}
}
function scheduleConnect(delay: number): void {
setState('reconnecting');
nextReconnectDelayMs = delay;
reconnectTimer = setTimeout(connect, delay);
unref(reconnectTimer);
}
function scheduleReconnectAfterFailure(closeCode: number | undefined): void {
if (!reconnectEnabled || attempt >= policy.maxAttempts) {
closedByUser = true;
budget.releaseClient(clientId, now());
registry.clear();
byChannel.clear();
setState('closed');
return;
}
const delay = reconnectDelay(attempt, closeCode, policy);
attempt++;
reconnects++;
scheduleConnect(delay);
}
function handleOpen(ws: WebSocketLike): void {
if (socket !== ws) return;
if (connectTimer) { clearTimeout(connectTimer); connectTimer = undefined; }
const t = now();
lastMessageAt = t;
lastPingAt = t;
socketOpened = true;
serverErrorSinceOpen = false;
capacityRejectsSinceOpen = 0;
pending.clear();
held.clear();
tickTimer = setInterval(() => tick(ws), tickMs);
unref(tickTimer);
for (const entry of registry.values()) {
if (suspects.has(entry.key)) {
pending.set(entry.key, undefined);
held.add(entry.key);
} else {
sendSubscribe(entry);
}
}
updateHealth();
}
function tick(ws: WebSocketLike): void {
if (socket !== ws) return;
const t = now();
if (t - lastMessageAt > silenceTimeoutMs) {
emit({ kind: 'silence', severity: 'error', message: `no frames for ${t - lastMessageAt} ms; terminating socket` });
terminate(ws);
return;
}
if (isOpen() && t - lastPingAt >= pingIntervalMs) {
lastPingAt = t;
enqueue('{"method":"ping"}');
}
}
function handleMessage(ws: WebSocketLike, event: { data?: unknown }): void {
if (socket !== ws) return;
// Liveness is updated before any parsing/dedupe: repeated snapshots must not count as silence.
lastMessageAt = now();
const text = decodeFrame(event?.data);
let frame: WsFrame | undefined;
if (text !== undefined) {
try {
const parsed: unknown = JSON.parse(text);
if (parsed && typeof parsed === 'object' && typeof (parsed as WsFrame).channel === 'string') frame = parsed as WsFrame;
} catch { /* reported below */ }
}
if (!frame) {
emit({ kind: 'parse', severity: 'warning', message: `unparseable frame skipped (${text?.slice(0, 120) ?? typeof event?.data})` });
return;
}
switch (frame.channel) {
case 'pong':
return;
case 'subscriptionResponse':
handleAck(frame.data);
return;
case 'error':
handleServerError(ws, frame.data);
return;
default:
dispatch(frame);
}
}
function pendingEntries(): Entry[] {
const out: Entry[] = [];
for (const key of pending.keys()) {
const entry = registry.get(key);
if (entry) out.push(entry);
}
return out;
}
function handleAck(data: unknown): void {
const d = data as { method?: unknown; subscription?: unknown } | undefined;
// Only an exact subscribe ack counts; unsubscribe acks must not make a partial set "healthy".
if (!d || d.method !== 'subscribe') return;
const key = matchEcho(pendingEntries().filter((e) => !held.has(e.key)), d.subscription);
if (key === undefined) return;
acknowledged(key);
updateHealth();
}
function handleServerError(ws: WebSocketLike, data: unknown): void {
const serverText = typeof data === 'string' ? data : JSON.stringify(data ?? null);
const info = classifyWsServerError(serverText);
const key = info.echo ? matchEcho([...registry.values()], info.echo) : undefined;
const entry = key !== undefined ? registry.get(key) : undefined;
const base = {
serverText,
...(entry ? { subscription: entry.subscription } : {}),
...(info.limit !== undefined ? { limit: info.limit } : {}),
};
if (info.kind === 'alreadySubscribed') {
if (entry && pending.has(entry.key) && !held.has(entry.key)) { acknowledged(entry.key); updateHealth(); }
emit({ ...base, kind: info.kind, severity: 'warning', message: 'server reports the subscription already exists' });
return;
}
if (info.kind === 'alreadyUnsubscribed') {
emit({ ...base, kind: info.kind, severity: 'warning', message: 'server reports the subscription did not exist' });
return;
}
serverErrorSinceOpen = true;
clearStable();
if (info.kind === 'userLimit' || info.kind === 'subscriptionLimit') {
trackingRejects++;
capacityRejectsSinceOpen++;
if (entry && pending.has(entry.key) && !held.has(entry.key)) { acknowledged(entry.key); capacityRejectsSinceOpen--; updateHealth(); }
emit({
...base,
kind: info.kind,
severity: 'error',
message: 'capacity refusal: the refused subscription looks accepted but will never push; socket kept (reconnect cannot add per-IP capacity)',
});
return;
}
if (info.kind === 'invalidSubscription' && entry) {
dropEntry(entry, false);
emit({ ...base, kind: info.kind, severity: 'error', message: 'server rejected the subscription; it was removed from the registry' });
updateHealth();
return;
}
emit({ ...base, kind: info.kind, severity: 'error', message: 'server error; terminating socket (HL may keep a failed socket open)' });
terminate(ws);
}
function dispatch(frame: WsFrame): void {
const entries = byChannel.get(frame.channel);
if (!entries) return;
for (const entry of [...entries]) {
if (!frameMatchesSubscription(entry.subscription, frame.channel, frame.data)) continue;
for (const handler of [...entry.handlers]) {
try {
handler(frame.data, frame);
} catch (cause) {
emit({ kind: 'handler', severity: 'error', message: `handler for ${entry.subscription.type} threw`, subscription: entry.subscription, cause });
}
}
}
}
function handleSocketError(ws: WebSocketLike, event: unknown): void {
if (socket !== ws) return;
const message = (event as { message?: unknown } | undefined)?.message;
emit({ kind: 'socket', severity: 'warning', message: typeof message === 'string' ? message : 'socket error', cause: event });
}
function teardownSocket(): void {
socket = null;
for (const timer of [connectTimer, confirmTimer, stableTimer, flushTimer]) if (timer) clearTimeout(timer);
if (tickTimer) clearInterval(tickTimer);
tickTimer = connectTimer = confirmTimer = stableTimer = flushTimer = undefined;
pending.clear();
held.clear();
outQueue = [];
budget.connectionClosed(clientId, now(), socketOpened);
socketOpened = false;
}
/**
* Called when the SERVER (or the network) closed an open socket: marks unacknowledged subscriptions
* as suspects and returns the one to quarantine, if any.
*/
function recordDropSuspects(code: number | undefined): Entry | undefined {
if (!socketOpened || code === WS_CLOSE_POLICY_VIOLATION) return undefined;
const unacked: Entry[] = [];
for (const [key, sentAt] of pending) {
const entry = registry.get(key);
if (entry && sentAt !== undefined && !held.has(key)) unacked.push(entry);
}
for (const entry of unacked) suspects.add(entry.key);
if (unacked.length !== 1) return undefined;
const only = unacked[0] as Entry;
only.strikes++;
return only.strikes >= POISON_STRIKES ? only : undefined;
}
function handleClose(ws: WebSocketLike, code: number | undefined, byServer = false): void {
if (socket !== ws) return;
const poison = byServer && !closedByUser ? recordDropSuspects(code) : undefined;
teardownSocket();
if (closedByUser) { setState('closed'); return; }
if (poison) {
const subscription = poison.subscription;
dropEntry(poison, false);
quarantined++;
emit({
kind: 'invalidSubscription',
severity: 'error',
subscription,
...(code !== undefined ? { closeCode: code } : {}),
message:
`server dropped the socket ${POISON_STRIKES} times right after this ${subscription.type} subscription without an ack ` +
'(unknown or wrong-case coin? names are case-sensitive: BTC, xyz:TSLA, @107); removed from the registry',
});
}
if (code === WS_CLOSE_POLICY_VIOLATION) {
emit({
kind: 'policyViolation',
severity: 'error',
closeCode: code,
message: `closed with 1008 (too many connections from this IP?); pausing >= ${policy.policyViolationPauseMs} ms`,
});
}
scheduleReconnectAfterFailure(code);
}
/** Drops the socket now (no closing handshake) and schedules a reconnect. */
function terminate(ws: WebSocketLike): void {
if (socket !== ws) return;
handleClose(ws, undefined);
try {
if (typeof ws.terminate === 'function') ws.terminate();
else ws.close();
} catch { /* the socket is already replaced; nothing else to do */ }
}
// ---- registry ------------------------------------------------------------
function dropEntry(entry: Entry, sendUnsubscribe: boolean): void {
if (registry.get(entry.key) !== entry) return;
registry.delete(entry.key);
for (const channel of entry.channels) {
const set = byChannel.get(channel);
set?.delete(entry);
if (set && set.size === 0) byChannel.delete(channel);
}
pending.delete(entry.key);
held.delete(entry.key);
suspects.delete(entry.key);
budget.removeSubscription(clientId, entry.key, entry.user, now());
if (sendUnsubscribe && isOpen()) {
enqueue(JSON.stringify({ method: 'unsubscribe', subscription: entry.subscription }));
}
}
const client: WsClient = {
subscribe(subscription, handler) {
if (closedByUser) throw new WsLimitError('closed', 'hl-kit ws: client is closed');
if (typeof handler !== 'function') throw new TypeError('hl-kit ws: handler must be a function');
const canonical = canonicalizeSubscription(subscription);
const key = JSON.stringify(canonical);
let entry = registry.get(key);
if (!entry) {
const user = trackedUserOf(canonical);
const t = now();
const check = budget.checkSubscribe(clientId, user, t);
if (!check.ok) {
throw new WsLimitError(
check.reason,
check.reason === 'maxUniqueUsers'
? `hl-kit ws: tracking one more user exceeds the per-IP cap (${check.used}/${check.limit} slots used, ghost slots included); poll it via REST or use another egress IP`
: `hl-kit ws: subscription cap reached (${check.used}/${check.limit})`,
check.limit,
check.used,
);
}
for (const warning of check.warnings) {
if (warning.kind === 'userLimitWarning') {
emit({
kind: 'userLimitWarning',
severity: 'warning',
limit: warning.limit,
subscription: canonical,
message: `${warning.used} tracked user slots exceed the documented ${warning.limit}; the surplus is best-effort`,
});
} else {
emit({
kind: 'migrationWarning',
severity: 'warning',
subscription: canonical,
message: 'user was released on another connection <60 s ago: its ghost slot is still held; resubscribe on the same connection instead',
});
}
}
entry = { key, subscription: canonical, user, channels: channelsOf(canonical), handlers: new Set(), strikes: 0 };
registry.set(key, entry);
for (const channel of entry.channels) {
let set = byChannel.get(channel);
if (!set) { set = new Set(); byChannel.set(channel, set); }
set.add(entry);
}
budget.addSubscription(clientId, key, user, t);
if (isOpen()) {
sendSubscribe(entry);
updateHealth();
}
}
const owned = entry;
const wrapped = (data: unknown, frame: WsFrame) => (handler as (d: unknown, f: WsFrame) => void)(data, frame);
owned.handlers.add(wrapped);
let active = true;
return () => {
if (!active) return;
active = false;
if (!owned.handlers.delete(wrapped)) return;
if (owned.handlers.size === 0) {
dropEntry(owned, true);
updateHealth();
}
};
},
resubscribe(subscription) {
const key = JSON.stringify(canonicalizeSubscription(subscription));
const entry = registry.get(key);
// A held suspect is sent by the probe sequence; sending it now would break the isolation.
if (!entry || !isOpen() || held.has(key)) return false;
enqueue(JSON.stringify({ method: 'unsubscribe', subscription: entry.subscription }));
sendSubscribe(entry);
updateHealth();
return true;
},
onError(listener) {
errorListeners.add(listener);
return () => { errorListeners.delete(listener); };
},
onStateChange(listener) {
stateListeners.add(listener);
return () => { stateListeners.delete(listener); };
},
close() {
if (closedByUser && state === 'closed') return;
closedByUser = true;
if (reconnectTimer) { clearTimeout(reconnectTimer); reconnectTimer = undefined; }
const ws = socket;
if (ws) teardownSocket();
budget.releaseClient(clientId, now());
registry.clear();
byChannel.clear();
setState('closed');
if (ws) {
try { ws.close(1000, 'client closed'); } catch { /* ignore */ }
}
},
get state() {
return state;
},
get isHealthy() {
return isOpen() && pending.size === 0;
},
stats() {
let handlers = 0;
const users = new Set<string>();
for (const entry of registry.values()) {
handlers += entry.handlers.size;
if (entry.user !== undefined) users.add(entry.user);
}
return {
state,
url: redactWsUrl(url),
subscriptions: registry.size,
handlers,
uniqueUsers: users.size,
pendingAcks: pending.size,
trackingRejects,
suspects: suspects.size,
quarantined,
reconnectAttempt: attempt,
reconnects,
nextReconnectDelayMs,
lastMessageAt,
budget: budget.stats(now()),
};
},
};
connect();
return client;
}
function positive(value: number | undefined, fallback: number, name: string): number {
if (value === undefined) return fallback;
if (!(value > 0) || !Number.isFinite(value)) throw new RangeError(`hl-kit ws: ${name} must be a positive number`);
return value;
}
function nonNegative(value: number | undefined, fallback: number, name: string): number {
if (value === undefined) return fallback;
if (!(value >= 0) || !Number.isFinite(value)) throw new RangeError(`hl-kit ws: ${name} must be a non-negative number`);
return value;
}