Skip to content
markpaper

src/ws/client.test.ts

v0.3.0 · 37 KB

Download file
import type WsPackageWebSocket from 'ws';
import { afterEach, beforeEach, describe, expect, expectTypeOf, it, vi } from 'vitest';
import { WS_URL } from '../transport/types.js';
import { createWsClient, NODE_WS_SOCKET_OPTIONS, redactWsUrl, type WebSocketConstructorLike, type WebSocketLike, type WsClientOptions } from './client.js';
import { WsLimitError, type WsClientError } from './errors.js';
import { createWsIpBudget } from './limits.js';
import { WsSubscriptionError } from './subscriptionKey.js';
import type { WsFrame, WsL2Book, WsOrderUpdate, WsUserFills } from './types.js';

// ---------------------------------------------------------------------------
// Fake WebSocket: no network. `open/receive/serverClose` drive the "server" side.
// ---------------------------------------------------------------------------

interface SentFrame { method: string; subscription?: Record<string, unknown> }

class FakeWebSocket implements WebSocketLike {
  static instances: FakeWebSocket[] = [];
  static last(): FakeWebSocket {
    const ws = FakeWebSocket.instances[FakeWebSocket.instances.length - 1];
    if (!ws) throw new Error('no socket created');
    return ws;
  }

  readyState = 0;
  readonly sent: SentFrame[] = [];
  readonly args: unknown[];
  terminated = false;
  closedWith: { code?: number; reason?: string } | undefined;
  onopen: ((event: unknown) => void) | null = null;
  onmessage: ((event: { data: unknown }) => void) | null = null;
  onclose: ((event: { code?: number; reason?: string }) => void) | null = null;
  onerror: ((event: unknown) => void) | null = null;

  constructor(readonly url: string, ...rest: unknown[]) {
    this.args = rest;
    FakeWebSocket.instances.push(this);
  }

  send(data: string): void {
    if (this.readyState !== 1) throw new Error('not open');
    this.sent.push(JSON.parse(data) as SentFrame);
  }

  close(code?: number, reason?: string): void {
    if (this.readyState === 3) return;
    this.closedWith = { code, reason };
    this.readyState = 3;
    this.onclose?.({ code: code ?? 1005, reason });
  }

  terminate(): void {
    this.terminated = true;
    this.readyState = 3;
    this.onclose?.({ code: 1006 });
  }

  // ---- server side ----
  open(): void {
    this.readyState = 1;
    this.onopen?.({});
  }

  receive(frame: unknown): void {
    this.onmessage?.({ data: typeof frame === 'string' || frame instanceof Uint8Array ? frame : JSON.stringify(frame) });
  }

  serverClose(code = 1006): void {
    this.readyState = 3;
    this.onclose?.({ code });
  }

  subscribes(): SentFrame[] {
    return this.sent.filter((f) => f.method === 'subscribe');
  }

  ackAll(): void {
    for (const f of this.subscribes()) this.receive({ channel: 'subscriptionResponse', data: { method: 'subscribe', subscription: f.subscription } });
  }
}

const U1 = '0x' + '1'.repeat(40);
const U2 = '0x' + '2'.repeat(40);
const U3 = '0x' + '3'.repeat(40);
const userN = (n: number) => '0x' + n.toString(16).padStart(40, '0');
const noop = () => {};

function makeClient(options: Partial<WsClientOptions> = {}) {
  const errors: WsClientError[] = [];
  const client = createWsClient({ WebSocket: FakeWebSocket, logger: false, ...options });
  client.onError((e) => errors.push(e));
  const ws = FakeWebSocket.last();
  return { client, errors, ws };
}

const kinds = (errors: WsClientError[]) => errors.map((e) => e.kind);

beforeEach(() => {
  FakeWebSocket.instances = [];
  vi.useFakeTimers();
  vi.setSystemTime(0);
});

afterEach(() => {
  vi.useRealTimers();
  vi.unstubAllGlobals();
});

describe('construction', () => {
  it('connects to the network URL immediately', () => {
    makeClient({ network: 'testnet' });
    expect(FakeWebSocket.last().url).toBe(WS_URL.testnet);
    makeClient();
    expect(FakeWebSocket.last().url).toBe(WS_URL.mainnet);
    makeClient({ url: 'wss://example.invalid/ws' });
    expect(FakeWebSocket.last().url).toBe('wss://example.invalid/ws');
  });

  it('requires a WebSocket implementation (Node 20 has no global one)', () => {
    vi.stubGlobal('WebSocket', undefined);
    expect(() => createWsClient({ logger: false })).toThrow(/ws/);
  });

  it('uses globalThis.WebSocket when no constructor is passed', () => {
    vi.stubGlobal('WebSocket', FakeWebSocket);
    createWsClient({ logger: false });
    expect(FakeWebSocket.instances).toHaveLength(1);
  });

  it('passes socketOptions as the second argument only when set', () => {
    makeClient();
    expect(FakeWebSocket.last().args).toEqual([]);
    makeClient({ socketOptions: NODE_WS_SOCKET_OPTIONS });
    expect(FakeWebSocket.last().args).toEqual([{ perMessageDeflate: false, handshakeTimeout: 15_000 }]);
  });

  it('validates timing options', () => {
    expect(() => makeClient({ pingIntervalMs: 0 })).toThrow(RangeError);
    expect(() => makeClient({ confirmTimeoutMs: -1 })).toThrow(RangeError);
  });

  it('accepts the ws package class and the WHATWG WebSocket at the type level', () => {
    expectTypeOf<typeof WsPackageWebSocket>().toMatchTypeOf<WebSocketConstructorLike>();
    expectTypeOf<typeof globalThis.WebSocket>().toMatchTypeOf<WebSocketConstructorLike>();
  });

  it('types handler data by subscription type', () => {
    const { client } = makeClient();
    client.subscribe({ type: 'l2Book', coin: 'BTC' }, (data) => expectTypeOf(data).toEqualTypeOf<WsL2Book>());
    client.subscribe({ type: 'userFills', user: U1 }, (data) => expectTypeOf(data).toEqualTypeOf<WsUserFills>());
    client.subscribe({ type: 'orderUpdates', user: U1 }, (data) => expectTypeOf(data).toEqualTypeOf<WsOrderUpdate[]>());
    client.subscribe({ type: 'somethingNew', foo: 1 }, (data) => expectTypeOf(data).toEqualTypeOf<unknown>());
  });

  it('redacts credentials and query from URLs for logs', () => {
    expect(redactWsUrl('wss://user:[email protected]/ws?token=1#x')).toBe('wss://host.invalid/ws'); // privacy-allow: placeholder credentials
    expect(redactWsUrl('not a url?secret=1')).toBe('not a url');
  });
});

describe('subscriptions and acks', () => {
  it('sends canonical subscriptions on open; connected only after an exact subscribe ack', () => {
    const { client, ws } = makeClient();
    client.subscribe({ type: 'userFills', user: '0x' + 'A'.repeat(40) }, noop);
    expect(ws.sent).toEqual([]); // not open yet: remembered, sent on open
    expect(client.state).toBe('connecting');

    ws.open();
    expect(ws.sent).toEqual([{ method: 'subscribe', subscription: { type: 'userFills', user: '0x' + 'a'.repeat(40) } }]);
    expect(client.state).toBe('subscribing');
    expect(client.isHealthy).toBe(false);

    // an unsubscribe ack or an unrelated data frame must not make the set healthy
    ws.receive({ channel: 'subscriptionResponse', data: { method: 'unsubscribe', subscription: { type: 'userFills', user: '0x' + 'a'.repeat(40) } } });
    ws.receive({ channel: 'userFills', data: { user: '0x' + 'a'.repeat(40), fills: [] } });
    expect(client.state).toBe('subscribing');

    // the server may echo the address in another case
    ws.receive({ channel: 'subscriptionResponse', data: { method: 'subscribe', subscription: { type: 'userFills', user: '0x' + 'A'.repeat(40) } } });
    expect(client.state).toBe('connected');
    expect(client.isHealthy).toBe(true);
  });

  it('an ack for one user does not confirm another user', () => {
    const { client, ws } = makeClient();
    ws.open();
    client.subscribe({ type: 'orderUpdates', user: U1 }, noop);
    ws.receive({ channel: 'subscriptionResponse', data: { method: 'subscribe', subscription: { type: 'orderUpdates', user: U2 } } });
    expect(client.stats().pendingAcks).toBe(1);
  });

  it('throws on an invalid address before anything is sent', () => {
    const { client, ws } = makeClient();
    ws.open();
    expect(() => client.subscribe({ type: 'orderUpdates', user: '0xYOUR_ADDRESS' }, noop)).toThrow(WsSubscriptionError);
    expect(ws.sent).toEqual([]);
    expect(client.stats().subscriptions).toBe(0);
  });

  it('refcounts handlers: one server subscription, unsubscribed on the same socket with the last handler', () => {
    const { client, ws } = makeClient();
    ws.open();
    const calls: string[] = [];
    const f = () => calls.push('f');
    const un1 = client.subscribe({ type: 'l2Book', coin: 'BTC' }, f);
    const un2 = client.subscribe({ coin: 'BTC', type: 'l2Book' }, f);
    expect(ws.subscribes()).toHaveLength(1);
    ws.receive({ channel: 'l2Book', data: { coin: 'BTC', time: 1, levels: [[], []] } });
    expect(calls).toEqual(['f', 'f']);

    un1();
    un1(); // idempotent
    expect(ws.sent.filter((s) => s.method === 'unsubscribe')).toHaveLength(0);
    ws.receive({ channel: 'l2Book', data: { coin: 'BTC', time: 2, levels: [[], []] } });
    expect(calls).toHaveLength(3);

    un2();
    expect(ws.sent.at(-1)).toEqual({ method: 'unsubscribe', subscription: { type: 'l2Book', coin: 'BTC' } });
    expect(FakeWebSocket.instances).toHaveLength(1);
    expect(client.stats().subscriptions).toBe(0);
  });

  it('routes frames to matching handlers only and filters pong/acks', () => {
    const { client, ws } = makeClient();
    ws.open();
    const got: Array<[string, unknown]> = [];
    const h = (name: string) => (data: unknown, frame: WsFrame) => got.push([name, frame.channel === 'pong' ? 'PONG' : data]);
    client.subscribe({ type: 'l2Book', coin: 'BTC' }, h('btc'));
    client.subscribe({ type: 'l2Book', coin: 'ETH' }, h('eth'));
    client.subscribe({ type: 'userFills', user: U1 }, h('fills'));
    client.subscribe({ type: 'userEvents', user: U1 }, h('events'));
    client.subscribe({ type: 'candle', coin: 'BTC', interval: '1m' }, h('candle'));
    ws.ackAll();

    ws.receive({ channel: 'pong' });
    ws.receive({ channel: 'l2Book', data: { coin: 'ETH', time: 1, levels: [[], []] } });
    const snapshot = { user: U1.toUpperCase().replace('0X', '0x'), isSnapshot: true, fills: [] };
    ws.receive({ channel: 'userFills', data: snapshot });
    ws.receive({ channel: 'user', data: { fills: [] } });
    ws.receive({ channel: 'candle', data: { s: 'BTC', i: '5m' } });
    ws.receive({ channel: 'candle', data: { s: 'BTC', i: '1m' } });
    ws.receive({ channel: 'unknownChannel', data: {} });

    expect(got).toEqual([
      ['eth', { coin: 'ETH', time: 1, levels: [[], []] }],
      ['fills', snapshot], // isSnapshot is passed through untouched
      ['events', { fills: [] }],
      ['candle', { s: 'BTC', i: '1m' }],
    ]);
  });

  it('decodes binary frames', () => {
    const { client, ws } = makeClient();
    ws.open();
    const got: unknown[] = [];
    client.subscribe({ type: 'bbo', coin: 'BTC' }, (d) => got.push(d));
    ws.receive(new TextEncoder().encode(JSON.stringify({ channel: 'bbo', data: { coin: 'BTC', time: 1, bbo: [null, null] } })));
    expect(got).toHaveLength(1);
  });

  it('bad JSON and a throwing handler do not close the socket', () => {
    const { client, ws, errors } = makeClient();
    ws.open();
    const good: unknown[] = [];
    client.subscribe({ type: 'bbo', coin: 'BTC' }, () => { throw new Error('boom'); });
    client.subscribe({ type: 'bbo', coin: 'BTC' }, (d) => good.push(d));
    ws.receive('not json');
    ws.receive({ no: 'channel' });
    ws.receive({ channel: 'bbo', data: { coin: 'BTC', time: 1, bbo: [null, null] } });
    expect(kinds(errors)).toEqual(['parse', 'parse', 'handler']);
    expect(good).toHaveLength(1);
    expect(ws.terminated).toBe(false);
    expect(FakeWebSocket.instances).toHaveLength(1);
  });

  it('resubscribe() re-sends unsubscribe + subscribe on the SAME socket', () => {
    const { client, ws } = makeClient();
    const sub = { type: 'allDexsClearinghouseState', user: U1 } as const;
    expect(client.resubscribe(sub)).toBe(false); // socket not open
    ws.open();
    client.subscribe(sub, noop);
    ws.ackAll();
    expect(client.state).toBe('connected');
    expect(client.resubscribe({ type: 'allDexsClearinghouseState', user: U1.toUpperCase().replace('0X', '0x') })).toBe(true);
    expect(ws.sent.slice(-2).map((s) => s.method)).toEqual(['unsubscribe', 'subscribe']);
    expect(client.state).toBe('subscribing');
    expect(FakeWebSocket.instances).toHaveLength(1);
    expect(client.resubscribe({ type: 'l2Book', coin: 'NOPE' })).toBe(false);
  });

  it('onStateChange reports transitions and can be removed', () => {
    const { client, ws } = makeClient();
    const seen: string[] = [];
    const off = client.onStateChange((s, p) => seen.push(`${p}->${s}`));
    ws.open();
    off();
    ws.serverClose();
    expect(seen).toEqual(['connecting->connected']);
  });

  it('falls back to the logger when nobody listens to errors (the error channel is never silent)', () => {
    const logger = { warn: vi.fn(), error: vi.fn() };
    createWsClient({ WebSocket: FakeWebSocket, logger });
    const ws = FakeWebSocket.last();
    ws.open();
    ws.receive('garbage');
    ws.receive({ channel: 'error', data: 'Cannot track more than 15 total users' });
    expect(logger.warn).toHaveBeenCalledTimes(1);
    expect(logger.error).toHaveBeenCalledTimes(1);
    expect(String(logger.error.mock.calls[0]?.[0])).toContain('userLimit');
  });
});

describe('ping and silence', () => {
  it('pings every 30 s', () => {
    const { ws } = makeClient();
    ws.open();
    vi.advanceTimersByTime(29_999);
    expect(ws.sent).toEqual([]);
    vi.advanceTimersByTime(1);
    expect(ws.sent).toEqual([{ method: 'ping' }]);
    vi.advanceTimersByTime(30_000);
    expect(ws.sent).toEqual([{ method: 'ping' }, { method: 'ping' }]);
  });

  it('terminates after 65 s without any frame and reconnects', () => {
    const { ws, errors, client } = makeClient();
    ws.open();
    vi.advanceTimersByTime(65_000);
    expect(ws.terminated).toBe(false);
    vi.advanceTimersByTime(5_000);
    expect(ws.terminated).toBe(true);
    expect(kinds(errors)).toContain('silence');
    expect(client.state).toBe('reconnecting');
    vi.advanceTimersByTime(1_000);
    expect(FakeWebSocket.instances).toHaveLength(2);
  });

  it('any frame, pong included, keeps the socket alive (l2Book silence is normal)', () => {
    const { ws } = makeClient();
    ws.open();
    vi.advanceTimersByTime(60_000);
    ws.receive({ channel: 'pong' });
    vi.advanceTimersByTime(60_000);
    expect(ws.terminated).toBe(false);
  });
});

describe('reconnect', () => {
  it('re-sends every subscription after a reconnect with backoff 1 s -> 2 s -> 4 s', () => {
    const { client, ws } = makeClient();
    client.subscribe({ type: 'l2Book', coin: 'BTC' }, noop);
    client.subscribe({ type: 'orderUpdates', user: U1 }, noop);
    ws.open();
    ws.ackAll();
    ws.serverClose(1006);
    expect(client.state).toBe('reconnecting');

    vi.advanceTimersByTime(999);
    expect(FakeWebSocket.instances).toHaveLength(1);
    vi.advanceTimersByTime(1);
    const ws2 = FakeWebSocket.last();
    expect(FakeWebSocket.instances).toHaveLength(2);
    ws2.open(); // open does NOT reset the backoff
    expect(ws2.subscribes().map((s) => s.subscription?.type)).toEqual(['l2Book', 'orderUpdates']);
    expect(client.state).toBe('subscribing');
    ws2.serverClose(1008 - 2);

    vi.advanceTimersByTime(1_999);
    expect(FakeWebSocket.instances).toHaveLength(2);
    vi.advanceTimersByTime(1);
    FakeWebSocket.last().serverClose();
    vi.advanceTimersByTime(3_999);
    expect(FakeWebSocket.instances).toHaveLength(3);
    vi.advanceTimersByTime(1);
    expect(FakeWebSocket.instances).toHaveLength(4);
    expect(client.stats().reconnects).toBe(3);
  });

  it('resets the backoff only after 8 s of a healthy socket with all acks', () => {
    const { client, ws } = makeClient();
    client.subscribe({ type: 'bbo', coin: 'BTC' }, noop);
    ws.open();
    ws.serverClose(); // attempt 1 -> 1 s
    vi.advanceTimersByTime(1_000);

    const ws2 = FakeWebSocket.last();
    ws2.open();
    ws2.ackAll();
    vi.advanceTimersByTime(7_999);
    ws2.serverClose(); // not stable long enough -> 2 s
    vi.advanceTimersByTime(1_999);
    expect(FakeWebSocket.instances).toHaveLength(2);
    vi.advanceTimersByTime(1);

    const ws3 = FakeWebSocket.last();
    ws3.open();
    ws3.ackAll();
    vi.advanceTimersByTime(8_000);
    expect(client.stats().reconnectAttempt).toBe(0);
    ws3.serverClose(); // reset -> 1 s
    vi.advanceTimersByTime(1_000);
    expect(FakeWebSocket.instances).toHaveLength(4);
  });

  it('does not reset the backoff when a server error arrived on the connection', () => {
    const { client, ws } = makeClient();
    client.subscribe({ type: 'userFills', user: U1 }, noop);
    ws.open();
    ws.serverClose();
    vi.advanceTimersByTime(1_000);
    const ws2 = FakeWebSocket.last();
    ws2.open();
    ws2.ackAll(); // ack of one subscription can arrive together with an error for another
    ws2.receive({ channel: 'error', data: 'Cannot track more than 15 total users' });
    vi.advanceTimersByTime(10_000);
    expect(client.stats().reconnectAttempt).toBe(1);
  });

  it('close code 1008 pauses at least 60 s', () => {
    const { ws, errors } = makeClient();
    ws.open();
    ws.serverClose(1008);
    expect(kinds(errors)).toContain('policyViolation');
    vi.advanceTimersByTime(59_999);
    expect(FakeWebSocket.instances).toHaveLength(1);
    vi.advanceTimersByTime(1);
    expect(FakeWebSocket.instances).toHaveLength(2);
  });

  it('ignores events of a replaced socket (no double reconnect)', () => {
    const { client, ws } = makeClient();
    const got: unknown[] = [];
    client.subscribe({ type: 'bbo', coin: 'BTC' }, (d) => got.push(d));
    ws.open();
    ws.receive({ channel: 'error', data: 'Unexpected failure' }); // terminate -> reconnect in 1 s
    vi.advanceTimersByTime(1_000);
    expect(FakeWebSocket.instances).toHaveLength(2);
    const ws2 = FakeWebSocket.last();
    ws2.open();
    ws2.ackAll();
    expect(client.state).toBe('connected');

    ws.onclose?.({ code: 1006 }); // late close of the old socket
    ws.onopen?.({});
    ws.onmessage?.({ data: JSON.stringify({ channel: 'bbo', data: { coin: 'BTC' } }) });
    vi.advanceTimersByTime(10_000);
    expect(FakeWebSocket.instances).toHaveLength(2);
    expect(client.state).toBe('connected');
    expect(got).toEqual([]);
  });

  it('terminates a hung handshake after 15 s', () => {
    const { ws, errors } = makeClient();
    vi.advanceTimersByTime(15_000);
    expect(ws.terminated).toBe(true);
    expect(kinds(errors)).toEqual(['connectTimeout']);
    vi.advanceTimersByTime(1_000);
    expect(FakeWebSocket.instances).toHaveLength(2);
  });

  it('reconnect: false closes for good after the first drop', () => {
    const { client, ws } = makeClient({ reconnect: false });
    ws.open();
    ws.serverClose();
    expect(client.state).toBe('closed');
    vi.advanceTimersByTime(120_000);
    expect(FakeWebSocket.instances).toHaveLength(1);
    expect(() => client.subscribe({ type: 'bbo', coin: 'BTC' }, noop)).toThrow(WsLimitError);
  });

  it('maxAttempts limits reconnects when set explicitly', () => {
    const { client, ws } = makeClient({ reconnect: { maxAttempts: 1 } });
    ws.serverClose();
    vi.advanceTimersByTime(1_000);
    FakeWebSocket.last().serverClose();
    expect(client.state).toBe('closed');
  });

  it('retries after the WebSocket constructor throws', () => {
    let fail = true;
    class Throwing extends FakeWebSocket {
      constructor(url: string) {
        super(url);
        if (fail) { fail = false; throw new Error('bad url'); }
      }
    }
    const errors: WsClientError[] = [];
    const client = createWsClient({ WebSocket: Throwing, logger: false });
    client.onError((e) => errors.push(e));
    expect(client.state).toBe('reconnecting');
    vi.advanceTimersByTime(1_000);
    expect(FakeWebSocket.instances).toHaveLength(2);
    expect(client.state).toBe('connecting');
  });
});

describe('server errors', () => {
  it('capacity refusal keeps the socket, counts the reject and does not reconnect', () => {
    const { client, ws, errors } = makeClient();
    ws.open();
    client.subscribe({ type: 'userFills', user: U1 }, noop);
    ws.receive({ channel: 'error', data: 'Cannot track more than 15 total users' });
    expect(ws.terminated).toBe(false);
    expect(client.stats().trackingRejects).toBe(1);
    expect(errors[0]).toMatchObject({ kind: 'userLimit', severity: 'error', limit: 15 });

    // the refused subscription is never acked: the confirm timeout keeps the socket
    vi.advanceTimersByTime(30_000);
    expect(ws.terminated).toBe(false);
    expect(kinds(errors)).toEqual(['userLimit', 'confirmTimeout']);
    expect(errors[1]?.severity).toBe('warning');
    expect(client.state).toBe('connected');
    expect(FakeWebSocket.instances).toHaveLength(1);
  });

  it('other server errors terminate the socket (HL may leave it open) and reconnect', () => {
    const { client, ws, errors } = makeClient();
    ws.open();
    ws.receive({ channel: 'error', data: 'Something failed' });
    expect(ws.terminated).toBe(true);
    expect(errors[0]).toMatchObject({ kind: 'other', serverText: 'Something failed' });
    expect(client.state).toBe('reconnecting');
    vi.advanceTimersByTime(1_000);
    expect(FakeWebSocket.instances).toHaveLength(2);
  });

  it('an invalid subscription echoed by the server is dropped instead of being resent forever', () => {
    const { client, ws, errors } = makeClient();
    ws.open();
    client.subscribe({ type: 'fooBar', coin: 'X' }, noop);
    ws.receive({ channel: 'error', data: 'Invalid subscription {"type":"fooBar","coin":"X"}' });
    expect(ws.terminated).toBe(false);
    expect(client.stats().subscriptions).toBe(0);
    expect(errors[0]).toMatchObject({ kind: 'invalidSubscription', subscription: { type: 'fooBar', coin: 'X' } });
    expect(client.state).toBe('connected');
  });

  it('"Already subscribed" counts as an ack; "Already unsubscribed" is only a warning', () => {
    const { client, ws, errors } = makeClient();
    ws.open();
    client.subscribe({ type: 'l2Book', coin: 'BTC' }, noop);
    ws.receive({ channel: 'error', data: 'Already subscribed: {"type":"l2Book","coin":"BTC","nSigFigs":null}' });
    expect(client.state).toBe('connected');
    ws.receive({ channel: 'error', data: 'Already unsubscribed: {"type":"trades","coin":"BTC"}' });
    expect(kinds(errors)).toEqual(['alreadySubscribed', 'alreadyUnsubscribed']);
    expect(ws.terminated).toBe(false);
  });

  it('terminates when subscriptions are not acknowledged within 30 s', () => {
    const { client, ws, errors } = makeClient();
    ws.open();
    client.subscribe({ type: 'bbo', coin: 'BTC' }, noop);
    vi.advanceTimersByTime(29_999);
    ws.receive({ channel: 'pong' });
    expect(ws.terminated).toBe(false);
    vi.advanceTimersByTime(1);
    expect(ws.terminated).toBe(true);
    expect(kinds(errors)).toEqual(['confirmTimeout']);
  });
});

describe('per-IP budget', () => {
  it('refuses a new tracked user above the cap BEFORE sending, across clients sharing one budget', () => {
    const budget = createWsIpBudget({ maxUniqueUsers: 2 });
    const a = makeClient({ budget });
    a.ws.open();
    a.client.subscribe({ type: 'userFills', user: U1 }, noop);
    a.client.subscribe({ type: 'orderUpdates', user: U2 }, noop);

    const b = makeClient({ budget });
    b.ws.open();
    let err: unknown;
    try { b.client.subscribe({ type: 'orderUpdates', user: U3 }, noop); } catch (e) { err = e; }
    expect(err).toBeInstanceOf(WsLimitError);
    expect((err as WsLimitError).reason).toBe('maxUniqueUsers');
    expect(b.ws.sent).toEqual([]);

    // same user on another connection and market data are fine
    b.client.subscribe({ type: 'orderUpdates', user: U1 }, noop);
    b.client.subscribe({ type: 'l2Book', coin: 'BTC' }, noop);
    expect(b.ws.subscribes()).toHaveLength(2);
    expect(b.client.stats().budget).toMatchObject({ uniqueUsers: 2, subscriptions: 4, connections: 2 });
  });

  it('warns when tracked users exceed the documented 10 after the cap was raised', () => {
    const { client, errors } = makeClient({ budgetLimits: { maxUniqueUsers: 20 } });
    for (let i = 1; i <= 11; i++) client.subscribe({ type: 'userFills', user: userN(i) }, noop);
    expect(kinds(errors)).toEqual(['userLimitWarning']);
    expect(() => {
      for (let i = 12; i <= 21; i++) client.subscribe({ type: 'userFills', user: userN(i) }, noop);
    }).toThrow(WsLimitError);
    expect(client.stats().uniqueUsers).toBe(20);
  });

  it('ghost slots of the old connection block new users for ~60 s after a reconnect', () => {
    const { client, ws } = makeClient({ budgetLimits: { maxUniqueUsers: 2 } });
    client.subscribe({ type: 'userFills', user: U1 }, noop);
    ws.open();
    ws.ackAll();
    ws.serverClose();
    vi.advanceTimersByTime(1_000);
    const ws2 = FakeWebSocket.last();
    ws2.open();
    ws2.ackAll();
    expect(() => client.subscribe({ type: 'userFills', user: U2 }, noop)).toThrow(WsLimitError);
    expect(client.stats().budget).toMatchObject({ uniqueUsers: 1, ghostSlots: 1 });
    ws2.receive({ channel: 'pong' });
    vi.advanceTimersByTime(59_001);
    expect(() => client.subscribe({ type: 'userFills', user: U2 }, noop)).not.toThrow();
  });

  it('warns when a user released on one connection is subscribed on another (<60 s)', () => {
    const budget = createWsIpBudget();
    const a = makeClient({ budget });
    a.ws.open();
    const un = a.client.subscribe({ type: 'userFills', user: U1 }, noop);
    un();
    const b = makeClient({ budget });
    b.ws.open();
    b.client.subscribe({ type: 'userFills', user: U1 }, noop);
    expect(kinds(b.errors)).toEqual(['migrationWarning']);
  });

  it('refuses subscriptions above the subscription cap', () => {
    const { client } = makeClient({ budgetLimits: { maxSubscriptions: 1 } });
    client.subscribe({ type: 'bbo', coin: 'BTC' }, noop);
    expect(() => client.subscribe({ type: 'bbo', coin: 'ETH' }, noop)).toThrow(expect.objectContaining({ reason: 'maxSubscriptions' }));
    // another handler for an existing subscription is not a new subscription
    expect(() => client.subscribe({ type: 'bbo', coin: 'BTC' }, noop)).not.toThrow();
  });

  it('a client above the per-IP connection cap waits instead of opening a socket', () => {
    const budget = createWsIpBudget({ maxConnections: 1 });
    const a = makeClient({ budget });
    const b = makeClient({ budget }); // no new socket: FakeWebSocket.last() is still a's
    expect(FakeWebSocket.instances).toHaveLength(1);
    expect(b.client.state).toBe('reconnecting');
    a.client.close();
    vi.advanceTimersByTime(60_000);
    expect(FakeWebSocket.instances).toHaveLength(2);
  });

  it('paces outgoing messages through the budget (experimental messages/min guard)', () => {
    const { client, ws } = makeClient({ budgetLimits: { maxMessagesPerMinute: 2 } });
    ws.open();
    client.subscribe({ type: 'bbo', coin: 'A' }, noop);
    client.subscribe({ type: 'bbo', coin: 'B' }, noop);
    client.subscribe({ type: 'bbo', coin: 'C' }, noop);
    expect(ws.subscribes()).toHaveLength(2);
    ws.ackAll();
    expect(client.stats().pendingAcks).toBe(1);
    vi.advanceTimersByTime(29_999);
    expect(ws.terminated).toBe(false); // a queued frame does not run the confirm timeout
    vi.advanceTimersByTime(30_001);
    expect(ws.subscribes()).toHaveLength(3);
  });

  it('close() stops for good: no reconnect, code 1000, budget released, subscribe throws', () => {
    const budget = createWsIpBudget();
    const { client, ws } = makeClient({ budget });
    client.subscribe({ type: 'bbo', coin: 'BTC' }, noop);
    ws.open();
    client.close();
    client.close();
    expect(ws.closedWith?.code).toBe(1000);
    expect(client.state).toBe('closed');
    vi.advanceTimersByTime(120_000);
    expect(FakeWebSocket.instances).toHaveLength(1);
    expect(budget.stats(Date.now())).toMatchObject({ connections: 0, subscriptions: 0 });
    expect(() => client.subscribe({ type: 'bbo', coin: 'BTC' }, noop)).toThrow(WsLimitError);
  });
});

// ---------------------------------------------------------------------------
// Behaviors observed on the live mainnet socket (2026-09-15), reproduced offline.
// ---------------------------------------------------------------------------

const ack = (ws: FakeWebSocket, subscription: Record<string, unknown>) =>
  ws.receive({ channel: 'subscriptionResponse', data: { method: 'subscribe', subscription } });

const coins = (ws: FakeWebSocket) => ws.subscribes().map((f) => f.subscription?.coin);

describe('poison subscriptions (server drops the socket without an error frame)', () => {
  it('quarantines an unknown-coin l2Book after two isolated drops; the rest keeps working', () => {
    const { client, errors } = makeClient();
    const good: WsL2Book[] = [];
    client.subscribe({ type: 'l2Book', coin: 'BTC' }, (d) => good.push(d));
    client.subscribe({ type: 'l2Book', coin: 'NOPE_COIN' }, noop);
    let ws = FakeWebSocket.last();
    ws.open();
    expect(coins(ws)).toEqual(['BTC', 'NOPE_COIN']);
    // Live: BTC is acked (with server-added fields), then the server drops the socket on NOPE_COIN.
    ack(ws, { type: 'l2Book', coin: 'BTC', nSigFigs: null, mantissa: null, fast: false });
    ws.serverClose(1006);

    vi.advanceTimersByTime(1_000);
    ws = FakeWebSocket.last();
    ws.open();
    // The suspect waits until everything else is acknowledged, then goes alone.
    expect(coins(ws)).toEqual(['BTC']);
    expect(client.state).toBe('subscribing');
    ack(ws, { type: 'l2Book', coin: 'BTC', nSigFigs: null, mantissa: null, fast: false });
    expect(coins(ws)).toEqual(['BTC', 'NOPE_COIN']);
    ws.serverClose(1006);
    expect(kinds(errors)).toEqual(['invalidSubscription']);
    expect(errors[0]?.subscription).toEqual({ type: 'l2Book', coin: 'NOPE_COIN' });
    expect(client.stats()).toMatchObject({ suspects: 0, quarantined: 1 });

    vi.advanceTimersByTime(2_000);
    ws = FakeWebSocket.last();
    ws.open();
    expect(coins(ws)).toEqual(['BTC']);
    ack(ws, { type: 'l2Book', coin: 'BTC' });
    expect(client.state).toBe('connected');
    ws.receive({ channel: 'l2Book', data: { coin: 'BTC', time: 1, levels: [[], []] } });
    expect(good).toHaveLength(1);
    expect(client.stats().subscriptions).toBe(1);
  });

  it('a burst drop marks every unacked subscription suspect; innocents are cleared by their ack', () => {
    const { client, errors } = makeClient();
    client.subscribe({ type: 'trades', coin: 'ETH' }, noop);
    client.subscribe({ type: 'trades', coin: 'btc' }, noop); // wrong case: live drop
    client.subscribe({ type: 'trades', coin: 'SOL' }, noop);
    let ws = FakeWebSocket.last();
    ws.open();
    ack(ws, { type: 'trades', coin: 'ETH' });
    ws.serverClose(1006); // btc and SOL unacked: both suspects, no strike yet

    vi.advanceTimersByTime(1_000);
    ws = FakeWebSocket.last();
    ws.open();
    expect(coins(ws)).toEqual(['ETH']);
    ack(ws, { type: 'trades', coin: 'ETH' });
    expect(coins(ws)).toEqual(['ETH', 'btc']);
    ws.serverClose(1006); // strike 1 for btc

    vi.advanceTimersByTime(2_000);
    ws = FakeWebSocket.last();
    ws.open();
    ack(ws, { type: 'trades', coin: 'ETH' });
    expect(coins(ws)).toEqual(['ETH', 'btc']);
    ws.serverClose(1006); // strike 2 -> quarantined
    expect(kinds(errors)).toEqual(['invalidSubscription']);

    vi.advanceTimersByTime(4_000);
    ws = FakeWebSocket.last();
    ws.open();
    ack(ws, { type: 'trades', coin: 'ETH' });
    expect(coins(ws)).toEqual(['ETH', 'SOL']);
    ack(ws, { type: 'trades', coin: 'SOL' });
    expect(client.state).toBe('connected');
    expect(client.stats().subscriptions).toBe(2);
  });

  it('one network drop right after a valid subscribe does not evict it (needs two isolated strikes)', () => {
    const { client, errors, ws } = makeClient();
    ws.open();
    client.subscribe({ type: 'l2Book', coin: 'ETH' }, noop);
    ws.serverClose(1006); // strike 1
    vi.advanceTimersByTime(1_000);
    const ws2 = FakeWebSocket.last();
    ws2.open();
    ack(ws2, { type: 'l2Book', coin: 'ETH' }); // acked: strikes reset
    ws2.serverClose(1006); // nothing unacked: no suspects
    vi.advanceTimersByTime(2_000);
    const ws3 = FakeWebSocket.last();
    ws3.open();
    expect(ws3.subscribes()).toHaveLength(1);
    ws3.serverClose(1006); // strike 1 again, not 2
    expect(kinds(errors)).toEqual([]);
    expect(client.stats().subscriptions).toBe(1);
  });

  it('close 1008 and client-side terminations never create suspects', () => {
    const { client, errors, ws } = makeClient();
    client.subscribe({ type: 'l2Book', coin: 'ETH' }, noop);
    client.subscribe({ type: 'l2Book', coin: 'SOL' }, noop);
    ws.open();
    ws.serverClose(1008);
    vi.advanceTimersByTime(60_000);
    let next = FakeWebSocket.last();
    next.open();
    expect(coins(next)).toEqual(['ETH', 'SOL']); // sent at once, nothing held
    vi.advanceTimersByTime(30_000); // confirmation timeout -> client-side terminate
    expect(next.terminated).toBe(true);
    vi.advanceTimersByTime(4_000);
    next = FakeWebSocket.last();
    next.open();
    expect(coins(next)).toEqual(['ETH', 'SOL']);
    expect(kinds(errors)).not.toContain('invalidSubscription');
  });

  it('a held suspect cannot be resubscribed out of turn; unsubscribing it releases the queue', () => {
    const { client } = makeClient();
    client.subscribe({ type: 'l2Book', coin: 'ETH' }, noop);
    const un = client.subscribe({ type: 'bbo', coin: 'NOPE_COIN' }, noop);
    let ws = FakeWebSocket.last();
    ws.open();
    ack(ws, { type: 'l2Book', coin: 'ETH' });
    ws.serverClose(1006);
    vi.advanceTimersByTime(1_000);
    ws = FakeWebSocket.last();
    ws.open();
    expect(coins(ws)).toEqual(['ETH']);
    expect(client.resubscribe({ type: 'bbo', coin: 'NOPE_COIN' })).toBe(false);
    un();
    ack(ws, { type: 'l2Book', coin: 'ETH' });
    expect(coins(ws)).toEqual(['ETH']);
    expect(client.state).toBe('connected');
  });
});

describe('live server error texts', () => {
  it('"Error parsing JSON into valid websocket request" with an echo drops the subscription without reconnecting', () => {
    const { client, ws, errors } = makeClient();
    client.subscribe({ type: 'candle', coin: 'BTC', interval: '7m' as '1m' }, noop);
    client.subscribe({ type: 'trades', coin: 'BTC' }, noop);
    ws.open();
    ack(ws, { type: 'trades', coin: 'BTC' });
    ws.receive({
      channel: 'error',
      data: 'Error parsing JSON into valid websocket request: {"method":"subscribe","subscription":{"coin":"BTC","interval":"7m","type":"candle"}}',
    });
    expect(ws.terminated).toBe(false);
    expect(kinds(errors)).toEqual(['invalidSubscription']);
    expect(client.stats().subscriptions).toBe(1);
    expect(client.state).toBe('connected');
  });

  it('"Invalid subscription {json}" for an unknown activeAssetCtx coin drops it', () => {
    const { client, ws, errors } = makeClient();
    client.subscribe({ type: 'activeAssetCtx', coin: 'NOPE_COIN' }, noop);
    ws.open();
    ws.receive({ channel: 'error', data: 'Invalid subscription {"type":"activeAssetCtx","coin":"NOPE_COIN"}' });
    expect(ws.terminated).toBe(false);
    expect(kinds(errors)).toEqual(['invalidSubscription']);
    expect(client.state).toBe('connected');
  });

  it('allMids with dex "" is the main-dex subscription (one server subscription; main frames carry no dex)', () => {
    const { client, ws, errors } = makeClient();
    const a: unknown[] = [];
    const b: unknown[] = [];
    const x: unknown[] = [];
    client.subscribe({ type: 'allMids' }, (d) => a.push(d));
    client.subscribe({ type: 'allMids', dex: '' }, (d) => b.push(d));
    client.subscribe({ type: 'allMids', dex: 'xyz' }, (d) => x.push(d));
    ws.open();
    expect(ws.subscribes().map((f) => f.subscription)).toEqual([{ type: 'allMids' }, { dex: 'xyz', type: 'allMids' }]);
    ack(ws, { type: 'allMids' });
    ack(ws, { type: 'allMids', dex: 'xyz' });
    expect(client.state).toBe('connected');
    ws.receive({ channel: 'allMids', data: { mids: { BTC: '1' } } });
    ws.receive({ channel: 'allMids', data: { dex: 'xyz', mids: { 'xyz:TSLA': '2' } } });
    expect([a.length, b.length, x.length]).toEqual([1, 1, 1]);
    expect(errors).toEqual([]);
  });

  it('routes a spot activeAssetCtx frame (channel activeSpotAssetCtx, coin @107) and HIP-3 perps', () => {
    const { client, ws } = makeClient();
    const spot: unknown[] = [];
    const hip3: unknown[] = [];
    client.subscribe({ type: 'activeAssetCtx', coin: '@107' }, (d) => spot.push(d));
    client.subscribe({ type: 'activeAssetCtx', coin: 'xyz:TSLA' }, (d) => hip3.push(d));
    ws.open();
    ws.receive({ channel: 'activeSpotAssetCtx', data: { coin: '@107', ctx: { markPx: '82.1', prevDayPx: '78.7', dayNtlVlm: '1' } } });
    ws.receive({ channel: 'activeAssetCtx', data: { coin: 'xyz:TSLA', ctx: { markPx: '359.6' } } });
    expect([spot.length, hip3.length]).toEqual([1, 1]);
  });
});

describe('budget accounting across failed connects', () => {
  it('handshakes that never open do not stack ghost slots', () => {
    const { client, ws } = makeClient({ budgetLimits: { maxUniqueUsers: 3 } });
    client.subscribe({ type: 'userFills', user: U1 }, noop);
    ws.open();
    ws.ackAll();
    ws.serverClose(1006);
    expect(client.stats().budget.ghostSlots).toBe(1);
    // Two reconnect attempts fail during an outage without ever opening.
    vi.advanceTimersByTime(1_000);
    FakeWebSocket.last().serverClose(1006);
    vi.advanceTimersByTime(2_000);
    FakeWebSocket.last().serverClose(1006);
    expect(client.stats().budget.ghostSlots).toBe(1);
    expect(client.stats().budget.usedUserSlots).toBe(2);
    expect(() => client.subscribe({ type: 'userFills', user: U2 }, noop)).not.toThrow();
  });
});
All files