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(); + }); +});