import { jetstreamEventSchema, type JetstreamEvent } from '@atcute/jetstream'; import * as v from 'valibot'; // jetstream only roughly orders events by time_us, so every connection starts // this far behind the newest event we've seen and `seen` drops the repeats export const REWIND_US = 10_000_000; export type Subscribe = ( url: string, cursor: number, collections: string[], onClose: () => void, ) => AsyncIterable; async function connect(url: URL): Promise { const res = await fetch(url, { headers: { Upgrade: 'websocket' } }); if (!res.webSocket) throw new Error(`jetstream upgrade failed: http ${res.status}`); res.webSocket.accept(); return res.webSocket; } function parseEvent(data: string | ArrayBuffer): JetstreamEvent | undefined { try { const raw = JSON.parse(typeof data === 'string' ? data : new TextDecoder().decode(data)); const result = v.safeParse(jetstreamEventSchema, raw); return result.success ? result.output : undefined; } catch { return undefined; } } // atcute's JetstreamSubscription rides on partysocket, which throws on every // message in workerd because it clones events with their own constructor. so // this is a plain socket using atcute's schema, and readJetstream reconnects export const subscribeJetstream: Subscribe = (url, cursor, collections, onClose) => ({ [Symbol.asyncIterator]() { const endpoint = new URL('/subscribe', url); endpoint.searchParams.set('cursor', String(cursor)); for (const collection of collections) endpoint.searchParams.append('wantedCollections', collection); const queue: JetstreamEvent[] = []; let ended = false; let wake: (() => void) | undefined; let socket: WebSocket | undefined; const end = () => { ended = true; wake?.(); }; // onClose means the socket is really gone, so it fires after return() too const gone = () => { end(); onClose(); }; const opened = connect(endpoint).then( (ws) => { socket = ws; ws.addEventListener('message', (msg) => { if (ended) return; const evt = parseEvent(msg.data); if (!evt) return; queue.push(evt); wake?.(); }); ws.addEventListener('close', gone); ws.addEventListener('error', gone); }, (err) => { console.warn('[sitemap] jetstream connect failed:', err); gone(); }, ); return { async next(): Promise> { await opened; while (queue.length === 0 && !ended) { await new Promise((resolve) => (wake = resolve)); } const evt = queue.shift(); return evt ? { value: evt, done: false } : { value: undefined, done: true }; }, async return(): Promise> { end(); socket?.close(1000, 'done'); return { value: undefined, done: true }; }, }; }, }); export interface ReadOptions { url: string; cursor: number; collections: string[]; maxEvents: number; // quiet for this long means the replay reached the live tip idleMs: number; // anything newer happened after the run started, so the backlog is done untilUs: number; // how long to wait for a closed socket to actually go closeMs: number; maxReconnects: number; seen: Set; } export interface ReadResult { events: JetstreamEvent[]; cursor: number; full: boolean; error?: string; } type Outcome = 'dropped' | 'idle' | 'caught-up' | 'full'; function eventKey(evt: JetstreamEvent): string { const commit = evt.kind === 'commit' ? `${evt.commit.collection}/${evt.commit.rkey}@${evt.commit.rev}` : ''; return `${evt.time_us}:${evt.did}:${evt.kind}:${commit}`; } export async function readJetstream(opts: ReadOptions, subscribe = subscribeJetstream): Promise { const events: JetstreamEvent[] = []; let cursor = opts.cursor; let drops = 0; while (true) { let dropped!: () => void; const drop = new Promise<'dropped'>((resolve) => (dropped = () => resolve('dropped'))); const from = Math.max(0, cursor - REWIND_US); const it = subscribe(opts.url, from, opts.collections, () => dropped())[Symbol.asyncIterator](); let outcome: Outcome; try { outcome = await drain(it, drop); } finally { // not awaited, a fake stream parked forever would never let return() finish it.return?.()?.catch(() => {}); } if (outcome !== 'dropped') { // a worker gets six open connections, and a socket we closed holds one // until the server's close frame lands behind all the backlog it already // sent. so wait, or the d1 queries after a few batches queue up forever await Promise.race([drop, new Promise((resolve) => setTimeout(resolve, opts.closeMs))]); return { events, cursor, full: outcome === 'full' }; } if (++drops > opts.maxReconnects) { return { events, cursor, full: false, error: `jetstream disconnected ${drops} times` }; } } async function drain(it: AsyncIterator, drop: Promise<'dropped'>): Promise { while (true) { let timer: ReturnType | null = null; const idle = new Promise<'idle'>((resolve) => (timer = setTimeout(() => resolve('idle'), opts.idleMs))); const next = await Promise.race([it.next(), drop, idle]).finally(() => clearTimeout(timer)); if (next === 'dropped' || next === 'idle') return next; if (next.done) return 'dropped'; const evt = next.value; const key = eventKey(evt); if (opts.seen.has(key)) continue; opts.seen.add(key); events.push(evt); cursor = Math.max(cursor, evt.time_us); if (evt.time_us >= opts.untilUs) return 'caught-up'; if (events.length >= opts.maxEvents) return 'full'; } } }