From 67de59c5a9f8bbe82bb701595a9d5371a6c83515 Mon Sep 17 00:00:00 2001 From: Steve Date: Wed, 14 Jan 2026 09:23:16 -0500 Subject: [PATCH] chore: adjustments to data flow --- packages/server/src/routes/webhook.ts | 27 ++++++++++++++++++++++----- packages/server/src/utils/document.ts | 2 +- packages/server/wrangler.toml | 2 +- 3 files changed, 24 insertions(+), 7 deletions(-) diff --git a/packages/server/src/routes/webhook.ts b/packages/server/src/routes/webhook.ts index 8a7221c..53cf507 100644 --- a/packages/server/src/routes/webhook.ts +++ b/packages/server/src/routes/webhook.ts @@ -1,6 +1,8 @@ import { Hono } from "hono"; import type { Bindings, TapEvent } from "../types"; -import { resolvePds, parseAtUri, resolveViewUrl } from "../utils"; +import { resolveViewUrl } from "../utils"; + +const STALE_OFFSET_HOURS = 24; const webhook = new Hono<{ Bindings: Bindings }>(); @@ -64,9 +66,9 @@ 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, stale_at) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, datetime('now'), datetime('now', '+12 hours')) + 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', '+12 hours')` + title = ?, path = ?, site = ?, content = ?, text_content = ?, published_at = ?, view_url = ?, resolved_at = datetime('now'), stale_at = datetime('now', '+${STALE_OFFSET_HOURS} hours')` ) .bind( uri, @@ -89,6 +91,13 @@ webhook.post("/tap", async (c) => { ) .run(); } + + // Queue for immediate full processing (verification, publication resolution, etc.) + await c.env.RESOLUTION_QUEUE.send({ + did: record.did, + collection: record.collection, + rkey: record.rkey, + }); } else if (record.action === "delete") { await db .prepare( @@ -189,9 +198,9 @@ 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, stale_at) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, datetime('now'), datetime('now', '+12 hours')) + 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', '+12 hours')` + title = ?, path = ?, site = ?, content = ?, text_content = ?, published_at = ?, view_url = ?, resolved_at = datetime('now'), stale_at = datetime('now', '+${STALE_OFFSET_HOURS} hours')` ) .bind( uri, @@ -214,6 +223,14 @@ webhook.post("/tap/batch", async (c) => { ) .run(); } + + // Queue for immediate full processing + await c.env.RESOLUTION_QUEUE.send({ + did: event.did, + collection: event.collection, + rkey: event.rkey, + }); + processed++; } else if ( event.type === "delete" && diff --git a/packages/server/src/utils/document.ts b/packages/server/src/utils/document.ts index 2fcba7a..c2eb613 100644 --- a/packages/server/src/utils/document.ts +++ b/packages/server/src/utils/document.ts @@ -204,7 +204,7 @@ export async function processDocument( const verified = await verifyDocumentRecord(pubUrl, site, viewUrl, uri); // 7. Insert/update resolved_documents - const STALE_OFFSET_HOURS = 12; + const STALE_OFFSET_HOURS = 24; await db .prepare( diff --git a/packages/server/wrangler.toml b/packages/server/wrangler.toml index 381ae3c..822c324 100644 --- a/packages/server/wrangler.toml +++ b/packages/server/wrangler.toml @@ -21,7 +21,7 @@ max_batch_timeout = 30 # Cron trigger to refresh stale documents [triggers] -crons = ["*/15 * * * *"] # Every 15 minutes +crons = ["0 * * * *"] # Every hour (at minute 0) # Environment variables (secrets should be set via wrangler secret) # TAP_WEBHOOK_SECRET - Optional secret for webhook authentication -- 2.51.2