src/ws/client.test.ts
v0.3.0 · 37 KB
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();
});
});