From 5882efc1098406ebd5efd57af76d5047e514dd85 Mon Sep 17 00:00:00 2001 From: troytt <47798984@qq.com> Date: Thu, 1 Oct 2026 12:01:47 +0000 Subject: [PATCH] fix(netproto): detect BNCS/MCP/D2GS session loss centrally in D2OnlineFlow D2OnlineFlow never subscribed to the BNCS/MCP socket close events and kept stale session references, so a dropped realm connection only surfaced later as "MCP realm session is not connected" on whatever screen issued the next command. The flow now watches every connection it opens and turns any unrequested close/error, request timeout, connect failure or command without its connection into exactly one onSessionLost event. Before emitting it tears down all remaining connections (references cleared first so the resulting close callbacks are not reported again), bumps an epoch so late-opening streams are disposed, and enters `closed`. A BNCS/MCP loss during a game is reported when the game ends. A failed D2GS join releases its D2GS session again. WsStream forwards the WebSocket close reason. Tests: in-memory BNCS/MCP/D2GS fake servers cover BNCS close on char select, MCP close on char create, realm logon / MCP_STARTUP timeouts, both sockets closing (one event), D2GS drop in game, deferred in-game BNCS loss and close() without an event. TAG=agy CONV=109a3012-103a-42f0-8ae1-bbc0857047ef --- src/netproto/flow/config.ts | 29 ++ src/netproto/flow/online-flow.ts | 476 ++++++++++++++---- src/netproto/index.ts | 9 +- src/netproto/transport/byte-stream.ts | 2 +- src/netproto/transport/ws-stream.ts | 7 +- tests/client/hud-session-play.test.ts | 2 + tests/netproto/fake-realm-servers.ts | 208 ++++++++ .../netproto/online-flow-session-lost.test.ts | 195 +++++++ 8 files changed, 817 insertions(+), 111 deletions(-) create mode 100644 tests/netproto/fake-realm-servers.ts create mode 100644 tests/netproto/online-flow-session-lost.test.ts diff --git a/src/netproto/flow/config.ts b/src/netproto/flow/config.ts index b24847c..8bca320 100644 --- a/src/netproto/flow/config.ts +++ b/src/netproto/flow/config.ts @@ -18,6 +18,35 @@ export type D2OnlineState = | 'ingame' | 'closed' +/** Which live connection of the flow was lost. */ +export type D2OnlineSessionKind = 'bncs' | 'mcp' | 'd2gs' + +/** + * Why the connection was lost: + * - `closed`: the socket closed without the flow asking for it (server/bridge close). + * - `error`: the transport reported an error (socket error / write on a dead socket). + * - `timeout`: a request/response exchange on that connection timed out. + * - `connect_failed`: opening the connection (or the realm MCP startup) failed. + * - `not_connected`: a command needed the connection but it is not connected. + */ +export type D2OnlineSessionLostCause = 'closed' | 'error' | 'timeout' | 'connect_failed' | 'not_connected' + +/** + * Emitted exactly once per `D2OnlineFlow` when a connection it owns is lost. By the time listeners run, + * the flow has already closed every remaining connection and is in state `closed`. + */ +export interface D2OnlineSessionLost { + readonly session: D2OnlineSessionKind + readonly cause: D2OnlineSessionLostCause + /** Concrete technical description (connection, close code/reason, timed-out packet, error text). */ + readonly message: string + readonly closeCode?: number | undefined + readonly closeReason?: string | undefined + readonly error?: Error | undefined + /** Flow state at the moment the loss was detected. */ + readonly flowState: D2OnlineState +} + export interface D2OnlineConfig { readonly resolver: EndpointResolver readonly bnetHost: string diff --git a/src/netproto/flow/online-flow.ts b/src/netproto/flow/online-flow.ts index deeb375..e96d1a0 100644 --- a/src/netproto/flow/online-flow.ts +++ b/src/netproto/flow/online-flow.ts @@ -1,12 +1,18 @@ /** * `D2OnlineFlow` orchestrating BNCS (port 6112) -> MCP (port 6113) -> D2GS (port 4000) * transitions and keeping BNCS/MCP alive while in-game. + * + * Connection loss is detected centrally here: every BNCS/MCP/D2GS connection the flow opens is + * watched, and any close the flow did not ask for, any request timeout, any connect failure and any + * command issued while its connection is gone produces exactly one `onSessionLost` event per flow. + * Before that event is emitted the flow closes every remaining connection (no leaked sockets, timers or + * reconnects) and enters state `closed`. */ import { BncsSession } from '../bncs/session.ts' import { defaultClock } from '../core/clock.ts' import { Emitter } from '../core/emitter.ts' -import { ProtocolError } from '../core/errors.ts' +import { ClosedError, ProtocolError, TimeoutError } from '../core/errors.ts' import { resolveTextCodec } from '../core/text-codec.ts' import { D2gsAdapter } from '../d2gs/adapter.ts' import { D2gsSession } from '../d2gs/session.ts' @@ -22,11 +28,22 @@ import type { RealmInfo, } from '../domain/lobby.ts' import { McpSession } from '../mcp/session.ts' -import type { D2OnlineConfig, D2OnlineState } from './config.ts' +import type { CloseReason } from '../transport/byte-stream.ts' +import type { + D2OnlineConfig, + D2OnlineSessionKind, + D2OnlineSessionLost, + D2OnlineSessionLostCause, + D2OnlineState, +} from './config.ts' export interface D2OnlineFlow { readonly state: D2OnlineState + /** The loss that ended this flow, or `null` while the flow is healthy / was closed on purpose. */ + readonly sessionLost: D2OnlineSessionLost | null onState(cb: (s: D2OnlineState) => void): () => void + /** Fires once when a BNCS/MCP/D2GS connection is lost (see `D2OnlineSessionLost`). */ + onSessionLost(cb: (ev: D2OnlineSessionLost) => void): () => void connectBnet(): Promise connect(): Promise createAccount(user: string, pass: string): Promise @@ -78,6 +95,59 @@ const FALLBACK_EMPTY_TABLES: ItemDataTables = { getItemMeta: () => undefined, } +const SESSION_LABEL: Readonly> = { + bncs: 'BNCS', + mcp: 'MCP realm', + d2gs: 'D2GS game', +} + +type SessionLostInfo = Omit + +/** Drop the `[ProtocolError conn=.. packet=.. offset=..]` prefix so user-facing reasons stay readable. */ +function stripProtocolPrefix(message: string): string { + return message.replace(/^\[ProtocolError [^\]]*\]\s*/, '') +} + +function errorText(err: unknown): string { + return stripProtocolPrefix(err instanceof Error ? err.message : String(err)) +} + +function describeClose(kind: D2OnlineSessionKind, reason: CloseReason): SessionLostInfo { + const label = SESSION_LABEL[kind] + if (reason.kind === 'error') { + return { + session: kind, + cause: 'error', + message: `${label} connection error: ${errorText(reason.error)}`, + error: reason.error, + } + } + const remoteReason = reason.kind === 'remote' && reason.reason ? reason.reason : undefined + const details: string[] = [] + if (reason.code !== undefined) details.push(`code ${reason.code}`) + if (remoteReason !== undefined) details.push(`reason "${remoteReason}"`) + return { + session: kind, + cause: 'closed', + message: `${label} connection closed${reason.kind === 'remote' ? ' by server' : ''}${ + details.length > 0 ? ` (${details.join(', ')})` : '' + }`, + ...(reason.code !== undefined ? { closeCode: reason.code } : {}), + ...(remoteReason !== undefined ? { closeReason: remoteReason } : {}), + } +} + +/** Transport-level failures of a request/response exchange mean the connection is unusable. */ +function classifyExchangeFailure(kind: D2OnlineSessionKind, err: unknown): SessionLostInfo | null { + if (err instanceof TimeoutError) { + return { session: kind, cause: 'timeout', message: `${SESSION_LABEL[kind]} ${errorText(err)}`, error: err } + } + if (err instanceof ClosedError) { + return { session: kind, cause: 'error', message: `${SESSION_LABEL[kind]} ${errorText(err)}`, error: err } + } + return null +} + export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { const bnetHost = config.bnetHost const bnetPort = config.bnetPort ?? 6112 @@ -90,6 +160,7 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { const stateEmitter = new Emitter() const chatEmitter = new Emitter() + const lostEmitter = new Emitter() let currentState: D2OnlineState = 'idle' let bncs: BncsSession | undefined @@ -99,6 +170,11 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { let selectedCharName = '' let selectedCharClass = 0 let knownChars: readonly CharSummary[] = [] + let lost: D2OnlineSessionLost | null = null + /** BNCS/MCP loss seen while a D2GS game runs: reported when that game ends instead of aborting it. */ + let pendingLost: D2OnlineSessionLost | null = null + /** Bumped by every teardown so connections that finish opening afterwards are disposed, not adopted. */ + let epoch = 0 function setState(next: D2OnlineState): void { if (currentState !== next) { @@ -107,18 +183,108 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { } } - function requireBncs(): BncsSession { - if (!bncs) { - throw new ProtocolError('BNCS session is not connected', { proto: 'bncs' }) + /** + * Close every connection the flow owns. References are cleared before closing, so the resulting + * close callbacks fail the identity checks in the watchers and are never reported as losses. + */ + function teardownConnections(): void { + epoch++ + const game = d2gsSession + d2gsSession = undefined + d2gsAdapter = undefined + const realm = mcp + mcp = undefined + const bnet = bncs + bncs = undefined + if (game) { + try { + game.leaveAndClose() + } catch { + // The socket is already gone; nothing left to release. + } } - return bncs + realm?.close() + bnet?.close() } - function requireMcp(): McpSession { - if (!mcp) { - throw new ProtocolError('MCP realm session is not connected', { proto: 'mcp' }) + function declareLost(info: SessionLostInfo): void { + if (lost) return + const ev: D2OnlineSessionLost = { ...info, flowState: currentState } + lost = ev + pendingLost = null + teardownConnections() + setState('closed') + lostEmitter.emit(ev) + } + + function onConnectionLost(info: SessionLostInfo): void { + if (lost) return + if (info.session !== 'd2gs' && d2gsSession !== undefined && currentState === 'ingame') { + pendingLost ??= { ...info, flowState: currentState } + if (info.session === 'bncs') { + bncs?.close() + bncs = undefined + } else { + mcp?.close() + mcp = undefined + } + return } - return mcp + declareLost(info) + } + + function watchLobbyConnection(kind: 'bncs' | 'mcp', session: BncsSession | McpSession): void { + session.onClose((reason) => { + const active = kind === 'bncs' ? bncs : mcp + if (lost || active !== session) return + onConnectionLost(describeClose(kind, reason)) + }) + } + + function watchGameConnection(session: D2gsSession): void { + session.onClose((reason) => { + if (lost || d2gsSession !== session) return + const info = describeClose('d2gs', reason) + const earlier = pendingLost + declareLost(earlier ? { ...info, message: `${info.message}; earlier: ${earlier.message}` } : info) + }) + } + + async function exchange( + kind: 'bncs' | 'mcp', + op: () => Promise, + otherFailureCause?: D2OnlineSessionLostCause, + ): Promise { + try { + return await op() + } catch (err) { + const info = + classifyExchangeFailure(kind, err) ?? + (otherFailureCause !== undefined + ? { session: kind, cause: otherFailureCause, message: `${SESSION_LABEL[kind]} ${errorText(err)}`, ...(err instanceof Error ? { error: err } : {}) } + : null) + if (info) onConnectionLost(info) + throw err + } + } + + function notConnected(kind: 'bncs' | 'mcp', op: string): never { + const message = `${SESSION_LABEL[kind]} session is not connected for ${op}` + const err = new ProtocolError(message, { proto: kind }) + onConnectionLost({ session: kind, cause: 'not_connected', message, error: err }) + throw err + } + + function requireBncs(op: string): BncsSession { + return bncs ?? notConnected('bncs', op) + } + + function requireMcp(op: string): McpSession { + return mcp ?? notConnected('mcp', op) + } + + function staleConnectionError(kind: 'bncs' | 'mcp'): ProtocolError { + return new ProtocolError(`D2OnlineFlow was closed while ${SESSION_LABEL[kind]} was connecting`, { proto: kind }) } const flow: D2OnlineFlow = { @@ -126,20 +292,52 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { return currentState }, + get sessionLost(): D2OnlineSessionLost | null { + return lost + }, + onState(cb: (s: D2OnlineState) => void): () => void { return stateEmitter.on(cb) }, + onSessionLost(cb: (ev: D2OnlineSessionLost) => void): () => void { + return lostEmitter.on(cb) + }, + async connectBnet(): Promise { + // A (re)connect starts a fresh lifecycle for this flow. + teardownConnections() + lost = null + pendingLost = null + const startEpoch = epoch setState('connecting_bnet') - const stream = await config.resolver.open('bnet', bnetHost, bnetPort) - bncs = new BncsSession(stream, { + let stream + try { + stream = await config.resolver.open('bnet', bnetHost, bnetPort) + } catch (err) { + if (startEpoch === epoch) { + declareLost({ + session: 'bncs', + cause: 'connect_failed', + message: `BNCS connect to ${bnetHost}:${bnetPort} failed: ${errorText(err)}`, + ...(err instanceof Error ? { error: err } : {}), + }) + } + throw err + } + if (startEpoch !== epoch) { + stream.close() + throw staleConnectionError('bncs') + } + const session = new BncsSession(stream, { clock, textCodec, packetTap: config.tap, }) - bncs.onChatEvent((ev) => chatEmitter.emit(ev)) - await bncs.handshake() + bncs = session + watchLobbyConnection('bncs', session) + session.onChatEvent((ev) => chatEmitter.emit(ev)) + await exchange('bncs', () => session.handshake()) }, async connect(): Promise { @@ -150,8 +348,8 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { if (!bncs) { await flow.connectBnet() } - const session = requireBncs() - const res = await session.createAccount(user, pass) + const session = requireBncs('createAccount') + const res = await exchange('bncs', () => session.createAccount(user, pass)) if (res.status !== 0) { throw new ProtocolError( `BNCS SID_CREATEACCOUNT2 failed for "${user}" with status=0x${res.status.toString(16)}`, @@ -172,18 +370,18 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { if (!bncs) { await flow.connectBnet() } - const session = requireBncs() + const session = requireBncs('login') if (opts?.register) { - const firstLogin = await session.login(user, pass) + const firstLogin = await exchange('bncs', () => session.login(user, pass)) if (firstLogin.status !== 0) { - const createRes = await session.createAccount(user, pass) + const createRes = await exchange('bncs', () => session.createAccount(user, pass)) if (createRes.status !== 0) { throw new ProtocolError( `BNCS SID_CREATEACCOUNT2 failed for "${user}" with status=0x${createRes.status.toString(16)}`, { proto: 'bncs', packetId: 0x3d }, ) } - const secondLogin = await session.login(user, pass) + const secondLogin = await exchange('bncs', () => session.login(user, pass)) if (secondLogin.status !== 0) { throw new ProtocolError( `BNCS SID_LOGONRESPONSE2 failed after registration for "${user}" with status=0x${secondLogin.status.toString(16)}`, @@ -192,7 +390,7 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { } } } else { - const loginRes = await session.login(user, pass) + const loginRes = await exchange('bncs', () => session.login(user, pass)) if (loginRes.status !== 0) { throw new ProtocolError( `BNCS SID_LOGONRESPONSE2 failed for "${user}" with status=0x${loginRes.status.toString(16)}`, @@ -204,28 +402,50 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { }, async listRealms(): Promise { - const session = requireBncs() - return await session.listRealms() + const session = requireBncs('listRealms') + return await exchange('bncs', () => session.listRealms()) }, async enterRealm(realm?: string): Promise { - const session = requireBncs() + const session = requireBncs('enterRealm') let targetRealm = realm ?? config.defaultRealm if (!targetRealm) { - const realms = await session.listRealms() + const realms = await exchange('bncs', () => session.listRealms()) targetRealm = realms[0]?.title ?? 'D2CS' } + const realmTitle = targetRealm setState('connecting_realm') - const realmData = await session.logonRealm(targetRealm) - const mcpStream = await config.resolver.open('realm', bnetHost, realmPort || realmData.mcpPort) - mcp = new McpSession(mcpStream, { + const realmData = await exchange('bncs', () => session.logonRealm(realmTitle)) + const port = realmPort || realmData.mcpPort + const startEpoch = epoch + let mcpStream + try { + mcpStream = await config.resolver.open('realm', bnetHost, port) + } catch (err) { + if (startEpoch === epoch) { + declareLost({ + session: 'mcp', + cause: 'connect_failed', + message: `MCP realm connect to ${bnetHost}:${port} failed: ${errorText(err)}`, + ...(err instanceof Error ? { error: err } : {}), + }) + } + throw err + } + if (startEpoch !== epoch) { + mcpStream.close() + throw staleConnectionError('mcp') + } + const realmSession = new McpSession(mcpStream, { clock, textCodec, packetTap: config.tap, }) - await mcp.startup(realmData) - knownChars = await mcp.listChars(8) + mcp = realmSession + watchLobbyConnection('mcp', realmSession) + await exchange('mcp', () => realmSession.startup(realmData), 'connect_failed') + knownChars = await exchange('mcp', () => realmSession.listChars(8)) setState('realm_ready') return knownChars }, @@ -235,8 +455,8 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { }, async listCharacters(): Promise { - const m = requireMcp() - knownChars = await m.listChars(8) + const m = requireMcp('listCharacters') + knownChars = await exchange('mcp', () => m.listChars(8)) return knownChars }, @@ -245,12 +465,14 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { }, async createCharacter(spec: CharCreateSpec): Promise { - const m = requireMcp() - const status = await m.createChar(spec.name, spec.classId, { - expansion: spec.expansion, - hardcore: spec.hardcore, - ladder: spec.ladder, - }) + const m = requireMcp('createCharacter') + const status = await exchange('mcp', () => + m.createChar(spec.name, spec.classId, { + expansion: spec.expansion, + hardcore: spec.hardcore, + ladder: spec.ladder, + }), + ) if (status !== 0) { throw new ProtocolError( `MCP_CHARCREATE failed for "${spec.name}" with status=0x${status.toString(16)}`, @@ -259,7 +481,7 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { } selectedCharName = spec.name selectedCharClass = spec.classId - knownChars = await m.listChars(8) + knownChars = await exchange('mcp', () => m.listChars(8)) }, async createChar( @@ -281,9 +503,9 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { }, async deleteCharacter(name: string): Promise { - const m = requireMcp() - await m.deleteChar(name) - knownChars = await m.listChars(8) + const m = requireMcp('deleteCharacter') + await exchange('mcp', () => m.deleteChar(name)) + knownChars = await exchange('mcp', () => m.listChars(8)) }, async deleteChar(name: string): Promise { @@ -291,9 +513,9 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { }, async upgradeCharacter(name: string): Promise { - const m = requireMcp() - await m.upgradeChar(name) - knownChars = await m.listChars(8) + const m = requireMcp('upgradeCharacter') + await exchange('mcp', () => m.upgradeChar(name)) + knownChars = await exchange('mcp', () => m.listChars(8)) }, async upgradeChar(name: string): Promise { @@ -301,8 +523,8 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { }, async selectCharacter(name: string): Promise { - const m = requireMcp() - const status = await m.charLogon(name) + const m = requireMcp('selectCharacter') + const status = await exchange('mcp', () => m.charLogon(name)) if (status !== 0) { throw new ProtocolError( `MCP_CHARLOGON failed for "${name}" with status=0x${status.toString(16)}`, @@ -321,18 +543,18 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { }, async listGames(filter = ''): Promise { - const m = requireMcp() - return await m.listGames(filter) + const m = requireMcp('listGames') + return await exchange('mcp', () => m.listGames(filter)) }, async getGameInfo(name: string): Promise { - const m = requireMcp() - return await m.getGameInfo(name) + const m = requireMcp('getGameInfo') + return await exchange('mcp', () => m.getGameInfo(name)) }, async createGame(spec: GameCreateSpec): Promise { - const m = requireMcp() - const res = await m.createGame(spec) + const m = requireMcp('createGame') + const res = await exchange('mcp', () => m.createGame(spec)) if (res.status !== 0 && res.status !== 0x1e) { throw new ProtocolError( `MCP_CREATEGAME failed for "${spec.name}" with status=0x${res.status.toString(16)}`, @@ -342,11 +564,11 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { }, async joinGame(name: string, password = ''): Promise { - const m = requireMcp() + const m = requireMcp('joinGame') if (!selectedCharName) { throw new ProtocolError('Cannot joinGame before selectCharacter()', { proto: 'mcp' }) } - const ticket = await m.joinGame(name, password) + const ticket = await exchange('mcp', () => m.joinGame(name, password)) if (ticket.status !== 0) { throw new ProtocolError( `MCP_JOINGAME failed for "${name}" with status=0x${ticket.status.toString(16)}`, @@ -355,8 +577,19 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { } setState('connecting_game') - const gameStream = await config.resolver.open('game', bnetHost, gamePort) - d2gsSession = new D2gsSession({ + const startEpoch = epoch + let gameStream + try { + gameStream = await config.resolver.open('game', bnetHost, gamePort) + } catch (err) { + if (startEpoch === epoch && !lost) setState(mcp ? 'realm_ready' : 'idle') + throw err + } + if (startEpoch !== epoch) { + gameStream.close() + throw new ProtocolError('D2OnlineFlow was closed while D2GS was connecting', { proto: 'd2gs' }) + } + const session = new D2gsSession({ stream: gameStream, logon: { gameHash: ticket.gameHash, @@ -367,58 +600,88 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { clock, tap: config.tap, }) - - d2gsAdapter = new D2gsAdapter({ - session: d2gsSession, + const adapter = new D2gsAdapter({ + session, tables, encoding, clock, }) + d2gsSession = session + d2gsAdapter = adapter - await new Promise((resolve, reject) => { - const cancelTimer = clock.setTimeout(() => { - unsubState() - unsubErr() - reject(new ProtocolError(`D2GS handshake timed out (state=${d2gsSession?.state})`, { proto: 'd2gs' })) - }, 12_000) - const unsubState = d2gsSession!.onState((s) => { - if (s === 'ingame') { + try { + await new Promise((resolve, reject) => { + const cancelTimer = clock.setTimeout(() => { + unsubState() + unsubErr() + reject(new ProtocolError(`D2GS handshake timed out (state=${session.state})`, { proto: 'd2gs' })) + }, 12_000) + const unsubState = session.onState((s) => { + if (s === 'ingame') { + cancelTimer() + unsubState() + unsubErr() + setState('ingame') + resolve() + } else if (s === 'closed') { + cancelTimer() + unsubState() + unsubErr() + reject(new ProtocolError('D2GS connection closed during handshake', { proto: 'd2gs' })) + } + }) + const unsubErr = session.onError((err) => { cancelTimer() unsubState() unsubErr() - setState('ingame') - resolve() - } else if (s === 'closed') { - cancelTimer() - unsubState() - unsubErr() - reject(new ProtocolError('D2GS connection closed during handshake', { proto: 'd2gs' })) + reject(err) + }) + }) + } catch (err) { + // A failed game join leaves BNCS/MCP usable (1.13c returns to the lobby); only the D2GS + // connection of this attempt is released. + if (d2gsSession === session) { + d2gsSession = undefined + d2gsAdapter = undefined + try { + session.leaveAndClose() + } catch { + // Already closed. } - }) - const unsubErr = d2gsSession!.onError((err) => { - cancelTimer() - unsubState() - unsubErr() - reject(err) - }) - }) + if (!lost) setState(mcp ? 'realm_ready' : 'idle') + } + throw err + } - return d2gsAdapter + watchGameConnection(session) + return adapter }, async enterChat(): Promise { - const session = requireBncs() - await session.enterChat(selectedCharName) + const session = requireBncs('enterChat') + await exchange('bncs', () => session.enterChat(selectedCharName)) }, async joinChannel(channel: string): Promise { - const session = requireBncs() - session.joinChannel(channel) + const session = requireBncs('joinChannel') + try { + session.joinChannel(channel) + } catch (err) { + const info = classifyExchangeFailure('bncs', err) + if (info) onConnectionLost(info) + throw err + } }, async sendChannelChat(text: string): Promise { - const session = requireBncs() - session.sendChat(text) + const session = requireBncs('sendChannelChat') + try { + session.sendChat(text) + } catch (err) { + const info = classifyExchangeFailure('bncs', err) + if (info) onConnectionLost(info) + throw err + } }, async sendChat(text: string): Promise { @@ -434,28 +697,27 @@ export function createD2OnlineFlow(config: D2OnlineConfig): D2OnlineFlow { }, async leaveToLobby(): Promise { - if (d2gsAdapter) { - await d2gsAdapter.close() - d2gsAdapter = undefined - d2gsSession = undefined + const game = d2gsSession + d2gsSession = undefined + d2gsAdapter = undefined + if (game) { + try { + game.leaveAndClose() + } catch { + // Already closed. + } + } + if (pendingLost) { + const { flowState: _flowState, ...info } = pendingLost + declareLost(info) + return } setState(mcp ? 'realm_ready' : 'idle') }, async close(): Promise { - if (d2gsAdapter) { - await d2gsAdapter.close() - d2gsAdapter = undefined - d2gsSession = undefined - } - if (mcp) { - mcp.close() - mcp = undefined - } - if (bncs) { - bncs.close() - bncs = undefined - } + pendingLost = null + teardownConnections() setState('closed') }, } diff --git a/src/netproto/index.ts b/src/netproto/index.ts index 1b4a089..c5611cc 100644 --- a/src/netproto/index.ts +++ b/src/netproto/index.ts @@ -14,7 +14,14 @@ export type { TcpResolverOptions, WsBridgeResolverOptions, } from './transport/endpoint.ts' -export type { D2OnlineConfig, D2OnlineConfig as OnlineConfig, D2OnlineState } from './flow/config.ts' +export type { + D2OnlineConfig, + D2OnlineConfig as OnlineConfig, + D2OnlineSessionKind, + D2OnlineSessionLost, + D2OnlineSessionLostCause, + D2OnlineState, +} from './flow/config.ts' export type { D2OnlineFlow } from './flow/online-flow.ts' export type { GameServerAdapter, GameServerStats } from './domain/game-server-adapter.ts' export type { ServerEvent } from './domain/server-event.ts' diff --git a/src/netproto/transport/byte-stream.ts b/src/netproto/transport/byte-stream.ts index 40d5aef..2d92b12 100644 --- a/src/netproto/transport/byte-stream.ts +++ b/src/netproto/transport/byte-stream.ts @@ -7,7 +7,7 @@ export type EndpointKind = 'bnet' | 'realm' | 'game' | 'ts' export type CloseReason = | { readonly kind: 'local'; readonly code?: number | undefined } - | { readonly kind: 'remote'; readonly code?: number | undefined } + | { readonly kind: 'remote'; readonly code?: number | undefined; readonly reason?: string | undefined } | { readonly kind: 'error'; readonly error: Error } export interface ByteStream { diff --git a/src/netproto/transport/ws-stream.ts b/src/netproto/transport/ws-stream.ts index 5acf5a9..59d5770 100644 --- a/src/netproto/transport/ws-stream.ts +++ b/src/netproto/transport/ws-stream.ts @@ -77,8 +77,11 @@ export class WsStream implements ByteStream { this.ws.onclose = (ev: { readonly code?: number; readonly reason?: string }) => { if (this.closed) return this.closed = true - const reason: CloseReason = - ev.code !== undefined ? { kind: 'remote', code: ev.code } : { kind: 'remote' } + const reason: CloseReason = { + kind: 'remote', + ...(ev.code !== undefined ? { code: ev.code } : {}), + ...(ev.reason ? { reason: ev.reason } : {}), + } this.closeEmitter.emit(reason) this.dataEmitter.clear() this.closeEmitter.clear() diff --git a/tests/client/hud-session-play.test.ts b/tests/client/hud-session-play.test.ts index 549f383..9b60e39 100644 --- a/tests/client/hud-session-play.test.ts +++ b/tests/client/hud-session-play.test.ts @@ -560,7 +560,9 @@ describe('Milestone M5 — HudModel, CommandMapper, and OnlineSession', () => { get state() { return flowState }, + sessionLost: null, onState: vi.fn(() => () => {}), + onSessionLost: vi.fn(() => () => {}), connectBnet: vi.fn(async () => {}), connect: vi.fn(async () => {}), login: vi.fn(async () => { diff --git a/tests/netproto/fake-realm-servers.ts b/tests/netproto/fake-realm-servers.ts new file mode 100644 index 0000000..e1c8ba5 --- /dev/null +++ b/tests/netproto/fake-realm-servers.ts @@ -0,0 +1,208 @@ +/** + * In-memory BNCS / MCP / D2GS fake servers for `createD2OnlineFlow` session-loss tests. + * + * Every `resolver.open()` creates a fresh `MemoryStream` pair (so a second login gets new + * connections) and records both ends, letting tests force-close the server side of a connection and + * assert that the flow disposed its client side. Responders follow the tier4 S6 harness. + */ + +import { + createMemoryStreamPair, + type ByteStream, + type EndpointResolver, + type EndpointRole, +} from '../../src/netproto/index.ts' +import { ByteReader } from '../../src/netproto/core/byte-reader.ts' +import { ByteWriter } from '../../src/netproto/core/byte-writer.ts' +import { BncsFramer, encodeBncsFrame } from '../../src/netproto/bncs/framing.ts' +import { BncsOpcode, encodeBncsPing } from '../../src/netproto/bncs/packets.ts' +import { encodeMcpFrame, McpFramer } from '../../src/netproto/mcp/framing.ts' +import { McpOpcode } from '../../src/netproto/mcp/packets.ts' +import { + defaultD2gsHuffmanCodec, + wrapD2gsCompressedBlock, +} from '../../src/netproto/d2gs/compression.ts' +import type { MemoryStream } from '../../src/netproto/transport/memory-stream.ts' + +export interface FakeConnection { + readonly role: EndpointRole + readonly client: MemoryStream + readonly server: MemoryStream +} + +export interface FakeRealmOptions { + /** BNCS request opcodes the fake server never answers (to provoke request timeouts). */ + readonly silentBncs?: readonly number[] | undefined + /** MCP request opcodes the fake server never answers. */ + readonly silentMcp?: readonly number[] | undefined + /** Characters in MCP_CHARLIST2 (0 = empty roster, i.e. the char_create screen). */ + readonly characters?: number | undefined +} + +export interface FakeRealm { + readonly resolver: EndpointResolver + readonly connections: FakeConnection[] + /** Most recent connection for `role`. */ + latest(role: EndpointRole): FakeConnection + /** Every client stream the resolver handed out is closed (nothing leaked). */ + allClientStreamsClosed(): boolean +} + +function charList(count: number): Uint8Array { + const w = new ByteWriter() + w.u16LE(8) + w.u32LE(count) + w.u16LE(count) + for (let i = 0; i < count; i++) { + w.u32LE(0x7fffffff) + w.cstring(`d2webbot_sorc${i === 0 ? '' : String(i)}`) + const stat = new Uint8Array(33).fill(0xff) + stat[0] = 0x84 + stat[1] = 0x80 + stat[13] = 0x02 + stat[25] = 80 + stat[26] = 0xa0 + stat[27] = 0x80 + w.bytes(stat) + w.u8(0) + } + return w.toUint8Array() +} + +function wireBncs(server: MemoryStream, silent: ReadonlySet): void { + const framer = new BncsFramer() + let initSeen = false + server.onData((chunk) => { + let data = chunk + if (!initSeen && data[0] === 0x01) { + initSeen = true + data = data.subarray(1) + } + for (const pkt of framer.push(data)) { + if (silent.has(pkt.id) || server.isClosed) continue + if (pkt.id === BncsOpcode.SID_AUTH_INFO) { + server.write(encodeBncsPing(0x11223344)) + const w = new ByteWriter() + w.u32LE(0) + w.u32LE(0xabcdef01) + w.u32LE(0) + w.u32LE(0) + w.u32LE(0) + w.cstring('ver-IX86-1.mpq') + w.cstring('A=1 B=2 C=3 4 A=A^S B=B^C C=C^A A=A^B') + server.write(encodeBncsFrame(BncsOpcode.SID_AUTH_INFO, w.toUint8Array())) + } else if (pkt.id === BncsOpcode.SID_AUTH_CHECK) { + server.write(encodeBncsFrame(BncsOpcode.SID_AUTH_CHECK, new Uint8Array([0, 0, 0, 0, 0]))) + } else if (pkt.id === BncsOpcode.SID_LOGONRESPONSE2) { + server.write(encodeBncsFrame(BncsOpcode.SID_LOGONRESPONSE2, new Uint8Array([0, 0, 0, 0]))) + } else if (pkt.id === BncsOpcode.SID_QUERYREALMS2) { + const w = new ByteWriter() + w.u32LE(0) + w.u32LE(1) + w.u32LE(1) + w.cstring('D2CS') + w.cstring('Local Test Realm') + server.write(encodeBncsFrame(BncsOpcode.SID_QUERYREALMS2, w.toUint8Array())) + } else if (pkt.id === BncsOpcode.SID_LOGONREALMEX) { + const w = new ByteWriter() + w.u32LE(0x11223344) + w.u32LE(0) + w.bytes(new Uint8Array([1, 2, 3, 4, 5, 6, 7, 8])) + w.bytes(new Uint8Array([127, 0, 0, 1])) + w.u16BE(6113) + w.u16LE(0) + w.bytes(new Uint8Array(48).fill(0x55)) + w.cstring('d2webbot_sorc') + server.write(encodeBncsFrame(BncsOpcode.SID_LOGONREALMEX, w.toUint8Array())) + } + } + }) +} + +function wireMcp(server: MemoryStream, silent: ReadonlySet, characters: number): void { + const framer = new McpFramer() + let initSeen = false + server.onData((chunk) => { + let data = chunk + if (!initSeen && data[0] === 0x01) { + initSeen = true + data = data.subarray(1) + } + for (const pkt of framer.push(data)) { + if (silent.has(pkt.id) || server.isClosed) continue + if (pkt.id === McpOpcode.MCP_STARTUP) { + server.write(encodeMcpFrame(McpOpcode.MCP_STARTUP, new Uint8Array([0, 0, 0, 0]))) + } else if (pkt.id === McpOpcode.MCP_CHARLIST2) { + server.write(encodeMcpFrame(McpOpcode.MCP_CHARLIST2, charList(characters))) + } else if (pkt.id === McpOpcode.MCP_CHARLOGON) { + server.write(encodeMcpFrame(McpOpcode.MCP_CHARLOGON, new Uint8Array([0, 0, 0, 0]))) + } else if (pkt.id === McpOpcode.MCP_CREATEGAME) { + const reqId = new ByteReader(pkt.payload).u16LE() + const w = new ByteWriter() + w.u16LE(reqId) + w.u16LE(1) + w.u16LE(0) + w.u32LE(0) + server.write(encodeMcpFrame(McpOpcode.MCP_CREATEGAME, w.toUint8Array())) + } else if (pkt.id === McpOpcode.MCP_JOINGAME) { + const reqId = new ByteReader(pkt.payload).u16LE() + const w = new ByteWriter() + w.u16LE(reqId) + w.u16LE(0x0042) + w.u16LE(0) + w.bytes(new Uint8Array([127, 0, 0, 1])) + w.u32LE(0xdeadbeef) + w.u32LE(0) + server.write(encodeMcpFrame(McpOpcode.MCP_JOINGAME, w.toUint8Array())) + } + } + }) +} + +function wireGame(server: MemoryStream): void { + server.onData((chunk) => { + if (chunk[0] === 0x68 && !server.isClosed) { + server.write( + wrapD2gsCompressedBlock( + defaultD2gsHuffmanCodec.compress(new Uint8Array([0x01, 0, 0, 0, 0, 0, 1, 1, 0x02, 0x04])), + ), + ) + } + }) +} + +export function createFakeRealm(options: FakeRealmOptions = {}): FakeRealm { + const silentBncs = new Set(options.silentBncs ?? []) + const silentMcp = new Set(options.silentMcp ?? []) + const characters = options.characters ?? 1 + const connections: FakeConnection[] = [] + + const resolver: EndpointResolver = { + open: async (role: EndpointRole): Promise => { + const [client, server] = createMemoryStreamPair(role) + connections.push({ role, client, server }) + if (role === 'bnet') wireBncs(server, silentBncs) + else if (role === 'realm') wireMcp(server, silentMcp, characters) + else { + wireGame(server) + setTimeout(() => { + if (!server.isClosed) server.write(new Uint8Array([0xaf, 0x01])) + }, 5) + } + return client + }, + } + + return { + resolver, + connections, + latest(role: EndpointRole): FakeConnection { + const found = [...connections].reverse().find((c) => c.role === role) + if (!found) throw new Error(`no ${role} connection was opened`) + return found + }, + allClientStreamsClosed(): boolean { + return connections.every((c) => c.client.isClosed) + }, + } +} diff --git a/tests/netproto/online-flow-session-lost.test.ts b/tests/netproto/online-flow-session-lost.test.ts new file mode 100644 index 0000000..51d03d7 --- /dev/null +++ b/tests/netproto/online-flow-session-lost.test.ts @@ -0,0 +1,195 @@ +/** + * `createD2OnlineFlow` connection-loss detection: every BNCS/MCP/D2GS loss (socket close, request + * timeout, command without a connection) produces exactly one `onSessionLost` event, the flow ends + * in state `closed`, and every connection it opened is closed. + */ + +import { describe, expect, it } from 'vitest' + +import { + createD2OnlineFlow, + type D2OnlineFlow, + type D2OnlineSessionLost, +} from '../../src/netproto/index.ts' +import { FakeClock } from '../../src/netproto/core/clock.ts' +import { TimeoutError } from '../../src/netproto/core/errors.ts' +import { BncsOpcode } from '../../src/netproto/bncs/packets.ts' +import { McpOpcode } from '../../src/netproto/mcp/packets.ts' +import { getCanonicalItemDataTables } from '../../src/client/world/item-tables-provider.ts' +import { createFakeRealm, type FakeRealm, type FakeRealmOptions } from './fake-realm-servers.ts' + +interface Harness { + readonly realm: FakeRealm + readonly clock: FakeClock + readonly flow: D2OnlineFlow + readonly lost: D2OnlineSessionLost[] +} + +function createHarness(options: FakeRealmOptions = {}): Harness { + const realm = createFakeRealm(options) + const clock = new FakeClock(1000) + const flow = createD2OnlineFlow({ + bnetHost: '127.0.0.1', + resolver: realm.resolver, + clock, + itemTables: getCanonicalItemDataTables(), + }) + const lost: D2OnlineSessionLost[] = [] + flow.onSessionLost((ev) => lost.push(ev)) + return { realm, clock, flow, lost } +} + +/** Log in and enter the realm: the char_select (1 char) / char_create (0 chars) screen. */ +async function toCharacterScreen(h: Harness): Promise { + await h.flow.connectBnet() + await h.flow.login('d2webbot2', 'secret') + await h.flow.enterRealm('D2CS') + expect(h.flow.state).toBe('realm_ready') +} + +async function flushMicrotasks(): Promise { + for (let i = 0; i < 10; i++) await Promise.resolve() +} + +/** Advance the fake clock in steps so promise continuations can arm their timers in between. */ +async function advance(clock: FakeClock, ms: number): Promise { + const step = 500 + for (let t = 0; t < ms; t += step) { + await flushMicrotasks() + clock.advance(Math.min(step, ms - t)) + } + await flushMicrotasks() +} + +describe('D2OnlineFlow session loss', () => { + it('BNCS server close on char_select: one loss event, flow closed, BNCS + MCP sockets disposed', async () => { + const h = createHarness({ characters: 1 }) + await toCharacterScreen(h) + + h.realm.latest('bnet').server.close() + + expect(h.lost).toHaveLength(1) + expect(h.lost[0]).toMatchObject({ session: 'bncs', cause: 'closed', flowState: 'realm_ready' }) + expect(h.lost[0]!.message).toContain('BNCS connection closed') + expect(h.flow.state).toBe('closed') + expect(h.flow.sessionLost).toBe(h.lost[0]) + expect(h.realm.latest('realm').client.isClosed).toBe(true) + expect(h.realm.allClientStreamsClosed()).toBe(true) + }) + + it('MCP server close on char_create: one loss event; a later createCharacter fails without a second event', async () => { + const h = createHarness({ characters: 0 }) + await toCharacterScreen(h) + + h.realm.latest('realm').server.close({ kind: 'error', error: new Error('ECONNRESET') }) + + expect(h.lost).toHaveLength(1) + expect(h.lost[0]).toMatchObject({ session: 'mcp', cause: 'error', flowState: 'realm_ready' }) + expect(h.lost[0]!.message).toContain('ECONNRESET') + expect(h.realm.allClientStreamsClosed()).toBe(true) + + await expect( + h.flow.createCharacter({ name: 'NewSorc', classId: 1, expansion: true }), + ).rejects.toThrow(/MCP realm session is not connected for createCharacter/) + expect(h.lost).toHaveLength(1) + expect(h.flow.state).toBe('closed') + }) + + it('realm logon timeout (SID_LOGONREALMEX unanswered): timeout loss, no reconnect, sockets disposed', async () => { + const h = createHarness({ silentBncs: [BncsOpcode.SID_LOGONREALMEX] }) + await h.flow.connectBnet() + await h.flow.login('d2webbot2', 'secret') + + const entering = h.flow.enterRealm('D2CS') + const outcome = entering.then( + () => null, + (err: unknown) => err, + ) + await advance(h.clock, 10_500) + + expect(await outcome).toBeInstanceOf(TimeoutError) + expect(h.lost).toHaveLength(1) + expect(h.lost[0]).toMatchObject({ session: 'bncs', cause: 'timeout', flowState: 'connecting_realm' }) + expect(h.flow.state).toBe('closed') + expect(h.realm.connections.filter((c) => c.role === 'realm')).toHaveLength(0) + expect(h.realm.connections.filter((c) => c.role === 'bnet')).toHaveLength(1) + expect(h.realm.allClientStreamsClosed()).toBe(true) + }) + + it('MCP_STARTUP timeout: realm timeout loss closes both the MCP and BNCS sockets', async () => { + const h = createHarness({ silentMcp: [McpOpcode.MCP_STARTUP] }) + await h.flow.connectBnet() + await h.flow.login('d2webbot2', 'secret') + + const outcome = h.flow.enterRealm('D2CS').then( + () => null, + (err: unknown) => err, + ) + await advance(h.clock, 10_500) + + expect(await outcome).toBeInstanceOf(TimeoutError) + expect(h.lost).toHaveLength(1) + expect(h.lost[0]).toMatchObject({ session: 'mcp', cause: 'timeout' }) + expect(h.realm.allClientStreamsClosed()).toBe(true) + }) + + it('BNCS and MCP closing at the same time are reported exactly once', async () => { + const h = createHarness({ characters: 1 }) + await toCharacterScreen(h) + + h.realm.latest('realm').server.close() + h.realm.latest('bnet').server.close() + await flushMicrotasks() + + expect(h.lost).toHaveLength(1) + expect(h.lost[0]!.session).toBe('mcp') + expect(h.flow.state).toBe('closed') + expect(h.realm.allClientStreamsClosed()).toBe(true) + }) + + it('D2GS drop in game: one d2gs loss event and BNCS/MCP/D2GS are all disposed', async () => { + const h = createHarness({ characters: 1 }) + await toCharacterScreen(h) + await h.flow.selectCharacter('d2webbot_sorc') + await h.flow.createGame({ name: 'loss-game', password: '', difficulty: 0 }) + await h.flow.joinGame('loss-game', '') + expect(h.flow.state).toBe('ingame') + + h.realm.latest('game').server.close() + + expect(h.lost).toHaveLength(1) + expect(h.lost[0]).toMatchObject({ session: 'd2gs', cause: 'closed', flowState: 'ingame' }) + expect(h.flow.state).toBe('closed') + expect(h.realm.allClientStreamsClosed()).toBe(true) + }) + + it('BNCS loss during a game is reported when the game ends, not by aborting the game', async () => { + const h = createHarness({ characters: 1 }) + await toCharacterScreen(h) + await h.flow.selectCharacter('d2webbot_sorc') + await h.flow.joinGame('loss-game', '') + + h.realm.latest('bnet').server.close() + expect(h.lost).toHaveLength(0) + expect(h.flow.state).toBe('ingame') + expect(h.realm.latest('game').client.isClosed).toBe(false) + + await h.flow.leaveToLobby() + expect(h.lost).toHaveLength(1) + expect(h.lost[0]).toMatchObject({ session: 'bncs', cause: 'closed', flowState: 'ingame' }) + expect(h.flow.state).toBe('closed') + expect(h.realm.allClientStreamsClosed()).toBe(true) + }) + + it('close() on purpose emits no loss event and disposes every socket', async () => { + const h = createHarness({ characters: 1 }) + await toCharacterScreen(h) + + await h.flow.close() + + expect(h.lost).toHaveLength(0) + expect(h.flow.sessionLost).toBeNull() + expect(h.flow.state).toBe('closed') + expect(h.realm.allClientStreamsClosed()).toBe(true) + }) +})