File
Blob: src/worker/queues/search-indexer.ts
| 1 | import { eq } from "drizzle-orm"; |
| 2 | import { createSessionDb } from "@/worker/db/d1/client"; |
| 3 | import { pages } from "@/worker/db/d1/schema"; |
| 4 | import { createLogger } from "@/worker/lib/logger"; |
| 5 | import { DEFAULT_PAGE_TITLE } from "@/worker/lib/constants"; |
| 6 | |
| 7 | const log = createLogger("search-indexer"); |
| 8 | |
| 9 | interface IndexPageMessage { |
| 10 | type: "index-page"; |
| 11 | pageId: string; |
| 12 | } |
| 13 | |
| 14 | export type SearchIndexResult = { kind: "ok" } | { kind: "retry"; delaySeconds: number; reason: string }; |
| 15 | |
| 16 | export async function handleSearchIndexMessage(msg: IndexPageMessage, env: Env): Promise<SearchIndexResult> { |
| 17 | const { pageId } = msg; |
| 18 | const { db } = createSessionDb(env.DB, "first-primary"); |
| 19 | |
| 20 | // Load page metadata from D1 (workspace_id for routing, archived_at for removal, kind for extraction) |
| 21 | const page = await db |
| 22 | .select({ |
| 23 | workspace_id: pages.workspace_id, |
| 24 | archived_at: pages.archived_at, |
| 25 | title: pages.title, |
| 26 | kind: pages.kind, |
| 27 | }) |
| 28 | .from(pages) |
| 29 | .where(eq(pages.id, pageId)) |
| 30 | .get(); |
| 31 | |
| 32 | // Page not visible yet from D1: most often the queue consumer raced ahead |
| 33 | // of the producer's bookmark. Retry so create-immediately-search works. |
| 34 | if (!page) { |
| 35 | log.info("fts_retry", { pageId, reason: "page_not_yet_visible" }); |
| 36 | return { kind: "retry", delaySeconds: 2, reason: "page_not_yet_visible" }; |
| 37 | } |
| 38 | |
| 39 | const indexer = env.WorkspaceIndexer.getByName(page.workspace_id); |
| 40 | |
| 41 | // Archived or deleted — remove from search index |
| 42 | if (page.archived_at) { |
| 43 | await indexer.removePage(pageId); |
| 44 | log.info("fts_removed", { pageId, reason: "archived" }); |
| 45 | return { kind: "ok" }; |
| 46 | } |
| 47 | |
| 48 | // Fetch indexable text from DocSync DO |
| 49 | const doc = env.DocSync.getByName(pageId); |
| 50 | const payload = await doc.getIndexPayload(pageId, page.kind); |
| 51 | |
| 52 | let title: string; |
| 53 | let bodyText: string; |
| 54 | |
| 55 | if (payload.kind === "found") { |
| 56 | title = payload.title; |
| 57 | bodyText = payload.bodyText; |
| 58 | } else { |
| 59 | // No snapshot yet — index by page title from D1 |
| 60 | title = page.title?.trim() || DEFAULT_PAGE_TITLE; |
| 61 | bodyText = ""; |
| 62 | } |
| 63 | |
| 64 | const result = await indexer.indexPage(pageId, title, bodyText); |
| 65 | if (result.kind === "error") { |
| 66 | throw new Error(`WorkspaceIndexer.indexPage failed for ${pageId}: ${result.message}`); |
| 67 | } |
| 68 | log.info("fts_indexed", { pageId, titleLen: title.length, bodyLen: bodyText.length }); |
| 69 | return { kind: "ok" }; |
| 70 | } |