Skip to content
595 linesCodeBlameRaw

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Docs index by meaning: passages of every page and project doc, embedded on save and recalled for agents; hybrid search for people1/**
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";
The docs service answers every artifacts call: docs can be made, listed, shared, moved, trashed, restored, searched, versioned and edited live in their own rooms, agents read, write and recall them only where their person and everyone in the conversation can, and folio events go out on the bus, while Docs' pages keep working as before.19import { scopesOf, type FolioRow, FOLIO_COLUMNS } from "./folios/access-store.ts";
20import { kindModel } from "./kinds/index.ts";
Docs index by meaning: passages of every page and project doc, embedded on save and recalled for agents; hybrid search for people21import { cloudflareEmbedder, cloudflareVectors, type Embedder, type VectorMetadata, type VectorStore } from "./vectors.ts";
22
23/** Passages embedded per workspace per hour at most. Logged when reached. */
24export const EMBED_PER_HOUR = 200;
25/** Documents one backfill job indexes before handing on to the next. */
26const BACKFILL_BATCH = 20;
27/** A backfill not finished after this long is taken to have died and may start again. */
28const BACKFILL_STALE_MS = 2 * 60 * 60 * 1000;
29/** How long a run waits when the hour's cap is reached. */
30const CAP_DELAY_SECONDS = 3600;
31
The docs service answers every artifacts call: docs can be made, listed, shared, moved, trashed, restored, searched, versioned and edited live in their own rooms, agents read, write and recall them only where their person and everyone in the conversation can, and folio events go out on the bus, while Docs' pages keep working as before.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 */
37export type DocsJob = { type: "docs.index"; workspace_id: string } | { type: "folios.reacl"; folio_id: string };
Docs index by meaning: passages of every page and project doc, embedded on save and recalled for agents; hybrid search for people38
The docs service answers every artifacts call: docs can be made, listed, shared, moved, trashed, restored, searched, versioned and edited live in their own rooms, agents read, write and recall them only where their person and everyone in the conversation can, and folio events go out on the bus, while Docs' pages keep working as before.39export 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};
Docs index by meaning: passages of every page and project doc, embedded on save and recalled for agents; hybrid search for people47
48/** The embedder and store this deployment has; null without them (words only). */
49export 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. */
54export 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
66type ChunkRow = { id: string; seq: number; space_id: string; hash: string; vector_hash: string | null };
67
68const hourOf = (now: Date) => now.toISOString().slice(0, 13);
69
70function 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. */
77async 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
82async 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 */
97export 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
159type 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). */
162async 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. */
223export 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
253type 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. */
256export 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
273type RepoFileForIndex = { space_id: string; path: string; title: string; markdown: string; workspace_id: string; repo_id: string };
274
275function 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 */
283export 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 */
318export 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
334async 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 */
347export async function ensureIndexed(env: IndexEnv, workspaceId: string): Promise<void> {
348 const row = await env.DB.prepare(
The docs service answers every artifacts call: docs can be made, listed, shared, moved, trashed, restored, searched, versioned and edited live in their own rooms, agents read, write and recall them only where their person and everyone in the conversation can, and folio events go out on the bus, while Docs' pages keep working as before.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",
Docs index by meaning: passages of every page and project doc, embedded on save and recalled for agents; hybrid search for people350 )
351 .bind(workspaceId)
The docs service answers every artifacts call: docs can be made, listed, shared, moved, trashed, restored, searched, versioned and edited live in their own rooms, agents read, write and recall them only where their person and everyone in the conversation can, and folio events go out on the bus, while Docs' pages keep working as before.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;
Docs index by meaning: passages of every page and project doc, embedded on save and recalled for agents; hybrid search for people354 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 */
362export 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;
The docs service answers every artifacts call: docs can be made, listed, shared, moved, trashed, restored, searched, versioned and edited live in their own rooms, agents read, write and recall them only where their person and everyone in the conversation can, and folio events go out on the bus, while Docs' pages keep working as before.371 if (cursor.startsWith("o:")) {
372 // Folios, after pages and before projects' docs.
Docs index by meaning: passages of every page and project doc, embedded on save and recalled for agents; hybrid search for people373 const rows = (
The docs service answers every artifacts call: docs can be made, listed, shared, moved, trashed, restored, searched, versioned and edited live in their own rooms, agents read, write and recall them only where their person and everyone in the conversation can, and folio events go out on the bus, while Docs' pages keep working as before.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 = (
Docs index by meaning: passages of every page and project doc, embedded on save and recalled for agents; hybrid search for people390 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 }
The docs service answers every artifacts call: docs can be made, listed, shared, moved, trashed, restored, searched, versioned and edited live in their own rooms, agents read, write and recall them only where their person and everyone in the conversation can, and folio events go out on the bus, while Docs' pages keep working as before.404 if (!capped && rows.length < BACKFILL_BATCH) next = "o:";
Docs index by meaning: passages of every page and project doc, embedded on save and recalled for agents; hybrid search for people405 } 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
434function 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}
The docs service answers every artifacts call: docs can be made, listed, shared, moved, trashed, restored, searched, versioned and edited live in their own rooms, agents read, write and recall them only where their person and everyone in the conversation can, and folio events go out on the bus, while Docs' pages keep working as before.439
440// ── Folios (Artifacts mode) ─────────────────────────────────────────────
441
442let 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 */
450export 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
459type 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 */
468export 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. */
569export 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}