File size: 4,652 Bytes
80d7f0c | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 | import { isJsonValue } from "@earendil-works/chord";
import { Check } from "typebox/value";
import { decodeCbor, encodeCbor } from "./cbor/index.ts";
import { DEFAULT_MAX_FRAME_LENGTH, encodeFrame, FrameDecoder, type FrameDecoderOptions } from "./framing.ts";
import {
type ClientMessage,
ClientMessageSchema,
PROTOCOL_VERSION,
type ServerMessage,
ServerMessageSchema,
} from "./protocol.ts";
export class ProtocolValidationError extends Error {
constructor(message: string) {
super(message);
this.name = "ProtocolValidationError";
}
}
export function parseClientMessage(value: unknown): ClientMessage {
if (!Check(ClientMessageSchema, value) || !isJsonValue(value)) {
throw new ProtocolValidationError("Invalid client protocol message");
}
return value;
}
export function parseServerMessage(value: unknown): ServerMessage {
if (!Check(ServerMessageSchema, value) || !isJsonValue(value)) {
throw new ProtocolValidationError("Invalid server protocol message");
}
return value;
}
function boundedErrorMessage(error: unknown): string {
if (!(error instanceof Error)) return "Unknown codec error";
return error.message.length <= 500 ? error.message : `${error.message.slice(0, 497)}...`;
}
function encodeProtocolMessage<T>(
value: T,
parse: (candidate: unknown) => T,
kind: string,
options?: FrameDecoderOptions,
): Uint8Array {
const validated = parse(value);
try {
const maxFrameLength = options?.maxFrameLength ?? DEFAULT_MAX_FRAME_LENGTH;
return encodeFrame(encodeCbor(validated, { maxByteLength: maxFrameLength }));
} catch (error) {
if (error instanceof ProtocolValidationError) throw error;
throw new ProtocolValidationError(`Unable to encode ${kind} protocol message: ${boundedErrorMessage(error)}`);
}
}
/** Validates and encodes one complete length-prefixed client message. */
export function encodeClientMessage(message: ClientMessage, options?: FrameDecoderOptions): Uint8Array {
return encodeProtocolMessage(message, parseClientMessage, "client", options);
}
/** Validates and encodes one complete length-prefixed server message. */
export function encodeServerMessage(message: ServerMessage, options?: FrameDecoderOptions): Uint8Array {
return encodeProtocolMessage(message, parseServerMessage, "server", options);
}
class ValidatedMessageDecoder<T> {
private failed = false;
private readonly frames: FrameDecoder;
private readonly kind: string;
private readonly maxFrameLength: number;
private readonly parse: (candidate: unknown) => T;
constructor(kind: string, parse: (candidate: unknown) => T, options?: FrameDecoderOptions) {
this.frames = new FrameDecoder(options);
this.kind = kind;
this.maxFrameLength = options?.maxFrameLength ?? DEFAULT_MAX_FRAME_LENGTH;
this.parse = parse;
}
push(chunk: Uint8Array): T[] {
if (this.failed) throw new ProtocolValidationError(`${this.kind} message decoder has failed`);
try {
const messages: T[] = [];
for (const frame of this.frames.push(chunk)) {
messages.push(this.parse(decodeCbor(frame, { maxByteLength: this.maxFrameLength })));
}
return messages;
} catch (error) {
this.failed = true;
if (error instanceof ProtocolValidationError) throw error;
throw new ProtocolValidationError(`Invalid ${this.kind} protocol frame: ${boundedErrorMessage(error)}`);
}
}
end(): void {
if (this.failed) throw new ProtocolValidationError(`${this.kind} message decoder has failed`);
try {
this.frames.end();
} catch (error) {
this.failed = true;
throw new ProtocolValidationError(`Invalid ${this.kind} protocol framing: ${boundedErrorMessage(error)}`);
}
}
}
/** Incrementally decodes and validates framed client messages. */
export class ClientMessageDecoder {
private readonly decoder: ValidatedMessageDecoder<ClientMessage>;
constructor(options?: FrameDecoderOptions) {
this.decoder = new ValidatedMessageDecoder("client", parseClientMessage, options);
}
push(chunk: Uint8Array): ClientMessage[] {
return this.decoder.push(chunk);
}
end(): void {
this.decoder.end();
}
}
/** Incrementally decodes and validates framed server messages. */
export class ServerMessageDecoder {
private readonly decoder: ValidatedMessageDecoder<ServerMessage>;
constructor(options?: FrameDecoderOptions) {
this.decoder = new ValidatedMessageDecoder("server", parseServerMessage, options);
}
push(chunk: Uint8Array): ServerMessage[] {
return this.decoder.push(chunk);
}
end(): void {
this.decoder.end();
}
}
export function isSupportedProtocolVersion(version: number): version is typeof PROTOCOL_VERSION {
return Number.isInteger(version) && version === PROTOCOL_VERSION;
}
|