Skip to content
markpaper

src/ws/client.ts

v0.3.0 · 35.3 KB

Download file
// 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;
}
All files