import { ArrayCodec, BigIntCodec, type Codec, type CodecType, MessageCodec, MutableArrayCodec, NumberCodec, OneOfCodec, OptionalCodec, StringCodec, UndefinedCodec, } from "./codec"; import { Cursor } from "./common"; import * as proto from "./proto"; /** Data finality. */ export const DataFinality: Codec< "finalized" | "accepted" | "pending" | "unknown", proto.stream.DataFinality > = { encode(x) { const enumMap = { finalized: proto.stream.DataFinality.FINALIZED, accepted: proto.stream.DataFinality.ACCEPTED, pending: proto.stream.DataFinality.PENDING, unknown: proto.stream.DataFinality.UNKNOWN, }; return enumMap[x] ?? proto.stream.DataFinality.UNKNOWN; }, decode(p) { const enumMap = { [proto.stream.DataFinality.FINALIZED]: "finalized", [proto.stream.DataFinality.ACCEPTED]: "accepted", [proto.stream.DataFinality.PENDING]: "pending", [proto.stream.DataFinality.UNKNOWN]: "unknown", [proto.stream.DataFinality.UNRECOGNIZED]: "unknown", } as const; return enumMap[p] ?? "unknown"; }, }; export type DataFinality = CodecType; /** Data production mode. */ export const DataProduction: Codec< "backfill" | "live" | "unknown", proto.stream.DataProduction > = { encode(x) { switch (x) { case "backfill": return proto.stream.DataProduction.BACKFILL; case "live": return proto.stream.DataProduction.LIVE; case "unknown": return proto.stream.DataProduction.UNKNOWN; default: return proto.stream.DataProduction.UNRECOGNIZED; } }, decode(p) { const enumMap = { [proto.stream.DataProduction.BACKFILL]: "backfill", [proto.stream.DataProduction.LIVE]: "live", [proto.stream.DataProduction.UNKNOWN]: "unknown", [proto.stream.DataProduction.UNRECOGNIZED]: "unknown", } as const; return enumMap[p] ?? "unknown"; }, }; export type DataProduction = CodecType; export const DurationCodec = MessageCodec({ seconds: BigIntCodec, nanos: NumberCodec, }); export type Duration = CodecType; /** Create a `StreamDataRequest` with the given filter schema. */ export const StreamDataRequest = (filter: Codec) => MessageCodec({ finality: OptionalCodec(DataFinality), startingCursor: OptionalCodec(Cursor), filter: MutableArrayCodec(filter), heartbeatInterval: OptionalCodec(DurationCodec), }); export type StreamDataRequest = CodecType< ReturnType> >; export const Invalidate = MessageCodec({ cursor: OptionalCodec(Cursor), }); export type Invalidate = CodecType; export const Finalize = MessageCodec({ cursor: OptionalCodec(Cursor), }); export type Finalize = CodecType; // TODO: Double check this; This is a hack to make the heartbeat variant undefined export const Heartbeat = UndefinedCodec; export type Heartbeat = CodecType; export const StdOut = StringCodec; export type StdOut = CodecType; export const StdErr = StringCodec; export type StdErr = CodecType; export const SystemMessage = MessageCodec({ output: OneOfCodec({ stdout: StdOut, stderr: StdErr, }), }); export type SystemMessage = CodecType; const _DataOrNull = ( schema: Codec, ): Codec => ({ encode(x) { if (x === null) { return new Uint8Array(); } return schema.encode(x); }, decode(p) { if (p.length === 0) { return null; } return schema.decode(p); }, }); export const Data = (schema: Codec) => MessageCodec({ cursor: OptionalCodec(Cursor), endCursor: Cursor, finality: DataFinality, production: DataProduction, data: ArrayCodec(_DataOrNull(schema)), }); export type Data = CodecType>>; export const StreamDataResponse = (schema: Codec) => OneOfCodec({ data: Data(schema), invalidate: Invalidate, finalize: Finalize, heartbeat: Heartbeat, systemMessage: SystemMessage, }); export const ResponseWithoutData = OneOfCodec({ invalidate: Invalidate, finalize: Finalize, heartbeat: Heartbeat, systemMessage: SystemMessage, }); export type ResponseWithoutData = CodecType; export type StreamDataResponse = CodecType< ReturnType> >;