parent
a987fd10ed
commit
1aa9c4dd30
@ -0,0 +1,65 @@
|
||||
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<TChannels extends ConnectionChannelMap> {
|
||||
readonly name: string;
|
||||
readonly raw: string | ArrayBuffer;
|
||||
readonly serializer: ConnectionSerializer;
|
||||
readonly state: ConnectionStateTracker;
|
||||
readonly heartbeat: ConnectionHeartbeat;
|
||||
readonly acks: ConnectionAckRegistry;
|
||||
readonly channels: ConnectionChannelRegistry<TChannels>;
|
||||
readonly events: ConnectionEventBus;
|
||||
readonly diagnostics: ConnectionDiagnostics;
|
||||
}
|
||||
|
||||
export function routeConnectionMessage<TChannels extends ConnectionChannelMap>(
|
||||
input: RouteConnectionMessageInput<TChannels>
|
||||
): 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);
|
||||
}
|
||||
Loading…
Reference in new issue