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) + }) +})