File
Blob: src/worker/durable-objects/workspace-indexer.ts
| 1 | import { DurableObject } from "cloudflare:workers"; |
| 2 | import { sql } from "drizzle-orm"; |
| 3 | import { drizzle, type DrizzleSqliteDODatabase } from "drizzle-orm/durable-sqlite"; |
| 4 | import { migrate } from "drizzle-orm/durable-sqlite/migrator"; |
| 5 | import * as indexerSchema from "@/worker/db/workspace-indexer/schema"; |
| 6 | import { createLogger, errorContext } from "@/worker/lib/logger"; |
| 7 | import workspaceIndexerMigrations from "../../../drizzle/workspace-indexer/migrations.js"; |
| 8 | |
| 9 | const log = createLogger("workspace-indexer"); |
| 10 | |
| 11 | type IndexerDb = DrizzleSqliteDODatabase<typeof indexerSchema>; |
| 12 | |
| 13 | export type IndexPageResult = { kind: "indexed" } | { kind: "error"; message: string }; |
| 14 | export type RemovePageResult = { kind: "removed" }; |
| 15 | export type SearchResult = { kind: "results"; items: { pageId: string; snippet: string }[] }; |
| 16 | export type ClearResult = { kind: "cleared" }; |
| 17 | |
| 18 | export class WorkspaceIndexer extends DurableObject<Env> { |
| 19 | private readonly db: IndexerDb; |
| 20 | |
| 21 | constructor(ctx: DurableObjectState, env: Env) { |
| 22 | super(ctx, env); |
| 23 | this.db = drizzle(ctx.storage, { schema: indexerSchema }); |
| 24 | |
| 25 | ctx.blockConcurrencyWhile(async () => { |
| 26 | await migrate(this.db, workspaceIndexerMigrations); |
| 27 | // FTS5 virtual tables are not managed by drizzle — create manually |
| 28 | this.db.run(sql`CREATE VIRTUAL TABLE IF NOT EXISTS pages_fts USING fts5( |
| 29 | page_id UNINDEXED, |
| 30 | title, |
| 31 | body_text, |
| 32 | tokenize='trigram' |
| 33 | )`); |
| 34 | }); |
| 35 | } |
| 36 | |
| 37 | async indexPage(pageId: string, title: string, bodyText: string): Promise<IndexPageResult> { |
| 38 | try { |
| 39 | // FTS5 doesn't support REPLACE INTO — delete then insert in a transaction |
| 40 | this.db.transaction((tx) => { |
| 41 | tx.run(sql`DELETE FROM pages_fts WHERE page_id = ${pageId}`); |
| 42 | tx.run(sql`INSERT INTO pages_fts (page_id, title, body_text) VALUES (${pageId}, ${title}, ${bodyText})`); |
| 43 | }); |
| 44 | log.debug("page_indexed", { pageId, titleLen: title.length, bodyLen: bodyText.length }); |
| 45 | return { kind: "indexed" }; |
| 46 | } catch (e) { |
| 47 | log.error("index_page_failed", { pageId, ...errorContext(e) }); |
| 48 | return { kind: "error", message: e instanceof Error ? e.message : String(e) }; |
| 49 | } |
| 50 | } |
| 51 | |
| 52 | async removePage(pageId: string): Promise<RemovePageResult> { |
| 53 | try { |
| 54 | this.db.run(sql`DELETE FROM pages_fts WHERE page_id = ${pageId}`); |
| 55 | log.debug("page_removed", { pageId }); |
| 56 | } catch (e) { |
| 57 | log.error("remove_page_failed", { pageId, ...errorContext(e) }); |
| 58 | } |
| 59 | return { kind: "removed" }; |
| 60 | } |
| 61 | |
| 62 | async search(query: string, limit: number): Promise<SearchResult> { |
| 63 | try { |
| 64 | // Double-quote wrapping escapes FTS5 operators in user input |
| 65 | const escaped = '"' + query.replace(/"/g, '""') + '"'; |
| 66 | const rows = this.db.all<{ page_id: string; snippet: string }>( |
| 67 | sql`SELECT page_id, |
| 68 | snippet(pages_fts, 2, '<mark>', '</mark>', '...', 32) as snippet |
| 69 | FROM pages_fts |
| 70 | WHERE pages_fts MATCH ${escaped} |
| 71 | LIMIT ${limit}`, |
| 72 | ); |
| 73 | return { |
| 74 | kind: "results", |
| 75 | items: rows.map((r) => ({ pageId: r.page_id, snippet: r.snippet })), |
| 76 | }; |
| 77 | } catch (e) { |
| 78 | log.error("search_failed", { query, ...errorContext(e) }); |
| 79 | return { kind: "results", items: [] }; |
| 80 | } |
| 81 | } |
| 82 | |
| 83 | async clear(): Promise<ClearResult> { |
| 84 | try { |
| 85 | this.db.run(sql`DELETE FROM pages_fts`); |
| 86 | log.info("index_cleared"); |
| 87 | } catch (e) { |
| 88 | log.error("clear_failed", errorContext(e)); |
| 89 | } |
| 90 | return { kind: "cleared" }; |
| 91 | } |
| 92 | } |