import { 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; 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; 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 { 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; }, ): Promise { if (options.signal?.aborted) return emptyConnection("aborted"); const socket = options.webSocketFactory(url); const queue = new AsyncQueue(); 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 { if (signal?.aborted) throw abortError(); await new Promise((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 { 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 { 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 { 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 { if (signal?.aborted) throw abortError(); await new Promise((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 { 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 { 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 { const value = this.values.shift(); if (value !== undefined) return value; return new Promise((resolve) => this.waiters.push(resolve)); } }