diff --git a/.gitignore b/.gitignore index cebc80d..2465a70 100644 --- a/.gitignore +++ b/.gitignore @@ -10,3 +10,7 @@ bun.lock .agents/ .claude/skills/ skills-lock.json + +# Explore worker +explore/node_modules/ +explore/.wrangler/ diff --git a/explore/Makefile b/explore/Makefile new file mode 100644 index 0000000..9af4e9d --- /dev/null +++ b/explore/Makefile @@ -0,0 +1,10 @@ +.PHONY: dev deploy schema + +dev: + npx wrangler dev + +deploy: + npx wrangler deploy + +schema: + npx wrangler d1 execute vit-explore --file=schema.sql diff --git a/explore/package.json b/explore/package.json new file mode 100644 index 0000000..9119ece --- /dev/null +++ b/explore/package.json @@ -0,0 +1,13 @@ +{ + "name": "vit-explore", + "version": "0.1.0", + "private": true, + "type": "module", + "scripts": { + "dev": "npx wrangler dev", + "deploy": "npx wrangler deploy" + }, + "devDependencies": { + "wrangler": "^4" + } +} diff --git a/explore/public/.gitkeep b/explore/public/.gitkeep new file mode 100644 index 0000000..8b13789 --- /dev/null +++ b/explore/public/.gitkeep @@ -0,0 +1 @@ + diff --git a/explore/schema.sql b/explore/schema.sql new file mode 100644 index 0000000..8e9c30d --- /dev/null +++ b/explore/schema.sql @@ -0,0 +1,50 @@ +CREATE TABLE IF NOT EXISTS caps ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + did TEXT NOT NULL, + rkey TEXT NOT NULL, + uri TEXT NOT NULL UNIQUE, + cid TEXT, + title TEXT NOT NULL, + description TEXT, + ref TEXT NOT NULL, + beacon TEXT, + record_json TEXT NOT NULL, + created_at TEXT NOT NULL, + indexed_at TEXT NOT NULL DEFAULT (datetime('now')), + UNIQUE(did, rkey) +); + +CREATE INDEX IF NOT EXISTS idx_caps_beacon ON caps(beacon); +CREATE INDEX IF NOT EXISTS idx_caps_created_at ON caps(created_at DESC); + +CREATE TABLE IF NOT EXISTS vouches ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + did TEXT NOT NULL, + rkey TEXT NOT NULL, + uri TEXT NOT NULL UNIQUE, + cid TEXT, + cap_uri TEXT NOT NULL, + ref TEXT NOT NULL, + beacon TEXT NOT NULL, + record_json TEXT NOT NULL, + created_at TEXT NOT NULL, + indexed_at TEXT NOT NULL DEFAULT (datetime('now')), + UNIQUE(did, rkey) +); + +CREATE INDEX IF NOT EXISTS idx_vouches_cap_uri ON vouches(cap_uri); +CREATE INDEX IF NOT EXISTS idx_vouches_beacon ON vouches(beacon); + +CREATE TABLE IF NOT EXISTS beacons ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + name TEXT NOT NULL UNIQUE, + cap_count INTEGER NOT NULL DEFAULT 0, + vouch_count INTEGER NOT NULL DEFAULT 0, + last_activity TEXT NOT NULL DEFAULT (datetime('now')) +); + +CREATE TABLE IF NOT EXISTS handles ( + did TEXT PRIMARY KEY, + handle TEXT NOT NULL, + fetched_at TEXT NOT NULL DEFAULT (datetime('now')) +); diff --git a/explore/src/api.js b/explore/src/api.js new file mode 100644 index 0000000..fa79ffa --- /dev/null +++ b/explore/src/api.js @@ -0,0 +1,119 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +const CORS_HEADERS = { + 'Access-Control-Allow-Origin': '*', + 'Access-Control-Allow-Methods': 'GET, OPTIONS', + 'Access-Control-Allow-Headers': 'Content-Type', +}; + +function json(data, status = 200) { + return new Response(JSON.stringify(data), { + status, + headers: { 'Content-Type': 'application/json', ...CORS_HEADERS }, + }); +} + +function parseLimit(value) { + const parsed = Number.parseInt(value ?? '50', 10); + if (!Number.isFinite(parsed) || parsed <= 0) { + return 50; + } + return Math.min(parsed, 100); +} + +function parseCursor(value) { + if (!value) { + return null; + } + + const parsed = Number.parseInt(value, 10); + return Number.isFinite(parsed) && parsed > 0 ? parsed : null; +} + +export async function handleRequest(request, env) { + if (request.method === 'OPTIONS') { + return new Response(null, { status: 204, headers: CORS_HEADERS }); + } + + const url = new URL(request.url); + const { pathname, searchParams } = url; + + if (request.method !== 'GET') { + return json({ error: 'method not allowed' }, 405); + } + + if (pathname === '/api/caps') { + const cursor = parseCursor(searchParams.get('cursor')); + const limit = parseLimit(searchParams.get('limit')); + const beacon = searchParams.get('beacon'); + + const conditions = []; + const bindings = []; + + if (beacon) { + conditions.push('c.beacon = ?'); + bindings.push(beacon); + } + + if (cursor) { + conditions.push('c.id < ?'); + bindings.push(cursor); + } + + let sql = 'SELECT c.*, h.handle FROM caps c LEFT JOIN handles h ON c.did = h.did'; + if (conditions.length > 0) { + sql += ` WHERE ${conditions.join(' AND ')}`; + } + sql += ' ORDER BY c.id DESC LIMIT ?'; + bindings.push(limit); + + const { results } = await env.DB.prepare(sql).bind(...bindings).all(); + return json({ + caps: results, + cursor: results.length > 0 ? results[results.length - 1].id : null, + }); + } + + if (pathname === '/api/vouches') { + const capUri = searchParams.get('cap_uri'); + if (!capUri) { + return json({ error: 'cap_uri is required' }, 400); + } + + const { results } = await env.DB.prepare( + `SELECT v.*, h.handle + FROM vouches v + LEFT JOIN handles h ON v.did = h.did + WHERE v.cap_uri = ? + ORDER BY v.id DESC`, + ) + .bind(capUri) + .all(); + + return json({ vouches: results }); + } + + if (pathname === '/api/beacons') { + const { results } = await env.DB.prepare('SELECT * FROM beacons ORDER BY last_activity DESC').all(); + return json({ beacons: results }); + } + + if (pathname === '/api/stats') { + const [caps, vouches, beacons, dids] = await env.DB.batch([ + env.DB.prepare('SELECT COUNT(*) as count FROM caps'), + env.DB.prepare('SELECT COUNT(*) as count FROM vouches'), + env.DB.prepare('SELECT COUNT(*) as count FROM beacons'), + env.DB.prepare('SELECT COUNT(DISTINCT did) as count FROM caps'), + ]); + + return json({ + total_caps: caps.results[0]?.count ?? 0, + total_vouches: vouches.results[0]?.count ?? 0, + total_beacons: beacons.results[0]?.count ?? 0, + active_dids: dids.results[0]?.count ?? 0, + }); + } + + return json({ error: 'not found' }, 404); +} diff --git a/explore/src/cursor.js b/explore/src/cursor.js new file mode 100644 index 0000000..33d9a15 --- /dev/null +++ b/explore/src/cursor.js @@ -0,0 +1,22 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +export class CursorStore { + constructor(state, env) { + this.state = state; + } + + async fetch(request) { + if (request.method === 'GET') { + return new Response((await this.state.storage.get('cursor')) || ''); + } + + if (request.method === 'PUT') { + const text = await request.text(); + await this.state.storage.put('cursor', text); + return new Response('ok'); + } + + return new Response('method not allowed', { status: 405 }); + } +} diff --git a/explore/src/index.js b/explore/src/index.js new file mode 100644 index 0000000..f6672e4 --- /dev/null +++ b/explore/src/index.js @@ -0,0 +1,39 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +import { handleRequest } from './api.js'; +import { streamEvents } from './jetstream.js'; + +export { CursorStore } from './cursor.js'; + +async function getCursor(env) { + const id = env.CURSOR_STORE.idFromName('cursor'); + const stub = env.CURSOR_STORE.get(id); + const res = await stub.fetch('http://cursor/'); + return (await res.text()) || null; +} + +async function saveCursor(env, cursor) { + const id = env.CURSOR_STORE.idFromName('cursor'); + const stub = env.CURSOR_STORE.get(id); + await stub.fetch('http://cursor/', { method: 'PUT', body: cursor }); +} + +export default { + async fetch(request, env) { + const url = new URL(request.url); + if (url.pathname.startsWith('/api/')) { + return handleRequest(request, env); + } + + return new Response('not found', { status: 404 }); + }, + + async scheduled(event, env, ctx) { + const cursor = await getCursor(env); + const result = await streamEvents(env, cursor); + if (result.latestCursor) { + await saveCursor(env, result.latestCursor); + } + }, +}; diff --git a/explore/src/jetstream.js b/explore/src/jetstream.js new file mode 100644 index 0000000..96787cf --- /dev/null +++ b/explore/src/jetstream.js @@ -0,0 +1,269 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +import { resolveHandles } from './resolve.js'; + +const CAP_COLLECTION = 'org.v-it.cap'; +const VOUCH_COLLECTION = 'org.v-it.vouch'; +const JETSTREAM_URL = 'wss://jetstream2.us-east.bsky.network/subscribe'; +const STREAM_DURATION_MS = 25_000; + +function beaconValue(value) { + return typeof value === 'string' && value.length > 0 ? value : null; +} + +function incrementCapBeaconStatements(env, beacon) { + return [ + env.DB.prepare( + `INSERT INTO beacons (name, cap_count, last_activity) + VALUES (?, 1, datetime('now')) + ON CONFLICT(name) DO UPDATE SET + cap_count = cap_count + 1, + last_activity = datetime('now')`, + ).bind(beacon), + ]; +} + +function incrementVouchBeaconStatements(env, beacon) { + return [ + env.DB.prepare( + `INSERT INTO beacons (name, vouch_count, last_activity) + VALUES (?, 1, datetime('now')) + ON CONFLICT(name) DO UPDATE SET + vouch_count = vouch_count + 1, + last_activity = datetime('now')`, + ).bind(beacon), + ]; +} + +function decrementCapBeaconStatement(env, beacon) { + return env.DB.prepare( + `UPDATE beacons + SET cap_count = MAX(0, cap_count - 1), + last_activity = datetime('now') + WHERE name = ?`, + ).bind(beacon); +} + +function decrementVouchBeaconStatement(env, beacon) { + return env.DB.prepare( + `UPDATE beacons + SET vouch_count = MAX(0, vouch_count - 1), + last_activity = datetime('now') + WHERE name = ?`, + ).bind(beacon); +} + +async function processCapEvent(env, did, commit) { + const { operation, rkey, record, cid } = commit; + const uri = `at://${did}/${CAP_COLLECTION}/${rkey}`; + + if (operation === 'create' || operation === 'update') { + const nextBeacon = beaconValue(record?.beacon); + const existing = await env.DB.prepare('SELECT beacon FROM caps WHERE did = ? AND rkey = ?') + .bind(did, rkey) + .first(); + const prevBeacon = beaconValue(existing?.beacon); + + const stmts = [ + env.DB.prepare( + `INSERT INTO caps (did, rkey, uri, cid, title, description, ref, beacon, record_json, created_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(did, rkey) DO UPDATE SET + cid = excluded.cid, + title = excluded.title, + description = excluded.description, + ref = excluded.ref, + beacon = excluded.beacon, + record_json = excluded.record_json, + created_at = excluded.created_at`, + ).bind( + did, + rkey, + uri, + cid ?? null, + record.title, + record.description || '', + record.ref, + nextBeacon, + JSON.stringify(record), + record.createdAt, + ), + ]; + + if (!existing && nextBeacon) { + stmts.push(...incrementCapBeaconStatements(env, nextBeacon)); + } else if (existing && prevBeacon !== nextBeacon) { + if (prevBeacon) { + stmts.push(decrementCapBeaconStatement(env, prevBeacon)); + } + if (nextBeacon) { + stmts.push(...incrementCapBeaconStatements(env, nextBeacon)); + } + } + + await env.DB.batch(stmts); + return; + } + + if (operation === 'delete') { + const existing = await env.DB.prepare('SELECT beacon FROM caps WHERE did = ? AND rkey = ?') + .bind(did, rkey) + .first(); + + const stmts = [ + env.DB.prepare('DELETE FROM caps WHERE did = ? AND rkey = ?').bind(did, rkey), + ]; + + const prevBeacon = beaconValue(existing?.beacon); + if (prevBeacon) { + stmts.unshift(decrementCapBeaconStatement(env, prevBeacon)); + } + + await env.DB.batch(stmts); + } +} + +async function processVouchEvent(env, did, commit) { + const { operation, rkey, record, cid } = commit; + const uri = `at://${did}/${VOUCH_COLLECTION}/${rkey}`; + + if (operation === 'create' || operation === 'update') { + const nextBeacon = beaconValue(record?.beacon); + const existing = await env.DB.prepare('SELECT beacon FROM vouches WHERE did = ? AND rkey = ?') + .bind(did, rkey) + .first(); + const prevBeacon = beaconValue(existing?.beacon); + + const stmts = [ + env.DB.prepare( + `INSERT INTO vouches (did, rkey, uri, cid, cap_uri, ref, beacon, record_json, created_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(did, rkey) DO UPDATE SET + cid = excluded.cid, + cap_uri = excluded.cap_uri, + ref = excluded.ref, + beacon = excluded.beacon, + record_json = excluded.record_json, + created_at = excluded.created_at`, + ).bind( + did, + rkey, + uri, + cid ?? null, + record.subject?.uri, + record.ref, + record.beacon, + JSON.stringify(record), + record.createdAt, + ), + ]; + + if (!existing && nextBeacon) { + stmts.push(...incrementVouchBeaconStatements(env, nextBeacon)); + } else if (existing && prevBeacon !== nextBeacon) { + if (prevBeacon) { + stmts.push(decrementVouchBeaconStatement(env, prevBeacon)); + } + if (nextBeacon) { + stmts.push(...incrementVouchBeaconStatements(env, nextBeacon)); + } + } + + await env.DB.batch(stmts); + return; + } + + if (operation === 'delete') { + const existing = await env.DB.prepare('SELECT beacon FROM vouches WHERE did = ? AND rkey = ?') + .bind(did, rkey) + .first(); + + const stmts = [ + env.DB.prepare('DELETE FROM vouches WHERE did = ? AND rkey = ?').bind(did, rkey), + ]; + + const prevBeacon = beaconValue(existing?.beacon); + if (prevBeacon) { + stmts.unshift(decrementVouchBeaconStatement(env, prevBeacon)); + } + + await env.DB.batch(stmts); + } +} + +export async function streamEvents(env, cursor) { + const url = new URL(JETSTREAM_URL); + url.searchParams.append('wantedCollections', CAP_COLLECTION); + url.searchParams.append('wantedCollections', VOUCH_COLLECTION); + if (cursor) { + url.searchParams.set('cursor', cursor); + } + + return await new Promise((resolve) => { + let latestCursor = cursor || null; + const newDids = new Set(); + const pending = new Set(); + const ws = new WebSocket(url.toString()); + + const timeout = setTimeout(() => { + ws.close(); + }, STREAM_DURATION_MS); + + const finish = async () => { + clearTimeout(timeout); + if (pending.size > 0) { + await Promise.allSettled([...pending]); + } + if (newDids.size > 0) { + await resolveHandles([...newDids], env); + } + resolve({ latestCursor }); + }; + + ws.addEventListener('message', (event) => { + const task = (async () => { + let msg; + try { + msg = JSON.parse(event.data); + } catch { + return; + } + + if (msg.kind !== 'commit') { + return; + } + + if (msg.time_us) { + latestCursor = String(msg.time_us); + } + + if (msg.did) { + newDids.add(msg.did); + } + + const commit = msg.commit; + if (!commit) { + return; + } + + if (commit.collection === CAP_COLLECTION) { + await processCapEvent(env, msg.did, commit); + } else if (commit.collection === VOUCH_COLLECTION) { + await processVouchEvent(env, msg.did, commit); + } + })(); + + pending.add(task); + task.finally(() => pending.delete(task)); + }); + + ws.addEventListener('close', () => { + void finish(); + }); + + ws.addEventListener('error', () => { + ws.close(); + }); + }); +} diff --git a/explore/src/resolve.js b/explore/src/resolve.js new file mode 100644 index 0000000..aba42c9 --- /dev/null +++ b/explore/src/resolve.js @@ -0,0 +1,57 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +export async function resolveHandle(did, env) { + try { + const cached = await env.DB.prepare( + "SELECT handle, fetched_at FROM handles WHERE did = ? AND fetched_at > datetime('now', '-24 hours')", + ) + .bind(did) + .first(); + + if (cached?.handle) { + return cached.handle; + } + + const res = await fetch(`https://plc.directory/${did}`); + if (!res.ok) { + return null; + } + + const data = await res.json(); + const handle = Array.isArray(data?.alsoKnownAs) + ? data.alsoKnownAs.find(value => typeof value === 'string' && value.startsWith('at://'))?.slice(5) ?? null + : null; + + if (!handle) { + return null; + } + + await env.DB.prepare( + `INSERT INTO handles (did, handle, fetched_at) + VALUES (?, ?, datetime('now')) + ON CONFLICT(did) DO UPDATE SET + handle = excluded.handle, + fetched_at = excluded.fetched_at`, + ) + .bind(did, handle) + .run(); + + return handle; + } catch { + return null; + } +} + +export async function resolveHandles(dids, env) { + const handles = new Map(); + + for (const did of dids) { + const handle = await resolveHandle(did, env); + if (handle) { + handles.set(did, handle); + } + } + + return handles; +} diff --git a/explore/wrangler.toml b/explore/wrangler.toml new file mode 100644 index 0000000..864ffd2 --- /dev/null +++ b/explore/wrangler.toml @@ -0,0 +1,28 @@ +name = "vit-explore" +main = "src/index.js" +compatibility_date = "2024-12-01" +account_id = "3f2c1528c7d4d9685819ea9e9e307c92" + +[assets] +directory = "public" + +[[d1_databases]] +binding = "DB" +database_name = "vit-explore" +database_id = "placeholder" + +[durable_objects] +bindings = [ + { name = "CURSOR_STORE", class_name = "CursorStore" } +] + +[[migrations]] +tag = "v1" +new_classes = ["CursorStore"] + +[triggers] +crons = ["*/1 * * * *"] + +[[routes]] +pattern = "explore.v-it.org" +custom_domain = true