import { CONNECTION_DIAGNOSTIC_EVENTS, CONNECTION_FRAME_TYPE_PONG } from './consts.ts'; import { emitConnectionDiagnostic, type ConnectionDiagnostics } from './diagnostics.ts'; import { assertConnectionFrame } from './serializer.ts'; import type { ConnectionAckRegistry } from './acks.ts'; import type { ConnectionChannelRegistry } from './channel-registry.ts'; import type { ConnectionEventBus } from './connection-events.ts'; import type { ConnectionStateTracker } from './connection-state.ts'; import type { ConnectionHeartbeat } from './heartbeat.ts'; import type { ConnectionChannelMap, ConnectionFrame, ConnectionMessageMeta, ConnectionSerializer } from './types.ts'; export interface RouteConnectionMessageInput { readonly name: string; readonly raw: string | ArrayBuffer; readonly serializer: ConnectionSerializer; readonly state: ConnectionStateTracker; readonly heartbeat: ConnectionHeartbeat; readonly acks: ConnectionAckRegistry; readonly channels: ConnectionChannelRegistry; readonly events: ConnectionEventBus; readonly diagnostics: ConnectionDiagnostics; } export function routeConnectionMessage( input: RouteConnectionMessageInput ): void { const receivedAt = input.state.touchMessage(); input.heartbeat.received(); let frame: ConnectionFrame; try { frame = input.serializer.decode(input.raw); assertConnectionFrame(frame); } catch (err) { input.state.setError(err); emitConnectionDiagnostic(input.diagnostics, CONNECTION_DIAGNOSTIC_EVENTS.FRAME_DECODE_FAILED, { error: err }); return; } if (frame.replyTo !== undefined) { input.acks.resolveFromFrame(frame); return; } if (frame.type === CONNECTION_FRAME_TYPE_PONG) return; const meta: ConnectionMessageMeta = { connection: input.name, topic: frame.topic, receivedAt, generation: input.state.generation }; if (frame.topic !== undefined) { input.channels.receive(frame.topic, frame, meta); } input.events.emitGlobal(frame, meta); }