import type { LoggerFriendly, LoggerFriendlyOptions } from "#Source/log/index.ts" import { Logger } from "#Source/log/index.ts" /** * @description 表示主动心跳轮次使用的消息处理器。 */ export interface ActiveHeartbeatMessageHandler { buildPingMessage: (clientId: string) => Message verifyPongMessage: (clientId: string, message: Message) => boolean } /** * @description 表示被动响应心跳时使用的消息处理器。 */ export interface PassiveHeartbeatMessageHandler { verifyPingMessage: (clientId: string, message: Message) => boolean buildPongMessage: (clientId: string, message: Message) => Message } /** * @description 表示主动心跳消息处理器的构造器。 */ export type ActiveHeartbeatMessageHandlerGenerator = () => ActiveHeartbeatMessageHandler /** * @description 表示被动心跳消息处理器的构造器。 */ export type PassiveHeartbeatMessageHandlerGenerator = () => PassiveHeartbeatMessageHandler /** * @description 表示 SocketUnitHeartbeat 的配置项。 */ export interface SocketUnitHeartbeatOptions extends LoggerFriendlyOptions { enableActiveHeartbeat?: boolean | undefined enablePassiveHeartbeat?: boolean | undefined heartbeatInterval?: number | undefined maxResponseTime?: number | undefined maxHeartbeatLostCount?: number | undefined activeHeartbeatMessageHandlerGenerator?: | ActiveHeartbeatMessageHandlerGenerator | undefined passiveHeartbeatMessageHandlerGenerator?: | PassiveHeartbeatMessageHandlerGenerator | undefined } /** * @description 验证心跳配置是否满足统一运行时的基本约束。 */ export const validateSocketUnitHeartbeatOptions = ( options: SocketUnitHeartbeatOptions, ): void => { if ( options.enableActiveHeartbeat === true && options.activeHeartbeatMessageHandlerGenerator === undefined ) { throw new Error( "Active heartbeat message handler is required when active heartbeat is enabled.", ) } if ( options.enablePassiveHeartbeat === true && options.passiveHeartbeatMessageHandlerGenerator === undefined ) { throw new Error( "Passive heartbeat message handler is required when passive heartbeat is enabled.", ) } if ( options.heartbeatInterval !== undefined && (!Number.isFinite(options.heartbeatInterval) || options.heartbeatInterval <= 0) ) { throw new Error("Heartbeat interval must be a positive finite number.") } if ( options.maxResponseTime !== undefined && (!Number.isFinite(options.maxResponseTime) || options.maxResponseTime <= 0) ) { throw new Error("Max response time must be a positive finite number.") } if ( options.maxHeartbeatLostCount !== undefined && (!Number.isFinite(options.maxHeartbeatLostCount) || options.maxHeartbeatLostCount < 1) ) { throw new Error("Max heartbeat lost count must be at least 1.") } } /** * @description 表示补齐默认值后的心跳配置。 */ export interface ResolvedSocketUnitHeartbeatOptions extends LoggerFriendlyOptions { enableActiveHeartbeat: boolean enablePassiveHeartbeat: boolean heartbeatInterval: number maxResponseTime: number maxHeartbeatLostCount: number activeHeartbeatMessageHandlerGenerator?: | ActiveHeartbeatMessageHandlerGenerator | undefined passiveHeartbeatMessageHandlerGenerator?: | PassiveHeartbeatMessageHandlerGenerator | undefined } /** * @description 为心跳配置补齐统一运行时所需的默认值。 */ export const resolveSocketUnitHeartbeatOptions = ( options: SocketUnitHeartbeatOptions, ): ResolvedSocketUnitHeartbeatOptions => { return { ...options, enableActiveHeartbeat: options.enableActiveHeartbeat ?? false, enablePassiveHeartbeat: options.enablePassiveHeartbeat ?? false, heartbeatInterval: options.heartbeatInterval ?? 30 * 1_000, maxResponseTime: options.maxResponseTime ?? 10 * 1_000, maxHeartbeatLostCount: options.maxHeartbeatLostCount ?? 3, activeHeartbeatMessageHandlerGenerator: options.activeHeartbeatMessageHandlerGenerator, passiveHeartbeatMessageHandlerGenerator: options.passiveHeartbeatMessageHandlerGenerator, } } /** * @description 表示心跳运行时与外部 SocketUnit 之间的回调边界。 */ export interface SocketUnitHeartbeatCallbacks { sendMessage: (message: { action: string; messageString: string }) => boolean close: () => void } /** * @description 表示单个心跳轮次的状态。 */ export type SocketUnitHeartbeatStateStatus = "waiting" | "success" | "fail" /** * @description 表示统一心跳运行时内部保留的一条轮次记录。 */ export interface SocketUnitHeartbeatState { activeHeartbeatMessageHandler: ActiveHeartbeatMessageHandler responseTimer: ReturnType | undefined status: SocketUnitHeartbeatStateStatus pingMessageString: string | undefined pongMessageString: string | undefined } /** * @description 表示单个心跳轮次对外暴露的只读快照。 */ export interface SocketUnitHeartbeatStateSnapshot { status: SocketUnitHeartbeatStateStatus hasPendingResponseTimer: boolean pingMessageString: string | undefined pongMessageString: string | undefined } /** * @description 表示统一心跳运行时的聚合快照。 */ export interface SocketUnitHeartbeatSnapshot { enabledActiveHeartbeat: boolean enabledPassiveHeartbeat: boolean isRunning: boolean waitingCount: number successCount: number failCount: number recentStates: SocketUnitHeartbeatStateSnapshot[] } /** * @description 管理 SocketUnit 的统一心跳运行时。 */ export class SocketUnitHeartbeat implements LoggerFriendly { protected readonly options: ResolvedSocketUnitHeartbeatOptions protected readonly callbacks: SocketUnitHeartbeatCallbacks readonly logger: Logger protected heartbeatTimer: ReturnType | undefined protected heartbeatStates: Array> constructor( options: SocketUnitHeartbeatOptions, callbacks: SocketUnitHeartbeatCallbacks, ) { validateSocketUnitHeartbeatOptions(options) this.options = resolveSocketUnitHeartbeatOptions(options) this.callbacks = callbacks this.logger = Logger.fromOptions(options).setDefaultName("SocketUnitHeartbeat") this.heartbeatTimer = undefined this.heartbeatStates = [] } /** * @description 返回统一心跳运行时的诊断快照。 */ getSnapshot(): SocketUnitHeartbeatSnapshot { let waitingCount = 0 let successCount = 0 let failCount = 0 const recentStates = this.heartbeatStates.map((heartbeatState) => { switch (heartbeatState.status) { case "waiting": { waitingCount = waitingCount + 1 break } case "success": { successCount = successCount + 1 break } case "fail": { failCount = failCount + 1 break } default: { throw new Error("Unexpected heartbeat state status.") } } return { status: heartbeatState.status, hasPendingResponseTimer: heartbeatState.responseTimer !== undefined, pingMessageString: heartbeatState.pingMessageString, pongMessageString: heartbeatState.pongMessageString, } }) return { enabledActiveHeartbeat: this.options.enableActiveHeartbeat, enabledPassiveHeartbeat: this.options.enablePassiveHeartbeat, isRunning: this.heartbeatTimer !== undefined, waitingCount, successCount, failCount, recentStates, } } /** * @description 在启用主动心跳时启动轮询调度。 */ start(clientId: string): void { if (this.options.enableActiveHeartbeat === false) { return } if (this.heartbeatTimer !== undefined) { return } this.heartbeatTimer = setInterval(() => { this.runActiveHeartbeat(clientId) }, this.options.heartbeatInterval) } /** * @description 停止心跳调度并清理全部轮次状态。 */ stop(): void { if (this.heartbeatTimer !== undefined) { clearInterval(this.heartbeatTimer) this.heartbeatTimer = undefined } this.heartbeatStates.forEach((heartbeatState) => { if (heartbeatState.responseTimer !== undefined) { clearTimeout(heartbeatState.responseTimer) } }) this.heartbeatStates = [] } protected runActiveHeartbeat(clientId: string): void { const activeHeartbeatMessageHandler = this.options.activeHeartbeatMessageHandlerGenerator!() const heartbeatState: SocketUnitHeartbeatState = { activeHeartbeatMessageHandler, responseTimer: undefined, status: "waiting", pingMessageString: undefined, pongMessageString: undefined, } this.heartbeatStates.push(heartbeatState) const pingMessage = activeHeartbeatMessageHandler.buildPingMessage(clientId) const pingMessageString = JSON.stringify(pingMessage) heartbeatState.pingMessageString = pingMessageString const sendResult = this.callbacks.sendMessage({ action: "Active heartbeat ping send", messageString: pingMessageString, }) if (sendResult === true) { this.logger.info(`Active heartbeat ping message sent: ${clientId}, ${pingMessageString}`) } else { heartbeatState.status = "fail" this.callbacks.close() return } heartbeatState.responseTimer = setTimeout(() => { if (heartbeatState.status !== "waiting") { return } heartbeatState.responseTimer = undefined heartbeatState.status = "fail" const timeoutCount = this.countConsecutiveHeartbeatTimeouts(this.heartbeatStates) if (timeoutCount >= this.options.maxHeartbeatLostCount) { this.logger.error("Heartbeat timeout limit reached. Closing WebSocket connection.") this.callbacks.close() } }, this.options.maxResponseTime) } handleMessage(clientId: string, message: Message): boolean { const activeHandleResult = this.handleActiveHeartbeatMessage(clientId, message) if (activeHandleResult === true) { return true } const passiveHandleResult = this.handlePassiveHeartbeatMessage(clientId, message) return passiveHandleResult } protected handleActiveHeartbeatMessage(clientId: string, message: Message): boolean { if (this.options.enableActiveHeartbeat === false) { return false } const heartbeatState = this.findMatchingWaitingHeartbeatState(clientId, message) if (heartbeatState === undefined) { return false } const messageString = JSON.stringify(message) this.logger.info(`Active heartbeat pong message received: ${clientId}, ${messageString}`) if (heartbeatState.responseTimer !== undefined) { clearTimeout(heartbeatState.responseTimer) heartbeatState.responseTimer = undefined } heartbeatState.pongMessageString = messageString heartbeatState.status = "success" this.cleanupHeartbeatStatesAfterSuccess(heartbeatState) return true } protected handlePassiveHeartbeatMessage(clientId: string, message: Message): boolean { if (this.options.enablePassiveHeartbeat === false) { return false } const passiveHeartbeatMessageHandler = this.options.passiveHeartbeatMessageHandlerGenerator!() if (passiveHeartbeatMessageHandler.verifyPingMessage(clientId, message) === false) { return false } const messageString = JSON.stringify(message) this.logger.info(`Passive heartbeat ping message received: ${clientId}, ${messageString}`) const pongMessage = passiveHeartbeatMessageHandler.buildPongMessage(clientId, message) const pongMessageString = JSON.stringify(pongMessage) const sendResult = this.callbacks.sendMessage({ action: "Passive heartbeat pong send", messageString: pongMessageString, }) if (sendResult === true) { this.logger.info(`Passive heartbeat pong message sent: ${clientId}, ${pongMessageString}`) } return true } protected findMatchingWaitingHeartbeatState( clientId: string, message: Message, ): SocketUnitHeartbeatState | undefined { for (let index = this.heartbeatStates.length - 1; index >= 0; index = index - 1) { const heartbeatState = this.heartbeatStates[index]! if (heartbeatState.status !== "waiting") { continue } if ( heartbeatState.activeHeartbeatMessageHandler.verifyPongMessage(clientId, message) === true ) { return heartbeatState } } return undefined } protected countConsecutiveHeartbeatTimeouts( heartbeatStates: Array>, ): number { let timeoutCount = 0 for (let index = heartbeatStates.length - 1; index >= 0; index = index - 1) { const heartbeatState = heartbeatStates[index]! if (heartbeatState.status === "fail") { timeoutCount = timeoutCount + 1 continue } break } return timeoutCount } protected cleanupHeartbeatStatesAfterSuccess( successState: SocketUnitHeartbeatState, ): void { const successStateIndex = this.heartbeatStates.lastIndexOf(successState) if (successStateIndex === -1) { return } this.heartbeatStates.splice(0, successStateIndex + 1) } }