feat(framework): net-client 编排(login流程/门控/TcpID/超时重连/服务器切换/事件)
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,149 @@
|
|||||||
|
import { EventBus } from '../core/events.ts';
|
||||||
|
import { Route, RECV_TIMEOUT_MS, RECONNECT_INTERVAL_MS, LOGIN_GUARD_MS, SERVER_ROTATE_EVERY } from '../core/constants.ts';
|
||||||
|
import type { Transport } from './transport.ts';
|
||||||
|
import { encodeOutbound, decodeFrame } from './envelope-codec.ts';
|
||||||
|
import { HeartbeatWatchdog, realClock, type Clock } from './heartbeat.ts';
|
||||||
|
import { ReconnectPolicy } from './reconnect.ts';
|
||||||
|
import type { InboundMessage } from '../core/types/envelope.ts';
|
||||||
|
import type { LoginRequestData } from '../core/types/login.ts';
|
||||||
|
|
||||||
|
export interface NetClientEvents extends Record<string, any[]> {
|
||||||
|
open: [];
|
||||||
|
login: [unknown]; // player_login 响应 data
|
||||||
|
message: [InboundMessage];
|
||||||
|
slow: [];
|
||||||
|
reconnecting: [string]; // 目标 server
|
||||||
|
serverSwitch: [string]; // 新 server 地址
|
||||||
|
kicked: [unknown];
|
||||||
|
close: [];
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface NetClientOptions {
|
||||||
|
servers: string | string[];
|
||||||
|
transportFactory: () => Transport;
|
||||||
|
clock?: Clock;
|
||||||
|
bus?: EventBus<NetClientEvents>;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** 网络客户端:编排 transport/codec/heartbeat/reconnect,事件驱动,忠实 docs/protocol/01。 */
|
||||||
|
export class NetClient {
|
||||||
|
readonly bus: EventBus<NetClientEvents>;
|
||||||
|
private transportFactory: () => Transport;
|
||||||
|
private clock: Clock;
|
||||||
|
private policy: ReconnectPolicy;
|
||||||
|
private watchdog: HeartbeatWatchdog;
|
||||||
|
|
||||||
|
private transport: Transport | null = null;
|
||||||
|
private tcpId = 0; // 每次连接自增;onMessage 闭包捕获本次 id 做去重
|
||||||
|
private identity: LoginRequestData | null = null;
|
||||||
|
private isSendLoginState = false;
|
||||||
|
private isLogin = false;
|
||||||
|
private loginGuard: unknown = null;
|
||||||
|
private reconnectTimer: unknown = null;
|
||||||
|
private stopped = false;
|
||||||
|
|
||||||
|
constructor(opts: NetClientOptions) {
|
||||||
|
this.bus = opts.bus ?? new EventBus<NetClientEvents>();
|
||||||
|
this.transportFactory = opts.transportFactory;
|
||||||
|
this.clock = opts.clock ?? realClock;
|
||||||
|
this.policy = new ReconnectPolicy(opts.servers, SERVER_ROTATE_EVERY);
|
||||||
|
this.watchdog = new HeartbeatWatchdog(RECV_TIMEOUT_MS, () => this.onRecvTimeout(), this.clock);
|
||||||
|
}
|
||||||
|
|
||||||
|
setIdentity(identity: LoginRequestData): void { this.identity = identity; }
|
||||||
|
on = <K extends keyof NetClientEvents & string>(t: K, fn: (...a: NetClientEvents[K]) => void) => this.bus.on(t, fn);
|
||||||
|
off = <K extends keyof NetClientEvents & string>(t: K, fn: (...a: NetClientEvents[K]) => void) => this.bus.off(t, fn);
|
||||||
|
|
||||||
|
start(): void { this.stopped = false; this.connect(this.policy.current()); }
|
||||||
|
|
||||||
|
stop(): void {
|
||||||
|
this.stopped = true;
|
||||||
|
this.watchdog.stop();
|
||||||
|
if (this.reconnectTimer != null) { this.clock.clearTimeout(this.reconnectTimer); this.reconnectTimer = null; }
|
||||||
|
if (this.loginGuard != null) { this.clock.clearTimeout(this.loginGuard); this.loginGuard = null; }
|
||||||
|
this.transport?.close();
|
||||||
|
}
|
||||||
|
|
||||||
|
/** 发业务包(单层信封)。 */
|
||||||
|
send(route: string, rpc: string, data: unknown): void {
|
||||||
|
this.transport?.send(encodeOutbound(route, rpc, data));
|
||||||
|
}
|
||||||
|
|
||||||
|
private connect(url: string): void {
|
||||||
|
const id = ++this.tcpId;
|
||||||
|
const t = this.transportFactory();
|
||||||
|
this.transport = t;
|
||||||
|
t.onOpen(() => { if (id === this.tcpId) this.onOpen(); });
|
||||||
|
t.onMessage((frame) => { if (id === this.tcpId) this.onMessage(frame); });
|
||||||
|
t.onClose(() => { if (id === this.tcpId) this.onClose(); });
|
||||||
|
t.connect(url);
|
||||||
|
}
|
||||||
|
|
||||||
|
private onOpen(): void {
|
||||||
|
this.bus.emit('open');
|
||||||
|
this.sendLogin();
|
||||||
|
}
|
||||||
|
|
||||||
|
private sendLogin(): void {
|
||||||
|
if (!this.identity) return; // 无身份不发(docs §6.1:首登需身份)
|
||||||
|
this.isSendLoginState = true;
|
||||||
|
this.transport?.send(encodeOutbound(Route.agent, 'player_login', this.identity));
|
||||||
|
if (this.loginGuard != null) this.clock.clearTimeout(this.loginGuard);
|
||||||
|
this.loginGuard = this.clock.setTimeout(() => { this.loginGuard = null; this.onLoginGuardTimeout(); }, LOGIN_GUARD_MS); // docs §6.2
|
||||||
|
}
|
||||||
|
|
||||||
|
private onMessage(frame: string): void {
|
||||||
|
const r = decodeFrame(frame);
|
||||||
|
if (r.kind === 'ignore') return;
|
||||||
|
this.watchdog.feed(); // 收任意有效帧重置看门狗 docs §4.1
|
||||||
|
if (r.kind === 'handshake' || r.kind === 'heartbeat') return; // 忽略,心跳不回包
|
||||||
|
if (r.kind === 'serverDown') { this.sendLogin(); return; } // docs §4.2(剥离 UI 判定,直接重发登录)
|
||||||
|
|
||||||
|
const msg = r.message;
|
||||||
|
if (msg.rpc === 'submit_error') return; // docs §3.2-9
|
||||||
|
if (this.isSendLoginState) { // docs §3.2-10 登录态门控
|
||||||
|
if (msg.rpc === 'player_login') {
|
||||||
|
this.isSendLoginState = false; this.isLogin = true;
|
||||||
|
if (this.loginGuard != null) { this.clock.clearTimeout(this.loginGuard); this.loginGuard = null; }
|
||||||
|
this.policy.reset();
|
||||||
|
this.bus.emit('login', msg.data);
|
||||||
|
} else if (msg.rpc === 'kick_server') {
|
||||||
|
this.bus.emit('kicked', msg.data);
|
||||||
|
}
|
||||||
|
return; // 其它一律丢弃
|
||||||
|
}
|
||||||
|
// 服务器切换指令 docs §7.4
|
||||||
|
if (msg.rpc === 'connect_roomserver' && (msg.data as any)?.roomserver) {
|
||||||
|
this.switchServer((msg.data as any).roomserver); return;
|
||||||
|
}
|
||||||
|
if (msg.rpc === 'connect_agentserver' && (msg.data as any)?.agentserver) {
|
||||||
|
this.switchServer((msg.data as any).agentserver); return;
|
||||||
|
}
|
||||||
|
this.bus.emit('message', msg);
|
||||||
|
}
|
||||||
|
|
||||||
|
private switchServer(addr: string): void {
|
||||||
|
this.policy.setCurrent(addr);
|
||||||
|
this.bus.emit('serverSwitch', addr);
|
||||||
|
this.reconnectNow(); // 关闭当前→重连到新地址→重登
|
||||||
|
}
|
||||||
|
|
||||||
|
private onRecvTimeout(): void { this.bus.emit('slow'); this.reconnectNow(); } // docs §4.1
|
||||||
|
private onLoginGuardTimeout(): void { this.reconnectNow(); } // docs §6.2
|
||||||
|
|
||||||
|
private onClose(): void {
|
||||||
|
this.watchdog.stop();
|
||||||
|
if (this.stopped) return;
|
||||||
|
const next = this.policy.onFailure();
|
||||||
|
this.bus.emit('reconnecting', next);
|
||||||
|
if (this.reconnectTimer != null) this.clock.clearTimeout(this.reconnectTimer);
|
||||||
|
this.reconnectTimer = this.clock.setTimeout(() => { this.reconnectTimer = null; this.connect(next); }, RECONNECT_INTERVAL_MS); // docs §7.2 间隔 10s
|
||||||
|
}
|
||||||
|
|
||||||
|
/** 主动关闭当前连接,交由 onClose 走重连。 */
|
||||||
|
private reconnectNow(): void {
|
||||||
|
this.isSendLoginState = false;
|
||||||
|
if (this.loginGuard != null) { this.clock.clearTimeout(this.loginGuard); this.loginGuard = null; }
|
||||||
|
this.transport?.close();
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,125 @@
|
|||||||
|
import { test } from 'node:test';
|
||||||
|
import assert from 'node:assert/strict';
|
||||||
|
import { NetClient } from '../../YouleNexus/assets/framework/net/net-client.ts';
|
||||||
|
import { FakeTransport } from '../helpers/fake-transport.ts';
|
||||||
|
import type { Clock } from '../../YouleNexus/assets/framework/net/heartbeat.ts';
|
||||||
|
|
||||||
|
function fakeClock() {
|
||||||
|
let seq = 0, now = 0;
|
||||||
|
const timers = new Map<number, { fn: () => void; at: number }>();
|
||||||
|
const clock: Clock = {
|
||||||
|
setTimeout: (fn, ms) => { const id = ++seq; timers.set(id, { fn, at: now + ms }); return id; },
|
||||||
|
clearTimeout: (id: any) => { timers.delete(id); },
|
||||||
|
};
|
||||||
|
const advance = (ms: number) => { now += ms; for (const [id, t] of [...timers]) if (t.at <= now) { timers.delete(id); t.fn(); } };
|
||||||
|
return { clock, advance };
|
||||||
|
}
|
||||||
|
|
||||||
|
const IDENTITY = { agentid: 1, gameid: 2, openid: 'o', nickname: 'n', avatar: 'a', sex: 0, province: '', city: '', unionid: 'u', version: '1', channelid: 0, marketid: 0 };
|
||||||
|
|
||||||
|
function makeClient(transports: FakeTransport[], clock: Clock) {
|
||||||
|
let i = 0;
|
||||||
|
const client = new NetClient({
|
||||||
|
servers: 'ws://a',
|
||||||
|
transportFactory: () => transports[i++] ?? transports[transports.length - 1],
|
||||||
|
clock,
|
||||||
|
});
|
||||||
|
client.setIdentity(IDENTITY);
|
||||||
|
return client;
|
||||||
|
}
|
||||||
|
|
||||||
|
test('onOpen 后自动发 player_login(单层信封)', async () => {
|
||||||
|
const t = new FakeTransport();
|
||||||
|
const { clock } = fakeClock();
|
||||||
|
const client = makeClient([t], clock);
|
||||||
|
client.start();
|
||||||
|
await t.flush();
|
||||||
|
assert.equal(t.sent.length, 1);
|
||||||
|
const env = JSON.parse(t.sent[0]);
|
||||||
|
assert.equal(env.app, 'youle');
|
||||||
|
assert.equal(env.route, 'agent');
|
||||||
|
assert.equal(env.rpc, 'player_login');
|
||||||
|
assert.equal(env.data.openid, 'o');
|
||||||
|
});
|
||||||
|
|
||||||
|
test('isSendLoginState 门控:login 响应前的其它包被丢弃,login 响应放行并 emit login', async () => {
|
||||||
|
const t = new FakeTransport();
|
||||||
|
const { clock } = fakeClock();
|
||||||
|
const client = makeClient([t], clock);
|
||||||
|
const messages: any[] = [];
|
||||||
|
let loginResp: any = null;
|
||||||
|
client.on('message', (m: any) => messages.push(m));
|
||||||
|
client.on('login', (d: any) => { loginResp = d; });
|
||||||
|
client.start();
|
||||||
|
await t.flush();
|
||||||
|
// 登录响应前推一个业务包 → 应被门控丢弃
|
||||||
|
t.serverPush(JSON.stringify({ data: JSON.stringify({ route: 'agent', rpc: 'update_bean', data: { bean: 5 } }) }));
|
||||||
|
await t.flush();
|
||||||
|
assert.equal(messages.length, 0);
|
||||||
|
// 推 login 响应 → 放行,emit login,清门控
|
||||||
|
t.serverPush(JSON.stringify({ data: JSON.stringify({ route: 'agent', rpc: 'player_login', data: { state: 0, playerid: 9 } }) }));
|
||||||
|
await t.flush();
|
||||||
|
assert.equal(loginResp.playerid, 9);
|
||||||
|
// 门控已清,后续业务包放行
|
||||||
|
t.serverPush(JSON.stringify({ data: JSON.stringify({ route: 'agent', rpc: 'update_bean', data: { bean: 5 } }) }));
|
||||||
|
await t.flush();
|
||||||
|
assert.equal(messages.length, 1);
|
||||||
|
assert.equal(messages[0].rpc, 'update_bean');
|
||||||
|
});
|
||||||
|
|
||||||
|
test('握手包与心跳包被忽略且不进 message', async () => {
|
||||||
|
const t = new FakeTransport();
|
||||||
|
const { clock } = fakeClock();
|
||||||
|
const client = makeClient([t], clock);
|
||||||
|
const messages: any[] = [];
|
||||||
|
client.on('message', (m: any) => messages.push(m));
|
||||||
|
client.start();
|
||||||
|
await t.flush();
|
||||||
|
// 先完成登录以清门控
|
||||||
|
t.serverPush(JSON.stringify({ data: JSON.stringify({ route: 'agent', rpc: 'player_login', data: { state: 0, playerid: 9 } }) }));
|
||||||
|
await t.flush();
|
||||||
|
t.serverPush(JSON.stringify({ data: '@toconconABC' })); // 握手
|
||||||
|
t.serverPush(JSON.stringify({ data: JSON.stringify({ com: '@serverheartbeat' }) })); // 心跳
|
||||||
|
await t.flush();
|
||||||
|
assert.equal(messages.length, 0);
|
||||||
|
});
|
||||||
|
|
||||||
|
test('收包超时 emit slow 并触发重连(onClose→新连接重发 login)', async () => {
|
||||||
|
const t1 = new FakeTransport(), t2 = new FakeTransport();
|
||||||
|
const { clock, advance } = fakeClock();
|
||||||
|
const client = makeClient([t1, t2], clock);
|
||||||
|
let slow = 0;
|
||||||
|
client.on('slow', () => { slow++; });
|
||||||
|
client.start();
|
||||||
|
await t1.flush();
|
||||||
|
// 先完成登录:清 4s 守护、feed 收包看门狗(否则 advance 会先撞上 loginGuard 而非 watchdog)
|
||||||
|
t1.serverPush(JSON.stringify({ data: JSON.stringify({ route: 'agent', rpc: 'player_login', data: { state: 0, playerid: 9 } }) }));
|
||||||
|
await t1.flush();
|
||||||
|
// 30s 无收包 → watchdog → slow + 关闭当前连接
|
||||||
|
advance(30000);
|
||||||
|
await t1.flush();
|
||||||
|
assert.equal(slow, 1);
|
||||||
|
// 重连定时器到点 → 新连接 t2 建立并重发 login
|
||||||
|
advance(10000);
|
||||||
|
await t2.flush();
|
||||||
|
assert.equal(t2.sent.length, 1);
|
||||||
|
assert.equal(JSON.parse(t2.sent[0]).rpc, 'player_login');
|
||||||
|
});
|
||||||
|
|
||||||
|
test('TcpID 去重:旧连接的残留包被丢弃', async () => {
|
||||||
|
const t1 = new FakeTransport(), t2 = new FakeTransport();
|
||||||
|
const { clock, advance } = fakeClock();
|
||||||
|
const client = makeClient([t1, t2], clock);
|
||||||
|
const messages: any[] = [];
|
||||||
|
client.on('message', (m: any) => messages.push(m));
|
||||||
|
client.start();
|
||||||
|
await t1.flush();
|
||||||
|
t1.serverPush(JSON.stringify({ data: JSON.stringify({ route: 'agent', rpc: 'player_login', data: { state: 0, playerid: 9 } }) }));
|
||||||
|
await t1.flush();
|
||||||
|
// 触发重连切到 t2
|
||||||
|
advance(30000); await t1.flush(); advance(10000); await t2.flush();
|
||||||
|
// 旧连接 t1 仍推包 → 应被 TcpID 去重丢弃
|
||||||
|
t1.serverPush(JSON.stringify({ data: JSON.stringify({ route: 'agent', rpc: 'update_bean', data: { bean: 1 } }) }));
|
||||||
|
await t1.flush();
|
||||||
|
assert.equal(messages.length, 0);
|
||||||
|
});
|
||||||
Reference in New Issue
Block a user