import { test } from 'node:test'; import assert from 'node:assert/strict'; import { WireClient, type WireEvent, } from '../../YouleNexus/assets/framework/net/wire-client.ts'; import type { OutboundEnvelope } from '../../YouleNexus/assets/framework/core/types/envelope.ts'; import type { Clock } from '../../YouleNexus/assets/framework/net/heartbeat.ts'; import { FakeTransport } from '../helpers/fake-transport.ts'; function fakeClock() { let sequence = 0; let now = 0; const timers = new Map void; readonly at: number }>(); const clearedTimers = new Map void>(); const clock: Clock = { setTimeout(fn, ms) { const id = ++sequence; timers.set(id, { fn, at: now + ms }); return id; }, clearTimeout(handle) { const id = handle as number; const timer = timers.get(id); if (timer) clearedTimers.set(id, timer.fn); timers.delete(id); }, }; const advance = (ms: number) => { now += ms; while (true) { const due = [...timers.entries()] .filter(([, timer]) => timer.at <= now) .sort(([leftId, left], [rightId, right]) => left.at - right.at || leftId - rightId)[0]; if (!due) return; timers.delete(due[0]); due[1].fn(); } }; const fireClearedTimers = () => { for (const [id, fn] of [...clearedTimers]) { clearedTimers.delete(id); fn(); } }; return { clock, advance, activeTimers: () => timers.size, fireClearedTimers }; } function makeHarness(servers: readonly string[] = ['ws://a']) { const transports: FakeTransport[] = []; const timer = fakeClock(); const client = new WireClient({ servers, clock: timer.clock, transportFactory: () => { const transport = new FakeTransport(); transports.push(transport); return transport; }, }); return { client, transports, ...timer }; } const envelope: OutboundEnvelope = { app: 'youle', route: 'agent', rpc: 'player_get_roomlist', data: { gameid: 101 }, }; test('debug packet logs preserve transmitted text and precede business dispatch',async(t)=>{ const logs:unknown[][]=[]; t.mock.method(console,'log',(...args:unknown[])=>logs.push(args)); const {client,transports}=makeHarness();client.start();await transports[0]!.flush(); const login:OutboundEnvelope={app:'youle',route:'agent',rpc:'player_login',data:{test:true}}; client.send(login); assert.ok(logs.some(row=>row[0]==='[收发包] 发送'&&row[2]===transports[0]!.sent[0])); const heartbeats=['{"com":"@serverheartbeat"}',JSON.stringify('{"com":"@serverheartbeat"}')]; for(const heartbeat of heartbeats)transports[0]!.serverPush(heartbeat); await transports[0]!.flush(); assert.equal(logs.some(row=>row[0]==='[收发包] 接收'&&heartbeats.includes(row[2] as string)),false); const raw='{"app":"youle","route":"agent","rpc":"player_login","data":{"roomtype":[0,1,0]}}'; let dispatched=false; client.subscribe(event=>{if(event.type==='message'){ assert.ok(logs.some(row=>row[0]==='[收发包] 接收'&&row[2]===raw));dispatched=true; }}); transports[0]!.serverPush(raw);await transports[0]!.flush();assert.ok(dispatched);client.stop(); }); test('start emits the opened server and never sends login automatically', async () => { const { client, transports } = makeHarness(); const events: WireEvent[] = []; client.subscribe((event) => events.push(event)); client.start(); await transports[0]!.flush(); assert.deepEqual(events, [{ type: 'open', server: 'ws://a' }]); assert.equal(transports[0]!.sent.length, 0, 'open 不应由 WireClient 自动发 login'); }); test('all decoded business RPCs are emitted exactly once without interception', async () => { const { client, transports } = makeHarness(); const events: WireEvent[] = []; client.subscribe((event) => events.push(event)); client.start(); const transport = transports[0]!; await transport.flush(); const messages = [ { route: 'agent', rpc: 'player_login', data: { state: 0 } }, { route: 'agent', rpc: 'kick_server', data: { msg: 'bye' } }, { route: 'room', rpc: 'connect_roomserver', data: { roomserver: 'ws://room' } }, { route: 'agent', rpc: 'connect_agentserver', data: { agentserver: 'ws://agent', opt: 'free_room' } }, ]; for (const message of messages) transport.serverPush(JSON.stringify(message)); await transport.flush(); assert.deepEqual(events.slice(1), messages.map((message) => ({ type: 'message', message }))); assert.equal(transport.sent.length, 0); }); test('handshake, heartbeat, server-down and invalid frames are not business messages', async () => { const { client, transports, advance } = makeHarness(); const events: WireEvent[] = []; client.subscribe((event) => events.push(event)); client.start(); const transport = transports[0]!; await transport.flush(); transport.serverPush('@toconconABC'); transport.serverPush(JSON.stringify({ com: '@serverheartbeat' })); transport.serverPush('webserve-服务器未工作'); transport.serverPush('not json'); await transport.flush(); assert.deepEqual(events, [{ type: 'open', server: 'ws://a' }]); assert.equal(transport.sent.length, 0); advance(29999); assert.equal(events.some((event) => event.type === 'slow'), false); advance(1); assert.equal(events.filter((event) => event.type === 'slow').length, 1); }); test('send requires an open connection and serializes the completed envelope once', async () => { const { client, transports } = makeHarness(); assert.throws(() => client.send(envelope), /connection.*open/i); client.start(); const transport = transports[0]!; assert.throws(() => client.send(envelope), /connection.*open/i); await transport.flush(); const completed = { data: { gameid: 101 }, rpc: 'player_get_roomlist', app: 'youle', route: 'agent', } as OutboundEnvelope; client.send(completed); assert.equal( transport.sent[0], JSON.stringify(completed), '调用方完成的信封必须保持属性和插入顺序原样序列化', ); client.stop(); assert.throws(() => client.send(envelope), /connection.*open/i); }); test('send rejects malformed completed envelopes before touching transport', async () => { const { client, transports } = makeHarness(); client.start(); const transport = transports[0]!; await transport.flush(); const invalid: unknown[] = [ { app: 'other', route: 'agent', rpc: 'x', data: {} }, { app: 'youle', route: '', rpc: 'x', data: {} }, { app: 'youle', route: 'agent', rpc: '', data: {} }, { app: 'youle', route: 'agent', rpc: 'x' }, { app: 'youle', route: 'agent', rpc: 'x', data: {}, extra: true }, { app: 'youle', route: 'agent', rpc: 'x', data: () => 1 }, null, ]; for (const value of invalid) { assert.throws(() => client.send(value as OutboundEnvelope), /outbound envelope/i); } assert.deepEqual(transport.sent, []); }); test('send accepts the complete protocol JSON value boundary without rebuilding bytes', async () => { const { client, transports } = makeHarness(); client.start(); const transport = transports[0]!; await transport.flush(); const nullPrototype = Object.assign(Object.create(null) as Record, { nested: [null, 'text', true, false, 0, -1.5], }); const valid = { rpc: 'json_boundary', data: nullPrototype, route: 'agent', app: 'youle', } as OutboundEnvelope; client.send(valid); assert.equal(transport.sent[0], JSON.stringify(valid)); }); test('send rejects every non-JSON data graph edge with its exact path', async () => { const { client, transports } = makeHarness(); client.start(); const transport = transports[0]!; await transport.flush(); const sparse: unknown[] = []; sparse.length = 1; const arrayWithExtra = [] as unknown[] & { extra?: number }; arrayWithExtra.extra = 1; const arrayWithSymbol = [] as unknown[]; Object.defineProperty(arrayWithSymbol, Symbol('extra'), { value: 1 }); const cyclic: Record = {}; cyclic.self = cyclic; const withAccessor: Record = {}; let getterCalls = 0; Object.defineProperty(withAccessor, 'secret', { enumerable: true, get() { getterCalls++; return 1; }, }); const withHidden: Record = {}; Object.defineProperty(withHidden, 'hidden', { enumerable: false, value: 1 }); const withSymbolKey: Record = { visible: true }; withSymbolKey[Symbol('hidden')] = 1; class Exotic { value = 1; } const cases: ReadonlyArray = [ [{ value: undefined }, /\$\.data\.value.*undefined/i], [{ value: () => 1 }, /\$\.data\.value.*function/i], [{ value: Symbol('x') }, /\$\.data\.value.*symbol/i], [{ value: 1n }, /\$\.data\.value.*bigint/i], [{ value: Number.NaN }, /\$\.data\.value.*finite number/i], [{ value: Number.POSITIVE_INFINITY }, /\$\.data\.value.*finite number/i], [{ value: -0 }, /\$\.data\.value.*negative zero/i], [{ nested: sparse }, /\$\.data\.nested\[0\].*array hole/i], [{ nested: arrayWithExtra }, /\$\.data\.nested\.extra.*extra array property/i], [{ nested: arrayWithSymbol }, /\$\.data\.nested.*symbol key/i], [{ nested: new Date(0) }, /\$\.data\.nested.*inherited toJSON/i], [{ nested: new Exotic() }, /\$\.data\.nested.*plain object/i], [{ nested: withAccessor }, /\$\.data\.nested\.secret.*accessor/i], [{ nested: withHidden }, /\$\.data\.nested\.hidden.*enumerable/i], [{ nested: withSymbolKey }, /\$\.data\.nested.*symbol key/i], [{ nested: { toJSON: () => ({ changed: true }) } }, /\$\.data\.nested\.toJSON.*function/i], [cyclic, /\$\.data\.self.*cycle/i], ]; for (const [data, expected] of cases) { assert.throws( () => client.send({ app: 'youle', route: 'agent', rpc: 'x', data }), expected, ); } assert.equal(getterCalls, 0, 'accessor validation must not invoke the getter'); assert.deepEqual(transport.sent, []); }); test('send preserves a dense non-enumerable array index in the validated wire bytes', async () => { const { client, transports } = makeHarness(); client.start(); const transport = transports[0]!; await transport.flush(); const values: unknown[] = []; Object.defineProperty(values, '0', { configurable: true, enumerable: false, value: 'descriptor value', writable: true, }); const completed = { app: 'youle', route: 'agent', rpc: 'array_descriptor', data: { values }, } as OutboundEnvelope; client.send(completed); assert.equal(transport.sent[0], JSON.stringify(completed)); assert.equal( transport.sent[0], '{"app":"youle","route":"agent","rpc":"array_descriptor","data":{"values":["descriptor value"]}}', ); }); test('send rejects inherited toJSON without invoking it', async () => { const { client, transports } = makeHarness(); client.start(); const transport = transports[0]!; await transport.flush(); const original = Object.getOwnPropertyDescriptor(Object.prototype, 'toJSON'); let calls = 0; Object.defineProperty(Object.prototype, 'toJSON', { configurable: true, value() { calls++; return { changed: true }; }, }); try { assert.throws( () => client.send({ app: 'youle', route: 'agent', rpc: 'x', data: { value: 1 } }), /\$.*inherited toJSON/i, ); assert.equal(calls, 0); assert.deepEqual(transport.sent, []); } finally { if (original) Object.defineProperty(Object.prototype, 'toJSON', original); else delete (Object.prototype as { toJSON?: unknown }).toJSON; } }); test('send serializes stable descriptor snapshots without invoking Proxy get traps', async () => { const { client, transports } = makeHarness(); client.start(); const transport = transports[0]!; await transport.flush(); const target = { value: 1 }; let getCalls = 0; const proxied = new Proxy(target, { get(source, key, receiver) { getCalls++; if (key === 'value') return 2; return Reflect.get(source, key, receiver); }, }); client.send({ app: 'youle', route: 'agent', rpc: 'proxy', data: { proxied } }); assert.equal(getCalls, 0); assert.equal( transport.sent[0], '{"app":"youle","route":"agent","rpc":"proxy","data":{"proxied":{"value":1}}}', ); }); test('send audits one shared source identity instead of accepting path-local stable reads', async () => { const { client, transports } = makeHarness(); client.start(); const transport = transports[0]!; await transport.flush(); const target = { value: 0 }; const reportedValues = [1, 2, 1, 2] as const; let descriptorReads = 0; const shared = new Proxy(target, { getOwnPropertyDescriptor(source, key) { const descriptor = Reflect.getOwnPropertyDescriptor(source, key); if (key !== 'value' || !descriptor) return descriptor; const value = reportedValues[descriptorReads]; descriptorReads++; if (value === undefined) throw new Error('shared source was audited more than twice'); return { ...descriptor, value }; }, }); assert.throws( () => client.send({ app: 'youle', route: 'agent', rpc: 'shared_unstable', data: { left: shared, right: shared }, }), /\$\.data\.left\.value.*descriptor.*changed/i, ); assert.equal(descriptorReads, 2, 'the shared source must have one capture and one verification'); assert.deepEqual(transport.sent, []); }); test('send expands an ordinary shared DAG reference at every wire path', async () => { const { client, transports } = makeHarness(); client.start(); const transport = transports[0]!; await transport.flush(); const shared = { score: 7 }; const completed = { app: 'youle', route: 'agent', rpc: 'shared_dag', data: { left: shared, right: shared }, } as OutboundEnvelope; client.send(completed); assert.equal( transport.sent[0], '{"app":"youle","route":"agent","rpc":"shared_dag","data":{"left":{"score":7},"right":{"score":7}}}', ); }); test('send rejects a cyclic prototype identity within a finite prototype walk', async () => { const { client, transports } = makeHarness(); client.start(); const transport = transports[0]!; await transport.flush(); const finiteWalkLimit = 4; let prototypeReads = 0; let getCalls = 0; let selfPrototype!: object; selfPrototype = new Proxy(Object.create(null) as object, { getPrototypeOf() { prototypeReads++; if (prototypeReads > finiteWalkLimit) { throw new Error('prototype walk exceeded the finite test bound'); } return selfPrototype; }, get(source, key, receiver) { getCalls++; return Reflect.get(source, key, receiver); }, }); assert.throws( () => client.send({ app: 'youle', route: 'agent', rpc: 'prototype_cycle', data: { loop: selfPrototype }, }), /\$\.data\.loop.*prototype.*cycle/i, ); assert.ok(prototypeReads <= finiteWalkLimit); assert.equal(getCalls, 0, 'prototype validation must not invoke Proxy get'); assert.deepEqual(transport.sent, []); }); test('send rejects source graphs whose descriptors, own-key order, or prototype change', async () => { const cases: ReadonlyArray object, RegExp]> = [ [ 'descriptor', () => { const target = { value: 1 }; let calls = 0; return new Proxy(target, { getOwnPropertyDescriptor(source, key) { const descriptor = Reflect.getOwnPropertyDescriptor(source, key); if (key !== 'value' || !descriptor) return descriptor; calls++; return { ...descriptor, value: calls === 1 ? 1 : 2 }; }, }); }, /\$\.data\.unstable\.value.*descriptor.*changed/i, ], [ 'own keys', () => { const target = { first: 1, second: 2 }; let calls = 0; return new Proxy(target, { ownKeys() { calls++; return calls === 1 ? ['first', 'second'] : ['second', 'first']; }, }); }, /\$\.data\.unstable.*own keys.*changed/i, ], [ 'prototype', () => { const target = { value: 1 }; let calls = 0; return new Proxy(target, { getPrototypeOf() { calls++; return calls === 1 ? Object.prototype : null; }, }); }, /\$\.data\.unstable.*prototype.*changed/i, ], ]; for (const [name, makeUnstable, expected] of cases) { const { client, transports } = makeHarness(); client.start(); const transport = transports[0]!; await transport.flush(); assert.throws( () => client.send({ app: 'youle', route: 'agent', rpc: name, data: { unstable: makeUnstable() } }), expected, ); assert.deepEqual(transport.sent, []); } }); test('send rejects top-level symbol, non-enumerable and accessor properties', async () => { const { client, transports } = makeHarness(); client.start(); const transport = transports[0]!; await transport.flush(); const symbolEnvelope = { app: 'youle', route: 'agent', rpc: 'x', data: {} }; Object.defineProperty(symbolEnvelope, Symbol('extra'), { enumerable: true, value: 1 }); const hiddenEnvelope = { app: 'youle', route: 'agent', rpc: 'x' } as Record; Object.defineProperty(hiddenEnvelope, 'data', { enumerable: false, value: {} }); let accessorCalls = 0; const accessorEnvelope = { app: 'youle', route: 'agent', rpc: 'x' } as Record; Object.defineProperty(accessorEnvelope, 'data', { enumerable: true, get() { accessorCalls++; return {}; }, }); assert.throws(() => client.send(symbolEnvelope as OutboundEnvelope), /top-level.*symbol/i); assert.throws(() => client.send(hiddenEnvelope as unknown as OutboundEnvelope), /\$\.data.*enumerable/i); assert.throws(() => client.send(accessorEnvelope as unknown as OutboundEnvelope), /\$\.data.*accessor/i); assert.equal(accessorCalls, 0); assert.deepEqual(transport.sent, []); }); test('constructor rejects an empty validated server list and start cannot run twice', () => { assert.throws( () => new WireClient({ servers: [], transportFactory: () => new FakeTransport() }), /servers.*non-empty/i, ); const { client } = makeHarness(); client.start(); assert.throws(() => client.start(), /already.*started/i); }); test('switchServer validates the target before closing and reconnects to it after policy delay', async () => { const { client, transports, advance } = makeHarness(['ws://a', 'ws://b']); const events: WireEvent[] = []; client.subscribe((event) => events.push(event)); client.start(); const first = transports[0]!; await first.flush(); assert.throws(() => client.switchServer('http://room'), /WebSocket.*URL/i); assert.throws(() => client.switchServer('ws://'), /WebSocket.*URL/i); assert.equal(first.closeCalls, 0); client.switchServer('ws://room'); assert.equal(first.closeCalls, 1); assert.deepEqual(events.slice(-2), [ { type: 'close' }, { type: 'reconnecting', server: 'ws://room' }, ]); advance(9999); assert.equal(transports.length, 1); advance(1); assert.equal(transports[1]!.url, 'ws://room'); await transports[1]!.flush(); assert.deepEqual(events[events.length - 1], { type: 'open', server: 'ws://room' }); }); test('reconnectCurrent closes once and reconnects to the current target after 10 seconds', async () => { const { client, transports, advance } = makeHarness(['ws://a', 'ws://b']); const events: WireEvent[] = []; client.subscribe((event) => events.push(event)); client.start(); await transports[0]!.flush(); client.reconnectCurrent(); assert.equal(transports[0]!.closeCalls, 1); assert.deepEqual(events.slice(-2), [ { type: 'close' }, { type: 'reconnecting', server: 'ws://a' }, ]); advance(10000); assert.equal(transports[1]!.url, 'ws://a'); }); test('duplicate transport close/error is handled once and failures rotate every three attempts', async () => { const { client, transports, advance, activeTimers } = makeHarness(['ws://a', 'ws://b']); const reconnectTargets: string[] = []; client.subscribe((event) => { if (event.type === 'reconnecting') reconnectTargets.push(event.server); }); client.start(); await transports[0]!.flush(); transports[0]!.serverClose(); await transports[0]!.flush(); assert.deepEqual(reconnectTargets, ['ws://a']); advance(10000); await transports[1]!.flush(); transports[1]!.serverError(); await transports[1]!.flush(); assert.deepEqual(reconnectTargets, ['ws://a', 'ws://a']); advance(10000); await transports[2]!.flush(); transports[2]!.serverError(); transports[2]!.serverClose(); await transports[2]!.flush(); assert.deepEqual(reconnectTargets, ['ws://a', 'ws://a', 'ws://b']); assert.equal(activeTimers(), 1); advance(10000); assert.equal(transports.length, 4); assert.equal(transports[3]!.url, 'ws://b'); }); test('a throwing close listener is reported without blocking later listeners or reconnect', async () => { const { client, transports, advance, activeTimers } = makeHarness(); const listenerError = new Error('close listener failed'); const reports: unknown[][] = []; let laterListenerCalls = 0; const originalConsoleError = console.error; console.error = (...args: unknown[]) => { reports.push(args); }; try { client.subscribe((event) => { if (event.type === 'close') throw listenerError; }); client.subscribe((event) => { if (event.type === 'close') laterListenerCalls++; }); client.start(); const first = transports[0]!; await first.flush(); first.serverClose(); await first.flush(); assert.equal(laterListenerCalls, 1); assert.equal(reports.length, 1); assert.equal(reports[0]!.includes(listenerError), true); assert.equal(activeTimers(), 1); advance(10000); await transports[1]!.flush(); assert.equal(transports[1]!.url, 'ws://a'); } finally { console.error = originalConsoleError; } }); test('a reconnecting listener can stop without leaving a timer or reviving a cleared callback', async () => { const { client, transports, advance, activeTimers, fireClearedTimers } = makeHarness(); client.subscribe((event) => { if (event.type === 'reconnecting' && event.server === 'ws://room') client.stop(); }); client.start(); const first = transports[0]!; await first.flush(); first.serverClose(); await first.flush(); assert.equal(activeTimers(), 1); client.switchServer('ws://room'); assert.equal(activeTimers(), 0); fireClearedTimers(); advance(60000); assert.equal(transports.length, 1); }); test('reconnectCurrent coalesces the same pending target before reconnecting listeners run', async () => { const { client, transports, advance, activeTimers } = makeHarness(); let reconnectEvents = 0; let reentered = false; const timerCountsDuringEvent: number[] = []; client.subscribe((event) => { if (event.type !== 'reconnecting') return; reconnectEvents++; timerCountsDuringEvent.push(activeTimers()); if (!reentered) { reentered = true; client.reconnectCurrent(); } }); client.start(); const first = transports[0]!; await first.flush(); first.serverClose(); await first.flush(); client.reconnectCurrent(); assert.equal(reconnectEvents, 1); assert.deepEqual(timerCountsDuringEvent, [1]); assert.equal(activeTimers(), 1); advance(10000); assert.equal(transports.length, 2); assert.equal(transports[1]!.url, 'ws://a'); }); test('slow listener stop prevents the timeout path from scheduling reconnect', async () => { const { client, transports, advance, activeTimers } = makeHarness(); client.subscribe((event) => { if (event.type === 'slow') client.stop(); }); client.start(); const first = transports[0]!; await first.flush(); first.serverPush('@toconconABC'); await first.flush(); advance(30000); await first.flush(); assert.equal(activeTimers(), 0); advance(10000); assert.equal(transports.length, 1); }); test('close listener stop prevents an unexpected close from scheduling reconnect', async () => { const { client, transports, advance, activeTimers } = makeHarness(); client.subscribe((event) => { if (event.type === 'close') client.stop(); }); client.start(); const first = transports[0]!; await first.flush(); first.serverClose(); await first.flush(); assert.equal(activeTimers(), 0); advance(10000); assert.equal(transports.length, 1); }); test('switch during pending reconnect invalidates the cleared callback and connects only the explicit target', async () => { const { client, transports, advance, fireClearedTimers, activeTimers } = makeHarness(); client.start(); const first = transports[0]!; await first.flush(); first.serverClose(); await first.flush(); assert.equal(activeTimers(), 1); client.switchServer('ws://room'); assert.equal(activeTimers(), 1); fireClearedTimers(); assert.equal(transports.length, 1, '已清除的 ws://a timer callback 不得重新建连'); advance(10000); assert.equal(transports.length, 2); assert.equal(transports[1]!.url, 'ws://room'); }); test('stop before pending open ignores the late open and message callbacks', async () => { const { client, transports } = makeHarness(); const events: WireEvent[] = []; client.subscribe((event) => events.push(event)); client.start(); const first = transports[0]!; client.stop(); first.serverPush(JSON.stringify({ route: 'agent', rpc: 'late', data: {} })); await first.flush(); assert.deepEqual(events, [{ type: 'close' }]); assert.deepEqual(first.sent, []); }); test('transportFactory synchronous failure is one failure and recovers after policy delay', async () => { const timer = fakeClock(); const recovered = new FakeTransport(); const events: WireEvent[] = []; let factoryCalls = 0; const client = new WireClient({ servers: ['ws://a', 'ws://b'], clock: timer.clock, transportFactory: () => { factoryCalls++; if (factoryCalls === 1) throw new Error('factory failed'); return recovered; }, }); client.subscribe((event) => events.push(event)); assert.doesNotThrow(() => client.start()); assert.deepEqual(events, [ { type: 'close' }, { type: 'reconnecting', server: 'ws://a' }, ]); assert.equal(timer.activeTimers(), 1); timer.advance(10000); await recovered.flush(); assert.equal(factoryCalls, 2); assert.deepEqual(events[events.length - 1], { type: 'open', server: 'ws://a' }); }); test('transport connect synchronous failure ignores duplicate close/error and counts once', async () => { const timer = fakeClock(); const broken = new FakeTransport(); broken.failNextConnect(new Error('connect failed')); const recovered = new FakeTransport(); const transports = [broken, recovered]; const reconnectTargets: string[] = []; let closeEvents = 0; let factoryIndex = 0; const client = new WireClient({ servers: ['ws://a', 'ws://b'], clock: timer.clock, transportFactory: () => transports[factoryIndex++]!, }); client.subscribe((event) => { if (event.type === 'reconnecting') reconnectTargets.push(event.server); if (event.type === 'close') closeEvents++; }); assert.doesNotThrow(() => client.start()); broken.serverError(); broken.serverClose(); await broken.flush(); assert.deepEqual(reconnectTargets, ['ws://a']); assert.equal(closeEvents, 1); assert.equal(timer.activeTimers(), 1); assert.equal(broken.closeCalls, 1); timer.advance(10000); await recovered.flush(); assert.equal(recovered.url, 'ws://a'); }); test('FakeTransport and WireClient preserve synchronous throw null as one connection failure', async () => { const timer = fakeClock(); const broken = new FakeTransport(); broken.failNextConnect(null); const recovered = new FakeTransport(); const transports = [broken, recovered]; let factoryIndex = 0; const events: WireEvent[] = []; const client = new WireClient({ servers: ['ws://a'], clock: timer.clock, transportFactory: () => transports[factoryIndex++]!, }); client.subscribe((event) => events.push(event)); client.start(); assert.deepEqual(events, [ { type: 'close' }, { type: 'reconnecting', server: 'ws://a' }, ]); assert.equal(broken.closeCalls, 1); assert.equal(timer.activeTimers(), 1); timer.advance(10000); await recovered.flush(); assert.deepEqual(events[events.length - 1], { type: 'open', server: 'ws://a' }); }); test('stale transport generations cannot emit messages or opens', async () => { const { client, transports, advance } = makeHarness(); const events: WireEvent[] = []; client.subscribe((event) => events.push(event)); client.start(); const first = transports[0]!; await first.flush(); first.serverClose(); await first.flush(); first.serverPush(JSON.stringify({ route: 'agent', rpc: 'stale_before_reconnect', data: {} })); await first.flush(); advance(10000); first.serverPush(JSON.stringify({ route: 'agent', rpc: 'stale_after_reconnect', data: {} })); await first.flush(); await transports[1]!.flush(); assert.deepEqual(events.filter((event) => event.type === 'message'), []); assert.equal(events.filter((event) => event.type === 'open').length, 2); }); test('unsubscribe removes the listener without affecting transport lifecycle', async () => { const { client, transports } = makeHarness(); const events: WireEvent[] = []; const unsubscribe = client.subscribe((event) => events.push(event)); client.start(); const transport = transports[0]!; await transport.flush(); unsubscribe(); unsubscribe(); transport.serverPush(JSON.stringify({ route: 'agent', rpc: 'update_bean', data: { bean: 1 } })); await transport.flush(); assert.deepEqual(events, [{ type: 'open', server: 'ws://a' }]); }); test('subscribe and unsubscribe during emit use a stable listener snapshot', async () => { const { client, transports } = makeHarness(); const calls: string[] = []; let unsubscribeSecond = (): void => undefined; const third = (event: WireEvent) => { if (event.type === 'message') calls.push(`third:${event.message.rpc}`); }; client.subscribe((event) => { if (event.type !== 'message') return; calls.push(`first:${event.message.rpc}`); unsubscribeSecond(); client.subscribe(third); }); unsubscribeSecond = client.subscribe((event) => { if (event.type === 'message') calls.push(`second:${event.message.rpc}`); }); client.start(); const transport = transports[0]!; await transport.flush(); transport.serverPush(JSON.stringify({ route: 'agent', rpc: 'one', data: {} })); await transport.flush(); transport.serverPush(JSON.stringify({ route: 'agent', rpc: 'two', data: {} })); await transport.flush(); assert.deepEqual(calls, ['first:one', 'second:one', 'first:two', 'third:two']); }); for (const [name, frame] of [ ['handshake', '@toconconABC'], ['heartbeat', JSON.stringify({ com: '@serverheartbeat' })], ] as const) { test(`${name} independently resets the receive watchdog sliding window`, async () => { const { client, transports, advance } = makeHarness(); const events: WireEvent[] = []; client.subscribe((event) => events.push(event)); client.start(); const transport = transports[0]!; await transport.flush(); transport.serverPush(frame); await transport.flush(); advance(20000); transport.serverPush(frame); await transport.flush(); advance(29999); assert.equal(events.some((event) => event.type === 'slow'), false); advance(1); assert.equal(events.filter((event) => event.type === 'slow').length, 1); }); } test('stop is idempotent, closes an active transport once and clears its watchdog', async () => { const { client, transports, advance, activeTimers } = makeHarness(); const events: WireEvent[] = []; client.subscribe((event) => events.push(event)); client.start(); const first = transports[0]!; await first.flush(); first.serverPush('@toconconABC'); await first.flush(); assert.equal(activeTimers(), 1); client.stop(); client.stop(); assert.equal(activeTimers(), 0); assert.equal(first.closeCalls, 1); advance(60000); await first.flush(); assert.equal(transports.length, 1); assert.equal(events.filter((event) => event.type === 'close').length, 1); assert.throws(() => client.send(envelope), /connection.*open/i); }); test('stop clears a pending reconnect timer and prevents a later transport generation', async () => { const { client, transports, advance, activeTimers, fireClearedTimers } = makeHarness(); client.start(); const first = transports[0]!; await first.flush(); client.reconnectCurrent(); assert.equal(first.closeCalls, 1); assert.equal(activeTimers(), 1); client.stop(); assert.equal(activeTimers(), 0); fireClearedTimers(); advance(60000); await first.flush(); assert.equal(transports.length, 1); });