Skip to content
408 linesCodeBlameRaw
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 */
18import { chunkId, chunkMarkdown, embedText, repoFileId, textHash } from "./chunks.ts";
19import { cloudflareEmbedder, cloudflareVectors, type Embedder, type VectorMetadata, type VectorStore } from "./vectors.ts";
20
21/** Passages embedded per workspace per hour at most. Logged when reached. */
22export const EMBED_PER_HOUR = 200;
23/** Documents one backfill job indexes before handing on to the next. */
24const BACKFILL_BATCH = 20;
25/** A backfill not finished after this long is taken to have died and may start again. */
26const BACKFILL_STALE_MS = 2 * 60 * 60 * 1000;
27/** How long a run waits when the hour's cap is reached. */
28const CAP_DELAY_SECONDS = 3600;
29
30/** A job on the queue: index more of a workspace's docs. */
31export type DocsJob = { type: "docs.index"; workspace_id: string };
32
33export 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). */
36export 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. */
41export 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
53type ChunkRow = { id: string; seq: number; space_id: string; hash: string; vector_hash: string | null };
54
55const hourOf = (now: Date) => now.toISOString().slice(0, 13);
56
57function 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. */
64async 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
69async 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 */
84export 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
146type 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). */
149async 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. */
210export 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
240type 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. */
243export 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
260type RepoFileForIndex = { space_id: string; path: string; title: string; markdown: string; workspace_id: string; repo_id: string };
261
262function 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 */
270export 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 */
305export 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
321async 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 */
334export 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 */
349export 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
404function 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}