From a39679ab0c566c4ce53fba0371750dc3641e8866 Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Mon, 6 Jul 2026 21:24:33 -0600 Subject: [PATCH] fix(explore): persist Jetstream cursor across cron windows The one-minute cron previously live-tailed Jetstream with no cursor, dropping any org.v-it.* commit that landed between cron windows. scheduled() now reads the cursor from the CursorStore Durable Object, starts from that cursor or a bounded 45-minute startup reconcile when unset, and persists the next cursor only after a fully successful stream window. The stored value advances to the latest observed event cursor, or to the window-open time when the window is quiet. Cursor read/write failures, malformed cursors, WebSocket and stream errors, and per-event D1 write failures now reject the run and leave the cursor untouched. There is no silent cursorless tailing fallback. Unit tests cover the scheduler with local fakes and verify replay idempotency with a bun:sqlite D1 shim over the real schema. They avoid live Jetstream, D1, and Cloudflare services. Co-Authored-By: Claude Opus 4.8 (1M context) --- explore/src/cursor.js | 23 +++ explore/src/index.js | 32 +++- explore/src/jetstream.js | 105 ++++++++++--- test/explore-cursor.test.js | 293 ++++++++++++++++++++++++++++++++++++ 4 files changed, 430 insertions(+), 23 deletions(-) create mode 100644 test/explore-cursor.test.js diff --git a/explore/src/cursor.js b/explore/src/cursor.js index e27d756..b97c12e 100644 --- a/explore/src/cursor.js +++ b/explore/src/cursor.js @@ -1,6 +1,29 @@ // SPDX-License-Identifier: MIT // Copyright (c) 2026 sol pbc +export const CURSOR_NAME = 'jetstream'; +const CURSOR_URL = 'https://cursor.internal/'; + +function cursorStub(env) { + const id = env.CURSOR_STORE.idFromName(CURSOR_NAME); + return env.CURSOR_STORE.get(id); +} + +export async function readCursor(env) { + const res = await cursorStub(env).fetch(CURSOR_URL, { method: 'GET' }); + if (!res.ok) { + throw new Error('cursor read failed: ' + res.status); + } + return await res.text(); +} + +export async function writeCursor(env, value) { + const res = await cursorStub(env).fetch(CURSOR_URL, { method: 'PUT', body: String(value) }); + if (!res.ok) { + throw new Error('cursor write failed: ' + res.status); + } +} + export class CursorStore { constructor(state, env) { this.state = state; diff --git a/explore/src/index.js b/explore/src/index.js index c918a05..317ae2a 100644 --- a/explore/src/index.js +++ b/explore/src/index.js @@ -2,10 +2,38 @@ // Copyright (c) 2026 sol pbc import { handleRequest } from './api.js'; +import { readCursor, writeCursor } from './cursor.js'; import { streamEvents } from './jetstream.js'; export { CursorStore } from './cursor.js'; +// 45 minutes in microseconds. +const STARTUP_REPLAY_US = 2_700_000_000; + +function validCursor(value) { + const parsed = Number(value); + return /^\d+$/.test(value) && parsed > 0 && Number.isSafeInteger(parsed); +} + +export async function runScheduled(env, { streamReader = streamEvents, now = Date.now } = {}) { + const windowOpen = now() * 1000; + const stored = await readCursor(env); + + let startCursor; + if (stored === '') { + startCursor = windowOpen - STARTUP_REPLAY_US; + } else if (validCursor(stored)) { + startCursor = Number(stored); + } else { + throw new Error('malformed cursor: ' + JSON.stringify(stored)); + } + + const { observedCursor } = await streamReader(env, startCursor); + const sawNewer = observedCursor != null && Number(observedCursor) > startCursor; + const nextCursor = sawNewer ? String(observedCursor) : String(windowOpen); + await writeCursor(env, nextCursor); +} + export default { async fetch(request, env) { const url = new URL(request.url); @@ -17,8 +45,6 @@ export default { }, async scheduled(event, env, ctx) { - // Always live-tail (no cursor) — Jetstream doesn't replay custom lexicon - // commits via cursor. D1 UNIQUE constraints handle deduplication. - const result = await streamEvents(env, null); + await runScheduled(env); }, }; diff --git a/explore/src/jetstream.js b/explore/src/jetstream.js index eabf3e4..72a6c1d 100644 --- a/explore/src/jetstream.js +++ b/explore/src/jetstream.js @@ -61,7 +61,7 @@ function decrementVouchBeaconStatement(env, beacon) { ).bind(beacon); } -async function processCapEvent(env, did, commit) { +export async function processCapEvent(env, did, commit) { const { operation, rkey, record, cid } = commit; const uri = `at://${did}/${CAP_COLLECTION}/${rkey}`; @@ -135,7 +135,7 @@ async function processCapEvent(env, did, commit) { } } -async function processVouchEvent(env, did, commit) { +export async function processVouchEvent(env, did, commit) { const { operation, rkey, record, cid } = commit; const uri = `at://${did}/${VOUCH_COLLECTION}/${rkey}`; @@ -207,7 +207,7 @@ async function processVouchEvent(env, did, commit) { } } -async function processSkillEvent(env, did, commit) { +export async function processSkillEvent(env, did, commit) { const { operation, rkey, record, cid } = commit; const uri = `at://${did}/${SKILL_COLLECTION}/${rkey}`; @@ -258,27 +258,88 @@ export async function streamEvents(env, cursor) { url.searchParams.set('cursor', cursor); } - return await new Promise((resolve) => { - let latestCursor = cursor || null; + return await new Promise((resolve, reject) => { + let observedCursor = null; const newDids = new Set(); const pending = new Set(); - const ws = new WebSocket(url.toString()); + let settled = false; + let ws; + let timeout; - const timeout = setTimeout(() => { - ws.close(); - }, STREAM_DURATION_MS); + const asError = (err) => { + if (err instanceof Error) { + return err; + } + return new Error(err?.message || 'WebSocket error'); + }; + + const clearWindow = () => { + if (timeout) { + clearTimeout(timeout); + timeout = null; + } + }; + + const fail = (err) => { + if (settled) { + return; + } + settled = true; + clearWindow(); + try { + ws?.close(); + } catch { + // Ignore close failures while rejecting the stream window. + } + reject(asError(err)); + }; + + const succeed = () => { + if (settled) { + return; + } + settled = true; + clearWindow(); + resolve({ observedCursor }); + }; const finish = async () => { - clearTimeout(timeout); - if (pending.size > 0) { - await Promise.allSettled([...pending]); + if (settled) { + return; } - if (newDids.size > 0) { - await resolveHandles([...newDids], env); + clearWindow(); + try { + if (pending.size > 0) { + await Promise.all([...pending]); + } + if (newDids.size > 0) { + try { + await resolveHandles([...newDids], env); + } catch { + // Handle resolution is best-effort and must not fail the window. + } + } + succeed(); + } catch (err) { + fail(err); } - resolve({ latestCursor }); }; + try { + ws = new WebSocket(url.toString()); + } catch (err) { + fail(err); + return; + } + + timeout = setTimeout(() => { + try { + ws.close(); + } catch (err) { + fail(err); + } + }, STREAM_DURATION_MS); + ws.addEventListener('message', (event) => { const task = (async () => { let msg; @@ -292,8 +353,11 @@ export async function streamEvents(env, cursor) { return; } - if (msg.time_us) { - latestCursor = String(msg.time_us); + if (msg.time_us != null) { + const timeUs = Number(msg.time_us); + if (Number.isFinite(timeUs) && (observedCursor === null || timeUs > Number(observedCursor))) { + observedCursor = String(msg.time_us); + } } if (msg.did) { @@ -315,15 +379,16 @@ export async function streamEvents(env, cursor) { })(); pending.add(task); - task.finally(() => pending.delete(task)); + task.catch(fail); + task.finally(() => pending.delete(task)).catch(() => {}); }); ws.addEventListener('close', () => { void finish(); }); - ws.addEventListener('error', () => { - ws.close(); + ws.addEventListener('error', (event) => { + fail(event?.error ?? event); }); }); } diff --git a/test/explore-cursor.test.js b/test/explore-cursor.test.js new file mode 100644 index 0000000..03deebd --- /dev/null +++ b/test/explore-cursor.test.js @@ -0,0 +1,293 @@ +// SPDX-License-Identifier: MIT +// Copyright (c) 2026 sol pbc + +import { describe, expect, test } from 'bun:test'; +import { readFileSync } from 'node:fs'; +import { join } from 'node:path'; +import { Database } from 'bun:sqlite'; +import { runScheduled } from '../explore/src/index.js'; +import { processCapEvent, processSkillEvent, processVouchEvent } from '../explore/src/jetstream.js'; + +function createCursorEnv({ + stored = '', + getStatus = 200, + putStatus = 200, + getThrows = null, + putThrows = null, +} = {}) { + const store = { + idNames: [], + ids: [], + gets: 0, + puts: 0, + putBodies: [], + idFromName(name) { + this.idNames.push(name); + return { name }; + }, + get(id) { + this.ids.push(id); + return { + fetch: async (url, { method, body } = {}) => { + if (method === 'GET') { + store.gets++; + if (getThrows) { + throw getThrows; + } + return new Response(stored, { status: getStatus }); + } + + if (method === 'PUT') { + store.puts++; + store.putBodies.push(String(body)); + if (putThrows) { + throw putThrows; + } + return new Response('ok', { status: putStatus }); + } + + return new Response('method not allowed', { status: 405 }); + }, + }; + }, + }; + + return { env: { CURSOR_STORE: store }, store }; +} + +function createStreamReader({ result = { observedCursor: null }, error = null } = {}) { + const calls = []; + const streamReader = async (env, startCursor) => { + calls.push({ env, startCursor }); + if (error) { + throw error; + } + return result; + }; + streamReader.calls = calls; + return streamReader; +} + +function d1Statement(db, sql, args = []) { + return { + sql, + args, + bind(...nextArgs) { + return d1Statement(db, sql, nextArgs); + }, + first() { + return db.query(sql).get(...args) ?? null; + }, + run() { + return db.query(sql).run(...args); + }, + all() { + return { results: db.query(sql).all(...args) }; + }, + }; +} + +function createD1(db) { + const executeBatch = db.transaction((statements) => statements.map((stmt) => { + if (/^\s*select\b/i.test(stmt.sql)) { + return stmt.all(); + } + return stmt.run(); + })); + + return { + prepare(sql) { + return d1Statement(db, sql); + }, + batch(statements) { + return executeBatch(statements); + }, + }; +} + +function createSqliteEnv() { + const db = new Database(':memory:'); + const schemaPath = join(import.meta.dir, '..', 'explore', 'schema.sql'); + db.exec(readFileSync(schemaPath, 'utf8')); + return { db, env: { DB: createD1(db) } }; +} + +describe('explore scheduled cursor', () => { + test('passes stored valid cursor and writes newer observed cursor', async () => { + const { env, store } = createCursorEnv({ stored: '12345' }); + const streamReader = createStreamReader({ result: { observedCursor: '12399' } }); + + await runScheduled(env, { streamReader, now: () => 999 }); + + expect(store.idNames).toEqual(['jetstream', 'jetstream']); + expect(streamReader.calls.length).toBe(1); + expect(streamReader.calls[0].env).toBe(env); + expect(streamReader.calls[0].startCursor).toBe(12345); + expect(store.putBodies).toEqual(['12399']); + }); + + test('uses startup replay window when no cursor is stored', async () => { + const withEvent = createCursorEnv({ stored: '' }); + const eventReader = createStreamReader({ result: { observedCursor: '300000123' } }); + + await runScheduled(withEvent.env, { streamReader: eventReader, now: () => 3_000_000 }); + + expect(eventReader.calls[0].startCursor).toBe(300_000_000); + expect(withEvent.store.putBodies).toEqual(['300000123']); + + const quiet = createCursorEnv({ stored: '' }); + const quietReader = createStreamReader({ result: { observedCursor: null } }); + + await runScheduled(quiet.env, { streamReader: quietReader, now: () => 3_000_000 }); + + expect(quietReader.calls[0].startCursor).toBe(300_000_000); + expect(quiet.store.putBodies).toEqual(['3000000000']); + }); + + test('writes window-open cursor for a quiet run', async () => { + const { env, store } = createCursorEnv({ stored: '1000' }); + const streamReader = createStreamReader({ result: { observedCursor: null } }); + + await runScheduled(env, { streamReader, now: () => 2 }); + + expect(streamReader.calls[0].startCursor).toBe(1000); + expect(store.putBodies).toEqual(['2000']); + }); + + test('cursor write failure rejects for event and quiet-window cursors', async () => { + const withEvent = createCursorEnv({ stored: '1000', putStatus: 500 }); + const eventReader = createStreamReader({ result: { observedCursor: '1500' } }); + + await expect(runScheduled(withEvent.env, { streamReader: eventReader, now: () => 2 })) + .rejects.toThrow('cursor write failed: 500'); + expect(withEvent.store.putBodies).toEqual(['1500']); + + const quiet = createCursorEnv({ stored: '1000', putStatus: 500 }); + const quietReader = createStreamReader({ result: { observedCursor: null } }); + + await expect(runScheduled(quiet.env, { streamReader: quietReader, now: () => 2 })) + .rejects.toThrow('cursor write failed: 500'); + expect(quiet.store.putBodies).toEqual(['2000']); + }); + + test('cursor read failure rejects before streaming', async () => { + const cases = [ + { fake: createCursorEnv({ stored: '1000', getStatus: 500 }), message: 'cursor read failed: 500' }, + { fake: createCursorEnv({ stored: '1000', getThrows: new Error('read exploded') }), message: 'read exploded' }, + ]; + + for (const { fake, message } of cases) { + const streamReader = createStreamReader(); + + await expect(runScheduled(fake.env, { streamReader, now: () => 2 })) + .rejects.toThrow(message); + expect(streamReader.calls.length).toBe(0); + expect(fake.store.putBodies).toEqual([]); + } + }); + + test('malformed stored cursors reject before streaming', async () => { + for (const stored of ['abc', '-5', '0', '1.5', ' ']) { + const { env, store } = createCursorEnv({ stored }); + const streamReader = createStreamReader(); + + await expect(runScheduled(env, { streamReader, now: () => 2 })) + .rejects.toThrow('malformed cursor: ' + JSON.stringify(stored)); + expect(streamReader.calls.length).toBe(0); + expect(store.putBodies).toEqual([]); + } + }); + + test('stream reader rejection propagates without writing cursor', async () => { + const { env, store } = createCursorEnv({ stored: '1000' }); + const streamReader = createStreamReader({ error: new Error('stream failed') }); + + await expect(runScheduled(env, { streamReader, now: () => 2 })) + .rejects.toThrow('stream failed'); + expect(streamReader.calls.length).toBe(1); + expect(store.putBodies).toEqual([]); + }); +}); + +describe('explore event idempotency', () => { + test('processes duplicate cap, vouch, and skill creates idempotently', async () => { + const { db, env } = createSqliteEnv(); + + const capDid = 'did:plc:capauthor'; + const capCommit = { + operation: 'create', + rkey: '3lcap', + cid: 'bafycap', + record: { + title: 'Idempotent Cap', + description: 'A cap inserted twice for testing.', + ref: 'idempotent-cap-test', + beacon: 'vit:example/repo', + kind: 'test', + createdAt: '2026-07-07T00:00:00.000Z', + }, + }; + + await processCapEvent(env, capDid, capCommit); + const capAfterOne = db + .query("SELECT (SELECT COUNT(*) FROM caps) AS rows, (SELECT cap_count FROM beacons WHERE name = ?) AS count") + .get('vit:example/repo'); + await processCapEvent(env, capDid, capCommit); + const capAfterTwo = db + .query("SELECT (SELECT COUNT(*) FROM caps) AS rows, (SELECT cap_count FROM beacons WHERE name = ?) AS count") + .get('vit:example/repo'); + + expect(capAfterOne).toEqual({ rows: 1, count: 1 }); + expect(capAfterTwo).toEqual(capAfterOne); + + const vouchDid = 'did:plc:vouchauthor'; + const vouchCommit = { + operation: 'create', + rkey: '3lvouch', + cid: 'bafyvouch', + record: { + subject: { uri: 'at://did:plc:capauthor/org.v-it.cap/3lcap' }, + ref: 'idempotent-cap-test', + beacon: 'vit:example/repo', + kind: 'want', + createdAt: '2026-07-07T00:00:01.000Z', + }, + }; + + await processVouchEvent(env, vouchDid, vouchCommit); + const vouchAfterOne = db + .query("SELECT (SELECT COUNT(*) FROM vouches) AS rows, (SELECT vouch_count FROM beacons WHERE name = ?) AS count") + .get('vit:example/repo'); + await processVouchEvent(env, vouchDid, vouchCommit); + const vouchAfterTwo = db + .query("SELECT (SELECT COUNT(*) FROM vouches) AS rows, (SELECT vouch_count FROM beacons WHERE name = ?) AS count") + .get('vit:example/repo'); + + expect(vouchAfterOne).toEqual({ rows: 1, count: 1 }); + expect(vouchAfterTwo).toEqual(vouchAfterOne); + + const skillDid = 'did:plc:skillauthor'; + const skillCommit = { + operation: 'create', + rkey: '3lskill', + cid: 'bafyskill', + record: { + name: 'idempotent-skill', + description: 'A skill inserted twice for testing.', + version: '1.0.0', + tags: ['test'], + createdAt: '2026-07-07T00:00:02.000Z', + }, + }; + + await processSkillEvent(env, skillDid, skillCommit); + const skillAfterOne = db.query('SELECT COUNT(*) AS rows FROM skills').get(); + await processSkillEvent(env, skillDid, skillCommit); + const skillAfterTwo = db.query('SELECT COUNT(*) AS rows FROM skills').get(); + + expect(skillAfterOne).toEqual({ rows: 1 }); + expect(skillAfterTwo).toEqual(skillAfterOne); + + db.close(); + }); +}); -- 2.51.2