From a8873b9a84ac4bcd4a2f7ee3ec538421bb211046 Mon Sep 17 00:00:00 2001 From: Joywayer Date: Sat, 5 Sep 2026 03:57:32 +0800 Subject: [PATCH] refactor(net): isolate transport and wire lifecycle --- .../assets/framework/net/wire-client.ts | 202 ++++++++++++ .../framework-tests/helpers/fake-transport.ts | 29 +- .../framework-tests/net/wire-client.test.ts | 289 ++++++++++++++++++ 3 files changed, 517 insertions(+), 3 deletions(-) create mode 100644 cocoscreator_projects/YouleNexus/assets/framework/net/wire-client.ts create mode 100644 cocoscreator_projects/framework-tests/net/wire-client.test.ts diff --git a/cocoscreator_projects/YouleNexus/assets/framework/net/wire-client.ts b/cocoscreator_projects/YouleNexus/assets/framework/net/wire-client.ts new file mode 100644 index 0000000..e02c31c --- /dev/null +++ b/cocoscreator_projects/YouleNexus/assets/framework/net/wire-client.ts @@ -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(); + + 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; + } +} diff --git a/cocoscreator_projects/framework-tests/helpers/fake-transport.ts b/cocoscreator_projects/framework-tests/helpers/fake-transport.ts index c8ad944..6f9eb6d 100644 --- a/cocoscreator_projects/framework-tests/helpers/fake-transport.ts +++ b/cocoscreator_projects/framework-tests/helpers/fake-transport.ts @@ -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 { diff --git a/cocoscreator_projects/framework-tests/net/wire-client.test.ts b/cocoscreator_projects/framework-tests/net/wire-client.test.ts new file mode 100644 index 0000000..07bd975 --- /dev/null +++ b/cocoscreator_projects/framework-tests/net/wire-client.test.ts @@ -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 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); +});