From fa68d0297f306668c11d74396ed7374387b9318c Mon Sep 17 00:00:00 2001 From: Steve Date: Sat, 11 Apr 2026 22:33:16 -0400 Subject: [PATCH] chore: refactor to jetstream durable object --- .../src/durable-objects/jetstream-consumer.ts | 205 ++++++++++++++++++ packages/server/src/index.ts | 59 ++++- packages/server/src/routes/index.ts | 1 + packages/server/src/routes/jetstream.ts | 29 +++ packages/server/src/types/index.ts | 41 ++++ packages/server/src/utils/index.ts | 1 + packages/server/src/utils/ingest.ts | 101 +++++++++ packages/server/worker-configuration.d.ts | 7 + packages/server/wrangler.toml | 10 + 9 files changed, 452 insertions(+), 2 deletions(-) create mode 100644 packages/server/src/durable-objects/jetstream-consumer.ts create mode 100644 packages/server/src/routes/jetstream.ts create mode 100644 packages/server/src/utils/ingest.ts create mode 100644 packages/server/worker-configuration.d.ts diff --git a/packages/server/src/durable-objects/jetstream-consumer.ts b/packages/server/src/durable-objects/jetstream-consumer.ts new file mode 100644 index 0000000..2908284 --- /dev/null +++ b/packages/server/src/durable-objects/jetstream-consumer.ts @@ -0,0 +1,205 @@ +import type { Bindings, JetstreamEvent } from "../types"; +import { resolvePds } from "../utils/resolver"; +import { ingestDocument, deleteDocument } from "../utils/ingest"; + +const JETSTREAM_URL = + "wss://jetstream2.us-east.bsky.network/subscribe?wantedCollections=site.standard.document"; +const ALARM_INTERVAL_MS = 30_000; +const CURSOR_SAVE_INTERVAL_MS = 10_000; +const CURSOR_SAVE_MESSAGE_COUNT = 100; + +export class JetstreamConsumer implements DurableObject { + private state: DurableObjectState; + private env: Bindings; + private ws: WebSocket | null = null; + private cursor: string | null = null; + private messageCount = 0; + private lastCursorSave = 0; + private pdsCache: Map = new Map(); + + constructor(state: DurableObjectState, env: Bindings) { + this.state = state; + this.env = env; + } + + async fetch(request: Request): Promise { + const url = new URL(request.url); + + switch (url.pathname) { + case "/start": + return this.handleStart(); + case "/status": + return this.handleStatus(); + case "/stop": + return this.handleStop(); + default: + return new Response("Not found", { status: 404 }); + } + } + + async alarm(): Promise { + if (!this.ws || this.ws.readyState !== WebSocket.OPEN) { + console.log("Alarm: WebSocket not connected, reconnecting..."); + await this.connect(); + } + // Reschedule alarm + await this.state.storage.setAlarm(Date.now() + ALARM_INTERVAL_MS); + } + + private async handleStart(): Promise { + if (this.ws && this.ws.readyState === WebSocket.OPEN) { + return Response.json({ + status: "already_connected", + cursor: this.cursor, + }); + } + + await this.connect(); + await this.state.storage.setAlarm(Date.now() + ALARM_INTERVAL_MS); + + return Response.json({ status: "started", cursor: this.cursor }); + } + + private handleStatus(): Response { + return Response.json({ + connected: this.ws?.readyState === WebSocket.OPEN, + cursor: this.cursor, + messageCount: this.messageCount, + }); + } + + private async handleStop(): Promise { + if (this.ws) { + this.ws.close(); + this.ws = null; + } + await this.state.storage.deleteAlarm(); + return Response.json({ status: "stopped" }); + } + + private async connect(): Promise { + // Close existing connection if any + if (this.ws) { + try { + this.ws.close(); + } catch {} + this.ws = null; + } + + // Load cursor from storage + if (!this.cursor) { + this.cursor = + (await this.state.storage.get("cursor")) || null; + } + + const url = this.cursor + ? `${JETSTREAM_URL}&cursor=${this.cursor}` + : JETSTREAM_URL; + + console.log(`Connecting to Jetstream: ${url}`); + + try { + const ws = new WebSocket(url); + + ws.addEventListener("message", (event) => { + this.onMessage(event.data as string); + }); + + ws.addEventListener("close", () => { + console.log("Jetstream WebSocket closed"); + this.ws = null; + // Alarm will handle reconnection + }); + + ws.addEventListener("error", (event) => { + console.error("Jetstream WebSocket error:", event); + try { + ws.close(); + } catch {} + this.ws = null; + }); + + this.ws = ws; + } catch (error) { + console.error("Failed to connect to Jetstream:", error); + this.ws = null; + } + } + + private async onMessage(data: string): Promise { + try { + const event = JSON.parse(data) as JetstreamEvent; + + // Update cursor from every event + this.cursor = String(event.time_us); + + // Only process commit events + if (event.kind !== "commit") return; + + const { commit } = event; + if (commit.collection !== "site.standard.document") return; + + // PDS filter: skip bridgy noise + const pds = await this.resolvePdsWithCache(event.did); + if (!pds || pds.includes("brid.gy")) return; + + if ( + commit.operation === "create" || + commit.operation === "update" + ) { + await ingestDocument(this.env.DB, this.env.RESOLUTION_QUEUE, { + did: event.did, + rkey: commit.rkey, + collection: commit.collection, + cid: commit.cid, + record: commit.record, + }); + } else if (commit.operation === "delete") { + await deleteDocument(this.env.DB, { + did: event.did, + collection: commit.collection, + rkey: commit.rkey, + }); + } + + this.messageCount++; + + // Periodically save cursor + await this.maybeSaveCursor(); + } catch (error) { + console.error("Error processing Jetstream message:", error); + } + } + + private async resolvePdsWithCache(did: string): Promise { + // Fast in-memory cache (cleared on DO eviction) + if (this.pdsCache.has(did)) { + return this.pdsCache.get(did)!; + } + + const pds = await resolvePds(this.env.DB, did); + this.pdsCache.set(did, pds); + + // Keep in-memory cache bounded + if (this.pdsCache.size > 10_000) { + const firstKey = this.pdsCache.keys().next().value; + if (firstKey) this.pdsCache.delete(firstKey); + } + + return pds; + } + + private async maybeSaveCursor(): Promise { + if (!this.cursor) return; + + const now = Date.now(); + const shouldSave = + this.messageCount % CURSOR_SAVE_MESSAGE_COUNT === 0 || + now - this.lastCursorSave > CURSOR_SAVE_INTERVAL_MS; + + if (shouldSave) { + await this.state.storage.put("cursor", this.cursor); + this.lastCursorSave = now; + } + } +} diff --git a/packages/server/src/index.ts b/packages/server/src/index.ts index 6554866..8350151 100644 --- a/packages/server/src/index.ts +++ b/packages/server/src/index.ts @@ -1,7 +1,16 @@ import { Hono } from "hono"; import { cors } from "hono/cors"; import type { Bindings } from "./types"; -import { health, webhook, feed, stats, records, admin, rss } from "./routes"; +import { + health, + webhook, + feed, + stats, + records, + admin, + rss, + jetstream, +} from "./routes"; import { processDocument } from "./utils"; const app = new Hono<{ Bindings: Bindings }>(); @@ -17,16 +26,20 @@ app.route("/stats", stats); app.route("/records", records); app.route("/admin", admin); app.route("/rss.xml", rss); +app.route("/jetstream", jetstream); // 404 handler app.notFound((c) => { return c.json({ error: "Not found" }, 404); }); +// Export Durable Object class +export { JetstreamConsumer } from "./durable-objects/jetstream-consumer"; + // Export for Cloudflare Workers export default { fetch: app.fetch, - async scheduled(event: ScheduledEvent, env: Bindings, ctx: ExecutionContext) { + async scheduled(_event: ScheduledEvent, env: Bindings, _ctx: ExecutionContext) { const batchSize = 50; // Select stale documents const { results } = await env.DB.prepare( @@ -50,6 +63,48 @@ export default { await env.RESOLUTION_QUEUE.sendBatch(messages); console.log(`Queued ${messages.length} documents for resolution`); } + + // Cleanup: keep only the 300 most recent verified documents, delete everything else + const deletedVerified = await env.DB.prepare( + `DELETE FROM resolved_documents WHERE verified = 1 AND uri NOT IN ( + SELECT uri FROM resolved_documents WHERE verified = 1 + ORDER BY published_at DESC LIMIT 300 + )`, + ).run(); + if (deletedVerified.meta.changes > 0) { + console.log(`Cleaned up ${deletedVerified.meta.changes} old verified documents`); + } + + // Delete unverified/stale documents older than 24 hours + const deletedUnverified = await env.DB.prepare( + `DELETE FROM resolved_documents WHERE (verified IS NULL OR verified = 0) + AND resolved_at < datetime('now', '-24 hours')`, + ).run(); + if (deletedUnverified.meta.changes > 0) { + console.log(`Cleaned up ${deletedUnverified.meta.changes} unverified documents`); + } + + // Clean up orphaned repo_records + if (deletedVerified.meta.changes > 0 || deletedUnverified.meta.changes > 0) { + const orphaned = await env.DB.prepare( + `DELETE FROM repo_records WHERE NOT EXISTS ( + SELECT 1 FROM resolved_documents WHERE resolved_documents.did = repo_records.did + AND resolved_documents.rkey = repo_records.rkey + )`, + ).run(); + if (orphaned.meta.changes > 0) { + console.log(`Cleaned up ${orphaned.meta.changes} orphaned repo_records`); + } + } + + // Ensure Jetstream consumer DO is alive + try { + const id = env.JETSTREAM_CONSUMER.idFromName("singleton"); + const stub = env.JETSTREAM_CONSUMER.get(id); + await stub.fetch(new Request("http://do/start")); + } catch (error) { + console.error("Failed to ping Jetstream consumer:", error); + } }, async queue(batch: MessageBatch, env: Bindings) { for (const message of batch.messages) { diff --git a/packages/server/src/routes/index.ts b/packages/server/src/routes/index.ts index c2e93a4..a071cc0 100644 --- a/packages/server/src/routes/index.ts +++ b/packages/server/src/routes/index.ts @@ -5,3 +5,4 @@ export { default as stats } from "./stats"; export { default as records } from "./records"; export { default as admin } from "./admin"; export { default as rss } from "./rss"; +export { default as jetstream } from "./jetstream"; diff --git a/packages/server/src/routes/jetstream.ts b/packages/server/src/routes/jetstream.ts new file mode 100644 index 0000000..f8543d0 --- /dev/null +++ b/packages/server/src/routes/jetstream.ts @@ -0,0 +1,29 @@ +import { Hono } from "hono"; +import type { Bindings } from "../types"; + +const jetstream = new Hono<{ Bindings: Bindings }>(); + +function getStub(env: Bindings) { + const id = env.JETSTREAM_CONSUMER.idFromName("singleton"); + return env.JETSTREAM_CONSUMER.get(id); +} + +jetstream.get("/start", async (c) => { + const stub = getStub(c.env); + const resp = await stub.fetch(new Request("http://do/start")); + return c.json(await resp.json()); +}); + +jetstream.get("/status", async (c) => { + const stub = getStub(c.env); + const resp = await stub.fetch(new Request("http://do/status")); + return c.json(await resp.json()); +}); + +jetstream.get("/stop", async (c) => { + const stub = getStub(c.env); + const resp = await stub.fetch(new Request("http://do/stop")); + return c.json(await resp.json()); +}); + +export default jetstream; diff --git a/packages/server/src/types/index.ts b/packages/server/src/types/index.ts index a712ada..f41dc79 100644 --- a/packages/server/src/types/index.ts +++ b/packages/server/src/types/index.ts @@ -2,8 +2,49 @@ export type Bindings = { DB: D1Database; RESOLUTION_QUEUE: Queue; TAP_WEBHOOK_SECRET?: string; + JETSTREAM_CONSUMER: DurableObjectNamespace; }; +// Jetstream WebSocket event types +export interface JetstreamCommitEvent { + did: string; + time_us: number; + kind: "commit"; + commit: { + rev: string; + operation: "create" | "update" | "delete"; + collection: string; + rkey: string; + record?: Record; + cid?: string; + }; +} + +export interface JetstreamIdentityEvent { + did: string; + time_us: number; + kind: "identity"; + identity: { did: string; seq: number; time: string }; +} + +export interface JetstreamAccountEvent { + did: string; + time_us: number; + kind: "account"; + account: { + active: boolean; + did: string; + seq: number; + time: string; + status?: string; + }; +} + +export type JetstreamEvent = + | JetstreamCommitEvent + | JetstreamIdentityEvent + | JetstreamAccountEvent; + export interface TapRecordEvent { id: number; type: "record"; diff --git a/packages/server/src/utils/index.ts b/packages/server/src/utils/index.ts index 1558992..9c8f189 100644 --- a/packages/server/src/utils/index.ts +++ b/packages/server/src/utils/index.ts @@ -3,3 +3,4 @@ export { resolvePds } from "./resolver"; export { resolveViewUrl, processDocument } from "./document"; export { buildBlobUrl, extractBlobCid } from "./blob"; export { verifyPublication, verifyDocument, verifyDocumentRecord } from "./verification"; +export { ingestDocument, deleteDocument } from "./ingest"; diff --git a/packages/server/src/utils/ingest.ts b/packages/server/src/utils/ingest.ts new file mode 100644 index 0000000..85e1dca --- /dev/null +++ b/packages/server/src/utils/ingest.ts @@ -0,0 +1,101 @@ +import { resolveViewUrl } from "./document"; + +const STALE_OFFSET_HOURS = 24; + +export async function ingestDocument( + db: D1Database, + queue: Queue, + params: { + did: string; + rkey: string; + collection: string; + cid?: string; + record?: Record; + }, +): Promise { + const { did, rkey, collection, cid, record } = params; + + // Upsert repo_records + await db + .prepare( + `INSERT INTO repo_records (did, rkey, collection, cid, synced_at) + VALUES (?, ?, ?, ?, datetime('now')) + ON CONFLICT(did, collection, rkey) DO UPDATE SET + cid = ?, + synced_at = datetime('now')`, + ) + .bind(did, rkey, collection, cid || null, cid || null) + .run(); + + // If we have the full record, upsert resolved_documents with initial data + if (record) { + const uri = `at://${did}/${collection}/${rkey}`; + const doc = record as { + title?: string; + path?: string; + site?: string; + content?: unknown; + textContent?: string; + publishedAt?: string; + coverImage?: unknown; + bskyPostRef?: { uri: string; cid: string }; + tags?: string[]; + }; + + let viewUrl: string | null = null; + if (doc.site && doc.path) { + viewUrl = await resolveViewUrl(db, doc.site, doc.path); + } + + await db + .prepare( + `INSERT INTO resolved_documents (uri, did, rkey, title, path, site, content, text_content, published_at, view_url, resolved_at, stale_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, datetime('now'), datetime('now', '+${STALE_OFFSET_HOURS} hours')) + ON CONFLICT(uri) DO UPDATE SET + title = ?, path = ?, site = ?, content = ?, text_content = ?, published_at = ?, view_url = ?, resolved_at = datetime('now'), stale_at = datetime('now', '+${STALE_OFFSET_HOURS} hours')`, + ) + .bind( + uri, + did, + rkey, + doc.title || null, + doc.path || null, + doc.site || null, + doc.content ? JSON.stringify(doc.content) : null, + doc.textContent || null, + doc.publishedAt || null, + viewUrl, + doc.title || null, + doc.path || null, + doc.site || null, + doc.content ? JSON.stringify(doc.content) : null, + doc.textContent || null, + doc.publishedAt || null, + viewUrl, + ) + .run(); + } + + // Queue for full resolution (verification, publication lookup, etc.) + await queue.send({ did, collection, rkey }); +} + +export async function deleteDocument( + db: D1Database, + params: { did: string; collection: string; rkey: string }, +): Promise { + const { did, collection, rkey } = params; + + await db + .prepare( + "DELETE FROM repo_records WHERE did = ? AND collection = ? AND rkey = ?", + ) + .bind(did, collection, rkey) + .run(); + + const uri = `at://${did}/${collection}/${rkey}`; + await db + .prepare("DELETE FROM resolved_documents WHERE uri = ?") + .bind(uri) + .run(); +} diff --git a/packages/server/worker-configuration.d.ts b/packages/server/worker-configuration.d.ts new file mode 100644 index 0000000..8913310 --- /dev/null +++ b/packages/server/worker-configuration.d.ts @@ -0,0 +1,7 @@ +// Generated by Wrangler by running `wrangler types` + +interface Env { + JETSTREAM_CONSUMER: DurableObjectNamespace; + DB: D1Database; + RESOLUTION_QUEUE: Queue; +} diff --git a/packages/server/wrangler.toml b/packages/server/wrangler.toml index 822c324..eab4745 100644 --- a/packages/server/wrangler.toml +++ b/packages/server/wrangler.toml @@ -19,6 +19,16 @@ max_batch_size = 10 max_batch_timeout = 30 +# Durable Object for Jetstream WebSocket consumer +[durable_objects] +bindings = [ + { name = "JETSTREAM_CONSUMER", class_name = "JetstreamConsumer" } +] + +[[migrations]] +tag = "v1" +new_classes = ["JetstreamConsumer"] + # Cron trigger to refresh stale documents [triggers] crons = ["0 * * * *"] # Every hour (at minute 0) -- 2.51.2