Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174import { 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 repeatsexport const REWIND_US = 10_000_000;
export type Subscribe = ( url: string, cursor: number, collections: string[], onClose: () => void,) => AsyncIterable<JetstreamEvent>;
async function connect(url: URL): Promise<WebSocket> { 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 reconnectsexport 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<IteratorResult<JetstreamEvent>> { await opened; while (queue.length === 0 && !ended) { await new Promise<void>((resolve) => (wake = resolve)); } const evt = queue.shift(); return evt ? { value: evt, done: false } : { value: undefined, done: true }; }, async return(): Promise<IteratorResult<JetstreamEvent>> { 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<string>;}
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<ReadResult> { 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<JetstreamEvent>, drop: Promise<'dropped'>): Promise<Outcome> { while (true) { let timer: ReturnType<typeof setTimeout> | 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'; } }}