import assert from 'node:assert/strict' import { it } from 'node:test' import { COLLECTIONS, MemoryRecordStore, collapse } from '../../core/dist/index.js' import { JetstreamSubscription, MemorySyncStateStore, RepoPoller } from '../dist/index.js' const DID = 'did:plc:agent' const OTHER = 'did:plc:human' const ENDPOINT = 'wss://jetstream.test/subscribe' /** A fake `WebSocket`: the four events the client listens for, driven by hand. */ class FakeSocket { static opened = [] listeners = new Map() closed = false constructor(url) { this.url = url FakeSocket.opened.push(this) } addEventListener(type, listener) { const list = this.listeners.get(type) ?? [] list.push(listener) this.listeners.set(type, list) } emit(type, event) { for (const listener of this.listeners.get(type) ?? []) listener(event) } open() { this.emit('open') } send(payload) { this.emit('message', { data: typeof payload === 'string' ? payload : JSON.stringify(payload) }) } close() { this.closed = true } /** What a real socket does when the far end goes away. */ drop() { this.closed = true this.emit('close') } } /** Deterministic timers: nothing in these tests waits on a real clock. */ function fakeTimers() { const pending = new Map() let next = 1 return { set(callback, milliseconds) { const handle = next++ pending.set(handle, { callback, milliseconds }) return handle }, clear(handle) { pending.delete(handle) }, /** Fire every scheduled callback once, oldest first. */ run() { const due = [...pending.entries()].sort(([a], [b]) => a - b) pending.clear() for (const [, entry] of due) entry.callback() return due.length }, get size() { return pending.size }, } } const goalRecord = (title = 'Ship it') => ({ $type: COLLECTIONS.goal, space: { uri: `at://${OTHER}/${COLLECTIONS.space}/space`, cid: 'cid-space' }, project: { uri: `at://${OTHER}/${COLLECTIONS.project}/project`, cid: 'cid-project' }, title, body: 'Body', createdAt: '2026-01-01T00:00:00Z', }) const commit = ({ commit: commitOverrides, ...overrides } = {}) => ({ did: DID, kind: 'commit', time_us: 1_700_000_000_000_000, ...overrides, commit: { rev: '3kaaaaaaaaa2a', operation: 'create', collection: COLLECTIONS.goal, rkey: 'goal-1', cid: 'cid-goal-1', record: goalRecord(), ...commitOverrides, }, }) function subscribe(options = {}) { FakeSocket.opened = [] const records = [] const cursors = [] const statuses = [] const timers = fakeTimers() const subscription = new JetstreamSubscription({ endpoint: ENDPOINT, dids: options.dids ?? (() => [DID]), onRecord: (record) => records.push(record), onCursor: (cursor) => cursors.push(cursor), onStatus: (status, detail) => statuses.push({ status, detail }), socket: (url) => new FakeSocket(url), timers, now: options.now, ...options.overrides, }) return { subscription, records, cursors, statuses, timers, socket: () => FakeSocket.opened.at(-1) } } it('subscribes with every Radial collection and the space DIDs, and stores a streamed record', () => { const s = subscribe({ now: () => Date.parse('2026-01-01T00:00:30Z') }) s.subscription.start() const url = new URL(s.socket().url) assert.deepEqual(url.searchParams.getAll('wantedDids'), [DID]) const wanted = url.searchParams.getAll('wantedCollections') assert.equal(wanted.includes(COLLECTIONS.goal), true) assert.equal(wanted.includes(COLLECTIONS.claim), true) assert.equal(wanted.length >= 20, true) assert.equal(url.searchParams.get('cursor'), null) s.socket().open() s.socket().send(commit()) assert.equal(s.records.length, 1) assert.deepEqual(s.records[0], { did: DID, collection: COLLECTIONS.goal, rkey: 'goal-1', uri: `at://${DID}/${COLLECTIONS.goal}/goal-1`, cid: 'cid-goal-1', rev: '3kaaaaaaaaa2a', // Stamped when the stream delivered it: observer-local, like the rev, and what claim liveness // measures a declared lease from. firstSeenAt: '2026-01-01T00:00:30.000Z', value: goalRecord(), }) assert.deepEqual(s.statuses.map((entry) => entry.status), ['connected']) s.subscription.stop() }) it('rejects a record that fails lexicon validation, and never stores it', () => { const s = subscribe() s.subscription.start() s.socket().open() s.socket().send(commit({ commit: { record: { ...goalRecord(), title: 42 } } })) assert.deepEqual(s.records, []) assert.equal(s.subscription.stats.rejected, 1) s.subscription.stop() }) it('counts a delete and never applies it', () => { const s = subscribe() s.subscription.start() s.socket().open() s.socket().send(commit({ commit: { operation: 'delete', record: undefined } })) assert.deepEqual(s.records, []) assert.equal(s.subscription.stats.deletes, 1) assert.equal(s.subscription.stats.rejected, 0) s.subscription.stop() }) it('ignores an unparseable frame, a non-commit event and a foreign collection', () => { const s = subscribe() s.subscription.start() s.socket().open() s.socket().send('{not json') s.socket().send({ kind: 'identity', did: DID, time_us: 1 }) s.socket().send(commit({ commit: { collection: 'app.bsky.feed.post' } })) assert.deepEqual(s.records, []) assert.equal(s.subscription.stats.ignored, 3) s.subscription.stop() }) it('persists the cursor for every event and replays from before it on reconnect', () => { const s = subscribe() s.subscription.start() s.socket().open() s.socket().send(commit({ time_us: 1_700_000_000_000_000 })) // The cursor advances even for events we ignore: a resume point that only moved on records we // liked would replay the whole quiet stretch after a restart. s.socket().send({ kind: 'identity', did: DID, time_us: 1_700_000_009_000_000 }) assert.deepEqual(s.cursors, ['1700000000000000', '1700000009000000']) assert.equal(s.subscription.cursor, '1700000009000000') s.socket().drop() assert.equal(s.timers.run(), 1) // the backoff reconnect const resumed = new URL(s.socket().url) // Rewound 5s (5_000_000µs), so the reconnect OVERLAPS. Duplicates are free — `collapse` dedupes // by cid — and a gap is not. assert.equal(resumed.searchParams.get('cursor'), '1700000004000000') s.subscription.stop() }) it('replaying a gap produces no duplicate records in the store', () => { const s = subscribe() const store = new MemoryRecordStore() s.subscription.start() s.socket().open() const event = commit() s.socket().send(event) store.put(s.records.at(-1)) s.socket().drop() s.timers.run() s.socket().open() s.socket().send(event) // the overlap re-delivers it store.put(s.records.at(-1)) assert.equal(s.records.length, 2) assert.equal(store.records().length, 1) assert.deepEqual(store.edits(), []) s.subscription.stop() }) it('reports degraded after a silent window, without dropping the socket', () => { let clock = 0 const s = subscribe({ now: () => clock, overrides: { degradedAfterMs: 10_000 } }) s.subscription.start() s.socket().open() clock = 30_000 s.timers.run() // the health check assert.deepEqual( s.statuses.map((entry) => entry.status), ['connected', 'degraded'], ) assert.equal(s.socket().closed, false) assert.equal(s.subscription.connected, false) // An event brings it back: degraded is "we have lost it", not a terminal state. s.socket().send(commit()) assert.equal(s.statuses.at(-1).status, 'connected') assert.equal(s.subscription.connected, true) s.subscription.stop() }) it('reconnects with a new DID filter when membership changes', () => { let dids = [DID] const s = subscribe({ dids: () => dids }) s.subscription.start() s.socket().open() const first = s.socket() s.subscription.refresh() assert.equal(s.socket(), first) // unchanged membership: no churn dids = [DID, OTHER] s.subscription.refresh() assert.notEqual(s.socket(), first) assert.equal(first.closed, true) assert.deepEqual(new URL(s.socket().url).searchParams.getAll('wantedDids'), [DID, OTHER].sort()) s.subscription.stop() }) it('leaks no socket when membership changes inside a reconnect backoff window', () => { let dids = [DID] const s = subscribe({ dids: () => dids }) s.subscription.start() s.socket().open() const first = s.socket() // The far end goes away: a reconnect is now armed and `#socket` is undefined. first.drop() assert.equal(s.timers.size, 1) // Membership changes before the backoff elapses. The refresh connects immediately; the armed // timer must NOT then open a second socket over it — nothing would ever close that one, and its // events would be silently dropped by the `this.#socket !== socket` guard every listener has. dids = [DID, OTHER] s.subscription.refresh() assert.equal(FakeSocket.opened.length, 2) assert.equal(s.timers.run(), 0) assert.equal(FakeSocket.opened.length, 2) // And every socket the client ever opened is either the live one or closed. s.subscription.stop() assert.deepEqual(FakeSocket.opened.map((socket) => socket.closed), [true, true]) }) it('stops cleanly: no socket, no timers, no further work', () => { const s = subscribe({ overrides: { degradedAfterMs: 1_000 } }) s.subscription.start() s.socket().open() s.subscription.stop() assert.equal(s.socket().closed, true) assert.equal(s.timers.size, 0) assert.equal(s.subscription.status, undefined) }) it('persists a cursor in its own SQLite table, separate from the repo checkpoints', async () => { const { SqliteSyncStateStore } = await import('../dist/node.js') const { mkdtemp, rm } = await import('node:fs/promises') const { tmpdir } = await import('node:os') const { join } = await import('node:path') const directory = await mkdtemp(join(tmpdir(), 'radial-cursor-')) try { const path = join(directory, 'sync.db') const key = `${ENDPOINT} at://did:plc:root/${COLLECTIONS.space}/space` const first = new SqliteSyncStateStore(path) assert.equal(first.getCursor(key), undefined) first.setCursor(key, '1700000000000000') first.set({ did: DID, rev: HEAD_REV, commitCid: 'head-cid', firstObservedRev: HEAD_REV, updatedAt: '2026-01-01T00:00:00Z', }) first.setCursor(key, '1700000009000000') first.close() const second = new SqliteSyncStateStore(path) assert.equal(second.getCursor(key), '1700000009000000') // A cursor is not a repo rev, and the two never share a row. assert.equal(second.get(DID).rev, HEAD_REV) assert.equal(second.getCursor('some other endpoint'), undefined) second.close() } finally { await rm(directory, { recursive: true, force: true }) } }) // --- D3: provenance reconciliation ------------------------------------------ // // The poller stamps every record with the repo HEAD rev at scan time; a Jetstream commit carries // the record's REAL rev, which is lower. `collapse` keeps the minimum rev per cid, so the stream's // truer provenance wins whichever path saw the record first. const HEAD_REV = '3kzzzzzzzzz2z' const COMMIT_REV = '3kaaaaaaaaa2a' class OneRecordTransport { constructor(record) { this.record = record } async getLatestCommit() { return { rev: HEAD_REV, commitCid: 'head-cid' } } async listRecords({ collection }) { if (collection !== COLLECTIONS.goal) return { records: [] } return { records: [{ uri: this.record.uri, cid: this.record.cid, value: this.record.value }] } } } it('reconciles poller and stream provenance to the same (lower) rev, either arrival order', async () => { // The stream saw it first, the poll a minute later: both observer-local stamps — the rev and the // first sighting — keep the value that came from the EARLIER sighting, so which path happened to // deliver a record cannot change what the fold reads off it. const streamed = { did: DID, collection: COLLECTIONS.goal, rkey: 'goal-1', uri: `at://${DID}/${COLLECTIONS.goal}/goal-1`, cid: 'cid-goal-1', rev: COMMIT_REV, firstSeenAt: '2026-01-01T00:00:00.000Z', value: goalRecord(), } const polled = { now: () => '2026-01-01T00:01:00.000Z' } const pollFirst = new MemoryRecordStore() await new RepoPoller(new OneRecordTransport(streamed), pollFirst, new MemorySyncStateStore(), polled).pollDid(DID) assert.equal(pollFirst.records()[0].rev, HEAD_REV) assert.equal(pollFirst.records()[0].firstSeenAt, '2026-01-01T00:01:00.000Z') pollFirst.put(streamed) const streamFirst = new MemoryRecordStore() streamFirst.put(streamed) await new RepoPoller(new OneRecordTransport(streamed), streamFirst, new MemorySyncStateStore(), polled).pollDid(DID) assert.equal(pollFirst.records().length, 1) assert.equal(streamFirst.records().length, 1) assert.equal(pollFirst.records()[0].rev, COMMIT_REV) assert.equal(pollFirst.records()[0].firstSeenAt, '2026-01-01T00:00:00.000Z') assert.deepEqual(streamFirst.records(), pollFirst.records()) // And neither path counts the other as an edit: same cid, so `collapse` merges rather than judges. assert.deepEqual(collapse([...pollFirst.records(), streamed]).edits, []) })