fix(net): harden wire lifecycle races
This commit is contained in:
@@ -1,10 +1,11 @@
|
||||
import {
|
||||
APP,
|
||||
RECONNECT_INTERVAL_MS,
|
||||
RECV_TIMEOUT_MS,
|
||||
SERVER_ROTATE_EVERY,
|
||||
} from '../core/constants.ts';
|
||||
import type { InboundMessage, OutboundEnvelope } from '../core/types/envelope.ts';
|
||||
import { decodeFrame, encodeOutbound } from './envelope-codec.ts';
|
||||
import { decodeFrame } from './envelope-codec.ts';
|
||||
import { HeartbeatWatchdog, realClock, type Clock } from './heartbeat.ts';
|
||||
import { ReconnectPolicy } from './reconnect.ts';
|
||||
import type { Transport } from './transport.ts';
|
||||
@@ -36,6 +37,25 @@ function assertSwitchTarget(value: string): void {
|
||||
}
|
||||
}
|
||||
|
||||
function serializeOutbound(envelope: OutboundEnvelope): string {
|
||||
if (typeof envelope !== 'object' || envelope === null || Array.isArray(envelope)) {
|
||||
throw new Error('invalid outbound envelope');
|
||||
}
|
||||
const value = envelope as unknown as Record<string, unknown>;
|
||||
const owns = (key: string) => Object.prototype.hasOwnProperty.call(value, key);
|
||||
if (
|
||||
!owns('app') || value.app !== APP
|
||||
|| !owns('route') || typeof value.route !== 'string' || value.route.length === 0
|
||||
|| !owns('rpc') || typeof value.rpc !== 'string' || value.rpc.length === 0
|
||||
|| !owns('data') || value.data === undefined
|
||||
) {
|
||||
throw new Error('invalid outbound envelope');
|
||||
}
|
||||
const frame = JSON.stringify(envelope);
|
||||
if (typeof frame !== 'string') throw new Error('invalid outbound envelope');
|
||||
return frame;
|
||||
}
|
||||
|
||||
/**
|
||||
* 只拥有 WebSocket 生命周期、帧编解码、收包看门狗和重连策略。
|
||||
* 登录、踢人、切服 RPC 等业务意图均作为 message 原样交给唯一上层编排。
|
||||
@@ -49,6 +69,7 @@ export class WireClient {
|
||||
|
||||
private transport: Transport | null = null;
|
||||
private reconnectTimer: unknown = null;
|
||||
private reconnectToken = 0;
|
||||
private generation = 0;
|
||||
private started = false;
|
||||
private stopped = false;
|
||||
@@ -84,7 +105,7 @@ export class WireClient {
|
||||
if (!this.open || this.stopped || this.transport === null) {
|
||||
throw new Error('connection is not open');
|
||||
}
|
||||
this.transport.send(encodeOutbound(envelope.route, envelope.rpc, envelope.data));
|
||||
this.transport.send(serializeOutbound(envelope));
|
||||
}
|
||||
|
||||
switchServer(url: string): void {
|
||||
@@ -119,16 +140,30 @@ export class WireClient {
|
||||
}
|
||||
|
||||
private emit(event: WireEvent): void {
|
||||
for (const listener of [...this.listeners]) listener(event);
|
||||
for (const listener of [...this.listeners]) {
|
||||
try {
|
||||
listener(event);
|
||||
} catch (error) {
|
||||
console.error('[WireClient] listener error:', error, event);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private connect(server: string): void {
|
||||
if (this.stopped) return;
|
||||
const generation = ++this.generation;
|
||||
const transport = this.transportFactory();
|
||||
this.transport = transport;
|
||||
this.transport = null;
|
||||
this.open = false;
|
||||
|
||||
let transport: Transport;
|
||||
try {
|
||||
transport = this.transportFactory();
|
||||
} catch {
|
||||
this.onConnectionFailure(generation, null, false);
|
||||
return;
|
||||
}
|
||||
this.transport = transport;
|
||||
|
||||
transport.onOpen(() => {
|
||||
if (!this.isCurrent(generation, transport)) return;
|
||||
this.open = true;
|
||||
@@ -140,9 +175,13 @@ export class WireClient {
|
||||
});
|
||||
transport.onClose(() => {
|
||||
if (!this.isCurrent(generation, transport)) return;
|
||||
this.onUnexpectedClose();
|
||||
this.onConnectionFailure(generation, transport, false);
|
||||
});
|
||||
try {
|
||||
transport.connect(server);
|
||||
} catch {
|
||||
this.onConnectionFailure(generation, transport, true);
|
||||
}
|
||||
}
|
||||
|
||||
private isCurrent(generation: number, transport: Transport): boolean {
|
||||
@@ -160,43 +199,72 @@ export class WireClient {
|
||||
private onReceiveTimeout(): void {
|
||||
if (!this.open || this.stopped) return;
|
||||
this.emit({ type: 'slow' });
|
||||
if (!this.open || this.stopped) return;
|
||||
this.disconnectAndSchedule(this.policy.onFailure());
|
||||
}
|
||||
|
||||
private onUnexpectedClose(): void {
|
||||
private onConnectionFailure(
|
||||
generation: number,
|
||||
transport: Transport | null,
|
||||
closeTransport: boolean,
|
||||
): void {
|
||||
if (this.stopped || generation !== this.generation || transport !== this.transport) return;
|
||||
this.transport = null;
|
||||
this.open = false;
|
||||
this.generation++;
|
||||
const failureGeneration = ++this.generation;
|
||||
this.watchdog.stop();
|
||||
if (closeTransport && transport) transport.close();
|
||||
this.emit({ type: 'close' });
|
||||
if (this.stopped || failureGeneration !== this.generation) return;
|
||||
this.scheduleReconnect(this.policy.onFailure());
|
||||
}
|
||||
|
||||
private disconnectAndSchedule(server: string): void {
|
||||
this.clearReconnectTimer();
|
||||
const transport = this.transport;
|
||||
const hadTransport = transport !== null;
|
||||
this.transport = null;
|
||||
this.open = false;
|
||||
this.generation++;
|
||||
const disconnectGeneration = ++this.generation;
|
||||
this.watchdog.stop();
|
||||
this.clearReconnectTimer();
|
||||
transport?.close();
|
||||
if (hadTransport) this.emit({ type: 'close' });
|
||||
if (
|
||||
this.stopped
|
||||
|| disconnectGeneration !== this.generation
|
||||
|| this.policy.current() !== server
|
||||
) return;
|
||||
this.scheduleReconnect(server);
|
||||
}
|
||||
|
||||
private scheduleReconnect(server: string): void {
|
||||
this.clearReconnectTimer();
|
||||
if (this.stopped || this.policy.current() !== server) return;
|
||||
const token = this.reconnectToken;
|
||||
this.emit({ type: 'reconnecting', server });
|
||||
this.reconnectTimer = this.clock.setTimeout(() => {
|
||||
if (
|
||||
this.stopped
|
||||
|| token !== this.reconnectToken
|
||||
|| this.policy.current() !== server
|
||||
) return;
|
||||
let handle: unknown = null;
|
||||
handle = this.clock.setTimeout(() => {
|
||||
if (
|
||||
this.stopped
|
||||
|| token !== this.reconnectToken
|
||||
|| this.policy.current() !== server
|
||||
|| this.reconnectTimer !== handle
|
||||
) return;
|
||||
this.reconnectTimer = null;
|
||||
this.reconnectToken++;
|
||||
this.connect(server);
|
||||
}, RECONNECT_INTERVAL_MS);
|
||||
this.reconnectTimer = handle;
|
||||
}
|
||||
|
||||
private clearReconnectTimer(): void {
|
||||
if (this.reconnectTimer === null) return;
|
||||
this.clock.clearTimeout(this.reconnectTimer);
|
||||
this.reconnectToken++;
|
||||
if (this.reconnectTimer !== null) this.clock.clearTimeout(this.reconnectTimer);
|
||||
this.reconnectTimer = null;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -11,10 +11,16 @@ export class FakeTransport implements Transport {
|
||||
private closeCb?: () => void;
|
||||
private pending: Array<() => void> = [];
|
||||
private closed = false;
|
||||
private connectFailure: unknown = null;
|
||||
|
||||
connect(url: string): void {
|
||||
this.connectCalls++;
|
||||
this.url = url;
|
||||
if (this.connectFailure !== null) {
|
||||
const error = this.connectFailure;
|
||||
this.connectFailure = null;
|
||||
throw error;
|
||||
}
|
||||
this.pending.push(() => this.openCb?.());
|
||||
}
|
||||
send(frame: string): void { if (!this.closed) this.sent.push(frame); }
|
||||
@@ -28,6 +34,9 @@ export class FakeTransport implements Transport {
|
||||
onMessage(cb: (f: string) => void): void { this.msgCb = cb; }
|
||||
onClose(cb: () => void): void { this.closeCb = cb; }
|
||||
|
||||
/** 测试侧:令下一次 connect 同步失败。 */
|
||||
failNextConnect(error: unknown): void { this.connectFailure = error; }
|
||||
|
||||
/** 测试侧:模拟服务器推一帧给客户端。 */
|
||||
serverPush(frame: string): void { this.pending.push(() => this.msgCb?.(frame)); }
|
||||
|
||||
|
||||
@@ -12,6 +12,7 @@ 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;
|
||||
@@ -19,7 +20,10 @@ function fakeClock() {
|
||||
return id;
|
||||
},
|
||||
clearTimeout(handle) {
|
||||
timers.delete(handle as number);
|
||||
const id = handle as number;
|
||||
const timer = timers.get(id);
|
||||
if (timer) clearedTimers.set(id, timer.fn);
|
||||
timers.delete(id);
|
||||
},
|
||||
};
|
||||
const advance = (ms: number) => {
|
||||
@@ -33,7 +37,13 @@ function fakeClock() {
|
||||
due[1].fn();
|
||||
}
|
||||
};
|
||||
return { clock, advance, activeTimers: () => timers.size };
|
||||
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']) {
|
||||
@@ -121,13 +131,42 @@ test('send requires an open connection and serializes the completed envelope onc
|
||||
const transport = transports[0]!;
|
||||
assert.throws(() => client.send(envelope), /connection.*open/i);
|
||||
await transport.flush();
|
||||
client.send(envelope);
|
||||
const completed = {
|
||||
data: { gameid: 101 },
|
||||
rpc: 'player_get_roomlist',
|
||||
app: 'youle',
|
||||
route: 'agent',
|
||||
} as OutboundEnvelope;
|
||||
client.send(completed);
|
||||
|
||||
assert.deepEqual(transport.sent.map((frame) => JSON.parse(frame)), [envelope]);
|
||||
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' },
|
||||
null,
|
||||
];
|
||||
for (const value of invalid) {
|
||||
assert.throws(() => client.send(value as OutboundEnvelope), /outbound envelope/i);
|
||||
}
|
||||
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() }),
|
||||
@@ -213,6 +252,185 @@ test('duplicate transport close/error is handled once and failures rotate every
|
||||
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('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('stale transport generations cannot emit messages or opens', async () => {
|
||||
const { client, transports, advance } = makeHarness();
|
||||
const events: WireEvent[] = [];
|
||||
@@ -249,6 +467,58 @@ test('unsubscribe removes the listener without affecting transport lifecycle', a
|
||||
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[] = [];
|
||||
@@ -272,7 +542,7 @@ test('stop is idempotent, closes an active transport once and clears its watchdo
|
||||
});
|
||||
|
||||
test('stop clears a pending reconnect timer and prevents a later transport generation', async () => {
|
||||
const { client, transports, advance, activeTimers } = makeHarness();
|
||||
const { client, transports, advance, activeTimers, fireClearedTimers } = makeHarness();
|
||||
client.start();
|
||||
const first = transports[0]!;
|
||||
await first.flush();
|
||||
@@ -282,6 +552,7 @@ test('stop clears a pending reconnect timer and prevents a later transport gener
|
||||
assert.equal(activeTimers(), 1);
|
||||
client.stop();
|
||||
assert.equal(activeTimers(), 0);
|
||||
fireClearedTimers();
|
||||
advance(60000);
|
||||
await first.flush();
|
||||
|
||||
|
||||
Reference in New Issue
Block a user