Complete room UI and protocol integration, move game definitions and resources behind bundle entries, publish authoritative version XML, and document single-game builds. Include all current resource changes and experiment artifacts.
970 lines
32 KiB
TypeScript
970 lines
32 KiB
TypeScript
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<number, { readonly fn: () => void; readonly at: number }>();
|
|
const clearedTimers = new Map<number, () => 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<string, unknown>, {
|
|
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<string, unknown> = {};
|
|
cyclic.self = cyclic;
|
|
const withAccessor: Record<string, unknown> = {};
|
|
let getterCalls = 0;
|
|
Object.defineProperty(withAccessor, 'secret', {
|
|
enumerable: true,
|
|
get() { getterCalls++; return 1; },
|
|
});
|
|
const withHidden: Record<string, unknown> = {};
|
|
Object.defineProperty(withHidden, 'hidden', { enumerable: false, value: 1 });
|
|
const withSymbolKey: Record<PropertyKey, unknown> = { visible: true };
|
|
withSymbolKey[Symbol('hidden')] = 1;
|
|
class Exotic { value = 1; }
|
|
|
|
const cases: ReadonlyArray<readonly [unknown, RegExp]> = [
|
|
[{ 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<readonly [string, () => 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<string, unknown>;
|
|
Object.defineProperty(hiddenEnvelope, 'data', { enumerable: false, value: {} });
|
|
let accessorCalls = 0;
|
|
const accessorEnvelope = { app: 'youle', route: 'agent', rpc: 'x' } as Record<string, unknown>;
|
|
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);
|
|
});
|