refactor(net): isolate transport and wire lifecycle

This commit is contained in:
2026-09-05 03:57:32 +08:00
parent 72227efd05
commit a8873b9a84
3 changed files with 517 additions and 3 deletions
@@ -0,0 +1,202 @@
import {
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 { HeartbeatWatchdog, realClock, type Clock } from './heartbeat.ts';
import { ReconnectPolicy } from './reconnect.ts';
import type { Transport } from './transport.ts';
export type WireEvent =
| { readonly type: 'open'; readonly server: string }
| { readonly type: 'message'; readonly message: InboundMessage }
| { readonly type: 'slow' }
| { readonly type: 'reconnecting'; readonly server: string }
| { readonly type: 'close' };
export interface WireClientOptions {
readonly servers: readonly string[];
readonly transportFactory: () => Transport;
readonly clock?: Clock;
}
type WireListener = (event: WireEvent) => void;
function assertSwitchTarget(value: string): void {
let parsed: URL;
try {
parsed = new URL(value);
} catch {
throw new Error('switchServer requires a valid WebSocket URL');
}
if ((parsed.protocol !== 'ws:' && parsed.protocol !== 'wss:') || parsed.host.length === 0) {
throw new Error('switchServer requires a valid WebSocket URL');
}
}
/**
* 只拥有 WebSocket 生命周期、帧编解码、收包看门狗和重连策略。
* 登录、踢人、切服 RPC 等业务意图均作为 message 原样交给唯一上层编排。
*/
export class WireClient {
private readonly transportFactory: () => Transport;
private readonly clock: Clock;
private readonly policy: ReconnectPolicy;
private readonly watchdog: HeartbeatWatchdog;
private readonly listeners = new Set<WireListener>();
private transport: Transport | null = null;
private reconnectTimer: unknown = null;
private generation = 0;
private started = false;
private stopped = false;
private open = false;
constructor(options: WireClientOptions) {
if (!Array.isArray(options.servers) || options.servers.length === 0) {
throw new Error('servers must be a non-empty validated list');
}
this.transportFactory = options.transportFactory;
this.clock = options.clock ?? realClock;
this.policy = new ReconnectPolicy([...options.servers], SERVER_ROTATE_EVERY);
this.watchdog = new HeartbeatWatchdog(
RECV_TIMEOUT_MS,
() => this.onReceiveTimeout(),
this.clock,
);
}
subscribe(listener: WireListener): () => void {
this.listeners.add(listener);
return () => { this.listeners.delete(listener); };
}
start(): void {
if (this.started) throw new Error('WireClient already started');
if (this.stopped) throw new Error('WireClient already stopped');
this.started = true;
this.connect(this.policy.current());
}
send(envelope: OutboundEnvelope): void {
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));
}
switchServer(url: string): void {
assertSwitchTarget(url);
this.assertRunning();
this.policy.setCurrent(url);
this.disconnectAndSchedule(this.policy.current());
}
reconnectCurrent(): void {
this.assertRunning();
this.disconnectAndSchedule(this.policy.current());
}
stop(): void {
if (this.stopped) return;
this.stopped = true;
this.watchdog.stop();
this.clearReconnectTimer();
const transport = this.transport;
const hadTransport = transport !== null;
this.transport = null;
this.open = false;
this.generation++;
transport?.close();
if (hadTransport) this.emit({ type: 'close' });
}
private assertRunning(): void {
if (!this.started || this.stopped) throw new Error('WireClient is not running');
}
private emit(event: WireEvent): void {
for (const listener of [...this.listeners]) listener(event);
}
private connect(server: string): void {
if (this.stopped) return;
const generation = ++this.generation;
const transport = this.transportFactory();
this.transport = transport;
this.open = false;
transport.onOpen(() => {
if (!this.isCurrent(generation, transport)) return;
this.open = true;
this.emit({ type: 'open', server });
});
transport.onMessage((frame) => {
if (!this.isCurrent(generation, transport) || !this.open) return;
this.onMessage(frame);
});
transport.onClose(() => {
if (!this.isCurrent(generation, transport)) return;
this.onUnexpectedClose();
});
transport.connect(server);
}
private isCurrent(generation: number, transport: Transport): boolean {
return !this.stopped && generation === this.generation && transport === this.transport;
}
private onMessage(frame: string): void {
const decoded = decodeFrame(frame);
if (decoded.kind === 'ignore') return;
this.watchdog.feed();
if (decoded.kind !== 'message') return;
this.emit({ type: 'message', message: decoded.message });
}
private onReceiveTimeout(): void {
if (!this.open || this.stopped) return;
this.emit({ type: 'slow' });
this.disconnectAndSchedule(this.policy.onFailure());
}
private onUnexpectedClose(): void {
this.transport = null;
this.open = false;
this.generation++;
this.watchdog.stop();
this.emit({ type: 'close' });
this.scheduleReconnect(this.policy.onFailure());
}
private disconnectAndSchedule(server: string): void {
const transport = this.transport;
const hadTransport = transport !== null;
this.transport = null;
this.open = false;
this.generation++;
this.watchdog.stop();
this.clearReconnectTimer();
transport?.close();
if (hadTransport) this.emit({ type: 'close' });
this.scheduleReconnect(server);
}
private scheduleReconnect(server: string): void {
this.clearReconnectTimer();
this.emit({ type: 'reconnecting', server });
this.reconnectTimer = this.clock.setTimeout(() => {
this.reconnectTimer = null;
this.connect(server);
}, RECONNECT_INTERVAL_MS);
}
private clearReconnectTimer(): void {
if (this.reconnectTimer === null) return;
this.clock.clearTimeout(this.reconnectTimer);
this.reconnectTimer = null;
}
}
@@ -4,21 +4,44 @@ import type { Transport } from '../../YouleNexus/assets/framework/net/transport.
export class FakeTransport implements Transport {
sent: string[] = [];
url = '';
connectCalls = 0;
closeCalls = 0;
private openCb?: () => void;
private msgCb?: (f: string) => void;
private closeCb?: () => void;
private pending: Array<() => void> = [];
private closed = false;
connect(url: string): void { this.url = url; this.pending.push(() => this.openCb?.()); }
connect(url: string): void {
this.connectCalls++;
this.url = url;
this.pending.push(() => this.openCb?.());
}
send(frame: string): void { if (!this.closed) this.sent.push(frame); }
close(): void { if (this.closed) return; this.closed = true; this.pending.push(() => this.closeCb?.()); }
close(): void {
this.closeCalls++;
if (this.closed) return;
this.closed = true;
this.pending.push(() => this.closeCb?.());
}
onOpen(cb: () => void): void { this.openCb = cb; }
onMessage(cb: (f: string) => void): void { this.msgCb = cb; }
onClose(cb: () => void): void { this.closeCb = cb; }
/** 测试侧:模拟服务器推一帧给客户端。 */
serverPush(frame: string): void { if (!this.closed) this.pending.push(() => this.msgCb?.(frame)); }
serverPush(frame: string): void { this.pending.push(() => this.msgCb?.(frame)); }
/** 测试侧:模拟远端 close;允许重复信号以验证连接代际去重。 */
serverClose(): void {
this.closed = true;
this.pending.push(() => this.closeCb?.());
}
/** 测试侧:生产 adapter 的 error 最终收敛到 close 生命周期。 */
serverError(): void {
this.closed = true;
this.pending.push(() => this.closeCb?.());
}
/** 冲刷所有挂起回调(解析微任务)。 */
async flush(): Promise<void> {
@@ -0,0 +1,289 @@
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 clock: Clock = {
setTimeout(fn, ms) {
const id = ++sequence;
timers.set(id, { fn, at: now + ms });
return id;
},
clearTimeout(handle) {
timers.delete(handle as number);
},
};
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();
}
};
return { clock, advance, activeTimers: () => timers.size };
}
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('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();
client.send(envelope);
assert.deepEqual(transport.sent.map((frame) => JSON.parse(frame)), [envelope]);
client.stop();
assert.throws(() => client.send(envelope), /connection.*open/i);
});
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('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('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 } = 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);
advance(60000);
await first.flush();
assert.equal(transports.length, 1);
});