Skip to content
File

Blob: src/worker/queues/search-indexer.ts

typescript71 lines
1import { eq } from "drizzle-orm";
2import { createSessionDb } from "@/worker/db/d1/client";
3import { pages } from "@/worker/db/d1/schema";
4import { createLogger } from "@/worker/lib/logger";
5import { DEFAULT_PAGE_TITLE } from "@/worker/lib/constants";
6 
7const log = createLogger("search-indexer");
8 
9interface IndexPageMessage {
10 type: "index-page";
11 pageId: string;
12}
13 
14export type SearchIndexResult = { kind: "ok" } | { kind: "retry"; delaySeconds: number; reason: string };
15 
16export 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}