| 1 | /** |
| 2 | * The semantic index: pages and projects' docs files as passages |
| 3 | * (src/chunks.ts) in D1 (`doc_chunks`, with full text in |
| 4 | * `doc_chunks_fts`) and as vectors in the vector store (src/vectors.ts). |
| 5 | * |
| 6 | * A page is indexed a little after its Markdown is saved (the room, |
| 7 | * src/room.ts, in the background, so editing never waits on it); a |
| 8 | * project's docs files after they are read (src/repo-spaces.ts). Only |
| 9 | * passages whose text changed are embedded; one that only moved keeps its |
| 10 | * vector. Embedding is capped per workspace and hour (EMBED_PER_HOUR): |
| 11 | * past the cap passages are kept for words, and a catch-up run on the |
| 12 | * queue embeds them later. A failure leaves them for the next save or run. |
| 13 | * Nothing here throws to its caller. |
| 14 | * |
| 15 | * The backfill indexes a workspace's existing pages and files in batches |
| 16 | * on the queue (`docs.index` jobs, on the artifacts service's own events queue). |
| 17 | */ |
| 18 | import { chunkId, chunkMarkdown, embedText, repoFileId, textHash } from "./chunks.ts"; |
| 19 | import { scopesOf, type FolioRow, FOLIO_COLUMNS } from "./folios/access-store.ts"; |
| 20 | import { kindModel } from "./kinds/index.ts"; |
| 21 | import { cloudflareEmbedder, cloudflareVectors, type Embedder, type VectorMetadata, type VectorStore } from "./vectors.ts"; |
| 22 | |
| 23 | /** Passages embedded per workspace per hour at most. Logged when reached. */ |
| 24 | export const EMBED_PER_HOUR = 200; |
| 25 | /** Documents one backfill job indexes before handing on to the next. */ |
| 26 | const BACKFILL_BATCH = 20; |
| 27 | /** A backfill not finished after this long is taken to have died and may start again. */ |
| 28 | const BACKFILL_STALE_MS = 2 * 60 * 60 * 1000; |
| 29 | /** How long a run waits when the hour's cap is reached. */ |
| 30 | const CAP_DELAY_SECONDS = 3600; |
| 31 | |
| 32 | /** |
| 33 | * A job on the queue: index more of a workspace's docs and folios; or |
| 34 | * bring a large folio subtree's access, rooms and index scope up to date |
| 35 | * after a move or a sharing change (src/folios/service.ts). |
| 36 | */ |
| 37 | export type DocsJob = { type: "docs.index"; workspace_id: string } | { type: "folios.reacl"; folio_id: string }; |
| 38 | |
| 39 | export type IndexEnv = { |
| 40 | DB: D1Database; |
| 41 | AI?: Ai; |
| 42 | VECTORS?: Vectorize; |
| 43 | /** Folios' semantic index, Vectorize `g1t-folios`. Optional: without it folios match words only. */ |
| 44 | FOLIO_VECTORS?: Vectorize; |
| 45 | JOBS?: Queue<DocsJob>; |
| 46 | }; |
| 47 | |
| 48 | /** The embedder and store this deployment has; null without them (words only). */ |
| 49 | export function adapters(env: IndexEnv): { embedder: Embedder | null; store: VectorStore | null } { |
| 50 | return env.AI && env.VECTORS ? { embedder: cloudflareEmbedder(env.AI), store: cloudflareVectors(env.VECTORS) } : { embedder: null, store: null }; |
| 51 | } |
| 52 | |
| 53 | /** One document to index: a page, or a project's docs file. */ |
| 54 | export type IndexDoc = { |
| 55 | kind: "page" | "repo_file"; |
| 56 | /** The page's id, or the file's `rf_` id. */ |
| 57 | doc_id: string; |
| 58 | workspace_id: string; |
| 59 | space_id: string; |
| 60 | title: string; |
| 61 | markdown: string; |
| 62 | repo_id?: string | null; |
| 63 | path?: string | null; |
| 64 | }; |
| 65 | |
| 66 | type ChunkRow = { id: string; seq: number; space_id: string; hash: string; vector_hash: string | null }; |
| 67 | |
| 68 | const hourOf = (now: Date) => now.toISOString().slice(0, 13); |
| 69 | |
| 70 | function metadataOf(doc: IndexDoc): VectorMetadata { |
| 71 | return doc.kind === "page" |
| 72 | ? { workspace_id: doc.workspace_id, space_id: doc.space_id, kind: "page", page_id: doc.doc_id } |
| 73 | : { workspace_id: doc.workspace_id, space_id: doc.space_id, kind: "repo_file", repo_file_id: doc.doc_id, repo_id: doc.repo_id ?? "" }; |
| 74 | } |
| 75 | |
| 76 | /** Passages the workspace may still embed this hour. */ |
| 77 | async function allowance(db: D1Database, workspaceId: string, now: Date): Promise<number> { |
| 78 | const used = await db.prepare("SELECT chunks FROM doc_embed_usage WHERE workspace_id = ? AND hour = ?").bind(workspaceId, hourOf(now)).first<{ chunks: number }>(); |
| 79 | return Math.max(0, EMBED_PER_HOUR - (used?.chunks ?? 0)); |
| 80 | } |
| 81 | |
| 82 | async function meter(db: D1Database, workspaceId: string, chunks: number, chars: number, now: Date): Promise<void> { |
| 83 | if (!chunks) return; |
| 84 | await db |
| 85 | .prepare( |
| 86 | "INSERT INTO doc_embed_usage (workspace_id, hour, chunks, tokens) VALUES (?1, ?2, ?3, ?4) ON CONFLICT (workspace_id, hour) DO UPDATE SET chunks = chunks + ?3, tokens = tokens + ?4", |
| 87 | ) |
| 88 | .bind(workspaceId, hourOf(now), chunks, Math.ceil(chars / 4)) |
| 89 | .run(); |
| 90 | } |
| 91 | |
| 92 | /** |
| 93 | * Brings one document's passages up to date: rows and full text in D1, |
| 94 | * vectors for those that changed (within the hour's cap). `capped` when |
| 95 | * some were left for later. |
| 96 | */ |
| 97 | export async function indexDoc( |
| 98 | db: D1Database, |
| 99 | embedder: Embedder | null, |
| 100 | store: VectorStore | null, |
| 101 | doc: IndexDoc, |
| 102 | now = new Date(), |
| 103 | ): Promise<{ chunks: number; embedded: number; capped: boolean }> { |
| 104 | const at = now.toISOString(); |
| 105 | const chunks = chunkMarkdown(doc.markdown, doc.title).map((c) => { |
| 106 | const embed = embedText(doc.title, c); |
| 107 | return { ...c, id: chunkId(doc.doc_id, c.seq), embed, hash: textHash(embed) }; |
| 108 | }); |
| 109 | const column = doc.kind === "page" ? "page_id" : "repo_file_id"; |
| 110 | const old = (await db.prepare(`SELECT id, seq, space_id, hash, vector_hash FROM doc_chunks WHERE ${column} = ?`).bind(doc.doc_id).all<ChunkRow>()).results; |
| 111 | const byId = new Map(old.map((r) => [r.id, r])); |
| 112 | const statements: D1PreparedStatement[] = []; |
| 113 | for (const c of chunks) { |
| 114 | const was = byId.get(c.id); |
| 115 | if (was && was.hash === c.hash && was.space_id === doc.space_id) continue; |
| 116 | statements.push( |
| 117 | db |
| 118 | .prepare( |
| 119 | `INSERT INTO doc_chunks (id, workspace_id, space_id, page_id, repo_file_id, repo_id, path, seq, heading, text, hash, vector_hash, updated_at) |
| 120 | VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, NULL, ?) |
| 121 | ON CONFLICT (id) DO UPDATE SET space_id = excluded.space_id, repo_id = excluded.repo_id, path = excluded.path, heading = excluded.heading, text = excluded.text, hash = excluded.hash, updated_at = excluded.updated_at`, |
| 122 | ) |
| 123 | .bind( |
| 124 | c.id, |
| 125 | doc.workspace_id, |
| 126 | doc.space_id, |
| 127 | doc.kind === "page" ? doc.doc_id : null, |
| 128 | doc.kind === "repo_file" ? doc.doc_id : null, |
| 129 | doc.repo_id ?? null, |
| 130 | doc.path ?? null, |
| 131 | c.seq, |
| 132 | c.heading, |
| 133 | c.text, |
| 134 | c.hash, |
| 135 | at, |
| 136 | ), |
| 137 | db.prepare("DELETE FROM doc_chunks_fts WHERE chunk_id = ?").bind(c.id), |
| 138 | db.prepare("INSERT INTO doc_chunks_fts (chunk_id, space_id, doc_id, heading, text) VALUES (?, ?, ?, ?, ?)").bind(c.id, doc.space_id, doc.doc_id, c.heading ?? "", c.text), |
| 139 | ); |
| 140 | } |
| 141 | const keep = new Set(chunks.map((c) => c.id)); |
| 142 | const gone = old.filter((r) => !keep.has(r.id)); |
| 143 | for (const r of gone) { |
| 144 | statements.push(db.prepare("DELETE FROM doc_chunks WHERE id = ?").bind(r.id), db.prepare("DELETE FROM doc_chunks_fts WHERE chunk_id = ?").bind(r.id)); |
| 145 | } |
| 146 | for (let i = 0; i < statements.length; i += 60) await db.batch(statements.slice(i, i + 60)); |
| 147 | |
| 148 | if (!store || !embedder) return { chunks: chunks.length, embedded: 0, capped: false }; |
| 149 | const result = await embedChanged(db, embedder, store, doc, chunks, old, byId, now); |
| 150 | // Last, so a passage that moved from a place now gone could take its vector first. |
| 151 | try { |
| 152 | if (gone.some((r) => r.vector_hash)) await store.delete(gone.filter((r) => r.vector_hash).map((r) => r.id)); |
| 153 | } catch (error) { |
| 154 | console.error("docs could not drop old passages from the index", doc.doc_id, String(error)); |
| 155 | } |
| 156 | return result; |
| 157 | } |
| 158 | |
| 159 | type IndexedChunk = { id: string; seq: number; heading: string | null; text: string; embed: string; hash: string }; |
| 160 | |
| 161 | /** Vectors for the passages that need one: reused when the index already holds their text, embedded otherwise (within the cap). */ |
| 162 | async function embedChanged( |
| 163 | db: D1Database, |
| 164 | embedder: Embedder, |
| 165 | store: VectorStore, |
| 166 | doc: IndexDoc, |
| 167 | chunks: IndexedChunk[], |
| 168 | old: ChunkRow[], |
| 169 | byId: Map<string, ChunkRow>, |
| 170 | now: Date, |
| 171 | ): Promise<{ chunks: number; embedded: number; capped: boolean }> { |
| 172 | // What needs a vector: a passage whose text the index doesn't hold under its id, or one whose space changed. |
| 173 | const stale = chunks.filter((c) => { |
| 174 | const was = byId.get(c.id); |
| 175 | return !was || was.vector_hash !== c.hash || was.space_id !== doc.space_id; |
| 176 | }); |
| 177 | if (!stale.length) return { chunks: chunks.length, embedded: 0, capped: false }; |
| 178 | // A passage that only moved (or changed space) has its vector already, under another id or this one. |
| 179 | const held = new Map<string, string>(); |
| 180 | for (const r of old) if (r.vector_hash) held.set(r.vector_hash, r.id); |
| 181 | const reuse = stale.filter((c) => held.has(c.hash)); |
| 182 | const fresh = stale.filter((c) => !held.has(c.hash)); |
| 183 | const budget = fresh.length ? await allowance(db, doc.workspace_id, now) : 0; |
| 184 | const embedNow = fresh.slice(0, budget); |
| 185 | const capped = embedNow.length < fresh.length; |
| 186 | if (capped) console.log("docs embedding cap reached", doc.workspace_id, `${fresh.length - embedNow.length} passages wait for the next hour`); |
| 187 | const metadata = metadataOf(doc); |
| 188 | const done: { id: string; hash: string }[] = []; |
| 189 | try { |
| 190 | const vectors: { id: string; values: number[]; metadata: VectorMetadata }[] = []; |
| 191 | if (reuse.length) { |
| 192 | const values = new Map((await store.get([...new Set(reuse.map((c) => held.get(c.hash)!))])).map((v) => [v.id, v.values])); |
| 193 | for (const c of reuse) { |
| 194 | const v = values.get(held.get(c.hash)!); |
| 195 | if (v) vectors.push({ id: c.id, values: v, metadata }); |
| 196 | } |
| 197 | } |
| 198 | if (embedNow.length) { |
| 199 | const embedded = await embedder.embed(embedNow.map((c) => c.embed)); |
| 200 | embedNow.forEach((c, i) => vectors.push({ id: c.id, values: embedded[i]!, metadata })); |
| 201 | await meter( |
| 202 | db, |
| 203 | doc.workspace_id, |
| 204 | embedNow.length, |
| 205 | embedNow.reduce((n, c) => n + c.embed.length, 0), |
| 206 | now, |
| 207 | ); |
| 208 | } |
| 209 | if (vectors.length) await store.upsert(vectors); |
| 210 | const hashOf = new Map(chunks.map((c) => [c.id, c.hash])); |
| 211 | for (const v of vectors) done.push({ id: v.id, hash: hashOf.get(v.id)! }); |
| 212 | } catch (error) { |
| 213 | console.error("docs could not embed passages; the next save tries again", doc.doc_id, String(error)); |
| 214 | } |
| 215 | if (done.length) { |
| 216 | const updates = done.map((d) => db.prepare("UPDATE doc_chunks SET vector_hash = ? WHERE id = ?").bind(d.hash, d.id)); |
| 217 | for (let i = 0; i < updates.length; i += 60) await db.batch(updates.slice(i, i + 60)); |
| 218 | } |
| 219 | return { chunks: chunks.length, embedded: embedNow.length, capped }; |
| 220 | } |
| 221 | |
| 222 | /** Takes passages out of D1 and the index: by page ids, file ids, or a whole space. */ |
| 223 | export async function forgetDocs(env: IndexEnv, which: { page_ids?: string[]; repo_file_ids?: string[]; space_id?: string }): Promise<void> { |
| 224 | const db = env.DB; |
| 225 | const { store } = adapters(env); |
| 226 | try { |
| 227 | const found: string[] = []; |
| 228 | const select = async (column: string, values: string[]) => { |
| 229 | for (let i = 0; i < values.length; i += 90) { |
| 230 | const part = values.slice(i, i + 90); |
| 231 | const rows = await db.prepare(`SELECT id FROM doc_chunks WHERE ${column} IN (${part.map(() => "?").join(",")})`).bind(...part).all<{ id: string }>(); |
| 232 | found.push(...rows.results.map((r) => r.id)); |
| 233 | } |
| 234 | }; |
| 235 | if (which.page_ids?.length) await select("page_id", which.page_ids); |
| 236 | if (which.repo_file_ids?.length) await select("repo_file_id", which.repo_file_ids); |
| 237 | if (which.space_id) await select("space_id", [which.space_id]); |
| 238 | if (!found.length) return; |
| 239 | if (store) { |
| 240 | try { |
| 241 | await store.delete(found); |
| 242 | } catch (error) { |
| 243 | console.error("docs could not drop passages from the index", found.length, String(error)); |
| 244 | } |
| 245 | } |
| 246 | const statements = found.flatMap((id) => [db.prepare("DELETE FROM doc_chunks WHERE id = ?").bind(id), db.prepare("DELETE FROM doc_chunks_fts WHERE chunk_id = ?").bind(id)]); |
| 247 | for (let i = 0; i < statements.length; i += 80) await db.batch(statements.slice(i, i + 80)); |
| 248 | } catch (error) { |
| 249 | console.error("docs could not forget passages", String(error)); |
| 250 | } |
| 251 | } |
| 252 | |
| 253 | type PageForIndex = { id: string; workspace_id: string; space_id: string; title: string; markdown: string; archived_at: string | null }; |
| 254 | |
| 255 | /** Indexes one page as it is saved now; forgets it when it is gone or in the trash. Never throws. */ |
| 256 | export async function indexPage(env: IndexEnv, pageId: string, now = new Date()): Promise<{ capped: boolean }> { |
| 257 | try { |
| 258 | const page = await env.DB.prepare("SELECT id, workspace_id, space_id, title, markdown, archived_at FROM pages WHERE id = ?").bind(pageId).first<PageForIndex>(); |
| 259 | if (!page || page.archived_at) { |
| 260 | await forgetDocs(env, { page_ids: [pageId] }); |
| 261 | return { capped: false }; |
| 262 | } |
| 263 | const { embedder, store } = adapters(env); |
| 264 | const result = await indexDoc(env.DB, embedder, store, { kind: "page", doc_id: page.id, workspace_id: page.workspace_id, space_id: page.space_id, title: page.title, markdown: page.markdown }, now); |
| 265 | if (result.capped) await startBackfill(env, page.workspace_id, { delaySeconds: CAP_DELAY_SECONDS }); |
| 266 | return { capped: result.capped }; |
| 267 | } catch (error) { |
| 268 | console.error("docs could not index a page", pageId, String(error)); |
| 269 | return { capped: false }; |
| 270 | } |
| 271 | } |
| 272 | |
| 273 | type RepoFileForIndex = { space_id: string; path: string; title: string; markdown: string; workspace_id: string; repo_id: string }; |
| 274 | |
| 275 | function repoDoc(row: RepoFileForIndex): IndexDoc { |
| 276 | return { kind: "repo_file", doc_id: repoFileId(row.space_id, row.path), workspace_id: row.workspace_id, space_id: row.space_id, title: row.title, markdown: row.markdown, repo_id: row.repo_id, path: row.path }; |
| 277 | } |
| 278 | |
| 279 | /** |
| 280 | * Indexes a project's docs files after they were read: `changed` paths |
| 281 | * again, `gone` paths forgotten. Never throws. |
| 282 | */ |
| 283 | export async function indexRepoFiles(env: IndexEnv, spaceId: string, changed: string[], gone: string[], now = new Date()): Promise<void> { |
| 284 | try { |
| 285 | if (gone.length) await forgetDocs(env, { repo_file_ids: gone.map((p) => repoFileId(spaceId, p)) }); |
| 286 | if (!changed.length) return; |
| 287 | const { embedder, store } = adapters(env); |
| 288 | let capped = false; |
| 289 | let workspaceId: string | null = null; |
| 290 | for (let i = 0; i < changed.length; i += 50) { |
| 291 | const part = changed.slice(i, i + 50); |
| 292 | const rows = ( |
| 293 | await env.DB.prepare( |
| 294 | `SELECT f.space_id, f.path, f.title, f.markdown, s.workspace_id, s.repo_id FROM repo_files f JOIN repo_spaces s ON s.id = f.space_id WHERE f.space_id = ? AND f.path IN (${part.map(() => "?").join(",")})`, |
| 295 | ) |
| 296 | .bind(spaceId, ...part) |
| 297 | .all<RepoFileForIndex>() |
| 298 | ).results; |
| 299 | for (const row of rows) { |
| 300 | workspaceId = row.workspace_id; |
| 301 | // Past the cap, the rest are only kept for words until the catch-up run. |
| 302 | const result = await indexDoc(env.DB, capped ? null : embedder, capped ? null : store, repoDoc(row), now); |
| 303 | capped ||= result.capped; |
| 304 | } |
| 305 | } |
| 306 | if (capped && workspaceId) await startBackfill(env, workspaceId, { delaySeconds: CAP_DELAY_SECONDS }); |
| 307 | } catch (error) { |
| 308 | console.error("docs could not index a project's docs", spaceId, String(error)); |
| 309 | } |
| 310 | } |
| 311 | |
| 312 | /** |
| 313 | * Starts (or, with `force`, restarts) indexing a workspace's existing |
| 314 | * pages and projects' docs on the queue. Does nothing while a run is |
| 315 | * going, unless it has been going so long it must have died. Without a |
| 316 | * queue, runs one batch now. True when a run was started. |
| 317 | */ |
| 318 | export async function startBackfill(env: IndexEnv, workspaceId: string, options: { force?: boolean; delaySeconds?: number } = {}): Promise<boolean> { |
| 319 | const at = new Date(); |
| 320 | const staleBefore = new Date(at.getTime() - BACKFILL_STALE_MS).toISOString(); |
| 321 | const started = await env.DB.prepare( |
| 322 | `INSERT INTO doc_index_runs (workspace_id, started_at, finished_at, cursor, pages, files) VALUES (?1, ?2, NULL, NULL, 0, 0) |
| 323 | ON CONFLICT (workspace_id) DO UPDATE SET started_at = ?2, finished_at = NULL, cursor = NULL, pages = 0, files = 0 |
| 324 | WHERE ?3 OR doc_index_runs.finished_at IS NOT NULL OR doc_index_runs.started_at < ?4 |
| 325 | RETURNING workspace_id`, |
| 326 | ) |
| 327 | .bind(workspaceId, at.toISOString(), options.force ? 1 : 0, staleBefore) |
| 328 | .first<{ workspace_id: string }>(); |
| 329 | if (!started) return false; |
| 330 | await enqueue(env, workspaceId, options.delaySeconds ?? 0); |
| 331 | return true; |
| 332 | } |
| 333 | |
| 334 | async function enqueue(env: IndexEnv, workspaceId: string, delaySeconds: number): Promise<void> { |
| 335 | if (env.JOBS) { |
| 336 | await env.JOBS.send({ type: "docs.index", workspace_id: workspaceId }, delaySeconds ? { delaySeconds } : undefined); |
| 337 | return; |
| 338 | } |
| 339 | // No queue (self-hosted without one): one batch now; saves and later recalls carry on from there. |
| 340 | if (!delaySeconds) await runBackfill(env, workspaceId, false); |
| 341 | } |
| 342 | |
| 343 | /** |
| 344 | * The first time anyone recalls from a workspace's docs: when it has |
| 345 | * pages or projects' docs and has never been indexed, start the backfill. |
| 346 | */ |
| 347 | export async function ensureIndexed(env: IndexEnv, workspaceId: string): Promise<void> { |
| 348 | const row = await env.DB.prepare( |
| 349 | "SELECT (SELECT 1 FROM doc_index_runs WHERE workspace_id = ?1) AS ran, (SELECT 1 FROM pages WHERE workspace_id = ?1 AND archived_at IS NULL LIMIT 1) AS page, (SELECT 1 FROM folios WHERE workspace_id = ?1 AND trashed_at IS NULL LIMIT 1) AS folio, (SELECT 1 FROM repo_spaces WHERE workspace_id = ?1 LIMIT 1) AS repo", |
| 350 | ) |
| 351 | .bind(workspaceId) |
| 352 | .first<{ ran: number | null; page: number | null; folio: number | null; repo: number | null }>(); |
| 353 | if (!row || row.ran || (!row.page && !row.folio && !row.repo)) return; |
| 354 | await startBackfill(env, workspaceId); |
| 355 | } |
| 356 | |
| 357 | /** |
| 358 | * One batch of a workspace's backfill: pages by id, then projects' docs |
| 359 | * files by space and path, from the run's cursor. Hands on to the next |
| 360 | * batch (or waits out the hour's cap) through the queue. |
| 361 | */ |
| 362 | export async function runBackfill(env: IndexEnv, workspaceId: string, chain = true): Promise<{ done: boolean }> { |
| 363 | const run = await env.DB.prepare("SELECT cursor, finished_at, pages, files FROM doc_index_runs WHERE workspace_id = ?").bind(workspaceId).first<{ cursor: string | null; finished_at: string | null; pages: number; files: number }>(); |
| 364 | if (!run || run.finished_at) return { done: true }; |
| 365 | const { embedder, store } = adapters(env); |
| 366 | const cursor = run.cursor ?? "p:"; |
| 367 | let next: string | null = cursor; |
| 368 | let pages = 0; |
| 369 | let files = 0; |
| 370 | let capped = false; |
| 371 | if (cursor.startsWith("o:")) { |
| 372 | // Folios, after pages and before projects' docs. |
| 373 | const rows = ( |
| 374 | await env.DB.prepare(`SELECT ${FOLIO_COLUMNS} FROM folios WHERE workspace_id = ? AND trashed_at IS NULL AND id > ? ORDER BY id LIMIT ?`) |
| 375 | .bind(workspaceId, cursor.slice(2), BACKFILL_BATCH) |
| 376 | .all<FolioRow>() |
| 377 | ).results; |
| 378 | for (const row of rows) { |
| 379 | const result = await indexFolio(env, row.id); |
| 380 | pages++; |
| 381 | if (result.capped) { |
| 382 | capped = true; |
| 383 | break; |
| 384 | } |
| 385 | next = `o:${row.id}`; |
| 386 | } |
| 387 | if (!capped && rows.length < BACKFILL_BATCH) next = "f:"; |
| 388 | } else if (cursor.startsWith("p:")) { |
| 389 | const rows = ( |
| 390 | await env.DB.prepare("SELECT id, workspace_id, space_id, title, markdown, archived_at FROM pages WHERE workspace_id = ? AND archived_at IS NULL AND id > ? ORDER BY id LIMIT ?") |
| 391 | .bind(workspaceId, cursor.slice(2), BACKFILL_BATCH) |
| 392 | .all<PageForIndex>() |
| 393 | ).results; |
| 394 | for (const page of rows) { |
| 395 | const result = await indexDoc(env.DB, embedder, store, { kind: "page", doc_id: page.id, workspace_id: page.workspace_id, space_id: page.space_id, title: page.title, markdown: page.markdown }); |
| 396 | pages++; |
| 397 | if (result.capped) { |
| 398 | // This page again next time, after the hour: the cursor stays before it. |
| 399 | capped = true; |
| 400 | break; |
| 401 | } |
| 402 | next = `p:${page.id}`; |
| 403 | } |
| 404 | if (!capped && rows.length < BACKFILL_BATCH) next = "o:"; |
| 405 | } else { |
| 406 | const [space, path] = splitFileCursor(cursor.slice(2)); |
| 407 | const rows = ( |
| 408 | await env.DB.prepare( |
| 409 | `SELECT f.space_id, f.path, f.title, f.markdown, s.workspace_id, s.repo_id FROM repo_files f JOIN repo_spaces s ON s.id = f.space_id |
| 410 | WHERE s.workspace_id = ? AND (f.space_id > ? OR (f.space_id = ? AND f.path > ?)) ORDER BY f.space_id, f.path LIMIT ?`, |
| 411 | ) |
| 412 | .bind(workspaceId, space, space, path, BACKFILL_BATCH) |
| 413 | .all<RepoFileForIndex>() |
| 414 | ).results; |
| 415 | for (const row of rows) { |
| 416 | const result = await indexDoc(env.DB, embedder, store, repoDoc(row)); |
| 417 | files++; |
| 418 | if (result.capped) { |
| 419 | capped = true; |
| 420 | break; |
| 421 | } |
| 422 | next = `f:${row.space_id}\n${row.path}`; |
| 423 | } |
| 424 | if (!capped && rows.length < BACKFILL_BATCH) next = null; |
| 425 | } |
| 426 | await env.DB.prepare("UPDATE doc_index_runs SET cursor = ?, pages = pages + ?, files = files + ?, finished_at = ? WHERE workspace_id = ?") |
| 427 | .bind(next, pages, files, next === null ? new Date().toISOString() : null, workspaceId) |
| 428 | .run(); |
| 429 | if (next === null) return { done: true }; |
| 430 | if (chain && env.JOBS) await enqueue(env, workspaceId, capped ? CAP_DELAY_SECONDS : 0); |
| 431 | return { done: false }; |
| 432 | } |
| 433 | |
| 434 | function splitFileCursor(cursor: string): [string, string] { |
| 435 | const at = cursor.indexOf("\n"); |
| 436 | // Empty: from the start. |
| 437 | return at < 0 ? ["", ""] : [cursor.slice(0, at), cursor.slice(at + 1)]; |
| 438 | } |
| 439 | |
| 440 | // ── Folios (Artifacts mode) ───────────────────────────────────────────── |
| 441 | |
| 442 | let saidNoFolioIndex = false; |
| 443 | |
| 444 | /** |
| 445 | * The embedder and folio index this deployment has. Without Workers AI or |
| 446 | * the `g1t-folios` index (FOLIO_VECTORS, which is made by hand before a |
| 447 | * deploy that uses it), folios keep their passages in D1 and recall and |
| 448 | * search match words; that is said once per isolate. |
| 449 | */ |
| 450 | export function folioAdapters(env: IndexEnv): { embedder: Embedder | null; store: VectorStore | null } { |
| 451 | if (env.AI && env.FOLIO_VECTORS) return { embedder: cloudflareEmbedder(env.AI), store: cloudflareVectors(env.FOLIO_VECTORS) }; |
| 452 | if (!saidNoFolioIndex) { |
| 453 | saidNoFolioIndex = true; |
| 454 | console.log(`folios: no ${env.AI ? "FOLIO_VECTORS (Vectorize g1t-folios)" : "AI"} binding, so artifacts are searched and recalled by their words only`); |
| 455 | } |
| 456 | return { embedder: null, store: null }; |
| 457 | } |
| 458 | |
| 459 | type FolioChunkRow = { id: string; seq: number; scope: string; hash: string; vector_hash: string | null }; |
| 460 | |
| 461 | /** |
| 462 | * Brings one folio's passages up to date, as `indexDoc` does for pages: |
| 463 | * rows and full text in D1, under the folio's index scope, and vectors |
| 464 | * for passages whose text or scope changed (a passage whose text the |
| 465 | * index already holds keeps its vector; only its metadata is written |
| 466 | * again). Forgets it when it is gone or in the trash. Never throws. |
| 467 | */ |
| 468 | export async function indexFolio(env: IndexEnv, folioId: string, now = new Date()): Promise<{ capped: boolean }> { |
| 469 | try { |
| 470 | const db = env.DB; |
| 471 | const row = await db.prepare(`SELECT ${FOLIO_COLUMNS.replace("'' AS text", "text")} FROM folios WHERE id = ?`).bind(folioId).first<FolioRow>(); |
| 472 | if (!row || row.trashed_at) { |
| 473 | await forgetFolios(env, [folioId]); |
| 474 | return { capped: false }; |
| 475 | } |
| 476 | const scope = (await scopesOf(db, [row])).get(row.id) ?? `folio:${row.acl_root}`; |
| 477 | const model = kindModel(row.kind); |
| 478 | const at = now.toISOString(); |
| 479 | const chunks = (model ? model.chunks(row.text, row.title) : chunkMarkdown(row.text, row.title)).map((c) => { |
| 480 | const embed = embedText(row.title, c); |
| 481 | return { ...c, id: chunkId(row.id, c.seq), embed, hash: textHash(embed) }; |
| 482 | }); |
| 483 | const old = (await db.prepare("SELECT id, seq, scope, hash, vector_hash FROM folio_chunks WHERE folio_id = ?").bind(row.id).all<FolioChunkRow>()).results; |
| 484 | const byId = new Map(old.map((r) => [r.id, r])); |
| 485 | const statements: D1PreparedStatement[] = []; |
| 486 | for (const c of chunks) { |
| 487 | const was = byId.get(c.id); |
| 488 | if (was && was.hash === c.hash && was.scope === scope) continue; |
| 489 | statements.push( |
| 490 | db |
| 491 | .prepare( |
| 492 | `INSERT INTO folio_chunks (id, workspace_id, folio_id, kind, scope, seq, heading, text, hash, vector_hash, updated_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, NULL, ?) |
| 493 | ON CONFLICT (id) DO UPDATE SET kind = excluded.kind, scope = excluded.scope, heading = excluded.heading, text = excluded.text, hash = excluded.hash, updated_at = excluded.updated_at`, |
| 494 | ) |
| 495 | .bind(c.id, row.workspace_id, row.id, row.kind, scope, c.seq, c.heading, c.text, c.hash, at), |
| 496 | db.prepare("DELETE FROM folio_chunks_fts WHERE chunk_id = ?").bind(c.id), |
| 497 | db.prepare("INSERT INTO folio_chunks_fts (chunk_id, scope, folio_id, heading, text) VALUES (?, ?, ?, ?, ?)").bind(c.id, scope, row.id, c.heading ?? "", c.text), |
| 498 | ); |
| 499 | } |
| 500 | const keep = new Set(chunks.map((c) => c.id)); |
| 501 | const gone = old.filter((r) => !keep.has(r.id)); |
| 502 | for (const r of gone) statements.push(db.prepare("DELETE FROM folio_chunks WHERE id = ?").bind(r.id), db.prepare("DELETE FROM folio_chunks_fts WHERE chunk_id = ?").bind(r.id)); |
| 503 | for (let i = 0; i < statements.length; i += 60) await db.batch(statements.slice(i, i + 60)); |
| 504 | |
| 505 | const { embedder, store } = folioAdapters(env); |
| 506 | if (!embedder || !store) return { capped: false }; |
| 507 | // What needs a vector written: text the index doesn't hold under this id, or a new scope. |
| 508 | const stale = chunks.filter((c) => { |
| 509 | const was = byId.get(c.id); |
| 510 | return !was || was.vector_hash !== c.hash || was.scope !== scope; |
| 511 | }); |
| 512 | let capped = false; |
| 513 | if (stale.length) { |
| 514 | const held = new Map<string, string>(); |
| 515 | for (const r of old) if (r.vector_hash) held.set(r.vector_hash, r.id); |
| 516 | const reuse = stale.filter((c) => held.has(c.hash)); |
| 517 | const fresh = stale.filter((c) => !held.has(c.hash)); |
| 518 | const budget = fresh.length ? await allowance(db, row.workspace_id, now) : 0; |
| 519 | const embedNow = fresh.slice(0, budget); |
| 520 | capped = embedNow.length < fresh.length; |
| 521 | if (capped) console.log("folios embedding cap reached", row.workspace_id, `${fresh.length - embedNow.length} passages wait for the next hour`); |
| 522 | const metadata: VectorMetadata = { workspace_id: row.workspace_id, scope, kind: row.kind, folio_id: row.id }; |
| 523 | const done: { id: string; hash: string }[] = []; |
| 524 | try { |
| 525 | const vectors: { id: string; values: number[]; metadata: VectorMetadata }[] = []; |
| 526 | if (reuse.length) { |
| 527 | const values = new Map((await store.get([...new Set(reuse.map((c) => held.get(c.hash)!))])).map((v) => [v.id, v.values])); |
| 528 | for (const c of reuse) { |
| 529 | const v = values.get(held.get(c.hash)!); |
| 530 | if (v) vectors.push({ id: c.id, values: v, metadata }); |
| 531 | } |
| 532 | } |
| 533 | if (embedNow.length) { |
| 534 | const embedded = await embedder.embed(embedNow.map((c) => c.embed)); |
| 535 | embedNow.forEach((c, i) => vectors.push({ id: c.id, values: embedded[i]!, metadata })); |
| 536 | await meter( |
| 537 | db, |
| 538 | row.workspace_id, |
| 539 | embedNow.length, |
| 540 | embedNow.reduce((n, c) => n + c.embed.length, 0), |
| 541 | now, |
| 542 | ); |
| 543 | } |
| 544 | if (vectors.length) await store.upsert(vectors); |
| 545 | const hashOf = new Map(chunks.map((c) => [c.id, c.hash])); |
| 546 | for (const v of vectors) done.push({ id: v.id, hash: hashOf.get(v.id)! }); |
| 547 | } catch (error) { |
| 548 | console.error("folios could not embed passages; the next save tries again", row.id, String(error)); |
| 549 | } |
| 550 | if (done.length) { |
| 551 | const updates = done.map((d) => db.prepare("UPDATE folio_chunks SET vector_hash = ? WHERE id = ?").bind(d.hash, d.id)); |
| 552 | for (let i = 0; i < updates.length; i += 60) await db.batch(updates.slice(i, i + 60)); |
| 553 | } |
| 554 | } |
| 555 | try { |
| 556 | if (gone.some((r) => r.vector_hash)) await store.delete(gone.filter((r) => r.vector_hash).map((r) => r.id)); |
| 557 | } catch (error) { |
| 558 | console.error("folios could not drop old passages from the index", row.id, String(error)); |
| 559 | } |
| 560 | if (capped) await startBackfill(env, row.workspace_id, { delaySeconds: CAP_DELAY_SECONDS }); |
| 561 | return { capped }; |
| 562 | } catch (error) { |
| 563 | console.error("folios could not index", folioId, String(error)); |
| 564 | return { capped: false }; |
| 565 | } |
| 566 | } |
| 567 | |
| 568 | /** Takes folios' passages out of D1 and the index. Never throws. */ |
| 569 | export async function forgetFolios(env: IndexEnv, ids: readonly string[]): Promise<void> { |
| 570 | if (!ids.length) return; |
| 571 | const db = env.DB; |
| 572 | try { |
| 573 | const found: string[] = []; |
| 574 | for (let i = 0; i < ids.length; i += 500) { |
| 575 | const rows = await db |
| 576 | .prepare("SELECT id FROM folio_chunks WHERE folio_id IN (SELECT value FROM json_each(?))") |
| 577 | .bind(JSON.stringify(ids.slice(i, i + 500))) |
| 578 | .all<{ id: string }>(); |
| 579 | found.push(...rows.results.map((r) => r.id)); |
| 580 | } |
| 581 | if (!found.length) return; |
| 582 | const { store } = folioAdapters(env); |
| 583 | if (store) { |
| 584 | try { |
| 585 | await store.delete(found); |
| 586 | } catch (error) { |
| 587 | console.error("folios could not drop passages from the index", found.length, String(error)); |
| 588 | } |
| 589 | } |
| 590 | const statements = found.flatMap((id) => [db.prepare("DELETE FROM folio_chunks WHERE id = ?").bind(id), db.prepare("DELETE FROM folio_chunks_fts WHERE chunk_id = ?").bind(id)]); |
| 591 | for (let i = 0; i < statements.length; i += 80) await db.batch(statements.slice(i, i + 80)); |
| 592 | } catch (error) { |
| 593 | console.error("folios could not forget passages", String(error)); |
| 594 | } |
| 595 | } |