Skip to content
markpaper

src/ws/errors.ts

v0.3.0 · 6.6 KB

Download file
// Classification of `channel:"error"` texts and the error objects the client emits.

import type { WsSubscription } from './types.js';

/** What a `channel:"error"` text means for the socket. */
export type WsServerErrorKind =
  /** `Cannot track more than N total users`: per-IP user-tracking capacity exhausted. */
  | 'userLimit'
  /** `Cannot subscribe to more than N ...`: subscription capacity exhausted. */
  | 'subscriptionLimit'
  /** `Already subscribed ...`: the server already serves this subscription. */
  | 'alreadySubscribed'
  /** `Already unsubscribed ...`: unsubscribe of a subscription the server does not hold. */
  | 'alreadyUnsubscribed'
  /**
   * The server refused the subscription itself. Live texts (checked 2026-09-15):
   * `Invalid subscription {json}` (e.g. `activeAssetCtx` with an unknown coin) and
   * `Error parsing JSON into valid websocket request: {request json}` (unknown `type`, bad candle
   * interval, malformed user). The client also uses this kind when it quarantines a subscription after
   * which the server repeatedly dropped the socket without any error frame.
   */
  | 'invalidSubscription'
  /** Anything else. */
  | 'other';

export interface WsServerErrorInfo {
  kind: WsServerErrorKind;
  /** Number parsed from `... more than N ...`, when present. */
  limit?: number;
  /** Subscription echoed inside the text (`{...}` JSON), when present. */
  echo?: Record<string, unknown>;
  /**
   * Whether the socket should be terminated and reconnected.
   *
   * - Capacity refusals: NO. The pool is per IP, a reconnect does not add capacity and itself leaves
   *   ~60 s ghost slots on the old connection. Keep the socket, count the refusal, shrink the set or
   *   move the overflow to REST polling. (Derived from the ghost-slot mechanics; not tested separately.)
   * - Other errors: YES. HL may leave the TCP connection open after an error, so `isConnected` "lies"
   *   forever and the rejected set never resubscribes; terminate (not graceful close, which waits
   *   for the unhealthy peer).
   */
  terminate: boolean;
}

function parseEcho(text: string): Record<string, unknown> | undefined {
  const match = /\{.*\}/s.exec(text);
  if (!match) return undefined;
  try {
    const parsed: unknown = JSON.parse(match[0]);
    if (!parsed || typeof parsed !== 'object' || Array.isArray(parsed)) return undefined;
    const envelope = parsed as Record<string, unknown>;
    const inner = envelope.subscription;
    if (inner && typeof inner === 'object' && !Array.isArray(inner)) return inner as Record<string, unknown>;
    return typeof envelope.type === 'string' ? envelope : undefined;
  } catch {
    return undefined;
  }
}

/**
 * Classifies the text of a `channel:"error"` frame (case-insensitive substring match).
 *
 * The capacity error is the one that matters most: a subscription refused for capacity still looks
 * accepted and simply never pushes, so without reading this channel the refusal is invisible and
 * the data silently falls back to REST, if there is a fallback at all.
 */
export function classifyWsServerError(text: unknown): WsServerErrorInfo {
  const raw = typeof text === 'string' ? text : JSON.stringify(text ?? '');
  const lower = raw.toLowerCase();
  const echo = parseEcho(raw);
  const withEcho = (info: WsServerErrorInfo): WsServerErrorInfo => (echo ? { ...info, echo } : info);

  const users = /cannot track more than (\d+)/.exec(lower);
  if (users) return withEcho({ kind: 'userLimit', limit: Number(users[1]), terminate: false });
  if (lower.includes('cannot track more than')) return withEcho({ kind: 'userLimit', terminate: false });

  // Text as composed by the SDK from the docs; the live server text is not recorded in the knowledge base.
  const subs = /cannot subscribe to more than (\d+)/.exec(lower);
  if (subs) return withEcho({ kind: 'subscriptionLimit', limit: Number(subs[1]), terminate: false });

  if (lower.includes('already subscribed')) return withEcho({ kind: 'alreadySubscribed', terminate: false });
  if (lower.includes('already unsubscribed')) return withEcho({ kind: 'alreadyUnsubscribed', terminate: false });
  // Terminating on an attributable invalid subscription would resend it after reconnect: endless loop.
  // The same holds for the request-parse refusal: an unknown `type` or a bad candle interval is echoed as
  // `Error parsing JSON into valid websocket request: {"method":"subscribe","subscription":{...}}`
  // (live, 2026-09-15) and the server keeps the socket open.
  if (lower.includes('invalid subscription') || lower.includes('valid websocket request')) {
    return withEcho({ kind: 'invalidSubscription', terminate: !echo });
  }
  return withEcho({ kind: 'other', terminate: true });
}

/** Kinds of problems reported through `client.onError`. */
export type WsClientErrorKind =
  | WsServerErrorKind
  /** Transport-level `error` event of the socket (a `close` always follows). */
  | 'socket'
  /** A frame was not valid JSON. The socket is kept: one bad payload must not cost every snapshot. */
  | 'parse'
  /** A subscription handler threw. Other handlers still run. */
  | 'handler'
  /** Close code 1008 (policy violation, in practice too many connections from the IP): >= 60 s pause. */
  | 'policyViolation'
  /** No frame (pong included) within the silence timeout: socket terminated. */
  | 'silence'
  /** Not all subscriptions were acknowledged in time: socket terminated. */
  | 'confirmTimeout'
  /** The socket did not open in time: socket terminated. */
  | 'connectTimeout'
  /** The per-IP connection budget refused to open a socket; retried later. */
  | 'connectionLimit'
  /** A new tracked user exceeds the DOCUMENTED guarantee (10); the extra is best-effort. */
  | 'userLimitWarning'
  /** A user released on another connection less than ~60 s ago: its ghost slot is still held. */
  | 'migrationWarning';

export interface WsClientError {
  kind: WsClientErrorKind;
  severity: 'warning' | 'error';
  message: string;
  /** Raw server text for server errors. Do not expose it on a public health endpoint. */
  serverText?: string;
  /** Subscription the problem is attributed to, when known. */
  subscription?: WsSubscription;
  closeCode?: number;
  limit?: number;
  cause?: unknown;
}

/** Why a subscribe was refused before anything was sent. */
export type WsLimitReason = 'maxSubscriptions' | 'maxUniqueUsers' | 'closed';

/** Thrown synchronously by `subscribe()` when a limit would be exceeded; nothing is sent. */
export class WsLimitError extends Error {
  override readonly name = 'WsLimitError';
  constructor(
    readonly reason: WsLimitReason,
    message: string,
    readonly limit?: number,
    readonly used?: number,
  ) {
    super(message);
  }
}
All files