open source is social v-it.org
0

Configure Feed

Select the types of activity you want to include in your feed.

Merge remote-tracking branch 'origin/main'

# Conflicts:
# explore/src/jetstream.js

+430 -23
+23
explore/src/cursor.js
··· 1 1 // SPDX-License-Identifier: MIT 2 2 // Copyright (c) 2026 sol pbc 3 3 4 + export const CURSOR_NAME = 'jetstream'; 5 + const CURSOR_URL = 'https://cursor.internal/'; 6 + 7 + function cursorStub(env) { 8 + const id = env.CURSOR_STORE.idFromName(CURSOR_NAME); 9 + return env.CURSOR_STORE.get(id); 10 + } 11 + 12 + export async function readCursor(env) { 13 + const res = await cursorStub(env).fetch(CURSOR_URL, { method: 'GET' }); 14 + if (!res.ok) { 15 + throw new Error('cursor read failed: ' + res.status); 16 + } 17 + return await res.text(); 18 + } 19 + 20 + export async function writeCursor(env, value) { 21 + const res = await cursorStub(env).fetch(CURSOR_URL, { method: 'PUT', body: String(value) }); 22 + if (!res.ok) { 23 + throw new Error('cursor write failed: ' + res.status); 24 + } 25 + } 26 + 4 27 export class CursorStore { 5 28 constructor(state, env) { 6 29 this.state = state;
+29 -3
explore/src/index.js
··· 2 2 // Copyright (c) 2026 sol pbc 3 3 4 4 import { handleRequest } from './api.js'; 5 + import { readCursor, writeCursor } from './cursor.js'; 5 6 import { streamEvents } from './jetstream.js'; 6 7 7 8 export { CursorStore } from './cursor.js'; 8 9 10 + // 45 minutes in microseconds. 11 + const STARTUP_REPLAY_US = 2_700_000_000; 12 + 13 + function validCursor(value) { 14 + const parsed = Number(value); 15 + return /^\d+$/.test(value) && parsed > 0 && Number.isSafeInteger(parsed); 16 + } 17 + 18 + export async function runScheduled(env, { streamReader = streamEvents, now = Date.now } = {}) { 19 + const windowOpen = now() * 1000; 20 + const stored = await readCursor(env); 21 + 22 + let startCursor; 23 + if (stored === '') { 24 + startCursor = windowOpen - STARTUP_REPLAY_US; 25 + } else if (validCursor(stored)) { 26 + startCursor = Number(stored); 27 + } else { 28 + throw new Error('malformed cursor: ' + JSON.stringify(stored)); 29 + } 30 + 31 + const { observedCursor } = await streamReader(env, startCursor); 32 + const sawNewer = observedCursor != null && Number(observedCursor) > startCursor; 33 + const nextCursor = sawNewer ? String(observedCursor) : String(windowOpen); 34 + await writeCursor(env, nextCursor); 35 + } 36 + 9 37 export default { 10 38 async fetch(request, env) { 11 39 const url = new URL(request.url); ··· 17 45 }, 18 46 19 47 async scheduled(event, env, ctx) { 20 - // Always live-tail (no cursor) — Jetstream doesn't replay custom lexicon 21 - // commits via cursor. D1 UNIQUE constraints handle deduplication. 22 - const result = await streamEvents(env, null); 48 + await runScheduled(env); 23 49 }, 24 50 };
+85 -20
explore/src/jetstream.js
··· 73 73 ).bind(beacon); 74 74 } 75 75 76 - async function processCapEvent(env, did, commit) { 76 + export async function processCapEvent(env, did, commit) { 77 77 const { operation, rkey, record, cid } = commit; 78 78 const uri = `at://${did}/${CAP_COLLECTION}/${rkey}`; 79 79 ··· 147 147 } 148 148 } 149 149 150 - async function processVouchEvent(env, did, commit) { 150 + export async function processVouchEvent(env, did, commit) { 151 151 const { operation, rkey, record, cid } = commit; 152 152 const uri = `at://${did}/${VOUCH_COLLECTION}/${rkey}`; 153 153 ··· 219 219 } 220 220 } 221 221 222 - async function processSkillEvent(env, did, commit) { 222 + export async function processSkillEvent(env, did, commit) { 223 223 const { operation, rkey, record, cid } = commit; 224 224 const uri = `at://${did}/${SKILL_COLLECTION}/${rkey}`; 225 225 ··· 270 270 url.searchParams.set('cursor', cursor); 271 271 } 272 272 273 - return await new Promise((resolve) => { 274 - let latestCursor = cursor || null; 273 + return await new Promise((resolve, reject) => { 274 + let observedCursor = null; 275 275 const newDids = new Set(); 276 276 const pending = new Set(); 277 277 const recordTasks = new Map(); 278 - const ws = new WebSocket(url.toString()); 278 + let settled = false; 279 + let ws; 280 + let timeout; 281 + 282 + const asError = (err) => { 283 + if (err instanceof Error) { 284 + return err; 285 + } 286 + return new Error(err?.message || 'WebSocket error'); 287 + }; 288 + 289 + const clearWindow = () => { 290 + if (timeout) { 291 + clearTimeout(timeout); 292 + timeout = null; 293 + } 294 + }; 279 295 280 - const timeout = setTimeout(() => { 281 - ws.close(); 282 - }, STREAM_DURATION_MS); 296 + const fail = (err) => { 297 + if (settled) { 298 + return; 299 + } 300 + settled = true; 301 + clearWindow(); 302 + try { 303 + ws?.close(); 304 + } catch { 305 + // Ignore close failures while rejecting the stream window. 306 + } 307 + reject(asError(err)); 308 + }; 309 + 310 + const succeed = () => { 311 + if (settled) { 312 + return; 313 + } 314 + settled = true; 315 + clearWindow(); 316 + resolve({ observedCursor }); 317 + }; 283 318 284 319 const finish = async () => { 285 - clearTimeout(timeout); 286 - if (pending.size > 0) { 287 - await Promise.allSettled([...pending]); 320 + if (settled) { 321 + return; 288 322 } 289 - if (newDids.size > 0) { 290 - await resolveHandles([...newDids], env); 323 + clearWindow(); 324 + try { 325 + if (pending.size > 0) { 326 + await Promise.all([...pending]); 327 + } 328 + if (newDids.size > 0) { 329 + try { 330 + await resolveHandles([...newDids], env); 331 + } catch { 332 + // Handle resolution is best-effort and must not fail the window. 333 + } 334 + } 335 + succeed(); 336 + } catch (err) { 337 + fail(err); 291 338 } 292 - resolve({ latestCursor }); 293 339 }; 294 340 341 + try { 342 + ws = new WebSocket(url.toString()); 343 + } catch (err) { 344 + fail(err); 345 + return; 346 + } 347 + 348 + timeout = setTimeout(() => { 349 + try { 350 + ws.close(); 351 + } catch (err) { 352 + fail(err); 353 + } 354 + }, STREAM_DURATION_MS); 355 + 295 356 ws.addEventListener('message', (event) => { 296 357 const task = (async () => { 297 358 let msg; ··· 305 366 return; 306 367 } 307 368 308 - if (msg.time_us) { 309 - latestCursor = String(msg.time_us); 369 + if (msg.time_us != null) { 370 + const timeUs = Number(msg.time_us); 371 + if (Number.isFinite(timeUs) && (observedCursor === null || timeUs > Number(observedCursor))) { 372 + observedCursor = String(msg.time_us); 373 + } 310 374 } 311 375 312 376 if (msg.did) { ··· 332 396 })(); 333 397 334 398 pending.add(task); 335 - task.finally(() => pending.delete(task)); 399 + task.catch(fail); 400 + task.finally(() => pending.delete(task)).catch(() => {}); 336 401 }); 337 402 338 403 ws.addEventListener('close', () => { 339 404 void finish(); 340 405 }); 341 406 342 - ws.addEventListener('error', () => { 343 - ws.close(); 407 + ws.addEventListener('error', (event) => { 408 + fail(event?.error ?? event); 344 409 }); 345 410 }); 346 411 }
+293
test/explore-cursor.test.js
··· 1 + // SPDX-License-Identifier: MIT 2 + // Copyright (c) 2026 sol pbc 3 + 4 + import { describe, expect, test } from 'bun:test'; 5 + import { readFileSync } from 'node:fs'; 6 + import { join } from 'node:path'; 7 + import { Database } from 'bun:sqlite'; 8 + import { runScheduled } from '../explore/src/index.js'; 9 + import { processCapEvent, processSkillEvent, processVouchEvent } from '../explore/src/jetstream.js'; 10 + 11 + function createCursorEnv({ 12 + stored = '', 13 + getStatus = 200, 14 + putStatus = 200, 15 + getThrows = null, 16 + putThrows = null, 17 + } = {}) { 18 + const store = { 19 + idNames: [], 20 + ids: [], 21 + gets: 0, 22 + puts: 0, 23 + putBodies: [], 24 + idFromName(name) { 25 + this.idNames.push(name); 26 + return { name }; 27 + }, 28 + get(id) { 29 + this.ids.push(id); 30 + return { 31 + fetch: async (url, { method, body } = {}) => { 32 + if (method === 'GET') { 33 + store.gets++; 34 + if (getThrows) { 35 + throw getThrows; 36 + } 37 + return new Response(stored, { status: getStatus }); 38 + } 39 + 40 + if (method === 'PUT') { 41 + store.puts++; 42 + store.putBodies.push(String(body)); 43 + if (putThrows) { 44 + throw putThrows; 45 + } 46 + return new Response('ok', { status: putStatus }); 47 + } 48 + 49 + return new Response('method not allowed', { status: 405 }); 50 + }, 51 + }; 52 + }, 53 + }; 54 + 55 + return { env: { CURSOR_STORE: store }, store }; 56 + } 57 + 58 + function createStreamReader({ result = { observedCursor: null }, error = null } = {}) { 59 + const calls = []; 60 + const streamReader = async (env, startCursor) => { 61 + calls.push({ env, startCursor }); 62 + if (error) { 63 + throw error; 64 + } 65 + return result; 66 + }; 67 + streamReader.calls = calls; 68 + return streamReader; 69 + } 70 + 71 + function d1Statement(db, sql, args = []) { 72 + return { 73 + sql, 74 + args, 75 + bind(...nextArgs) { 76 + return d1Statement(db, sql, nextArgs); 77 + }, 78 + first() { 79 + return db.query(sql).get(...args) ?? null; 80 + }, 81 + run() { 82 + return db.query(sql).run(...args); 83 + }, 84 + all() { 85 + return { results: db.query(sql).all(...args) }; 86 + }, 87 + }; 88 + } 89 + 90 + function createD1(db) { 91 + const executeBatch = db.transaction((statements) => statements.map((stmt) => { 92 + if (/^\s*select\b/i.test(stmt.sql)) { 93 + return stmt.all(); 94 + } 95 + return stmt.run(); 96 + })); 97 + 98 + return { 99 + prepare(sql) { 100 + return d1Statement(db, sql); 101 + }, 102 + batch(statements) { 103 + return executeBatch(statements); 104 + }, 105 + }; 106 + } 107 + 108 + function createSqliteEnv() { 109 + const db = new Database(':memory:'); 110 + const schemaPath = join(import.meta.dir, '..', 'explore', 'schema.sql'); 111 + db.exec(readFileSync(schemaPath, 'utf8')); 112 + return { db, env: { DB: createD1(db) } }; 113 + } 114 + 115 + describe('explore scheduled cursor', () => { 116 + test('passes stored valid cursor and writes newer observed cursor', async () => { 117 + const { env, store } = createCursorEnv({ stored: '12345' }); 118 + const streamReader = createStreamReader({ result: { observedCursor: '12399' } }); 119 + 120 + await runScheduled(env, { streamReader, now: () => 999 }); 121 + 122 + expect(store.idNames).toEqual(['jetstream', 'jetstream']); 123 + expect(streamReader.calls.length).toBe(1); 124 + expect(streamReader.calls[0].env).toBe(env); 125 + expect(streamReader.calls[0].startCursor).toBe(12345); 126 + expect(store.putBodies).toEqual(['12399']); 127 + }); 128 + 129 + test('uses startup replay window when no cursor is stored', async () => { 130 + const withEvent = createCursorEnv({ stored: '' }); 131 + const eventReader = createStreamReader({ result: { observedCursor: '300000123' } }); 132 + 133 + await runScheduled(withEvent.env, { streamReader: eventReader, now: () => 3_000_000 }); 134 + 135 + expect(eventReader.calls[0].startCursor).toBe(300_000_000); 136 + expect(withEvent.store.putBodies).toEqual(['300000123']); 137 + 138 + const quiet = createCursorEnv({ stored: '' }); 139 + const quietReader = createStreamReader({ result: { observedCursor: null } }); 140 + 141 + await runScheduled(quiet.env, { streamReader: quietReader, now: () => 3_000_000 }); 142 + 143 + expect(quietReader.calls[0].startCursor).toBe(300_000_000); 144 + expect(quiet.store.putBodies).toEqual(['3000000000']); 145 + }); 146 + 147 + test('writes window-open cursor for a quiet run', async () => { 148 + const { env, store } = createCursorEnv({ stored: '1000' }); 149 + const streamReader = createStreamReader({ result: { observedCursor: null } }); 150 + 151 + await runScheduled(env, { streamReader, now: () => 2 }); 152 + 153 + expect(streamReader.calls[0].startCursor).toBe(1000); 154 + expect(store.putBodies).toEqual(['2000']); 155 + }); 156 + 157 + test('cursor write failure rejects for event and quiet-window cursors', async () => { 158 + const withEvent = createCursorEnv({ stored: '1000', putStatus: 500 }); 159 + const eventReader = createStreamReader({ result: { observedCursor: '1500' } }); 160 + 161 + await expect(runScheduled(withEvent.env, { streamReader: eventReader, now: () => 2 })) 162 + .rejects.toThrow('cursor write failed: 500'); 163 + expect(withEvent.store.putBodies).toEqual(['1500']); 164 + 165 + const quiet = createCursorEnv({ stored: '1000', putStatus: 500 }); 166 + const quietReader = createStreamReader({ result: { observedCursor: null } }); 167 + 168 + await expect(runScheduled(quiet.env, { streamReader: quietReader, now: () => 2 })) 169 + .rejects.toThrow('cursor write failed: 500'); 170 + expect(quiet.store.putBodies).toEqual(['2000']); 171 + }); 172 + 173 + test('cursor read failure rejects before streaming', async () => { 174 + const cases = [ 175 + { fake: createCursorEnv({ stored: '1000', getStatus: 500 }), message: 'cursor read failed: 500' }, 176 + { fake: createCursorEnv({ stored: '1000', getThrows: new Error('read exploded') }), message: 'read exploded' }, 177 + ]; 178 + 179 + for (const { fake, message } of cases) { 180 + const streamReader = createStreamReader(); 181 + 182 + await expect(runScheduled(fake.env, { streamReader, now: () => 2 })) 183 + .rejects.toThrow(message); 184 + expect(streamReader.calls.length).toBe(0); 185 + expect(fake.store.putBodies).toEqual([]); 186 + } 187 + }); 188 + 189 + test('malformed stored cursors reject before streaming', async () => { 190 + for (const stored of ['abc', '-5', '0', '1.5', ' ']) { 191 + const { env, store } = createCursorEnv({ stored }); 192 + const streamReader = createStreamReader(); 193 + 194 + await expect(runScheduled(env, { streamReader, now: () => 2 })) 195 + .rejects.toThrow('malformed cursor: ' + JSON.stringify(stored)); 196 + expect(streamReader.calls.length).toBe(0); 197 + expect(store.putBodies).toEqual([]); 198 + } 199 + }); 200 + 201 + test('stream reader rejection propagates without writing cursor', async () => { 202 + const { env, store } = createCursorEnv({ stored: '1000' }); 203 + const streamReader = createStreamReader({ error: new Error('stream failed') }); 204 + 205 + await expect(runScheduled(env, { streamReader, now: () => 2 })) 206 + .rejects.toThrow('stream failed'); 207 + expect(streamReader.calls.length).toBe(1); 208 + expect(store.putBodies).toEqual([]); 209 + }); 210 + }); 211 + 212 + describe('explore event idempotency', () => { 213 + test('processes duplicate cap, vouch, and skill creates idempotently', async () => { 214 + const { db, env } = createSqliteEnv(); 215 + 216 + const capDid = 'did:plc:capauthor'; 217 + const capCommit = { 218 + operation: 'create', 219 + rkey: '3lcap', 220 + cid: 'bafycap', 221 + record: { 222 + title: 'Idempotent Cap', 223 + description: 'A cap inserted twice for testing.', 224 + ref: 'idempotent-cap-test', 225 + beacon: 'vit:example/repo', 226 + kind: 'test', 227 + createdAt: '2026-07-07T00:00:00.000Z', 228 + }, 229 + }; 230 + 231 + await processCapEvent(env, capDid, capCommit); 232 + const capAfterOne = db 233 + .query("SELECT (SELECT COUNT(*) FROM caps) AS rows, (SELECT cap_count FROM beacons WHERE name = ?) AS count") 234 + .get('vit:example/repo'); 235 + await processCapEvent(env, capDid, capCommit); 236 + const capAfterTwo = db 237 + .query("SELECT (SELECT COUNT(*) FROM caps) AS rows, (SELECT cap_count FROM beacons WHERE name = ?) AS count") 238 + .get('vit:example/repo'); 239 + 240 + expect(capAfterOne).toEqual({ rows: 1, count: 1 }); 241 + expect(capAfterTwo).toEqual(capAfterOne); 242 + 243 + const vouchDid = 'did:plc:vouchauthor'; 244 + const vouchCommit = { 245 + operation: 'create', 246 + rkey: '3lvouch', 247 + cid: 'bafyvouch', 248 + record: { 249 + subject: { uri: 'at://did:plc:capauthor/org.v-it.cap/3lcap' }, 250 + ref: 'idempotent-cap-test', 251 + beacon: 'vit:example/repo', 252 + kind: 'want', 253 + createdAt: '2026-07-07T00:00:01.000Z', 254 + }, 255 + }; 256 + 257 + await processVouchEvent(env, vouchDid, vouchCommit); 258 + const vouchAfterOne = db 259 + .query("SELECT (SELECT COUNT(*) FROM vouches) AS rows, (SELECT vouch_count FROM beacons WHERE name = ?) AS count") 260 + .get('vit:example/repo'); 261 + await processVouchEvent(env, vouchDid, vouchCommit); 262 + const vouchAfterTwo = db 263 + .query("SELECT (SELECT COUNT(*) FROM vouches) AS rows, (SELECT vouch_count FROM beacons WHERE name = ?) AS count") 264 + .get('vit:example/repo'); 265 + 266 + expect(vouchAfterOne).toEqual({ rows: 1, count: 1 }); 267 + expect(vouchAfterTwo).toEqual(vouchAfterOne); 268 + 269 + const skillDid = 'did:plc:skillauthor'; 270 + const skillCommit = { 271 + operation: 'create', 272 + rkey: '3lskill', 273 + cid: 'bafyskill', 274 + record: { 275 + name: 'idempotent-skill', 276 + description: 'A skill inserted twice for testing.', 277 + version: '1.0.0', 278 + tags: ['test'], 279 + createdAt: '2026-07-07T00:00:02.000Z', 280 + }, 281 + }; 282 + 283 + await processSkillEvent(env, skillDid, skillCommit); 284 + const skillAfterOne = db.query('SELECT COUNT(*) AS rows FROM skills').get(); 285 + await processSkillEvent(env, skillDid, skillCommit); 286 + const skillAfterTwo = db.query('SELECT COUNT(*) AS rows FROM skills').get(); 287 + 288 + expect(skillAfterOne).toEqual({ rows: 1 }); 289 + expect(skillAfterTwo).toEqual(skillAfterOne); 290 + 291 + db.close(); 292 + }); 293 + });