import type { JetstreamEvent } from '@atcute/jetstream'; import type { Subscribe } from '../src/jetstream'; const SECOND = 1_000_000; export function commit( seconds: number, did: string, collection: string, rkey: string, operation: 'create' | 'update' | 'delete' = 'create', record: Record = {}, ): JetstreamEvent { const base = { rev: `rev-${seconds}`, collection, rkey, operation }; return { did, time_us: seconds * SECOND, kind: 'commit', commit: operation === 'delete' ? base : { ...base, cid: 'bafy', record }, } as JetstreamEvent; } export function identity(seconds: number, did: string, handle: string): JetstreamEvent { return { did, time_us: seconds * SECOND, kind: 'identity', identity: { did, handle, seq: seconds, time: new Date(0).toISOString() }, } as JetstreamEvent; } const parked = () => new Promise(() => {}); // replays whatever source() holds from the cursor on, then goes quiet like a // caught up jetstream. dropAfter[n] closes the nth connection after that many events export function fakeJetstream(source: () => JetstreamEvent[], dropAfter: (number | undefined)[] = []) { const cursors: number[] = []; const subscribe: Subscribe = (_url, cursor, _collections, onClose) => { const connection = cursors.push(cursor) - 1; const limit = dropAfter[connection]; const pending = source().filter((e) => e.time_us >= cursor); return { async *[Symbol.asyncIterator]() { for (const [i, evt] of pending.entries()) { if (i === limit) break; yield evt; } if (limit !== undefined) onClose(); await parked(); }, }; }; return { subscribe, cursors }; } // a jetstream that refuses every connection export const deadJetstream: Subscribe = (_url, _cursor, _collections, onClose) => ({ async *[Symbol.asyncIterator]() { queueMicrotask(onClose); await parked(); }, });