import type { IncomingMessage } from "node:http" import type { WebSocket } from "ws" import type { BuildEvents } from "#Source/event/index.ts" import { EventManager } from "#Source/event/index.ts" import { generateUuidV4FromUrl } from "#Source/identifier/index.ts" import type { LoggerFriendly, LoggerFriendlyOptions } from "#Source/log/index.ts" import { Logger } from "#Source/log/index.ts" import type { SocketUnitHeartbeatOptions } from "../common/index.ts" import type { WebSocketServer, WebSocketServerStartResult } from "./server.ts" import type { InitialMessageBuilder, SocketUnitSnapshot } from "./socket-unit.ts" import { SocketUnit } from "./socket-unit.ts" /** * @description 表示服务端 Socket 接入成功时的事件负载。 */ export interface SocketConnectedEvent { clientId: string } /** * @description 表示服务端 Socket 连接关闭时的事件负载。 */ export interface SocketClosedEvent { clientId: string } /** * @description 表示服务端 Socket 派发业务消息时的事件负载。 */ export interface SocketMessageEvent { clientId: string message: Message } /** * @description 表示服务端 Socket 对外派发的事件表。 */ export type SocketEvents = BuildEvents<{ connect: (payload: SocketConnectedEvent) => void close: (payload: SocketClosedEvent) => void message: (payload: SocketMessageEvent) => void }> /** * @description 表示服务端 Socket 的诊断快照。 */ export interface SocketSnapshot { isStarted: boolean isRunning: boolean managesWebSocketServerLifecycle: boolean activeSocketUnitCount: number clientIds: string[] socketUnitSnapshots: SocketUnitSnapshot[] } /** * @description 表示服务端 Socket 的构造参数。 */ export interface SocketOptions extends LoggerFriendlyOptions { webSocketServer: WebSocketServer enableInitialMessage?: boolean | undefined initialMessageBuilder?: InitialMessageBuilder | undefined heartbeat?: SocketUnitHeartbeatOptions | undefined onMessage?: ((payload: SocketMessageEvent) => void) | undefined } /** * @description 管理服务端接入的多个 SocketUnit,并对外提供统一事件面。 */ export class Socket implements LoggerFriendly { protected readonly options: SocketOptions readonly logger: Logger eventManager: EventManager> protected isStarted: boolean protected managesWebSocketServerLifecycle: boolean protected socketUnitMap: Map> protected connectionListener: | ((webSocket: WebSocket, incomingMessage: IncomingMessage) => void) | undefined constructor(options: SocketOptions) { this.options = options this.logger = Logger.fromOptions(options).setDefaultName("Socket") this.eventManager = new EventManager>() this.isStarted = false this.managesWebSocketServerLifecycle = false this.socketUnitMap = new Map() this.connectionListener = undefined } /** * @description 将当前 Socket 挂接到一个已在运行的 WebSocketServer 上。 */ attach(): WebSocketServerStartResult { if (this.isStarted === true) { throw new Error("Socket has already attached.") } const webSocketServer = this.options.webSocketServer const server = webSocketServer.getServer() const startResult = webSocketServer.getStartResult() const connectionListener = (webSocket: WebSocket, incomingMessage: IncomingMessage): void => { this.handleConnection(webSocket, incomingMessage) } server.on("connection", connectionListener) this.connectionListener = connectionListener this.isStarted = true return startResult } /** * @description 启动服务端 Socket,并在必要时托管底层服务生命周期。 */ async open(): Promise { if (this.isStarted === true) { throw new Error("Socket has already opened.") } const webSocketServer = this.options.webSocketServer if (webSocketServer.isRunning() === false) { await webSocketServer.start() this.managesWebSocketServerLifecycle = true } else { this.managesWebSocketServerLifecycle = false } try { return this.attach() } catch (exception) { if (this.managesWebSocketServerLifecycle === true) { this.managesWebSocketServerLifecycle = false await webSocketServer.close() } throw exception } } /** * @description 关闭后重新挂接到服务上。 */ async reopen(): Promise { await this.close() await this.open() } /** * @description 关闭当前 Socket 及其托管的全部连接。 */ async close(): Promise { const closeTasks: Array> = [] if (this.isStarted === true) { this.removeConnectionListener() this.isStarted = false } const socketUnits = [...this.socketUnitMap.values()] closeTasks.push( ...socketUnits.map(async (socketUnit) => { await socketUnit.close() }), ) if (closeTasks.length !== 0) { await Promise.all(closeTasks) } if (this.managesWebSocketServerLifecycle === true) { this.managesWebSocketServerLifecycle = false await this.options.webSocketServer.close() } } /** * @description 返回当前服务端 Socket 的聚合快照。 */ getSnapshot(): SocketSnapshot { const socketUnitSnapshots = [...this.socketUnitMap.values()].map((item) => { return item.getSnapshot() }) return { isStarted: this.isStarted, isRunning: this.options.webSocketServer.isRunning(), managesWebSocketServerLifecycle: this.managesWebSocketServerLifecycle, activeSocketUnitCount: socketUnitSnapshots.length, clientIds: socketUnitSnapshots.map((item) => { return item.clientId }), socketUnitSnapshots, } } /** * @description 返回指定 clientId 对应的 SocketUnit 快照。 */ getSocketUnitSnapshot(clientId: string): SocketUnitSnapshot | undefined { return this.socketUnitMap.get(clientId)?.getSnapshot() } /** * @description 向指定连接发送业务消息。 */ sendMessage(clientId: string, message: Message): void { const socketUnit = this.socketUnitMap.get(clientId) if (socketUnit === undefined) { this.logger.warn(`Business message send skipped because client was not found: ${clientId}`) return } socketUnit.sendMessage(message) } protected handleConnection(webSocket: WebSocket, incomingMessage: IncomingMessage): void { const clientId: string = generateUuidV4FromUrl() const socketUnit = new SocketUnit({ logger: this.logger.derive().setName(`SocketUnit-${clientId}`), clientId, webSocket, incomingMessage, enableInitialMessage: this.options.enableInitialMessage, initialMessageBuilder: this.options.initialMessageBuilder, heartbeat: this.options.heartbeat, onMessage: (payload) => { this.eventManager.emit("message", payload) this.options.onMessage?.(payload) }, onClose: (closedClientId) => { this.socketUnitMap.delete(closedClientId) this.eventManager.emit("close", { clientId: closedClientId }) }, }) this.socketUnitMap.set(clientId, socketUnit) this.logger.log(`WebSocket connected: ${clientId}`) socketUnit.start() this.eventManager.emit("connect", { clientId }) } protected removeConnectionListener(): void { const connectionListener = this.connectionListener this.connectionListener = undefined if (connectionListener === undefined) { return } if (this.options.webSocketServer.isRunning() === false) { return } const server = this.options.webSocketServer.getServer() server.off("connection", connectionListener) } }