| 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 docs service's own events queue). |
| 17 | */ |
| 18 | import { chunkId, chunkMarkdown, embedText, repoFileId, textHash } from "./chunks.ts"; |
| 19 | import { cloudflareEmbedder, cloudflareVectors, type Embedder, type VectorMetadata, type VectorStore } from "./vectors.ts"; |
| 20 | |
| 21 | /** Passages embedded per workspace per hour at most. Logged when reached. */ |
| 22 | export const EMBED_PER_HOUR = 200; |
| 23 | /** Documents one backfill job indexes before handing on to the next. */ |
| 24 | const BACKFILL_BATCH = 20; |
| 25 | /** A backfill not finished after this long is taken to have died and may start again. */ |
| 26 | const BACKFILL_STALE_MS = 2 * 60 * 60 * 1000; |
| 27 | /** How long a run waits when the hour's cap is reached. */ |
| 28 | const CAP_DELAY_SECONDS = 3600; |
| 29 | |
| 30 | /** A job on the queue: index more of a workspace's docs. */ |
| 31 | export type DocsJob = { type: "docs.index"; workspace_id: string }; |
| 32 | |
| 33 | export type IndexEnv = { DB: D1Database; AI?: Ai; VECTORS?: Vectorize; JOBS?: Queue<DocsJob> }; |
| 34 | |
| 35 | /** The embedder and store this deployment has; null without them (words only). */ |
| 36 | export function adapters(env: IndexEnv): { embedder: Embedder | null; store: VectorStore | null } { |
| 37 | return env.AI && env.VECTORS ? { embedder: cloudflareEmbedder(env.AI), store: cloudflareVectors(env.VECTORS) } : { embedder: null, store: null }; |
| 38 | } |
| 39 | |
| 40 | /** One document to index: a page, or a project's docs file. */ |
| 41 | export type IndexDoc = { |
| 42 | kind: "page" | "repo_file"; |
| 43 | /** The page's id, or the file's `rf_` id. */ |
| 44 | doc_id: string; |
| 45 | workspace_id: string; |
| 46 | space_id: string; |
| 47 | title: string; |
| 48 | markdown: string; |
| 49 | repo_id?: string | null; |
| 50 | path?: string | null; |
| 51 | }; |
| 52 | |
| 53 | type ChunkRow = { id: string; seq: number; space_id: string; hash: string; vector_hash: string | null }; |
| 54 | |
| 55 | const hourOf = (now: Date) => now.toISOString().slice(0, 13); |
| 56 | |
| 57 | function metadataOf(doc: IndexDoc): VectorMetadata { |
| 58 | return doc.kind === "page" |
| 59 | ? { workspace_id: doc.workspace_id, space_id: doc.space_id, kind: "page", page_id: doc.doc_id } |
| 60 | : { workspace_id: doc.workspace_id, space_id: doc.space_id, kind: "repo_file", repo_file_id: doc.doc_id, repo_id: doc.repo_id ?? "" }; |
| 61 | } |
| 62 | |
| 63 | /** Passages the workspace may still embed this hour. */ |
| 64 | async function allowance(db: D1Database, workspaceId: string, now: Date): Promise<number> { |
| 65 | const used = await db.prepare("SELECT chunks FROM doc_embed_usage WHERE workspace_id = ? AND hour = ?").bind(workspaceId, hourOf(now)).first<{ chunks: number }>(); |
| 66 | return Math.max(0, EMBED_PER_HOUR - (used?.chunks ?? 0)); |
| 67 | } |
| 68 | |
| 69 | async function meter(db: D1Database, workspaceId: string, chunks: number, chars: number, now: Date): Promise<void> { |
| 70 | if (!chunks) return; |
| 71 | await db |
| 72 | .prepare( |
| 73 | "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", |
| 74 | ) |
| 75 | .bind(workspaceId, hourOf(now), chunks, Math.ceil(chars / 4)) |
| 76 | .run(); |
| 77 | } |
| 78 | |
| 79 | /** |
| 80 | * Brings one document's passages up to date: rows and full text in D1, |
| 81 | * vectors for those that changed (within the hour's cap). `capped` when |
| 82 | * some were left for later. |
| 83 | */ |
| 84 | export async function indexDoc( |
| 85 | db: D1Database, |
| 86 | embedder: Embedder | null, |
| 87 | store: VectorStore | null, |
| 88 | doc: IndexDoc, |
| 89 | now = new Date(), |
| 90 | ): Promise<{ chunks: number; embedded: number; capped: boolean }> { |
| 91 | const at = now.toISOString(); |
| 92 | const chunks = chunkMarkdown(doc.markdown, doc.title).map((c) => { |
| 93 | const embed = embedText(doc.title, c); |
| 94 | return { ...c, id: chunkId(doc.doc_id, c.seq), embed, hash: textHash(embed) }; |
| 95 | }); |
| 96 | const column = doc.kind === "page" ? "page_id" : "repo_file_id"; |
| 97 | 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; |
| 98 | const byId = new Map(old.map((r) => [r.id, r])); |
| 99 | const statements: D1PreparedStatement[] = []; |
| 100 | for (const c of chunks) { |
| 101 | const was = byId.get(c.id); |
| 102 | if (was && was.hash === c.hash && was.space_id === doc.space_id) continue; |
| 103 | statements.push( |
| 104 | db |
| 105 | .prepare( |
| 106 | `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) |
| 107 | VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, NULL, ?) |
| 108 | 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`, |
| 109 | ) |
| 110 | .bind( |
| 111 | c.id, |
| 112 | doc.workspace_id, |
| 113 | doc.space_id, |
| 114 | doc.kind === "page" ? doc.doc_id : null, |
| 115 | doc.kind === "repo_file" ? doc.doc_id : null, |
| 116 | doc.repo_id ?? null, |
| 117 | doc.path ?? null, |
| 118 | c.seq, |
| 119 | c.heading, |
| 120 | c.text, |
| 121 | c.hash, |
| 122 | at, |
| 123 | ), |
| 124 | db.prepare("DELETE FROM doc_chunks_fts WHERE chunk_id = ?").bind(c.id), |
| 125 | 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), |
| 126 | ); |
| 127 | } |
| 128 | const keep = new Set(chunks.map((c) => c.id)); |
| 129 | const gone = old.filter((r) => !keep.has(r.id)); |
| 130 | for (const r of gone) { |
| 131 | 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)); |
| 132 | } |
| 133 | for (let i = 0; i < statements.length; i += 60) await db.batch(statements.slice(i, i + 60)); |
| 134 | |
| 135 | if (!store || !embedder) return { chunks: chunks.length, embedded: 0, capped: false }; |
| 136 | const result = await embedChanged(db, embedder, store, doc, chunks, old, byId, now); |
| 137 | // Last, so a passage that moved from a place now gone could take its vector first. |
| 138 | try { |
| 139 | if (gone.some((r) => r.vector_hash)) await store.delete(gone.filter((r) => r.vector_hash).map((r) => r.id)); |
| 140 | } catch (error) { |
| 141 | console.error("docs could not drop old passages from the index", doc.doc_id, String(error)); |
| 142 | } |
| 143 | return result; |
| 144 | } |
| 145 | |
| 146 | type IndexedChunk = { id: string; seq: number; heading: string | null; text: string; embed: string; hash: string }; |
| 147 | |
| 148 | /** Vectors for the passages that need one: reused when the index already holds their text, embedded otherwise (within the cap). */ |
| 149 | async function embedChanged( |
| 150 | db: D1Database, |
| 151 | embedder: Embedder, |
| 152 | store: VectorStore, |
| 153 | doc: IndexDoc, |
| 154 | chunks: IndexedChunk[], |
| 155 | old: ChunkRow[], |
| 156 | byId: Map<string, ChunkRow>, |
| 157 | now: Date, |
| 158 | ): Promise<{ chunks: number; embedded: number; capped: boolean }> { |
| 159 | // What needs a vector: a passage whose text the index doesn't hold under its id, or one whose space changed. |
| 160 | const stale = chunks.filter((c) => { |
| 161 | const was = byId.get(c.id); |
| 162 | return !was || was.vector_hash !== c.hash || was.space_id !== doc.space_id; |
| 163 | }); |
| 164 | if (!stale.length) return { chunks: chunks.length, embedded: 0, capped: false }; |
| 165 | // A passage that only moved (or changed space) has its vector already, under another id or this one. |
| 166 | const held = new Map<string, string>(); |
| 167 | for (const r of old) if (r.vector_hash) held.set(r.vector_hash, r.id); |
| 168 | const reuse = stale.filter((c) => held.has(c.hash)); |
| 169 | const fresh = stale.filter((c) => !held.has(c.hash)); |
| 170 | const budget = fresh.length ? await allowance(db, doc.workspace_id, now) : 0; |
| 171 | const embedNow = fresh.slice(0, budget); |
| 172 | const capped = embedNow.length < fresh.length; |
| 173 | if (capped) console.log("docs embedding cap reached", doc.workspace_id, `${fresh.length - embedNow.length} passages wait for the next hour`); |
| 174 | const metadata = metadataOf(doc); |
| 175 | const done: { id: string; hash: string }[] = []; |
| 176 | try { |
| 177 | const vectors: { id: string; values: number[]; metadata: VectorMetadata }[] = []; |
| 178 | if (reuse.length) { |
| 179 | const values = new Map((await store.get([...new Set(reuse.map((c) => held.get(c.hash)!))])).map((v) => [v.id, v.values])); |
| 180 | for (const c of reuse) { |
| 181 | const v = values.get(held.get(c.hash)!); |
| 182 | if (v) vectors.push({ id: c.id, values: v, metadata }); |
| 183 | } |
| 184 | } |
| 185 | if (embedNow.length) { |
| 186 | const embedded = await embedder.embed(embedNow.map((c) => c.embed)); |
| 187 | embedNow.forEach((c, i) => vectors.push({ id: c.id, values: embedded[i]!, metadata })); |
| 188 | await meter( |
| 189 | db, |
| 190 | doc.workspace_id, |
| 191 | embedNow.length, |
| 192 | embedNow.reduce((n, c) => n + c.embed.length, 0), |
| 193 | now, |
| 194 | ); |
| 195 | } |
| 196 | if (vectors.length) await store.upsert(vectors); |
| 197 | const hashOf = new Map(chunks.map((c) => [c.id, c.hash])); |
| 198 | for (const v of vectors) done.push({ id: v.id, hash: hashOf.get(v.id)! }); |
| 199 | } catch (error) { |
| 200 | console.error("docs could not embed passages; the next save tries again", doc.doc_id, String(error)); |
| 201 | } |
| 202 | if (done.length) { |
| 203 | const updates = done.map((d) => db.prepare("UPDATE doc_chunks SET vector_hash = ? WHERE id = ?").bind(d.hash, d.id)); |
| 204 | for (let i = 0; i < updates.length; i += 60) await db.batch(updates.slice(i, i + 60)); |
| 205 | } |
| 206 | return { chunks: chunks.length, embedded: embedNow.length, capped }; |
| 207 | } |
| 208 | |
| 209 | /** Takes passages out of D1 and the index: by page ids, file ids, or a whole space. */ |
| 210 | export async function forgetDocs(env: IndexEnv, which: { page_ids?: string[]; repo_file_ids?: string[]; space_id?: string }): Promise<void> { |
| 211 | const db = env.DB; |
| 212 | const { store } = adapters(env); |
| 213 | try { |
| 214 | const found: string[] = []; |
| 215 | const select = async (column: string, values: string[]) => { |
| 216 | for (let i = 0; i < values.length; i += 90) { |
| 217 | const part = values.slice(i, i + 90); |
| 218 | const rows = await db.prepare(`SELECT id FROM doc_chunks WHERE ${column} IN (${part.map(() => "?").join(",")})`).bind(...part).all<{ id: string }>(); |
| 219 | found.push(...rows.results.map((r) => r.id)); |
| 220 | } |
| 221 | }; |
| 222 | if (which.page_ids?.length) await select("page_id", which.page_ids); |
| 223 | if (which.repo_file_ids?.length) await select("repo_file_id", which.repo_file_ids); |
| 224 | if (which.space_id) await select("space_id", [which.space_id]); |
| 225 | if (!found.length) return; |
| 226 | if (store) { |
| 227 | try { |
| 228 | await store.delete(found); |
| 229 | } catch (error) { |
| 230 | console.error("docs could not drop passages from the index", found.length, String(error)); |
| 231 | } |
| 232 | } |
| 233 | 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)]); |
| 234 | for (let i = 0; i < statements.length; i += 80) await db.batch(statements.slice(i, i + 80)); |
| 235 | } catch (error) { |
| 236 | console.error("docs could not forget passages", String(error)); |
| 237 | } |
| 238 | } |
| 239 | |
| 240 | type PageForIndex = { id: string; workspace_id: string; space_id: string; title: string; markdown: string; archived_at: string | null }; |
| 241 | |
| 242 | /** Indexes one page as it is saved now; forgets it when it is gone or in the trash. Never throws. */ |
| 243 | export async function indexPage(env: IndexEnv, pageId: string, now = new Date()): Promise<{ capped: boolean }> { |
| 244 | try { |
| 245 | const page = await env.DB.prepare("SELECT id, workspace_id, space_id, title, markdown, archived_at FROM pages WHERE id = ?").bind(pageId).first<PageForIndex>(); |
| 246 | if (!page || page.archived_at) { |
| 247 | await forgetDocs(env, { page_ids: [pageId] }); |
| 248 | return { capped: false }; |
| 249 | } |
| 250 | const { embedder, store } = adapters(env); |
| 251 | 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); |
| 252 | if (result.capped) await startBackfill(env, page.workspace_id, { delaySeconds: CAP_DELAY_SECONDS }); |
| 253 | return { capped: result.capped }; |
| 254 | } catch (error) { |
| 255 | console.error("docs could not index a page", pageId, String(error)); |
| 256 | return { capped: false }; |
| 257 | } |
| 258 | } |
| 259 | |
| 260 | type RepoFileForIndex = { space_id: string; path: string; title: string; markdown: string; workspace_id: string; repo_id: string }; |
| 261 | |
| 262 | function repoDoc(row: RepoFileForIndex): IndexDoc { |
| 263 | 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 }; |
| 264 | } |
| 265 | |
| 266 | /** |
| 267 | * Indexes a project's docs files after they were read: `changed` paths |
| 268 | * again, `gone` paths forgotten. Never throws. |
| 269 | */ |
| 270 | export async function indexRepoFiles(env: IndexEnv, spaceId: string, changed: string[], gone: string[], now = new Date()): Promise<void> { |
| 271 | try { |
| 272 | if (gone.length) await forgetDocs(env, { repo_file_ids: gone.map((p) => repoFileId(spaceId, p)) }); |
| 273 | if (!changed.length) return; |
| 274 | const { embedder, store } = adapters(env); |
| 275 | let capped = false; |
| 276 | let workspaceId: string | null = null; |
| 277 | for (let i = 0; i < changed.length; i += 50) { |
| 278 | const part = changed.slice(i, i + 50); |
| 279 | const rows = ( |
| 280 | await env.DB.prepare( |
| 281 | `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(",")})`, |
| 282 | ) |
| 283 | .bind(spaceId, ...part) |
| 284 | .all<RepoFileForIndex>() |
| 285 | ).results; |
| 286 | for (const row of rows) { |
| 287 | workspaceId = row.workspace_id; |
| 288 | // Past the cap, the rest are only kept for words until the catch-up run. |
| 289 | const result = await indexDoc(env.DB, capped ? null : embedder, capped ? null : store, repoDoc(row), now); |
| 290 | capped ||= result.capped; |
| 291 | } |
| 292 | } |
| 293 | if (capped && workspaceId) await startBackfill(env, workspaceId, { delaySeconds: CAP_DELAY_SECONDS }); |
| 294 | } catch (error) { |
| 295 | console.error("docs could not index a project's docs", spaceId, String(error)); |
| 296 | } |
| 297 | } |
| 298 | |
| 299 | /** |
| 300 | * Starts (or, with `force`, restarts) indexing a workspace's existing |
| 301 | * pages and projects' docs on the queue. Does nothing while a run is |
| 302 | * going, unless it has been going so long it must have died. Without a |
| 303 | * queue, runs one batch now. True when a run was started. |
| 304 | */ |
| 305 | export async function startBackfill(env: IndexEnv, workspaceId: string, options: { force?: boolean; delaySeconds?: number } = {}): Promise<boolean> { |
| 306 | const at = new Date(); |
| 307 | const staleBefore = new Date(at.getTime() - BACKFILL_STALE_MS).toISOString(); |
| 308 | const started = await env.DB.prepare( |
| 309 | `INSERT INTO doc_index_runs (workspace_id, started_at, finished_at, cursor, pages, files) VALUES (?1, ?2, NULL, NULL, 0, 0) |
| 310 | ON CONFLICT (workspace_id) DO UPDATE SET started_at = ?2, finished_at = NULL, cursor = NULL, pages = 0, files = 0 |
| 311 | WHERE ?3 OR doc_index_runs.finished_at IS NOT NULL OR doc_index_runs.started_at < ?4 |
| 312 | RETURNING workspace_id`, |
| 313 | ) |
| 314 | .bind(workspaceId, at.toISOString(), options.force ? 1 : 0, staleBefore) |
| 315 | .first<{ workspace_id: string }>(); |
| 316 | if (!started) return false; |
| 317 | await enqueue(env, workspaceId, options.delaySeconds ?? 0); |
| 318 | return true; |
| 319 | } |
| 320 | |
| 321 | async function enqueue(env: IndexEnv, workspaceId: string, delaySeconds: number): Promise<void> { |
| 322 | if (env.JOBS) { |
| 323 | await env.JOBS.send({ type: "docs.index", workspace_id: workspaceId }, delaySeconds ? { delaySeconds } : undefined); |
| 324 | return; |
| 325 | } |
| 326 | // No queue (self-hosted without one): one batch now; saves and later recalls carry on from there. |
| 327 | if (!delaySeconds) await runBackfill(env, workspaceId, false); |
| 328 | } |
| 329 | |
| 330 | /** |
| 331 | * The first time anyone recalls from a workspace's docs: when it has |
| 332 | * pages or projects' docs and has never been indexed, start the backfill. |
| 333 | */ |
| 334 | export async function ensureIndexed(env: IndexEnv, workspaceId: string): Promise<void> { |
| 335 | const row = await env.DB.prepare( |
| 336 | "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 repo_spaces WHERE workspace_id = ?1 LIMIT 1) AS repo", |
| 337 | ) |
| 338 | .bind(workspaceId) |
| 339 | .first<{ ran: number | null; page: number | null; repo: number | null }>(); |
| 340 | if (!row || row.ran || (!row.page && !row.repo)) return; |
| 341 | await startBackfill(env, workspaceId); |
| 342 | } |
| 343 | |
| 344 | /** |
| 345 | * One batch of a workspace's backfill: pages by id, then projects' docs |
| 346 | * files by space and path, from the run's cursor. Hands on to the next |
| 347 | * batch (or waits out the hour's cap) through the queue. |
| 348 | */ |
| 349 | export async function runBackfill(env: IndexEnv, workspaceId: string, chain = true): Promise<{ done: boolean }> { |
| 350 | 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 }>(); |
| 351 | if (!run || run.finished_at) return { done: true }; |
| 352 | const { embedder, store } = adapters(env); |
| 353 | const cursor = run.cursor ?? "p:"; |
| 354 | let next: string | null = cursor; |
| 355 | let pages = 0; |
| 356 | let files = 0; |
| 357 | let capped = false; |
| 358 | if (cursor.startsWith("p:")) { |
| 359 | const rows = ( |
| 360 | 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 ?") |
| 361 | .bind(workspaceId, cursor.slice(2), BACKFILL_BATCH) |
| 362 | .all<PageForIndex>() |
| 363 | ).results; |
| 364 | for (const page of rows) { |
| 365 | 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 }); |
| 366 | pages++; |
| 367 | if (result.capped) { |
| 368 | // This page again next time, after the hour: the cursor stays before it. |
| 369 | capped = true; |
| 370 | break; |
| 371 | } |
| 372 | next = `p:${page.id}`; |
| 373 | } |
| 374 | if (!capped && rows.length < BACKFILL_BATCH) next = "f:"; |
| 375 | } else { |
| 376 | const [space, path] = splitFileCursor(cursor.slice(2)); |
| 377 | const rows = ( |
| 378 | await env.DB.prepare( |
| 379 | `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 |
| 380 | WHERE s.workspace_id = ? AND (f.space_id > ? OR (f.space_id = ? AND f.path > ?)) ORDER BY f.space_id, f.path LIMIT ?`, |
| 381 | ) |
| 382 | .bind(workspaceId, space, space, path, BACKFILL_BATCH) |
| 383 | .all<RepoFileForIndex>() |
| 384 | ).results; |
| 385 | for (const row of rows) { |
| 386 | const result = await indexDoc(env.DB, embedder, store, repoDoc(row)); |
| 387 | files++; |
| 388 | if (result.capped) { |
| 389 | capped = true; |
| 390 | break; |
| 391 | } |
| 392 | next = `f:${row.space_id}\n${row.path}`; |
| 393 | } |
| 394 | if (!capped && rows.length < BACKFILL_BATCH) next = null; |
| 395 | } |
| 396 | await env.DB.prepare("UPDATE doc_index_runs SET cursor = ?, pages = pages + ?, files = files + ?, finished_at = ? WHERE workspace_id = ?") |
| 397 | .bind(next, pages, files, next === null ? new Date().toISOString() : null, workspaceId) |
| 398 | .run(); |
| 399 | if (next === null) return { done: true }; |
| 400 | if (chain && env.JOBS) await enqueue(env, workspaceId, capped ? CAP_DELAY_SECONDS : 0); |
| 401 | return { done: false }; |
| 402 | } |
| 403 | |
| 404 | function splitFileCursor(cursor: string): [string, string] { |
| 405 | const at = cursor.indexOf("\n"); |
| 406 | // Empty: from the start. |
| 407 | return at < 0 ? ["", ""] : [cursor.slice(0, at), cursor.slice(at + 1)]; |
| 408 | } |