You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
100 lines
3.1 KiB
100 lines
3.1 KiB
import { isPromiseLike } from '$libs/standard-schema';
|
|
import {
|
|
CONNECTION_DIAGNOSTIC_EVENTS,
|
|
CONNECTION_SEND_REASON_BUFFER_FULL,
|
|
CONNECTION_SEND_REASON_CLOSED,
|
|
CONNECTION_SEND_REASON_INVALID_FRAME,
|
|
CONNECTION_SEND_REASON_SEND_NOT_SUPPORTED,
|
|
CONNECTION_SEND_REASON_SERIALIZE_FAILED,
|
|
CONNECTION_SEND_REASON_TRANSPORT_ERROR
|
|
} from './consts.ts';
|
|
import { emitConnectionDiagnostic, type ConnectionDiagnostics } from './diagnostics.ts';
|
|
import type { ConnectionFrameBuffer } from './frame-buffer.ts';
|
|
import { createFrame } from './serializer.ts';
|
|
import type {
|
|
ConnectionFrame,
|
|
ConnectionSendResult,
|
|
ConnectionSerializer,
|
|
ConnectionTransport
|
|
} from './types.ts';
|
|
|
|
interface ConnectionSenderRuntime {
|
|
readonly serializer: ConnectionSerializer;
|
|
readonly frameBuffer: ConnectionFrameBuffer;
|
|
isConnected(): boolean;
|
|
getTransport(): ConnectionTransport | null;
|
|
setError(error: unknown): void;
|
|
readonly diagnostics: ConnectionDiagnostics;
|
|
}
|
|
|
|
export interface ConnectionSender {
|
|
sendFrame(frame: ConnectionFrame, allowBuffer: boolean): Promise<ConnectionSendResult>;
|
|
flushBuffer(): Promise<void>;
|
|
}
|
|
|
|
export function createConnectionSender(runtime: ConnectionSenderRuntime): ConnectionSender {
|
|
async function sendFrame(
|
|
frame: ConnectionFrame,
|
|
allowBuffer: boolean
|
|
): Promise<ConnectionSendResult> {
|
|
try {
|
|
createFrame(frame);
|
|
} catch (err) {
|
|
return { ok: false, reason: CONNECTION_SEND_REASON_INVALID_FRAME, error: err };
|
|
}
|
|
|
|
if (!runtime.isConnected()) {
|
|
if (runtime.frameBuffer.canBuffer(frame, allowBuffer)) return runtime.frameBuffer.push(frame);
|
|
if (runtime.frameBuffer.shouldDropClosedFrame()) return { ok: true, id: frame.id };
|
|
return { ok: false, reason: CONNECTION_SEND_REASON_CLOSED };
|
|
}
|
|
|
|
const currentTransport = runtime.getTransport();
|
|
if (currentTransport === null || !currentTransport.canSend) {
|
|
return { ok: false, reason: CONNECTION_SEND_REASON_SEND_NOT_SUPPORTED };
|
|
}
|
|
if (currentTransport.bufferedAmount > runtime.frameBuffer.maxBytes) {
|
|
return { ok: false, reason: CONNECTION_SEND_REASON_BUFFER_FULL };
|
|
}
|
|
|
|
let encoded: string | ArrayBuffer;
|
|
try {
|
|
encoded = runtime.serializer.encode(frame);
|
|
} catch (err) {
|
|
emitConnectionDiagnostic(
|
|
runtime.diagnostics,
|
|
CONNECTION_DIAGNOSTIC_EVENTS.FRAME_ENCODE_FAILED,
|
|
{ error: err }
|
|
);
|
|
return { ok: false, reason: CONNECTION_SEND_REASON_SERIALIZE_FAILED, error: err };
|
|
}
|
|
|
|
try {
|
|
const maybe = currentTransport.send(encoded);
|
|
if (isPromiseLike(maybe)) await maybe;
|
|
return { ok: true, id: frame.id };
|
|
} catch (err) {
|
|
runtime.setError(err);
|
|
emitConnectionDiagnostic(runtime.diagnostics, CONNECTION_DIAGNOSTIC_EVENTS.SEND_FAILED, {
|
|
error: err
|
|
});
|
|
return { ok: false, reason: CONNECTION_SEND_REASON_TRANSPORT_ERROR, error: err };
|
|
}
|
|
}
|
|
|
|
return {
|
|
sendFrame,
|
|
async flushBuffer(): Promise<void> {
|
|
while (runtime.isConnected() && runtime.frameBuffer.length > 0) {
|
|
const frame = runtime.frameBuffer.shift();
|
|
if (frame === undefined) return;
|
|
const result = await sendFrame(frame, false);
|
|
if (!result.ok) {
|
|
runtime.frameBuffer.unshift(frame);
|
|
return;
|
|
}
|
|
}
|
|
}
|
|
};
|
|
}
|