From 2f03c5c3549c643773e9a18096e17c7c1e60e410 Mon Sep 17 00:00:00 2001 From: Steve Date: Mon, 12 Jan 2026 22:57:00 -0500 Subject: [PATCH] feat: added queues and updated auth --- packages/server/src/index.ts | 42 ++++++++ packages/server/src/routes/webhook.ts | 57 ++++------- packages/server/src/types/index.ts | 1 + packages/server/src/utils/document.ts | 142 ++++++++++++++++++++++++++ packages/server/src/utils/index.ts | 1 + packages/server/tables.csv | 51 +++++++++ packages/server/wrangler.toml | 15 +-- 7 files changed, 264 insertions(+), 45 deletions(-) create mode 100644 packages/server/src/utils/document.ts create mode 100644 packages/server/tables.csv diff --git a/packages/server/src/index.ts b/packages/server/src/index.ts index 140808e..04f4463 100644 --- a/packages/server/src/index.ts +++ b/packages/server/src/index.ts @@ -2,6 +2,7 @@ import { Hono } from "hono"; import { cors } from "hono/cors"; import type { Bindings } from "./types"; import { health, webhook, feed, stats, records } from "./routes"; +import { processDocument } from "./utils"; const app = new Hono<{ Bindings: Bindings }>(); @@ -47,4 +48,45 @@ app.notFound((c) => { // Export for Cloudflare Workers export default { fetch: app.fetch, + async scheduled( + event: ScheduledEvent, + env: Bindings, + ctx: ExecutionContext, + ) { + const batchSize = 50; + // Select stale documents + const { results } = await env.DB.prepare( + `SELECT did, rkey FROM resolved_documents + WHERE stale_at < datetime('now') OR stale_at IS NULL + LIMIT ?`, + ) + .bind(batchSize) + .all<{ did: string; rkey: string }>(); + + if (results && results.length > 0) { + const messages = results.map((row) => ({ + body: { + did: row.did, + collection: "site.standard.document", + rkey: row.rkey, + }, + })); + + // Send to queue + await env.RESOLUTION_QUEUE.sendBatch(messages); + console.log(`Queued ${messages.length} documents for resolution`); + } + }, + async queue(batch: MessageBatch, env: Bindings) { + for (const message of batch.messages) { + try { + const { did, collection, rkey } = message.body; + await processDocument(env.DB, did, collection, rkey); + message.ack(); + } catch (error) { + console.error("Queue processing error:", error); + message.retry(); + } + } + }, }; diff --git a/packages/server/src/routes/webhook.ts b/packages/server/src/routes/webhook.ts index d0a8ecc..8a7221c 100644 --- a/packages/server/src/routes/webhook.ts +++ b/packages/server/src/routes/webhook.ts @@ -1,38 +1,9 @@ import { Hono } from "hono"; import type { Bindings, TapEvent } from "../types"; -import { resolvePds, parseAtUri } from "../utils"; +import { resolvePds, parseAtUri, resolveViewUrl } from "../utils"; const webhook = new Hono<{ Bindings: Bindings }>(); -async function resolveViewUrl( - db: D1Database, - siteUri: string, - path: string -): Promise { - const parsed = parseAtUri(siteUri); - if (!parsed) return null; - - try { - const pds = await resolvePds(db, parsed.did); - if (!pds) return null; - - const url = `${pds}/xrpc/com.atproto.repo.getRecord?repo=${encodeURIComponent(parsed.did)}&collection=${encodeURIComponent(parsed.collection)}&rkey=${encodeURIComponent(parsed.rkey)}`; - const response = await fetch(url); - if (!response.ok) return null; - - const data = (await response.json()) as { - value?: { url?: string; domain?: string }; - }; - const siteUrl = data.value?.url || data.value?.domain; - if (!siteUrl) return null; - - const baseUrl = siteUrl.startsWith("http") ? siteUrl : `https://${siteUrl}`; - return `${baseUrl}${path}`; - } catch { - return null; - } -} - webhook.post("/tap", async (c) => { try { const db = c.env.DB; @@ -40,7 +11,12 @@ webhook.post("/tap", async (c) => { const secret = c.env.TAP_WEBHOOK_SECRET; if (secret) { const auth = c.req.header("Authorization"); - if (auth !== `Bearer ${secret}`) { + // Support both Bearer token (legacy) and Basic Auth (Tap default) + // Tap sends Basic Auth as base64("admin:password") + const expectedBasic = `Basic ${btoa(`admin:${secret}`)}`; + const expectedBearer = `Bearer ${secret}`; + + if (auth !== expectedBasic && auth !== expectedBearer) { return c.json({ error: "Unauthorized" }, 401); } } @@ -87,10 +63,10 @@ webhook.post("/tap", async (c) => { await db .prepare( - `INSERT INTO resolved_documents (uri, did, rkey, title, path, site, content, text_content, published_at, view_url, resolved_at) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, datetime('now')) + `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', '+12 hours')) ON CONFLICT(uri) DO UPDATE SET - title = ?, path = ?, site = ?, content = ?, text_content = ?, published_at = ?, view_url = ?, resolved_at = datetime('now')` + title = ?, path = ?, site = ?, content = ?, text_content = ?, published_at = ?, view_url = ?, resolved_at = datetime('now'), stale_at = datetime('now', '+12 hours')` ) .bind( uri, @@ -147,7 +123,12 @@ webhook.post("/tap/batch", async (c) => { const secret = c.env.TAP_WEBHOOK_SECRET; if (secret) { const auth = c.req.header("Authorization"); - if (auth !== `Bearer ${secret}`) { + // Support both Bearer token (legacy) and Basic Auth (Tap default) + // Tap sends Basic Auth as base64("admin:password") + const expectedBasic = `Basic ${btoa(`admin:${secret}`)}`; + const expectedBearer = `Bearer ${secret}`; + + if (auth !== expectedBasic && auth !== expectedBearer) { return c.json({ error: "Unauthorized" }, 401); } } @@ -207,10 +188,10 @@ webhook.post("/tap/batch", async (c) => { await db .prepare( - `INSERT INTO resolved_documents (uri, did, rkey, title, path, site, content, text_content, published_at, view_url, resolved_at) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, datetime('now')) + `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', '+12 hours')) ON CONFLICT(uri) DO UPDATE SET - title = ?, path = ?, site = ?, content = ?, text_content = ?, published_at = ?, view_url = ?, resolved_at = datetime('now')` + title = ?, path = ?, site = ?, content = ?, text_content = ?, published_at = ?, view_url = ?, resolved_at = datetime('now'), stale_at = datetime('now', '+12 hours')` ) .bind( uri, diff --git a/packages/server/src/types/index.ts b/packages/server/src/types/index.ts index 9ce4e8f..9f21719 100644 --- a/packages/server/src/types/index.ts +++ b/packages/server/src/types/index.ts @@ -1,5 +1,6 @@ export type Bindings = { DB: D1Database; + RESOLUTION_QUEUE: Queue; TAP_WEBHOOK_SECRET?: string; }; diff --git a/packages/server/src/utils/document.ts b/packages/server/src/utils/document.ts new file mode 100644 index 0000000..a0613eb --- /dev/null +++ b/packages/server/src/utils/document.ts @@ -0,0 +1,142 @@ +import { resolvePds } from "./resolver"; +import { parseAtUri } from "./at-uri"; + +export async function resolveViewUrl( + db: D1Database, + siteUri: string, + path: string +): Promise { + const parsed = parseAtUri(siteUri); + if (!parsed) return null; + + try { + const pds = await resolvePds(db, parsed.did); + if (!pds) return null; + + const url = `${pds}/xrpc/com.atproto.repo.getRecord?repo=${encodeURIComponent( + parsed.did + )}&collection=${encodeURIComponent(parsed.collection)}&rkey=${encodeURIComponent( + parsed.rkey + )}`; + const response = await fetch(url); + if (!response.ok) return null; + + const data = (await response.json()) as { + value?: { url?: string; domain?: string }; + }; + const siteUrl = data.value?.url || data.value?.domain; + if (!siteUrl) return null; + + const baseUrl = siteUrl.startsWith("http") ? siteUrl : `https://${siteUrl}`; + return `${baseUrl}${path}`; + } catch { + return null; + } +} + +export async function processDocument( + db: D1Database, + did: string, + collection: string, + rkey: string +) { + try { + // 1. Resolve PDS + const pds = await resolvePds(db, did); + if (!pds) { + console.warn(`Could not resolve PDS for ${did}`); + return; + } + + // 2. Fetch Record + const url = `${pds}/xrpc/com.atproto.repo.getRecord?repo=${encodeURIComponent( + did + )}&collection=${encodeURIComponent(collection)}&rkey=${encodeURIComponent(rkey)}`; + + const response = await fetch(url); + if (!response.ok) { + if (response.status === 404) { + // Record deleted? + console.warn(`Record not found: ${did}/${collection}/${rkey}`); + } + return; + } + + const data = (await response.json()) as { + uri: string; + cid?: string; + value: { + title?: string; + path?: string; + site?: string; + content?: unknown; + textContent?: string; + publishedAt?: string; + [key: string]: unknown; + }; + }; + + const { value, cid } = data; + + // 3. Update 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(); + + // 4. Resolve View URL and Update resolved_documents + const uri = `at://${did}/${collection}/${rkey}`; + let viewUrl: string | null = null; + if (value.site && value.path) { + viewUrl = await resolveViewUrl(db, value.site, value.path); + } + + // Set stale_at to 12 hours from now + const STALE_OFFSET_HOURS = 12; + + 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, + value.title || null, + value.path || null, + value.site || null, + value.content ? JSON.stringify(value.content) : null, + value.textContent || null, + value.publishedAt || null, + viewUrl, + // Update bindings + value.title || null, + value.path || null, + value.site || null, + value.content ? JSON.stringify(value.content) : null, + value.textContent || null, + value.publishedAt || null, + viewUrl + ) + .run(); + } catch (error) { + console.error(`Error processing document ${did}/${collection}/${rkey}:`, error); + } +} diff --git a/packages/server/src/utils/index.ts b/packages/server/src/utils/index.ts index 1b28b86..5e48159 100644 --- a/packages/server/src/utils/index.ts +++ b/packages/server/src/utils/index.ts @@ -1,2 +1,3 @@ export { parseAtUri, buildAtUri, type AtUriComponents } from "./at-uri"; export { resolvePds } from "./resolver"; +export { resolveViewUrl, processDocument } from "./document"; diff --git a/packages/server/tables.csv b/packages/server/tables.csv new file mode 100644 index 0000000..0813b62 --- /dev/null +++ b/packages/server/tables.csv @@ -0,0 +1,51 @@ +name,sql +_cf_KV,"CREATE TABLE _cf_KV ( + key TEXT PRIMARY KEY, + value BLOB + ) WITHOUT ROWID" +repo_records,"CREATE TABLE repo_records ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + did TEXT NOT NULL, + rkey TEXT NOT NULL, + collection TEXT NOT NULL, + cid TEXT, + synced_at TEXT DEFAULT (datetime('now')), + UNIQUE(did, collection, rkey) +)" +pds_cache,"CREATE TABLE pds_cache ( + did TEXT PRIMARY KEY, + pds_endpoint TEXT NOT NULL, + cached_at TEXT DEFAULT (datetime('now')) +)" +record_cache,"CREATE TABLE record_cache ( + uri TEXT PRIMARY KEY, + did TEXT NOT NULL, + collection TEXT NOT NULL, + rkey TEXT NOT NULL, + record_data TEXT NOT NULL, -- JSON blob + cached_at TEXT DEFAULT (datetime('now')) +)" +publication_cache,"CREATE TABLE publication_cache ( + at_uri TEXT PRIMARY KEY, + base_url TEXT NOT NULL, + cached_at TEXT DEFAULT (datetime('now')) +)" +sync_metadata,"CREATE TABLE sync_metadata ( + key TEXT PRIMARY KEY, + value TEXT NOT NULL, + updated_at TEXT DEFAULT (datetime('now')) +)" +resolved_documents,"CREATE TABLE resolved_documents ( + uri TEXT PRIMARY KEY, + did TEXT NOT NULL, + rkey TEXT NOT NULL, + title TEXT, + path TEXT, + site TEXT, + content TEXT, -- JSON blob + text_content TEXT, + published_at TEXT, + view_url TEXT, + resolved_at TEXT DEFAULT (datetime('now')), + stale_at TEXT -- When this record should be re-resolved +)" diff --git a/packages/server/wrangler.toml b/packages/server/wrangler.toml index a9cfe22..381ae3c 100644 --- a/packages/server/wrangler.toml +++ b/packages/server/wrangler.toml @@ -9,14 +9,15 @@ database_name = "atfeeds-db" database_id = "bfbb9955-1496-47e9-9602-e32c9b1fa7b2" # Queue for processing document resolution -# [[queues.producers]] -# queue = "document-resolution" -# binding = "RESOLUTION_QUEUE" +[[queues.producers]] +queue = "document-resolution" +binding = "RESOLUTION_QUEUE" + +[[queues.consumers]] +queue = "document-resolution" +max_batch_size = 10 +max_batch_timeout = 30 -# [[queues.consumers]] -# queue = "document-resolution" -# max_batch_size = 10 -# max_batch_timeout = 30 # Cron trigger to refresh stale documents [triggers] -- 2.51.2