Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
21 kB · 595 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596import { newId } from "../core/ids.js";import type { JsonObject } from "../core/json.js";import type { JazzThoughtStore } from "../jazz/store.js";import type { SourceCursor } from "../store/types.js";import { cursorTimeUs, JetstreamConnector } from "./jetstream.js";
const DEFAULT_ENDPOINT = "wss://jetstream2.us-west.bsky.network/subscribe";
export interface JetstreamWebSocket { readonly readyState: number; addEventListener(type: string, listener: (event: any) => void, options?: { once?: boolean }): void; removeEventListener(type: string, listener: (event: any) => void): void; close(code?: number, reason?: string): void;}
export type JetstreamWebSocketFactory = (url: string) => JetstreamWebSocket;
export interface JetstreamSubscriptionOptions { endpoint?: string; signal?: AbortSignal; rewindUs?: number; maxMessages?: number; maxReconnects?: number; initialBackoffMs?: number; maxBackoffMs?: number; maxMessageSizeBytes?: number; pendingMessageLimit?: number; webSocketFactory?: JetstreamWebSocketFactory; sleep?: (milliseconds: number, signal?: AbortSignal) => Promise<void>; random?: () => number;}
export interface JetstreamSubscriptionResult { source: string; reason: "aborted" | "message-limit" | "closed"; messages: number; inserted: number; unchanged: number; ignored: number; stale: number; reconnects: number; cursor?: SourceCursor;}
interface ConnectionResult { reason: "aborted" | "message-limit" | "closed"; messages: number; inserted: number; unchanged: number; ignored: number; stale: number; closeCode?: number; closeReason?: string; overflowed: boolean; transportError?: string;}
type ConnectionProgress = Pick<ConnectionResult, "messages" | "inserted" | "unchanged" | "ignored" | "stale">;
interface SocketMessage { type: "message"; data: unknown;}
interface SocketClosed { type: "closed"; code?: number; reason?: string;}
type SocketItem = SocketMessage | SocketClosed;
class ConnectionProgressError extends Error { constructor( readonly cause: unknown, readonly progress: ConnectionProgress, readonly phase: "live-subscribe", readonly failureRecorded: boolean, ) { super(cause instanceof Error ? cause.message : String(cause)); this.name = "ConnectionProgressError"; }}
export async function subscribeJetstream( store: JazzThoughtStore, connector: JetstreamConnector, options: JetstreamSubscriptionOptions = {},): Promise<JetstreamSubscriptionResult> { const config = normalizedOptions(options); const subscriptionId = newId("jetstream_subscription"); const startedAt = new Date().toISOString(); await appendLifecycle(store, connector, "stream.thought.connector.subscription.started", subscriptionId, startedAt, { status: "started", endpoint: redactedEndpoint(config.endpoint), rewindUs: config.rewindUs, maxMessages: Number.isFinite(config.maxMessages) ? config.maxMessages : "unbounded", });
let messages = 0; let inserted = 0; let unchanged = 0; let ignored = 0; let stale = 0; let reconnects = 0; let recoveredAfterFailure = false; let finalReason: JetstreamSubscriptionResult["reason"] = "closed"; let terminalError: unknown;
try { while (!config.signal?.aborted && messages < config.maxMessages) { const prior = await connector.readCursor(store); const durableTimeUs = cursorTimeUs(prior); const replayFromTimeUs = Math.max(0, durableTimeUs - config.rewindUs); const connectionId = newId("jetstream_connection"); const url = buildJetstreamSubscriptionUrl(connector, { endpoint: config.endpoint, ...(durableTimeUs > 0 ? { cursor: replayFromTimeUs } : {}), maxMessageSizeBytes: config.maxMessageSizeBytes, });
try { const connection = await consumeConnection(store, connector, url, connectionId, replayFromTimeUs, { ...(config.signal ? { signal: config.signal } : {}), remainingMessages: config.maxMessages - messages, pendingMessageLimit: config.pendingMessageLimit, maxMessageSizeBytes: config.maxMessageSizeBytes, webSocketFactory: config.webSocketFactory, onOpen: async () => { const connectedAt = new Date().toISOString(); await appendLifecycle(store, connector, "stream.thought.connector.subscription.connected", connectionId, connectedAt, { status: "connected", subscriptionId, reconnect: reconnects > 0, resumeCursor: durableTimeUs > 0 ? replayFromTimeUs : null, }); if (recoveredAfterFailure) { await recordRecovery(store, connector, connectionId, connectedAt, prior); recoveredAfterFailure = false; } }, });
messages += connection.messages; inserted += connection.inserted; unchanged += connection.unchanged; ignored += connection.ignored; stale += connection.stale;
if (connection.reason === "aborted" || connection.reason === "message-limit") { finalReason = connection.reason; break; }
const closeDescription = connection.transportError ?? (connection.overflowed ? `Jetstream pending-message limit ${config.pendingMessageLimit} exceeded` : `Jetstream socket closed${connection.closeCode === undefined ? "" : ` with code ${connection.closeCode}`}${connection.closeReason ? `: ${connection.closeReason}` : ""}`); await recordFailure(store, connector, connectionId, "live-subscribe", closeDescription); recoveredAfterFailure = true; } catch (error) { if (error instanceof ConnectionProgressError) { messages += error.progress.messages; inserted += error.progress.inserted; unchanged += error.progress.unchanged; ignored += error.progress.ignored; stale += error.progress.stale; } if (config.signal?.aborted) { finalReason = "aborted"; break; } if (!(error instanceof ConnectionProgressError && error.failureRecorded)) { await recordFailure( store, connector, connectionId, error instanceof ConnectionProgressError ? error.phase : "live-subscribe", error instanceof ConnectionProgressError ? error.cause : error, ); } recoveredAfterFailure = true; }
if (messages >= config.maxMessages) { finalReason = "message-limit"; break; } if (reconnects >= config.maxReconnects) { throw new Error(`Jetstream reconnect limit reached after ${reconnects} reconnects`); } const delayMs = reconnectDelay(reconnects, config.initialBackoffMs, config.maxBackoffMs, config.random); reconnects += 1; try { await config.sleep(delayMs, config.signal); } catch (error) { if (config.signal?.aborted) { finalReason = "aborted"; break; } throw error; } }
if (config.signal?.aborted) finalReason = "aborted"; else if (messages >= config.maxMessages) finalReason = "message-limit"; } catch (error) { terminalError = error; }
const stoppedAt = new Date().toISOString(); await appendLifecycle(store, connector, "stream.thought.connector.subscription.stopped", subscriptionId, stoppedAt, { status: terminalError ? "failed" : "stopped", reason: terminalError ? (terminalError instanceof Error ? terminalError.message : String(terminalError)) : finalReason, messages, inserted, unchanged, ignored, stale, reconnects, }); if (terminalError) throw terminalError;
const cursor = await connector.readCursor(store); return { source: connector.id, reason: finalReason, messages, inserted, unchanged, ignored, stale, reconnects, ...(cursor ? { cursor } : {}), };}
export function buildJetstreamSubscriptionUrl( connector: JetstreamConnector, options: { endpoint?: string; cursor?: number; maxMessageSizeBytes?: number } = {},): string { const endpoint = new URL(options.endpoint ?? DEFAULT_ENDPOINT); if (endpoint.protocol !== "ws:" && endpoint.protocol !== "wss:") { throw new Error("Jetstream endpoint must use ws: or wss:"); } for (const collection of connector.collections) endpoint.searchParams.append("wantedCollections", collection); for (const did of connector.dids) endpoint.searchParams.append("wantedDids", did); if (options.cursor !== undefined) { if (!Number.isSafeInteger(options.cursor) || options.cursor < 0) throw new Error("Jetstream cursor must be a nonnegative safe integer"); endpoint.searchParams.set("cursor", String(options.cursor)); } if (options.maxMessageSizeBytes !== undefined) { endpoint.searchParams.set("maxMessageSizeBytes", String(options.maxMessageSizeBytes)); } return endpoint.toString();}
async function consumeConnection( store: JazzThoughtStore, connector: JetstreamConnector, url: string, connectionId: string, replayFromTimeUs: number, options: { signal?: AbortSignal; remainingMessages: number; pendingMessageLimit: number; maxMessageSizeBytes: number; webSocketFactory: JetstreamWebSocketFactory; onOpen: () => Promise<void>; },): Promise<ConnectionResult> { if (options.signal?.aborted) return emptyConnection("aborted"); const socket = options.webSocketFactory(url); const queue = new AsyncQueue<SocketItem>(); let accepting = true; let overflowed = false; let transportError: string | undefined;
const onMessage = (event: { data?: unknown }) => { if (!accepting) return; if (queue.size >= options.pendingMessageLimit) { overflowed = true; accepting = false; socket.close(1013, "local backpressure"); queue.push({ type: "closed", code: 1013, reason: "local backpressure" }); return; } queue.push({ type: "message", data: event.data }); }; const onClose = (event: { code?: number; reason?: string }) => { accepting = false; queue.push({ type: "closed", ...(event.code === undefined ? {} : { code: event.code }), ...(event.reason ? { reason: event.reason } : {}) }); }; const onError = (event: { error?: unknown; message?: string }) => { transportError = event.error instanceof Error ? event.error.message : event.message ?? "Jetstream WebSocket error"; }; socket.addEventListener("message", onMessage); socket.addEventListener("close", onClose); socket.addEventListener("error", onError);
const onAbort = () => { if (!accepting) return; accepting = false; socket.close(1000, "aborted"); queue.push({ type: "closed", code: 1000, reason: "aborted" }); }; options.signal?.addEventListener("abort", onAbort, { once: true });
try { await waitForOpen(socket, options.signal); await options.onOpen(); const totals = { messages: 0, inserted: 0, unchanged: 0, ignored: 0, stale: 0 };
while (true) { const item = await queue.take(); if (item.type === "closed") { return { reason: options.signal?.aborted ? "aborted" : "closed", ...totals, ...(item.code === undefined ? {} : { closeCode: item.code }), ...(item.reason ? { closeReason: item.reason } : {}), overflowed, ...(transportError ? { transportError } : {}), }; }
totals.messages += 1; let batch; let raw: unknown; try { raw = JSON.parse(await socketMessageText(item.data, options.maxMessageSizeBytes)) as unknown; } catch (error) { accepting = false; socket.close(1003, "invalid message"); throw new ConnectionProgressError(error, { ...totals }, "live-subscribe", false); } try { batch = await connector.ingestBatch(store, [raw], { replayFromTimeUs }); } catch (error) { accepting = false; socket.close(1003, "invalid message"); throw new ConnectionProgressError(error, { ...totals }, "live-subscribe", true); } totals.inserted += batch.inserted; totals.unchanged += batch.unchanged; totals.ignored += batch.ignored; totals.stale += batch.stale; if (totals.messages >= options.remainingMessages) { accepting = false; socket.close(1000, "message limit reached"); return { reason: "message-limit", ...totals, overflowed }; } if (options.signal?.aborted) { accepting = false; socket.close(1000, "aborted"); return { reason: "aborted", ...totals, overflowed }; } } } finally { accepting = false; if (socket.readyState < 2) socket.close(1011, "consumer stopped"); options.signal?.removeEventListener("abort", onAbort); socket.removeEventListener("message", onMessage); socket.removeEventListener("close", onClose); socket.removeEventListener("error", onError); }}
async function waitForOpen(socket: JetstreamWebSocket, signal?: AbortSignal): Promise<void> { if (signal?.aborted) throw abortError(); await new Promise<void>((resolve, reject) => { const cleanup = () => { socket.removeEventListener("open", onOpen); socket.removeEventListener("error", onError); socket.removeEventListener("close", onClose); signal?.removeEventListener("abort", onAbort); }; const onOpen = () => { cleanup(); resolve(); }; const onError = (event: { error?: unknown; message?: string }) => { cleanup(); reject(event.error instanceof Error ? event.error : new Error(event.message ?? "Jetstream WebSocket connection failed")); }; const onClose = (event: { code?: number; reason?: string }) => { cleanup(); reject(new Error(`Jetstream WebSocket closed before opening${event.code === undefined ? "" : ` (${event.code})`}${event.reason ? `: ${event.reason}` : ""}`)); }; const onAbort = () => { cleanup(); socket.close(1000, "aborted"); reject(abortError()); }; socket.addEventListener("open", onOpen, { once: true }); socket.addEventListener("error", onError, { once: true }); socket.addEventListener("close", onClose, { once: true }); signal?.addEventListener("abort", onAbort, { once: true }); });}
async function appendLifecycle( store: JazzThoughtStore, connector: JetstreamConnector, type: "stream.thought.connector.subscription.started" | "stream.thought.connector.subscription.connected" | "stream.thought.connector.subscription.stopped", id: string, at: string, payload: JsonObject,): Promise<void> { await store.appendEvent({ type, schemaVersion: 1, source: connector.id, sourceKind: connector.kind, externalId: id, idempotencyKey: `${id}:${type}`, occurredAt: at, actor: connector.id, correlationId: id, privacy: "public-source", payload, });}
async function recordFailure( store: JazzThoughtStore, connector: JetstreamConnector, connectionId: string, phase: string, error: unknown,): Promise<void> { const at = new Date().toISOString(); const message = error instanceof Error ? error.message : String(error); const prior = await store.getSourceCursor(`cursor:${connector.id}`); await store.appendEvent({ type: "stream.thought.connector.failed", schemaVersion: 1, source: connector.id, sourceKind: connector.kind, externalId: connectionId, idempotencyKey: `${connectionId}:failed:${phase}`, occurredAt: at, actor: connector.id, correlationId: connectionId, privacy: "public-source", payload: { status: "failed", phase, error: message }, }); await store.upsertSourceCursor({ id: `cursor:${connector.id}`, source: connector.id, cursor: prior?.cursor ?? { filterRevision: connector.filterRevision }, ...(prior?.lastSuccessAt ? { lastSuccessAt: prior.lastSuccessAt } : {}), lastFailureAt: at, lastError: message, updatedAt: at, });}
async function recordRecovery( store: JazzThoughtStore, connector: JetstreamConnector, connectionId: string, at: string, prior: SourceCursor | undefined,): Promise<void> { await store.appendEvent({ type: "stream.thought.connector.recovered", schemaVersion: 1, source: connector.id, sourceKind: connector.kind, externalId: connectionId, idempotencyKey: `${connectionId}:recovered`, occurredAt: at, actor: connector.id, correlationId: connectionId, privacy: "public-source", payload: { status: "recovered", phase: "live-subscribe" }, }); await store.upsertSourceCursor({ id: `cursor:${connector.id}`, source: connector.id, cursor: prior?.cursor ?? { filterRevision: connector.filterRevision }, lastSuccessAt: at, ...(prior?.lastFailureAt ? { lastFailureAt: prior.lastFailureAt } : {}), updatedAt: at, });}
function normalizedOptions(options: JetstreamSubscriptionOptions) { const endpoint = options.endpoint ?? DEFAULT_ENDPOINT; const rewindUs = boundedInteger(options.rewindUs ?? 2_000_000, "rewindUs", 0); const maxMessages = options.maxMessages ?? Number.POSITIVE_INFINITY; if (!(maxMessages === Number.POSITIVE_INFINITY || (Number.isSafeInteger(maxMessages) && maxMessages > 0))) { throw new Error("maxMessages must be a positive safe integer or Infinity"); } const maxReconnects = options.maxReconnects ?? Number.POSITIVE_INFINITY; if (!(maxReconnects === Number.POSITIVE_INFINITY || (Number.isSafeInteger(maxReconnects) && maxReconnects >= 0))) { throw new Error("maxReconnects must be a nonnegative safe integer or Infinity"); } const initialBackoffMs = boundedInteger(options.initialBackoffMs ?? 1_000, "initialBackoffMs", 0); const maxBackoffMs = boundedInteger(options.maxBackoffMs ?? 30_000, "maxBackoffMs", initialBackoffMs); const maxMessageSizeBytes = boundedInteger(options.maxMessageSizeBytes ?? 1_000_000, "maxMessageSizeBytes", 1); const pendingMessageLimit = boundedInteger(options.pendingMessageLimit ?? 1_000, "pendingMessageLimit", 1); return { endpoint, rewindUs, maxMessages, maxReconnects, initialBackoffMs, maxBackoffMs, maxMessageSizeBytes, pendingMessageLimit, signal: options.signal, webSocketFactory: options.webSocketFactory ?? defaultWebSocketFactory, sleep: options.sleep ?? sleep, random: options.random ?? Math.random, };}
function defaultWebSocketFactory(url: string): JetstreamWebSocket { if (typeof WebSocket === "undefined") throw new Error("This runtime does not provide WebSocket"); return new WebSocket(url) as unknown as JetstreamWebSocket;}
function reconnectDelay(attempt: number, initialMs: number, maxMs: number, random: () => number): number { const exponential = Math.min(maxMs, initialMs * (2 ** Math.min(attempt, 30))); return Math.round(exponential * (0.8 + Math.min(1, Math.max(0, random())) * 0.4));}
async function sleep(milliseconds: number, signal?: AbortSignal): Promise<void> { if (signal?.aborted) throw abortError(); await new Promise<void>((resolve, reject) => { const timer = setTimeout(() => { cleanup(); resolve(); }, milliseconds); const onAbort = () => { cleanup(); reject(abortError()); }; const cleanup = () => { clearTimeout(timer); signal?.removeEventListener("abort", onAbort); }; signal?.addEventListener("abort", onAbort, { once: true }); });}
async function socketMessageText(data: unknown, maxBytes: number): Promise<string> { let value: string; if (typeof data === "string") value = data; else if (data instanceof ArrayBuffer) value = Buffer.from(data).toString("utf8"); else if (ArrayBuffer.isView(data)) value = Buffer.from(data.buffer, data.byteOffset, data.byteLength).toString("utf8"); else if (typeof Blob !== "undefined" && data instanceof Blob) value = await data.text(); else throw new Error("Jetstream WebSocket message must be text or binary JSON"); if (Buffer.byteLength(value, "utf8") > maxBytes) throw new Error(`Jetstream message exceeds ${maxBytes} bytes`); return value;}
function redactedEndpoint(value: string): string { const url = new URL(value); url.username = ""; url.password = ""; url.search = ""; url.hash = ""; return url.toString();}
function boundedInteger(value: number, label: string, minimum: number): number { if (!Number.isSafeInteger(value) || value < minimum) throw new Error(`${label} must be a safe integer >= ${minimum}`); return value;}
function emptyConnection(reason: "aborted"): ConnectionResult { return { reason, messages: 0, inserted: 0, unchanged: 0, ignored: 0, stale: 0, overflowed: false };}
function abortError(): Error { const error = new Error("Operation aborted"); error.name = "AbortError"; return error;}
class AsyncQueue<T> { private readonly values: T[] = []; private readonly waiters: Array<(value: T) => void> = [];
get size(): number { return this.values.length; }
push(value: T): void { const waiter = this.waiters.shift(); if (waiter) waiter(value); else this.values.push(value); }
async take(): Promise<T> { const value = this.values.shift(); if (value !== undefined) return value; return new Promise<T>((resolve) => this.waiters.push(resolve)); }}