Skip to content

Commit

Docs index by meaning: passages of every page and project doc, embedded on save and recalled for agents; hybrid search for people

- Pages and repository docs are split into passages under their headings and embedded (Workers AI, Vectorize g1t-docs) after they settle, only what changed, at most 200 passages a workspace an hour; words (FTS5) always - recall_for_agent: the passages closest in meaning, only from spaces the person and everyone reading can read (repository docs only where they all can read), required reading first, two passages a page at most, words filling in, stale pages marked - Existing pages are indexed the first time a workspace is recalled from; owners can reindex - Docs search mixes words and meaning and shows the matching section - Embedder and VectorStore adapters; self-hosted installs search by words

syntaqxcommitted Parentd7ec3a0Browse files
20 files+1978−620/20 viewed
+37−3
200200 thread. The page is yours: you're its owner, and it links back to where it
201201 came from.
202202
203+### How agents find your docs
204+
205+You don't have to point an agent at the right page. Before an agent answers
206+in chat, and as it works through each step of a session, it recalls the
207+passages of your Docs closest in meaning to what it's been asked: a few
208+sections of pages (and of projects' docs), each with the page and heading it
209+came from, so it can follow your runbooks and decisions and link you to them.
210+When nothing in Docs is about the question, it recalls nothing.
211+
212+- **Only what everyone in the conversation can read.** In a DM or a private
213+ channel, that's the spaces every person there can read; in a public
214+ channel, only spaces the whole workspace can read. Private spaces stay
215+ private: their pages never reach a conversation that includes someone
216+ outside them. Projects' docs come only from repositories everyone there
217+ can read (in a public channel, only public repositories).
218+- **Required reading first.** An agent can be given spaces to read first.
219+ Recall looks there before the rest of the workspace, so a support agent
220+ leans on the support runbooks.
221+- **Current, not cached.** A page is indexed within about half a minute of
222+ a change, without slowing the editor, and pages in the trash, or in an archived space, are never recalled. A page
223+ marked possibly out of date is recalled with that warning, so the agent
224+ says so instead of trusting it.
225+- **Meaning, then words.** Recall matches by meaning, so "how do we roll
226+ back?" finds a section titled "Reverting a deploy". When meaning finds
227+ too little, passages with the question's words fill in.
228+
229+An agent can still read a whole page, search, or list a space's pages when
230+recall isn't enough; recall is the head start.
231+
203232 ### Write a thread up
204233
205234 In chat, **⋯ → Write this up in Docs** on a message or an open thread asks
322351 ## Search and projects
323352
324353 The search box at the top of the sidebar, and **Search docs** on Docs' home,
325−look through the titles and text of every page you can read. Narrow it to one
326−space or one project.
354+look through every page you can read, and the projects' docs shown in Docs, by
355+their words and by what they mean. Type words, or ask a question: "how do we
356+rotate the API key?" finds the section that explains it even when it's
357+worded differently. Each result shows the passage that matched, under its
358+heading. Narrow it to one space or one project.
359+
360+Linking to a page from the editor searches titles and text as you type.
327361
328362 A page belongs to a project when the page, or its space, is linked to the
329363 project. On Docs' home, **All projects** filters the recently edited pages and
330364 the spaces to one project.
331365
332366 Agents search the same way, and only find what the people they answer can
333−read.
367+read. See [How agents find your docs](#how-agents-find-your-docs).
334368
335369 ## Projects' docs
336370
+12−4
2424 const space = url.searchParams.get("space") || null;
2525 const project = url.searchParams.get("project") || null;
2626 const found = q.trim()
27− ? await docs.search(params.owner.toLowerCase(), viewer, { query: q, space_id: space, project, limit: 50 }).catch(() => null)
27+ ? await docs.search(params.owner.toLowerCase(), viewer, { query: q, space_id: space, project, limit: 50, mode: "hybrid" }).catch(() => null)
2828 : null;
2929 return { q, space, project, hits: found?.ok ? found.value : q.trim() ? null : ([] as DocSearchHit[]) };
3030 }
3131
32−/** Full text over every page the viewer can read; by space and by project. */
32+/**
33+ * Every page the viewer can read, and the projects' docs shown in Docs,
34+ * by words and by meaning at once; by space and by project. Each hit shows
35+ * the passage that matched, under its heading.
36+ */
3337 export default function DocsSearch({ loaderData, params }: Route.ComponentProps) {
3438 const slug = params.owner.toLowerCase();
3539 const { q, space, project, hits } = loaderData;
4448 <Form method="get" id="docs-search" className="mt-4 flex flex-wrap gap-2">
4549 <label className="flex h-10 min-w-0 grow items-center gap-2 rounded-md border border-line bg-surface px-3 focus-within:border-accent/60">
4650 <Search size={16} className="text-faint" aria-hidden="true" />
47− <input name="q" defaultValue={q} autoFocus placeholder="Words in a title or a page" aria-label="Search" className="min-w-0 grow bg-transparent text-sm outline-none placeholder:text-faint" />
51+ <input name="q" defaultValue={q} autoFocus placeholder="Words, or a question" aria-label="Search" className="min-w-0 grow bg-transparent text-sm outline-none placeholder:text-faint" />
4852 </label>
4953 <SelectField
5054 name="space"
6266 {hits === null ? (
6367 <EmptyState title="Search didn't answer">Try again in a moment.</EmptyState>
6468 ) : !q.trim() ? (
65− <p className="text-sm text-muted">Search titles and text across every space you can read, and the projects&apos; docs shown in Docs. Agents search the same way, and only see what the people they answer can see.</p>
69+ <p className="text-sm text-muted">
70+ Search every space you can read, and the projects&apos; docs shown in Docs, by their words and by what they mean: ask a question and the passage that answers it comes up even when it words it differently. Agents recall from
71+ Docs the same way, and only from what everyone they answer can read.
72+ </p>
6673 ) : hits.length === 0 ? (
6774 <EmptyState title="Nothing found">No page you can read matches &ldquo;{q}&rdquo;.</EmptyState>
6875 ) : (
7784 {hit.space_name} · <TimeAgo at={hit.updated_at} />
7885 </span>
7986 </span>
87+ {hit.heading && <span className="mt-0.5 block truncate text-xs font-medium text-muted">{hit.heading}</span>}
8088 <span className="mt-1 line-clamp-2 block text-xs leading-relaxed text-muted">
8189 {snippetParts(hit.snippet).map((part, i) =>
8290 part.match ? (
+2−1
139139 "setup": [
140140 "The D1 database, before the first deploy: npx wrangler d1 create g1t-docs, then put its id in services/docs/wrangler.jsonc",
141141 "The R2 bucket for files in pages: npx wrangler r2 bucket create g1t-docs-files",
142− "The queue the events service sends it merges and pushes on (pages whose cited code changed, projects' docs): npx wrangler queues create g1t-events-docs"
142+ "The queue the events service sends it merges and pushes on (pages whose cited code changed, projects' docs): npx wrangler queues create g1t-events-docs. The service also sends its own backfill jobs to it (JOBS)",
143+ "The Vectorize index agents recall Docs from: npx wrangler vectorize create g1t-docs --dimensions=768 --metric=cosine, with string metadata indexes on workspace_id and space_id (npx wrangler vectorize create-metadata-index g1t-docs --property-name=<name> --type=string)"
143144 ],
144145 "self_host": "run"
145146 },
+12−0
315315 - Off: the service already degrades to keyword search when `AI` or
316316 `VECTORS` is missing (`index.ts:919`), so phase 3 can first run context
317317 with neither bound.
318+- **Docs' semantic index** (`services/docs`, what agents recall from)
319+ talks only to an `Embedder` and a `VectorStore`
320+ (`services/docs/src/vectors.ts`). Only Cloudflare's are built today:
321+ Workers AI (`AI`, `@cf/baai/bge-base-en-v1.5`) and Vectorize (`VECTORS`,
322+ the `g1t-docs` index, filtered by `workspace_id` and `space_id`).
323+ Another model or vector database fits behind the same two interfaces
324+ (sqlite-vec with an OpenAI-compatible embeddings endpoint, as for the
325+ context hub, is the likely default), but none is written yet. Without
326+ them Docs still keeps every passage in D1 with full text, so agents'
327+ recall and the Docs search page match words instead of meaning. The
328+ backfill and catch-up jobs ride the docs service's events queue
329+ (`JOBS`); without a queue one batch runs at a time as recall asks.
318330
319331 ### Files in Docs pages
320332
+66−18
461461
462462 This is where Docs earns its place:
463463
464−- **Agents read it.** Pages are indexed by `services/context`, so the
465− knowledge reaches every reply, session and plan. A space can be pinned to
466− an agent as required reading.
464+- **Agents read it.** Pages and projects' docs are split into passages
465+ and indexed by meaning in Docs' own semantic index, and agents recall
466+ the most relevant passages on every reply and session step, only from
467+ spaces everyone in the conversation can read. A space can be pinned to
468+ an agent as required reading, which recall looks in first.
467469 - **Agents write it.** An agent edits pages directly where the space lets
468470 agents edit and the person it acts for can edit. Otherwise the edit
469471 becomes a **suggestion**: tracked changes a person accepts or rejects
583585 Not built yet: the documenter agent (it reads `stalePagesForAgent` and
584586 updates with `marks_current`; the routine and its `doc.page.stale` trigger
585587 are the agents service's), editing a project's docs from Docs as a pull
586−request, and indexing pages in `services/context`.
588+request.
587589
588−**Pages in `services/context`: what it needs.** The context hub indexes
589−items per project, and decides who sees one by the project's privacy
590−(`private` and membership). A Docs page's access is per space (private and
591−team spaces, listed members), which that model can't express, so indexing
592−pages there as they are would show private spaces' pages to every member.
593−Doing it right needs: an ingest RPC on context for documents with an
594−access key (`docs:<space id>`), a check at query time that asks the docs
595−service which spaces the reader (and audience) can read
596−(`spacesForAgent` already answers that), and a `doc.page.updated`
597−consumer that re-embeds the page's Markdown (`SUBSCRIBER_CONTEXT` would
598−route `doc.page.*`). Until then agents reach pages through the docs tools
599−(`searchForAgent`, `pageMarkdown`), which apply exactly those rules.
600−Projects' docs folders are already in context: it reads each project's
601−`docs/` on every push.
590+**The semantic index: how agents find what Docs say.** Docs keeps its
591+own index rather than putting pages in `services/context`, because a
592+page's access is per space (private and team spaces, listed members),
593+which the context hub's per-project privacy can't express. Built
594+(`services/docs` src/chunks.ts, src/indexer.ts, src/recall.ts,
595+src/vectors.ts):
596+
597+- **Passages.** A page's derived Markdown, and each projects' docs file,
598+ is split by heading into passages of about 300 to 1,500 characters
599+ (long sections split on paragraph boundaries, a code block kept whole,
600+ tiny sections joined to the next), each with its heading path
601+ ("Runbook › Rollback"). They live in D1 (`doc_chunks`, with FTS5 over
602+ heading and text in `doc_chunks_fts`), and as vectors in Vectorize
603+ (`g1t-docs`, 768 dimensions, cosine), embedded by Workers AI
604+ (`@cf/baai/bge-base-en-v1.5`) with the title and heading in front.
605+ Vector ids are `<page or file id>:<seq>`; metadata is `workspace_id`
606+ and `space_id` (both indexed), `kind`, and `page_id`, or `repo_file_id`
607+ and `repo_id`. A projects' docs file's space is its repo space, and its
608+ id `rf_<hash of space and path>`.
609+- **Indexing never slows editing.** The page's room indexes it half a
610+ minute after its Markdown first changes, on its own alarm. Only
611+ passages whose text changed are embedded; one that only moved keeps its
612+ vector. Creating, renaming, moving and restoring a page index it; the
613+ trash, deleting, removing a project's docs and a purged repository take
614+ passages out. A project's docs files are indexed after each read (on
615+ adding them, and after each push to the default branch).
616+- **A cap, and catching up.** At most 200 passages are embedded per
617+ workspace per hour (`doc_embed_usage`, which also counts characters as
618+ an estimate of tokens); past it the rest are kept for words and a
619+ catch-up run embeds them the next hour. Any failure leaves passages for
620+ the next save or run; saving never waits on, or fails for, the index.
621+- **Backfill.** A run (`doc_index_runs`) indexes a workspace's existing
622+ pages and files 20 at a time, as `docs.index` jobs on the docs
623+ service's own events queue. It starts by itself the first time an agent
624+ recalls from a workspace that has pages but was never indexed, and an
625+ owner can start it again (`reindexDocs`).
626+- **Recall** (`recallForAgent`, RPC `recall_for_agent`). The spaces an
627+ agent may read for the viewer and audience come from the same check as
628+ every other agent read (`agentSpaces`, behind `spacesForAgent` and
629+ `searchForAgent`); projects' docs only from repositories the viewer can
630+ read and, for a DM or private channel, every person in it (a public
631+ channel, or more than 20 people: public repositories only). The query is
632+ embedded once (kept a minute per isolate) and the index asked for the
633+ 24 nearest passages within those spaces (by a `$in` filter on
634+ `space_id`, or for more than 40 spaces by workspace, 50 of them,
635+ filtered after). Passages below a cosine similarity of 0.6 are dropped;
636+ at most two per page; required spaces (`spaces`) are asked separately
637+ and come first; when meaning finds fewer than `limit` (5 by default, 10
638+ at most), passages matching any of the query's words fill in (score
639+ 0.5). Each passage comes with its page or file, heading, and whether the
640+ page is possibly out of date. Archived pages and spaces are never
641+ returned: what comes back is checked against D1 as it is now.
642+- **People's search** asks the same index: the Docs search page is hybrid
643+ (`mode: "hybrid"`), word hits and meaning hits fused by reciprocal rank,
644+ each with the passage and heading that matched. Search as you type (page
645+ links) stays words only.
646+- **Adapters.** Docs talks to an `Embedder` and a `VectorStore`
647+ (src/vectors.ts), so a self-hosted g1t could use another model or vector
648+ database. Only Cloudflare's are built. Without them, passages are still
649+ kept and recall matches words.
602650
603651 ## Chat
604652
+17−0
384384 projects: string[];
385385 /** Set for a file from a project's docs folder (then `id` is `repo:<space>:<path>` and `path` its address in Docs). */
386386 repo_file?: { repo: string; path: string } | null;
387+ /** The heading of the passage that matched, when search found one (`mode: "hybrid"`). */
388+ heading?: string | null;
389+ /** How it was found: by its words, by meaning (the semantic index), or both. */
390+ matched?: "words" | "meaning" | "both" | null;
387391 };
388392
389393 export type DocTemplate = {
486490 /** `owner/name`: pages linked to it, or in a space linked to it. */
487491 project?: string | null;
488492 limit?: number | null;
493+ /**
494+ * `words` (the default): titles and text by their words, as you type.
495+ * `hybrid`: words and meaning (the semantic index) together, each hit
496+ * with the passage and heading that matched; the Docs search page.
497+ */
498+ mode?: "words" | "hybrid" | null;
489499 };
490500
491501 /**
631641 removeRepoSpace(workspace: string, viewer: User, id: string): Promise<Result<boolean>>;
632642 /** One file of a project's docs, for a viewer who can read the repository; not found otherwise. */
633643 repoPage(workspace: string, viewer: User, repo: string, path: string): Promise<Result<DocRepoPage>>;
644+ /**
645+ * Indexes the workspace's pages and projects' docs for agents' recall
646+ * again, in batches in the background. Workspace owners. True when a run
647+ * started (false: one is already going).
648+ */
649+ reindexDocs(workspace: string, viewer: User): Promise<Result<boolean>>;
634650
635651 // ── Agents (services/agents) ─────────────────────────────────────────
636652 //
757773 addRepoSpace: (workspace, viewer, repo) => call("add_repo_space", { workspace, viewer, repo }),
758774 removeRepoSpace: (workspace, viewer, id) => call("remove_repo_space", { workspace, viewer, id }),
759775 repoPage: (workspace, viewer, repo, path) => call("repo_page", { workspace, viewer, repo, path }),
776+ reindexDocs: (workspace, viewer) => call("reindex_docs", { workspace, viewer }),
760777 spacesForAgent: (workspace, agentId, viewer, audience) => call("spaces_for_agent", { workspace, agent_id: agentId, viewer, audience: audience ?? null }),
761778 pageMarkdown: (workspace, agentId, viewer, pageId, audience) =>
762779 call("page_markdown", { workspace, agent_id: agentId, viewer, page_id: pageId, audience: audience ?? null }),
+58−0
1+-- Passages: pages and projects' docs files split by heading, for the
2+-- semantic index (Vectorize `g1t-docs`) and agents' recall
3+-- (docs/WORKSPACE.md, "Agents and docs"; services/docs src/chunks.ts,
4+-- src/indexer.ts).
5+
6+-- One passage. `id` is `<page or file id>:<seq>`, the same as its
7+-- vector's. `hash` is of what is embedded (title, heading, text);
8+-- `vector_hash` is the hash the index holds a vector for, NULL when it
9+-- holds none yet: the two differ until the passage is embedded, which a
10+-- later save or the backfill retries. For a project's docs file,
11+-- `space_id` is the repo space's id, `repo_file_id` is `rf_<hash of space
12+-- and path>`, and `path` the file's.
13+CREATE TABLE doc_chunks (
14+ id TEXT PRIMARY KEY,
15+ workspace_id TEXT NOT NULL,
16+ space_id TEXT NOT NULL,
17+ page_id TEXT,
18+ repo_file_id TEXT,
19+ repo_id TEXT,
20+ path TEXT,
21+ seq INTEGER NOT NULL,
22+ heading TEXT,
23+ text TEXT NOT NULL,
24+ hash TEXT NOT NULL,
25+ vector_hash TEXT,
26+ updated_at TEXT NOT NULL,
27+ CHECK ((page_id IS NULL) <> (repo_file_id IS NULL))
28+);
29+CREATE INDEX doc_chunks_page ON doc_chunks (page_id, seq) WHERE page_id IS NOT NULL;
30+CREATE INDEX doc_chunks_file ON doc_chunks (repo_file_id, seq) WHERE repo_file_id IS NOT NULL;
31+CREATE INDEX doc_chunks_space ON doc_chunks (space_id);
32+CREATE INDEX doc_chunks_workspace ON doc_chunks (workspace_id);
33+CREATE INDEX doc_chunks_unembedded ON doc_chunks (workspace_id) WHERE vector_hash IS NULL OR vector_hash <> hash;
34+
35+-- Full text over passages: recall's word fallback, and the heading a
36+-- search hit sits under.
37+CREATE VIRTUAL TABLE doc_chunks_fts USING fts5 (chunk_id UNINDEXED, space_id UNINDEXED, doc_id UNINDEXED, heading, text, tokenize = 'unicode61 remove_diacritics 2');
38+
39+-- Passages embedded per workspace and hour: the cap (src/indexer.ts
40+-- EMBED_PER_HOUR) and what it cost (`tokens`, estimated as characters / 4).
41+CREATE TABLE doc_embed_usage (
42+ workspace_id TEXT NOT NULL,
43+ hour TEXT NOT NULL,
44+ chunks INTEGER NOT NULL DEFAULT 0,
45+ tokens INTEGER NOT NULL DEFAULT 0,
46+ PRIMARY KEY (workspace_id, hour)
47+);
48+
49+-- A workspace's backfill: (re)indexing its existing pages and projects'
50+-- docs in batches on the queue. `cursor` is where the next batch starts.
51+CREATE TABLE doc_index_runs (
52+ workspace_id TEXT PRIMARY KEY,
53+ started_at TEXT NOT NULL,
54+ finished_at TEXT,
55+ cursor TEXT,
56+ pages INTEGER NOT NULL DEFAULT 0,
57+ files INTEGER NOT NULL DEFAULT 0
58+);
+117−0
1+import assert from "node:assert/strict";
2+import { test } from "node:test";
3+
4+import { EMBED_CHARS, MAX_CHARS, MAX_CHUNKS, MIN_CHARS, chunkId, chunkMarkdown, embedText, repoFileId, textHash } from "./chunks.ts";
5+
6+const para = (words: number, word = "rollback") => Array.from({ length: words }, (_, i) => `${word}${i % 7}`).join(" ") + ".";
7+
8+test("passages follow headings, with the heading path", () => {
9+ const md = [
10+ "Intro words that set the scene for the runbook. " + para(40, "intro"),
11+ "",
12+ "# Runbook",
13+ "",
14+ "## Deploy",
15+ "",
16+ para(60, "deploy"),
17+ "",
18+ "## Rollback",
19+ "",
20+ para(60, "rollback"),
21+ "",
22+ "### Database",
23+ "",
24+ para(60, "database"),
25+ ].join("\n");
26+ const chunks = chunkMarkdown(md, "Operations");
27+ assert.deepEqual(
28+ chunks.map((c) => c.heading),
29+ [null, "Runbook › Deploy", "Runbook › Rollback", "Runbook › Rollback › Database"],
30+ );
31+ assert.deepEqual(
32+ chunks.map((c) => c.seq),
33+ [0, 1, 2, 3],
34+ );
35+ assert.match(chunks[2]!.text, /^rollback0/);
36+ // The heading isn't repeated in the text.
37+ assert.ok(!chunks[2]!.text.includes("## Rollback"));
38+});
39+
40+test("long sections split on paragraphs, never past the most", () => {
41+ const md = ["## Big", "", ...Array.from({ length: 12 }, () => para(40)).flatMap((p) => [p, ""])].join("\n");
42+ const chunks = chunkMarkdown(md, "T");
43+ assert.ok(chunks.length > 1);
44+ for (const c of chunks) {
45+ assert.ok(c.text.length <= MAX_CHARS, `${c.text.length}`);
46+ assert.equal(c.heading, "Big");
47+ }
48+ // Nothing lost.
49+ const words = (t: string) => t.split(/\s+/).filter(Boolean).length;
50+ assert.equal(words(chunks.map((c) => c.text).join("\n\n")), words(md) - 2);
51+});
52+
53+test("one huge paragraph is cut at sentences", () => {
54+ const huge = Array.from({ length: 80 }, (_, i) => `Sentence number ${i} explains a part of the system.`).join(" ");
55+ const chunks = chunkMarkdown(huge, "T");
56+ assert.ok(chunks.length >= 3);
57+ for (const c of chunks) assert.ok(c.text.length <= MAX_CHARS);
58+ assert.ok(chunks[0]!.text.endsWith("."));
59+});
60+
61+test("tiny sections join the next, naming its heading", () => {
62+ const md = ["## A", "", "Short.", "", "## B", "", "Also short.", "", "## C", "", para(80)].join("\n");
63+ const chunks = chunkMarkdown(md, "T");
64+ assert.equal(chunks[0]!.heading, "A");
65+ assert.match(chunks[0]!.text, /Short\.\n\n\*\*B\*\*\n\nAlso short\./);
66+ for (const c of chunks.slice(0, -1)) assert.ok(c.text.length >= MIN_CHARS || chunks.length === 1);
67+});
68+
69+test("a tiny last section joins the one before", () => {
70+ const md = ["## A", "", para(60), "", "## B", "", "The end."].join("\n");
71+ const chunks = chunkMarkdown(md, "T");
72+ assert.equal(chunks.length, 1);
73+ assert.match(chunks[0]!.text, /\*\*B\*\*\n\nThe end\.$/);
74+});
75+
76+test("headings inside code fences are code, and fences stay whole", () => {
77+ const code = ["```sh", "# not a heading", ...Array.from({ length: 10 }, (_, i) => `echo step ${i}`), "```"].join("\n");
78+ const md = ["## Script", "", para(30), "", code, "", para(30)].join("\n");
79+ const chunks = chunkMarkdown(md, "T");
80+ assert.ok(chunks.every((c) => c.heading === "Script"));
81+ const withCode = chunks.find((c) => c.text.includes("```sh"))!;
82+ assert.ok(withCode.text.includes("# not a heading"));
83+ assert.equal((withCode.text.match(/```/g) ?? []).length, 2);
84+});
85+
86+test("a file's front matter and leading title are not passages", () => {
87+ const md = ["---", "title: Setup", "---", "# Setup", "", para(60, "install")].join("\n");
88+ const chunks = chunkMarkdown(md, "Setup");
89+ assert.equal(chunks.length, 1);
90+ assert.equal(chunks[0]!.heading, null);
91+ assert.ok(!chunks[0]!.text.includes("title:"));
92+});
93+
94+test("empty documents have no passages; huge ones stop at the most", () => {
95+ assert.deepEqual(chunkMarkdown("", "T"), []);
96+ assert.deepEqual(chunkMarkdown("\n\n \n", "T"), []);
97+ const many = Array.from({ length: MAX_CHUNKS + 30 }, (_, i) => `## S${i}\n\n${para(70)}`).join("\n\n");
98+ assert.equal(chunkMarkdown(many, "T").length, MAX_CHUNKS);
99+});
100+
101+test("what is embedded leads with the title and heading", () => {
102+ assert.equal(embedText("Ops", { heading: "Runbook › Rollback", text: "Do this." }), "Ops › Runbook › Rollback\n\nDo this.");
103+ assert.equal(embedText("", { heading: null, text: "Just text." }), "Just text.");
104+ assert.equal(embedText("T", { heading: null, text: "x".repeat(5000) }).length, EMBED_CHARS);
105+});
106+
107+test("hashes and ids are stable and short", () => {
108+ assert.equal(textHash("a"), textHash("a"));
109+ assert.notEqual(textHash("a"), textHash("b"));
110+ assert.match(textHash("anything"), /^[0-9a-f]{16}$/);
111+ const id = repoFileId("rds_1", "docs/a/very/long/path/that/goes/on/and/on/README.md");
112+ assert.equal(id, repoFileId("rds_1", "docs/a/very/long/path/that/goes/on/and/on/README.md"));
113+ assert.notEqual(id, repoFileId("rds_2", "docs/a/very/long/path/that/goes/on/and/on/README.md"));
114+ // A vector id is at most 64 bytes.
115+ assert.ok(chunkId(id, 149).length <= 64);
116+ assert.equal(chunkId("pag_1", 3), "pag_1:3");
117+});
+219−0
1+/**
2+ * Passages: a page's Markdown (or a repository's docs file) split into
3+ * sections an agent can be handed whole, for the semantic index
4+ * (src/indexer.ts) and recall. Pure.
5+ *
6+ * A passage is one section under its heading path ("Runbook › Rollback"),
7+ * between MIN_CHARS and MAX_CHARS where it can be: long sections split on
8+ * paragraph boundaries (a fenced code block stays one paragraph), tiny ones
9+ * join the next. What is embedded is the passage with its document's title
10+ * and heading in front (`embedText`); what an agent reads is the passage.
11+ */
12+
13+/** A passage shorter than this joins the next one. */
14+export const MIN_CHARS = 300;
15+/** A passage longer than this splits on paragraphs. */
16+export const MAX_CHARS = 1500;
17+/** The most passages one document keeps; past it the rest goes unindexed (words still find the page). */
18+export const MAX_CHUNKS = 150;
19+/** The most text embedded for one passage: the model reads 512 tokens. */
20+export const EMBED_CHARS = 2000;
21+/** Between heading levels in a passage's heading path. */
22+export const HEADING_SEPARATOR = " › ";
23+
24+export type Chunk = {
25+ seq: number;
26+ /** The heading path the passage sits under, or null before the first heading. */
27+ heading: string | null;
28+ /** The passage as Markdown. */
29+ text: string;
30+};
31+
32+type Section = { heading: string | null; lines: string[] };
33+
34+const FENCE = /^\s{0,3}(`{3,}|~{3,})/;
35+const HEADING = /^\s{0,3}(#{1,6})\s+(.+?)\s*#*\s*$/;
36+
37+/** A heading's text as plain words: no emphasis, code ticks or links. */
38+function headingText(raw: string): string {
39+ return raw
40+ .replace(/!\[([^\]]*)\]\([^)]*\)/g, "$1")
41+ .replace(/\[([^\]]*)\]\([^)]*\)/g, "$1")
42+ .replace(/[*_`~]+/g, "")
43+ .replace(/\s+/g, " ")
44+ .trim();
45+}
46+
47+/** Drops YAML front matter, as repository docs files often have. */
48+function withoutFrontMatter(markdown: string): string {
49+ const front = /^---\r?\n[\s\S]*?\r?\n---[ \t]*(\r?\n|$)/.exec(markdown);
50+ return front ? markdown.slice(front[0].length) : markdown;
51+}
52+
53+/** The document split at its headings, each section with its heading path. */
54+function sections(markdown: string, title: string): Section[] {
55+ const lines = withoutFrontMatter(markdown).replace(/\r\n?/g, "\n").split("\n");
56+ const out: Section[] = [{ heading: null, lines: [] }];
57+ const path: { level: number; text: string }[] = [];
58+ let fence: string | null = null;
59+ let first = true;
60+ const wantTitle = headingText(title).toLowerCase();
61+ for (const line of lines) {
62+ const f = FENCE.exec(line);
63+ if (fence) {
64+ if (f && f[1]![0] === fence[0] && f[1]!.length >= fence.length && line.trim() === f[1]) fence = null;
65+ out[out.length - 1]!.lines.push(line);
66+ continue;
67+ }
68+ if (f) {
69+ fence = f[1]!;
70+ out[out.length - 1]!.lines.push(line);
71+ continue;
72+ }
73+ const h = HEADING.exec(line);
74+ if (!h) {
75+ if (line.trim()) first = false;
76+ out[out.length - 1]!.lines.push(line);
77+ continue;
78+ }
79+ const level = h[1]!.length;
80+ const text = headingText(h[2]!);
81+ // A file's leading `# Title` is its title, not a section of it.
82+ if (first && level === 1 && text.toLowerCase() === wantTitle) {
83+ first = false;
84+ continue;
85+ }
86+ first = false;
87+ while (path.length && path[path.length - 1]!.level >= level) path.pop();
88+ path.push({ level, text });
89+ out.push({ heading: path.map((p) => p.text).filter(Boolean).join(HEADING_SEPARATOR) || null, lines: [] });
90+ }
91+ return out;
92+}
93+
94+/** Paragraphs: runs of lines between blank lines, with a fenced block kept whole. */
95+function paragraphs(lines: string[]): string[] {
96+ const out: string[] = [];
97+ let current: string[] = [];
98+ let fence: string | null = null;
99+ const flush = () => {
100+ const text = current.join("\n").trim();
101+ if (text) out.push(text);
102+ current = [];
103+ };
104+ for (const line of lines) {
105+ const f = FENCE.exec(line);
106+ if (fence) {
107+ current.push(line);
108+ if (f && f[1]![0] === fence[0] && f[1]!.length >= fence.length && line.trim() === f[1]) fence = null;
109+ continue;
110+ }
111+ if (f) {
112+ fence = f[1]!;
113+ current.push(line);
114+ continue;
115+ }
116+ if (!line.trim()) flush();
117+ else current.push(line);
118+ }
119+ flush();
120+ return out;
121+}
122+
123+/** A paragraph longer than MAX_CHARS, cut at sentence ends or spaces. */
124+function cut(text: string): string[] {
125+ const out: string[] = [];
126+ let rest = text;
127+ while (rest.length > MAX_CHARS) {
128+ const window = rest.slice(0, MAX_CHARS);
129+ let at = Math.max(window.lastIndexOf(". "), window.lastIndexOf(".\n"), window.lastIndexOf("\n"));
130+ if (at < MAX_CHARS / 2) at = window.lastIndexOf(" ");
131+ if (at < MAX_CHARS / 2) at = MAX_CHARS - 1;
132+ out.push(rest.slice(0, at + 1).trim());
133+ rest = rest.slice(at + 1).trim();
134+ }
135+ if (rest) out.push(rest);
136+ return out;
137+}
138+
139+/** One section's passages: its paragraphs packed up to MAX_CHARS. */
140+function pack(paras: string[]): string[] {
141+ const out: string[] = [];
142+ let current = "";
143+ for (const para of paras.flatMap(cut)) {
144+ if (!current) current = para;
145+ else if (current.length + 2 + para.length <= MAX_CHARS) current += `\n\n${para}`;
146+ else {
147+ out.push(current);
148+ current = para;
149+ }
150+ }
151+ if (current) out.push(current);
152+ return out;
153+}
154+
155+/**
156+ * A document's passages, in order. `title` is the page's or file's title:
157+ * not part of any passage, but a file's leading `# Title` is dropped as it.
158+ */
159+export function chunkMarkdown(markdown: string, title = ""): Chunk[] {
160+ const pieces: { heading: string | null; text: string }[] = [];
161+ for (const section of sections(String(markdown ?? ""), title)) {
162+ for (const text of pack(paragraphs(section.lines))) pieces.push({ heading: section.heading, text });
163+ // A heading with nothing under it but subsections still names them: nothing to keep alone.
164+ }
165+ // Tiny passages join the next one (its heading written in), while that fits.
166+ const merged: { heading: string | null; text: string }[] = [];
167+ for (const piece of pieces) {
168+ const last = merged[merged.length - 1];
169+ if (last && last.text.length < MIN_CHARS) {
170+ const joined = piece.heading && piece.heading !== last.heading ? `${last.text}\n\n**${piece.heading.split(HEADING_SEPARATOR).pop()}**\n\n${piece.text}` : `${last.text}\n\n${piece.text}`;
171+ if (joined.length <= MAX_CHARS) {
172+ last.text = joined;
173+ if (!last.heading) last.heading = piece.heading;
174+ continue;
175+ }
176+ }
177+ merged.push({ ...piece });
178+ }
179+ // A tiny last passage joins the one before it, when that fits.
180+ if (merged.length > 1) {
181+ const last = merged[merged.length - 1]!;
182+ const before = merged[merged.length - 2]!;
183+ if (last.text.length < MIN_CHARS && before.text.length + last.text.length + 40 <= MAX_CHARS) {
184+ before.text = last.heading && last.heading !== before.heading ? `${before.text}\n\n**${last.heading.split(HEADING_SEPARATOR).pop()}**\n\n${last.text}` : `${before.text}\n\n${last.text}`;
185+ merged.pop();
186+ }
187+ }
188+ return merged.slice(0, MAX_CHUNKS).map((piece, seq) => ({ seq, heading: piece.heading, text: piece.text }));
189+}
190+
191+/** What is embedded for a passage: its document's title and heading in front, so a passage alone still says what it is about. */
192+export function embedText(title: string, chunk: Pick<Chunk, "heading" | "text">): string {
193+ const head = [title.trim(), chunk.heading ?? ""].filter(Boolean).join(HEADING_SEPARATOR);
194+ return (head ? `${head}\n\n${chunk.text}` : chunk.text).slice(0, EMBED_CHARS);
195+}
196+
197+/** A short, stable hash (hex) of a string: whether a passage changed since it was embedded. Not for security. */
198+export function textHash(text: string): string {
199+ let h1 = 0xdeadbeef;
200+ let h2 = 0x41c6ce57;
201+ for (let i = 0; i < text.length; i++) {
202+ const c = text.charCodeAt(i);
203+ h1 = Math.imul(h1 ^ c, 2654435761);
204+ h2 = Math.imul(h2 ^ c, 1597334677);
205+ }
206+ h1 = Math.imul(h1 ^ (h1 >>> 16), 2246822507) ^ Math.imul(h2 ^ (h2 >>> 13), 3266489909);
207+ h2 = Math.imul(h2 ^ (h2 >>> 16), 2246822507) ^ Math.imul(h1 ^ (h1 >>> 13), 3266489909);
208+ return (h2 >>> 0).toString(16).padStart(8, "0") + (h1 >>> 0).toString(16).padStart(8, "0");
209+}
210+
211+/** A repository docs file's id in the index (`rf_` and 32 hex): short enough for a vector id, the same each time. */
212+export function repoFileId(spaceId: string, path: string): string {
213+ return `rf_${textHash(`${spaceId}\n${path}`)}${textHash(`${path}\n${spaceId}`)}`;
214+}
215+
216+/** A passage's id, in D1 and in the vector index: `<page or file id>:<seq>`. */
217+export function chunkId(docId: string, seq: number): string {
218+ return `${docId}:${seq}`;
219+}
+453−24
4747 type DocPageChange,
4848 type DocPageDetail,
4949 type DocPageRef,
50+ type DocPassage,
5051 type DocRole,
5152 type DocSearchHit,
5253 type DocSearchQuery,
8283 import { cleanDescribes } from "./citations.ts";
8384 import { publishDocEvent } from "./events.ts";
8485 import { fileStore, safeName, servedType, type FileStoreEnv } from "./files.ts";
86+import { repoFileId } from "./chunks.ts";
87+import { adapters, ensureIndexed, forgetDocs, indexPage, indexRepoFiles, runBackfill, startBackfill, type DocsJob } from "./indexer.ts";
8588 import { excerpt, searchText } from "./markdown.ts";
89+import { QueryCache, fuseRanks, pickPassages, queryKey, recallLimit, requiredSpaces, vectorQueryPlan, MEANING_FLOOR, WORDS_SCORE, type Candidate } from "./recall.ts";
8690 import { ROOM_MEMBER_HEADER, type Origin, type PageRoom, type RoomMember } from "./room.ts";
8791 import { indexRepoSpace, reindexRepo, type RepoSpaceRow } from "./repo-spaces.ts";
88−import { ftsQuery, inProject, projectRef, searchSpaces } from "./search.ts";
92+import { ftsAnyQuery, ftsQuery, inProject, projectRef, searchSpaces } from "./search.ts";
8993 import { freeSlug, pageSlug, validSpaceSlug } from "./slugs.ts";
9094 import { BUILTIN_TEMPLATES, builtinTemplate } from "./templates.ts";
9195 import type { ThreadResult } from "./threads.ts";
106110 /** The bus: `doc.page.*` events (src/events.ts). */
107111 EVENTS?: ServiceBinding;
108112 PAGES: DurableObjectNamespace<PageRoom>;
113+ /** Workers AI: embeds passages and queries for the semantic index (src/vectors.ts). Without it, words only. */
114+ AI?: Ai;
115+ /** The semantic index, Vectorize `g1t-docs` (src/vectors.ts, src/indexer.ts). */
116+ VECTORS?: Vectorize;
117+ /** The docs service's own events queue, also carrying its backfill jobs (`docs.index`, src/indexer.ts). */
118+ JOBS?: Queue<DocsJob>;
109119 };
110120
121+/** Queries' embeddings, a minute per isolate (src/recall.ts). */
122+const queryVectors = new QueryCache();
123+
111124 type SpaceRow = {
112125 id: string;
113126 workspace_id: string;
174187
175188 type VersionRow = { id: string; page_id: string; created_at: string; kind: DocVersion["kind"]; authors: string; note: string | null; markdown: string; state: ArrayBuffer | null };
176189
190+/** A passage as recall and search read it back (`passages`): its page's or file's title, and the page's space now. */
191+type PassageRow = {
192+ id: string;
193+ page_id: string | null;
194+ repo_file_id: string | null;
195+ path: string | null;
196+ heading: string | null;
197+ text: string;
198+ updated_at: string;
199+ space_id: string;
200+ title: string | null;
201+ icon: string | null;
202+ page_updated_at: string | null;
203+};
204+
177205 /** A space, with who is in it and the viewer's role. */
178206 type Space = { row: SpaceRow; members: { principal: string; role: DocRole }[]; projects: string[]; role: DocRole | null };
179207
673701 .first<{ id: string }>();
674702 if (!inserted) return fail("conflict", `${ref}'s docs are already in Docs.`);
675703 try {
676− await indexRepoSpace({ DB: this.db, REPOS: this.env.REPOS }, row);
704+ const read = await indexRepoSpace({ DB: this.db, REPOS: this.env.REPOS }, row);
705+ this.defer(indexRepoFiles(this.env, row.id, read.changed, read.gone));
677706 } catch (error) {
678707 console.error("docs could not read a project's docs", row.repo, String(error));
679708 }
689718 if (!row) return fail("not_found", "No such project's docs.");
690719 if (row.added_by !== this.userKey(a.viewer!) && !this.viewerOwner(a.viewer!, a.workspace)) return fail("forbidden", "Only whoever added a project's docs, or an owner, can remove them.");
691720 await this.db.batch([this.db.prepare("DELETE FROM repo_files_fts WHERE space_id = ?").bind(row.id), this.db.prepare("DELETE FROM repo_spaces WHERE id = ?").bind(row.id)]);
721+ this.defer(forgetDocs(this.env, { space_id: row.id }));
692722 return ok(true);
693723 }
694724
12311261 ]);
12321262 await this.room(id).ensure({ page_id: id, workspace_slug: workspace.slug, markdown: input.markdown, state: input.state ?? null });
12331263 this.defer(publishDocEvent(this.env.EVENTS, "doc.page.created", this.eventData(workspace, space, row), author));
1264+ // Its passages, for agents' recall (src/indexer.ts); later edits are indexed by its room.
1265+ if (input.markdown.trim()) this.defer(indexPage(this.env, id));
12341266 return row;
12351267 }
12361268
13091341 const after = await this.db.prepare("SELECT * FROM pages WHERE id = ?").bind(page.id).first<PageRow>();
13101342 const [detail] = await this.toPages(workspace, new Map([[space.row.id, space.row]]), [after!]);
13111343 this.tell(page.id, { type: "page.updated", page: detail! });
1344+ // The title is part of what each passage is embedded with.
1345+ if (c.title !== undefined && cleanTitle(c.title) !== page.title) this.defer(indexPage(this.env, page.id));
13121346 return ok(detail!);
13131347 }
13141348
13301364 const placed = placeBefore(targetRows.results, page.id, parent, move.before_id ?? null);
13311365 const statements: D1PreparedStatement[] = [this.db.prepare("UPDATE pages SET parent_id = ?, position = ? WHERE id = ?").bind(parent, placed.position, page.id)];
13321366 for (const [id, position] of placed.renumber) statements.push(this.db.prepare("UPDATE pages SET position = ? WHERE id = ?").bind(position, id));
1367+ let moved: string[] = [];
13331368 if (target.row.id !== page.space_id) {
13341369 // The page and everything under it move to the other space.
13351370 const all = (await this.db.prepare("SELECT id, parent_id, position FROM pages WHERE space_id = ?").bind(page.space_id).all<PageRow>()).results;
1336− for (const id of descendants(all, page.id)) statements.push(this.db.prepare("UPDATE pages SET space_id = ? WHERE id = ?").bind(target.row.id, id));
1371+ moved = descendants(all, page.id);
1372+ for (const id of moved) statements.push(this.db.prepare("UPDATE pages SET space_id = ? WHERE id = ?").bind(target.row.id, id));
13371373 }
13381374 await this.db.batch(statements);
1375+ // Their passages are filed under the new space (no new embeddings: they only moved).
1376+ if (moved.length) this.defer(this.reindexPages(moved));
13391377 const after = await this.db.prepare("SELECT * FROM pages WHERE id = ?").bind(page.id).first<PageRow>();
13401378 const [detail] = await this.toPages(workspace, new Map(spaces.map((s) => [s.row.id, s.row])), [after!]);
13411379 this.tell(page.id, { type: "page.updated", page: detail! });
13721410 const at = now();
13731411 await this.db.batch(ids.map((id) => this.db.prepare("UPDATE pages SET archived_at = ?, archived_by = ? WHERE id = ? AND archived_at IS NULL").bind(at, this.userKey(a.viewer!), id)));
13741412 for (const id of ids) this.defer(this.room(id).closeAll("Moved to the trash").catch(() => undefined));
1413+ // Out of agents' recall while in the trash; restoring indexes them again.
1414+ this.defer(forgetDocs(this.env, { page_ids: ids }));
13751415 this.defer(publishDocEvent(this.env.EVENTS, "doc.page.archived", this.eventData(workspace, space.row, page), this.userKey(a.viewer!)));
13761416 const after = await this.db.prepare("SELECT * FROM pages WHERE id = ?").bind(page.id).first<PageRow>();
13771417 const [detail] = await this.toPages(workspace, new Map([[space.row.id, space.row]]), [after!]);
13911431 const statements = ids.map((id) => this.db.prepare("UPDATE pages SET archived_at = NULL, archived_by = NULL WHERE id = ?").bind(id));
13921432 if (parentGone) statements.push(this.db.prepare("UPDATE pages SET parent_id = NULL WHERE id = ?").bind(page.id));
13931433 await this.db.batch(statements);
1434+ this.defer(this.reindexPages(ids));
13941435 const after = await this.db.prepare("SELECT * FROM pages WHERE id = ?").bind(page.id).first<PageRow>();
13951436 const [detail] = await this.toPages(workspace, new Map([[space.row.id, space.row]]), [after!]);
13961437 return ok(detail!);
14041445 const all = (await this.db.prepare("SELECT id, parent_id, position FROM pages WHERE space_id = ?").bind(page.space_id).all<PageRow>()).results;
14051446 const ids = descendants(all, page.id);
14061447 await this.db.batch(ids.flatMap((id) => [this.db.prepare("DELETE FROM pages_fts WHERE page_id = ?").bind(id), this.db.prepare("DELETE FROM pages WHERE id = ?").bind(id)]));
1448+ this.defer(forgetDocs(this.env, { page_ids: ids }));
14071449 return ok(true);
14081450 }
14091451
15021544 const workspace = found.value;
15031545 const query = a.query ?? { query: "" };
15041546 const spaces = (await this.spacesFor(workspace, a.viewer!)).filter((s) => s.role);
1505− const [pages, files] = await Promise.all([
1547+ const hybrid = query.mode === "hybrid" && !!ftsQuery(query.query);
1548+ // A project's docs, when the search isn't narrowed to one of the workspace's spaces.
1549+ const repoSpaces = query.space_id || !ftsQuery(query.query)
1550+ ? []
1551+ : await this.repoSpacesMatching(workspace, a.viewer!, query).catch((error: unknown) => {
1552+ console.error("docs could not list projects' docs for search", String(error));
1553+ return [] as { row: RepoSpaceRow; repo: Repo }[];
1554+ });
1555+ const [pages, files, meaning] = await Promise.all([
15061556 this.searchIn(workspace, spaces, query),
1507− // A project's docs, when the search isn't narrowed to one of the workspace's spaces.
1508− query.space_id
1509− ? Promise.resolve([] as DocSearchHit[])
1510− : this.searchRepoFiles(workspace, a.viewer!, query).catch((error: unknown) => {
1511− console.error("docs could not search projects' docs", String(error));
1512− return [] as DocSearchHit[];
1513− }),
1557+ this.searchRepoFiles(workspace, repoSpaces, query).catch((error: unknown) => {
1558+ console.error("docs could not search projects' docs", String(error));
1559+ return [] as DocSearchHit[];
1560+ }),
1561+ hybrid
1562+ ? this.meaningHits(workspace, spaces, repoSpaces, query).catch((error: unknown) => {
1563+ console.error("docs could not search by meaning", String(error));
1564+ return null;
1565+ })
1566+ : Promise.resolve(null),
15141567 ]);
15151568 const limit = Math.min(Math.max(Number(query.limit) || 20, 1), 50);
1516− // Pages first, then files, as many as asked for.
1517− return ok([...pages, ...files].slice(0, limit));
1569+ // Words only: pages first, then files, as many as asked for.
1570+ if (!hybrid) return ok([...pages, ...files].slice(0, limit));
1571+ return ok(await this.fuseHits(workspace, spaces, repoSpaces, [...pages, ...files], meaning ?? [], query, limit));
1572+ }
1573+
1574+ /** The projects' docs the viewer can read, narrowed to the search's project. */
1575+ private async repoSpacesMatching(workspace: Workspace, viewer: User, query: DocSearchQuery): Promise<{ row: RepoSpaceRow; repo: Repo }[]> {
1576+ const spaces = await this.readableRepoSpaces(workspace, viewer);
1577+ const project = query.project ? projectRef(query.project) : null;
1578+ return project ? spaces.filter((s) => `${s.repo.namespace}/${s.repo.name}`.toLowerCase() === project) : spaces;
1579+ }
1580+
1581+ /** A search hit's key: a page's id, or `repo:<space>:<path>` for a project's docs file. */
1582+ private hitKey(row: Pick<PassageRow, "page_id" | "space_id" | "path">): string {
1583+ return row.page_id ?? `repo:${row.space_id}:${row.path}`;
1584+ }
1585+
1586+ /** By meaning: each page's or file's closest passage above the floor, closest first. */
1587+ private async meaningHits(workspace: Workspace, spaces: Space[], repoSpaces: { row: RepoSpaceRow }[], query: DocSearchQuery): Promise<{ key: string; row: PassageRow; score: number }[]> {
1588+ const allowed = [...searchSpaces(spaces.map((s) => s.row.id), query.space_id ?? null), ...repoSpaces.map((r) => r.row.id)];
1589+ if (!allowed.length) return [];
1590+ const vector = await this.queryVector(query.query);
1591+ if (!vector) return [];
1592+ const matches = (await this.meaningMatches(workspace.id, allowed, vector)).filter((m) => m.score >= MEANING_FLOOR);
1593+ const rows = await this.passages(workspace.id, matches.map((m) => m.id));
1594+ const may = new Set(allowed);
1595+ const best = new Map<string, { key: string; row: PassageRow; score: number }>();
1596+ for (const m of matches) {
1597+ const row = rows.get(m.id);
1598+ if (!row || !may.has(row.space_id)) continue;
1599+ const key = this.hitKey(row);
1600+ if ((best.get(key)?.score ?? -1) < m.score) best.set(key, { key, row, score: m.score });
1601+ }
1602+ return [...best.values()].sort((a, b) => b.score - a.score);
15181603 }
15191604
1520− /** Full text over the projects' docs the viewer can read. */
1521− private async searchRepoFiles(workspace: Workspace, viewer: User, query: DocSearchQuery): Promise<DocSearchHit[]> {
1605+ /**
1606+ * Hybrid search's answer: word hits and meaning hits fused by rank, each
1607+ * with the passage that matched and its heading. Word hits get theirs
1608+ * from the passages' full text; meaning-only hits show their passage.
1609+ */
1610+ private async fuseHits(
1611+ workspace: Workspace,
1612+ spaces: Space[],
1613+ repoSpaces: { row: RepoSpaceRow; repo: Repo }[],
1614+ words: DocSearchHit[],
1615+ meaning: { key: string; row: PassageRow; score: number }[],
1616+ query: DocSearchQuery,
1617+ limit: number,
1618+ ): Promise<DocSearchHit[]> {
1619+ const order = fuseRanks(
1620+ words.map((h) => h.id),
1621+ meaning.map((m) => m.key),
1622+ );
1623+ const byWords = new Map(words.map((h) => [h.id, h]));
1624+ const byMeaning = new Map(meaning.map((m) => [m.key, m]));
1625+ // The passage each word hit matched in, for its heading and a closer snippet.
1626+ const docIds = new Map(words.map((h) => [h.repo_file ? repoFileId(h.space_id, h.repo_file.path) : h.id, h.id]));
1627+ const passageOf = new Map<string, { heading: string; snippet: string }>();
15221628 const q = ftsQuery(query.query);
1523− if (!q) return [];
1524− let spaces = await this.readableRepoSpaces(workspace, viewer);
1629+ if (q && docIds.size) {
1630+ // The best-ranked 90, within D1's bound parameters.
1631+ const ids = [...docIds.keys()].slice(0, 90);
1632+ const found = await this.db
1633+ .prepare(
1634+ `SELECT doc_id, heading, snippet(doc_chunks_fts, 4, '[[', ']]', '…', 16) AS snippet FROM doc_chunks_fts
1635+ WHERE doc_chunks_fts MATCH ? AND doc_id IN (${ids.map(() => "?").join(",")}) ORDER BY bm25(doc_chunks_fts, 0, 0, 0, 4.0, 1.0) LIMIT 200`,
1636+ )
1637+ .bind(q, ...ids)
1638+ .all<{ doc_id: string; heading: string; snippet: string }>()
1639+ .catch(() => ({ results: [] as { doc_id: string; heading: string; snippet: string }[] }));
1640+ for (const r of found.results) {
1641+ const key = docIds.get(r.doc_id);
1642+ if (key && !passageOf.has(key)) passageOf.set(key, { heading: r.heading, snippet: r.snippet });
1643+ }
1644+ }
1645+ // Meaning-only pages: their projects, for the project filter and the hit.
1646+ const onlyMeaning = meaning.filter((m) => !byWords.has(m.key) && m.row.page_id);
1647+ const pageIds = onlyMeaning.map((m) => m.row.page_id!);
1648+ const projects = pageIds.length
1649+ ? (
1650+ await this.db
1651+ .prepare(`SELECT page_id, repo FROM page_projects WHERE page_id IN (${pageIds.map(() => "?").join(",")})`)
1652+ .bind(...pageIds)
1653+ .all<{ page_id: string; repo: string }>()
1654+ ).results
1655+ : [];
15251656 const project = query.project ? projectRef(query.project) : null;
1526− if (project) spaces = spaces.filter((s) => `${s.repo.namespace}/${s.repo.name}`.toLowerCase() === project);
1657+ const bySpace = new Map(spaces.map((s) => [s.row.id, s]));
1658+ const byRepo = new Map(repoSpaces.map((r) => [r.row.id, r]));
1659+ const out: DocSearchHit[] = [];
1660+ for (const key of order) {
1661+ if (out.length >= limit) break;
1662+ const w = byWords.get(key);
1663+ const m = byMeaning.get(key);
1664+ if (w) {
1665+ const passage = passageOf.get(key);
1666+ out.push({
1667+ ...w,
1668+ snippet: passage?.snippet || w.snippet,
1669+ heading: (passage ? passage.heading || null : null) ?? m?.row.heading ?? null,
1670+ matched: m ? "both" : "words",
1671+ });
1672+ continue;
1673+ }
1674+ if (!m) continue;
1675+ const row = m.row;
1676+ const snippet = excerpt(row.text, 200);
1677+ if (row.page_id) {
1678+ const space = bySpace.get(row.space_id);
1679+ if (!space) continue;
1680+ const own = projects.filter((p) => p.page_id === row.page_id).map((p) => p.repo);
1681+ if (!inProject(project, own, space.projects)) continue;
1682+ out.push({
1683+ ...this.ref(workspace.slug, space.row, { id: row.page_id, title: row.title ?? "", icon: row.icon }),
1684+ space_name: space.row.name,
1685+ snippet,
1686+ updated_at: row.page_updated_at ?? row.updated_at,
1687+ projects: [...new Set([...own, ...space.projects])],
1688+ heading: row.heading,
1689+ matched: "meaning",
1690+ });
1691+ } else {
1692+ const r = byRepo.get(row.space_id);
1693+ if (!r || !row.path) continue;
1694+ const repo = `${r.repo.namespace}/${r.repo.name}`;
1695+ out.push({
1696+ id: key,
1697+ space_id: row.space_id,
1698+ space_slug: "repo",
1699+ title: row.title ?? row.path,
1700+ icon: null,
1701+ slug: row.path,
1702+ path: `/${workspace.slug}/-/docs/repo/${repo}/${row.path.split("/").map(encodeURIComponent).join("/")}`,
1703+ space_name: repo,
1704+ snippet,
1705+ updated_at: r.row.indexed_at ?? r.row.added_at,
1706+ projects: [repo.toLowerCase()],
1707+ repo_file: { repo, path: row.path },
1708+ heading: row.heading,
1709+ matched: "meaning",
1710+ });
1711+ }
1712+ }
1713+ return out;
1714+ }
1715+
1716+ /** Full text over these projects' docs (the viewer's to read). */
1717+ private async searchRepoFiles(workspace: Workspace, spaces: { row: RepoSpaceRow; repo: Repo }[], query: DocSearchQuery): Promise<DocSearchHit[]> {
1718+ const q = ftsQuery(query.query);
1719+ if (!q) return [];
15271720 if (!spaces.length) return [];
15281721 const limit = Math.min(Math.max(Number(query.limit) || 20, 1), 50);
15291722 const rows = (
19472140 return ok(await this.searchIn(found.value.workspace, found.value.spaces, { ...(a.query ?? { query: "" }), limit: Math.min(Number(a.query?.limit) || 10, 20) }));
19482141 }
19492142
2143+ // ── Recall: the semantic index ──────────────────────────────────────────
2144+
2145+ /** Pages indexed again one after another (moved, restored). */
2146+ private async reindexPages(ids: string[]): Promise<void> {
2147+ for (const id of ids.slice(0, 500)) await indexPage(this.env, id);
2148+ }
2149+
2150+ /** A query's embedding, kept a minute; null without an embedder or when it fails (then words only). */
2151+ private async queryVector(query: string): Promise<number[] | null> {
2152+ const { embedder } = adapters(this.env);
2153+ const key = queryKey(query);
2154+ if (!embedder || !key) return null;
2155+ const cached = queryVectors.get(key);
2156+ if (cached) return cached;
2157+ try {
2158+ const [vector] = await embedder.embed([key]);
2159+ if (vector) queryVectors.set(key, vector);
2160+ return vector ?? null;
2161+ } catch (error) {
2162+ console.error("docs could not embed a query; matching words instead", String(error));
2163+ return null;
2164+ }
2165+ }
2166+
2167+ /** The passages nearest a vector, only from `allowed` spaces (by the index's filter, or after). */
2168+ private async meaningMatches(workspaceId: string, allowed: string[], vector: number[]): Promise<{ id: string; score: number }[]> {
2169+ const { store } = adapters(this.env);
2170+ const plan = vectorQueryPlan(workspaceId, allowed);
2171+ if (!store || !plan) return [];
2172+ try {
2173+ return await store.query(vector, { topK: plan.topK, filter: plan.filter });
2174+ } catch (error) {
2175+ console.error("docs semantic query failed; matching words instead", String(error));
2176+ return [];
2177+ }
2178+ }
2179+
2180+ /** Passages by their words (any of them), best first, from `allowed` spaces. */
2181+ private async wordMatches(workspaceId: string, allowed: string[], fts: string, limit: number): Promise<string[]> {
2182+ if (!allowed.length) return [];
2183+ const named = allowed.length <= 80;
2184+ const rows = await this.db
2185+ .prepare(
2186+ `SELECT doc_chunks_fts.chunk_id AS id FROM doc_chunks_fts JOIN doc_chunks c ON c.id = doc_chunks_fts.chunk_id
2187+ WHERE doc_chunks_fts MATCH ? AND c.workspace_id = ? ${named ? `AND doc_chunks_fts.space_id IN (${allowed.map(() => "?").join(",")})` : ""}
2188+ ORDER BY bm25(doc_chunks_fts, 0, 0, 0, 4.0, 1.0) LIMIT ?`,
2189+ )
2190+ .bind(fts, workspaceId, ...(named ? allowed : []), limit)
2191+ .all<{ id: string }>()
2192+ .catch((error: unknown) => {
2193+ console.error("docs word recall failed", String(error));
2194+ return { results: [] as { id: string }[] };
2195+ });
2196+ return rows.results.map((r) => r.id);
2197+ }
2198+
2199+ /**
2200+ * Passages by id as they read now, with their page or file: only those
2201+ * whose page is still out of the trash and whose file is still there.
2202+ * A page's passages count as in the page's space now, whatever the index says.
2203+ */
2204+ private async passages(workspaceId: string, ids: string[]): Promise<Map<string, PassageRow>> {
2205+ const out = new Map<string, PassageRow>();
2206+ const unique = [...new Set(ids)];
2207+ for (let i = 0; i < unique.length; i += 90) {
2208+ const part = unique.slice(i, i + 90);
2209+ const rows = await this.db
2210+ .prepare(
2211+ `SELECT c.id, c.page_id, c.repo_file_id, c.path, c.heading, c.text, c.updated_at,
2212+ CASE WHEN c.page_id IS NOT NULL THEN p.space_id ELSE c.space_id END AS space_id,
2213+ COALESCE(p.title, f.title) AS title, p.icon AS icon, p.updated_at AS page_updated_at
2214+ FROM doc_chunks c
2215+ LEFT JOIN pages p ON p.id = c.page_id
2216+ LEFT JOIN repo_files f ON f.space_id = c.space_id AND f.path = c.path
2217+ WHERE c.workspace_id = ? AND c.id IN (${part.map(() => "?").join(",")})
2218+ AND ((c.page_id IS NOT NULL AND p.id IS NOT NULL AND p.archived_at IS NULL) OR (c.repo_file_id IS NOT NULL AND f.path IS NOT NULL))`,
2219+ )
2220+ .bind(workspaceId, ...part)
2221+ .all<PassageRow>();
2222+ for (const r of rows.results) out.set(r.id, r);
2223+ }
2224+ return out;
2225+ }
2226+
2227+ /**
2228+ * Projects' docs an agent may recall from: those the viewer can read
2229+ * and, with an audience, everyone in it. A workspace-wide audience, or
2230+ * one too large to ask about person by person, gets public repositories
2231+ * only. Never wider than the viewer.
2232+ */
2233+ private async repoSpacesForAudience(workspace: Workspace, viewer: User, audience: DocAudience | null): Promise<{ row: RepoSpaceRow; repo: Repo }[]> {
2234+ const mine = await this.readableRepoSpaces(workspace, viewer);
2235+ if (!mine.length || !audience) return mine;
2236+ const publicOnly = () => mine.filter((s) => !s.repo.isPrivate);
2237+ if (audience.kind === "workspace") return publicOnly();
2238+ if (audience.kind !== "people" || !Array.isArray(audience.user_ids)) return mine;
2239+ const others = [...new Set(audience.user_ids.map(String))].filter((id) => id !== viewer.id);
2240+ if (!others.length) return mine;
2241+ if (others.length > 20 || !this.env.REPOS) return publicOnly();
2242+ await this.nameUsers(others);
2243+ const members = await this.members(workspace);
2244+ let keep = new Set(mine.map((s) => s.row.repo_id));
2245+ for (const id of others) {
2246+ const username = this.usernames.get(id)?.toLowerCase();
2247+ if (!username) {
2248+ const open = new Set(publicOnly().map((s) => s.row.repo_id));
2249+ keep = new Set([...keep].filter((r) => open.has(r)));
2250+ continue;
2251+ }
2252+ const member = members.get(username);
2253+ // As repos sees them: their membership here and no direct grants, so never wider than they are.
2254+ const person: User = { id, username, verified: true, workspaces: member ? [{ slug: workspace.slug, role: member.role }] : [] };
2255+ const readable = await reposClient(this.env.REPOS)
2256+ .readable([...keep], person)
2257+ .catch(() => [] as Repo[]);
2258+ keep = new Set(readable.map((r) => r.id));
2259+ if (!keep.size) break;
2260+ }
2261+ return mine.filter((s) => keep.has(s.row.repo_id));
2262+ }
2263+
2264+ /**
2265+ * What the workspace's Docs say about a query, for an agent about to
2266+ * answer (DocsApi.recallForAgent): the closest passages by meaning above
2267+ * MEANING_FLOOR, required spaces first, at most two per page, filled with
2268+ * passages matching its words when meaning finds too few. Only from what
2269+ * the viewer and audience can all read, by the same rules as every other
2270+ * agent read (`agentSpaces`).
2271+ */
2272+ async recallForAgent(a: {
2273+ workspace: string;
2274+ agent_id: string;
2275+ viewer: Viewer;
2276+ query: string;
2277+ limit?: number | null;
2278+ spaces?: string[] | null;
2279+ audience: DocAudience | null;
2280+ }): Promise<Result<DocPassage[]>> {
2281+ const found = await this.agentSpaces(a.workspace, a.agent_id, a.viewer, a.audience ?? null);
2282+ if (!found.ok) return found;
2283+ const { workspace, spaces } = found.value;
2284+ this.defer(ensureIndexed(this.env, workspace.id).catch((error: unknown) => console.error("docs could not start indexing", workspace.id, String(error))));
2285+ const query = String(a.query ?? "").trim().slice(0, 2000);
2286+ if (!query) return ok([]);
2287+ const limit = recallLimit(a.limit);
2288+ const repoSpaces = await this.repoSpacesForAudience(workspace, a.viewer!, a.audience ?? null).catch((error: unknown) => {
2289+ console.error("docs could not check projects' docs for recall", String(error));
2290+ return [] as { row: RepoSpaceRow; repo: Repo }[];
2291+ });
2292+ const allowed = [...spaces.map((s) => s.row.id), ...repoSpaces.map((r) => r.row.id)];
2293+ if (!allowed.length) return ok([]);
2294+ const required = requiredSpaces(allowed, a.spaces);
2295+ const fts = ftsAnyQuery(query);
2296+ const vector = await this.queryVector(query);
2297+ const [meaning, requiredMeaning, words] = await Promise.all([
2298+ vector ? this.meaningMatches(workspace.id, allowed, vector) : Promise.resolve([]),
2299+ // Required reading asked on its own too, so the rest of the workspace can't crowd it out.
2300+ vector && required.length && required.length < allowed.length ? this.meaningMatches(workspace.id, required, vector) : Promise.resolve([]),
2301+ fts ? this.wordMatches(workspace.id, allowed, fts, 30) : Promise.resolve([] as string[]),
2302+ ]);
2303+ const scores = new Map<string, number>();
2304+ for (const m of [...meaning, ...requiredMeaning]) scores.set(m.id, Math.max(scores.get(m.id) ?? 0, m.score));
2305+ const rows = await this.passages(workspace.id, [...scores.keys(), ...words]);
2306+ const candidates: (Candidate & { row: PassageRow })[] = [];
2307+ for (const [id, score] of scores) {
2308+ const row = rows.get(id);
2309+ if (row) candidates.push({ id, doc_id: row.page_id ?? row.repo_file_id!, space_id: row.space_id, score, by: "meaning", row });
2310+ }
2311+ for (const id of words) {
2312+ const row = rows.get(id);
2313+ if (row) candidates.push({ id, doc_id: row.page_id ?? row.repo_file_id!, space_id: row.space_id, score: WORDS_SCORE, by: "words", row });
2314+ }
2315+ const picked = pickPassages(candidates, { allowed: new Set(allowed), required, limit });
2316+ const stale = await this.staleIds([...new Set(picked.map((c) => c.row.page_id).filter((id): id is string => !!id))]);
2317+ const bySpace = new Map(spaces.map((s) => [s.row.id, s.row]));
2318+ const byRepo = new Map(repoSpaces.map((r) => [r.row.id, r]));
2319+ const out: DocPassage[] = [];
2320+ for (const c of picked) {
2321+ const row = c.row;
2322+ const score = Math.round(c.score * 1000) / 1000;
2323+ if (row.page_id) {
2324+ const space = bySpace.get(row.space_id);
2325+ if (!space) continue;
2326+ out.push({
2327+ page: this.ref(workspace.slug, space, { id: row.page_id, title: row.title ?? "", icon: row.icon }),
2328+ repo_file: null,
2329+ space_name: space.name,
2330+ heading: row.heading,
2331+ text: row.text,
2332+ score,
2333+ updated_at: row.page_updated_at ?? row.updated_at,
2334+ stale: stale.has(row.page_id),
2335+ });
2336+ } else {
2337+ const repo = byRepo.get(row.space_id);
2338+ if (!repo || !row.path) continue;
2339+ const name = `${repo.repo.namespace}/${repo.repo.name}`;
2340+ out.push({
2341+ page: null,
2342+ repo_file: { repo: name, path: row.path, href: `/${workspace.slug}/-/docs/repo/${name}/${row.path.split("/").map(encodeURIComponent).join("/")}` },
2343+ space_name: name,
2344+ heading: row.heading,
2345+ text: row.text,
2346+ score,
2347+ updated_at: repo.row.indexed_at ?? row.updated_at,
2348+ stale: false,
2349+ });
2350+ }
2351+ }
2352+ return ok(out);
2353+ }
2354+
2355+ /** Indexes the workspace's pages and projects' docs again, on the queue. Owners only. */
2356+ async reindexDocs(a: { workspace: string; viewer: Viewer }): Promise<Result<boolean>> {
2357+ const found = await this.viewerWorkspace(a.workspace, a.viewer);
2358+ if (!found.ok) return found;
2359+ if (!this.viewerOwner(a.viewer!, a.workspace)) return fail("forbidden", "Only an owner can index the workspace's docs again.");
2360+ return ok(await startBackfill(this.env, found.value.id, { force: true }));
2361+ }
2362+
19502363 private async fileSuggestion(
19512364 workspace: Workspace,
19522365 page: PageRow,
22722685 return Response.json(await service.removeRepoSpace(args));
22732686 case "repo_page":
22742687 return Response.json(await service.repoPage(args));
2688+ case "recall_for_agent":
2689+ return Response.json(await service.recallForAgent(args));
2690+ case "reindex_docs":
2691+ return Response.json(await service.reindexDocs(args));
22752692 default:
22762693 return new Response("Unknown method\n", { status: 404 });
22772694 }
23022719 /**
23032720 * Events from the events service (SUBSCRIBER_DOCS): pages whose cited
23042721 * code changed become possibly out of date, and projects' docs are read
2305− * again after a push (src/staleness.ts). One failing event is retried on
2306− * its own.
2722+ * again after a push (src/staleness.ts) and their passages indexed
2723+ * (src/indexer.ts). The same queue carries this service's own backfill
2724+ * jobs (`docs.index`). One failing message is retried on its own.
23072725 */
2308− async queue(batch: MessageBatch<G1tEvent>, env: Env): Promise<void> {
2726+ async queue(batch: MessageBatch<G1tEvent | DocsJob>, env: Env): Promise<void> {
23092727 const reindex = async (repoId: string) => {
2310− if (env.REPOS) await reindexRepo({ DB: env.DB, REPOS: env.REPOS }, repoId);
2728+ if (env.REPOS) await reindexRepo({ DB: env.DB, REPOS: env.REPOS }, repoId, (spaceId, changed, gone) => indexRepoFiles(env, spaceId, changed, gone));
23112729 };
23122730 for (const message of batch.messages) {
23132731 try {
2314− await onEvent(env, message.body, reindex);
2732+ const body = message.body;
2733+ if (body.type === "docs.index") {
2734+ await runBackfill(env, (body as DocsJob).workspace_id);
2735+ message.ack();
2736+ continue;
2737+ }
2738+ if (body.type === "repo.purged") {
2739+ // Its docs leave Docs (src/staleness.ts); their passages leave the index first.
2740+ const gone = await env.DB.prepare("SELECT id FROM repo_spaces WHERE repo_id = ?").bind((body as G1tEvent<"repo.purged">).data.repoId).all<{ id: string }>();
2741+ for (const space of gone.results) await forgetDocs(env, { space_id: space.id });
2742+ }
2743+ await onEvent(env, body as G1tEvent, reindex);
23152744 message.ack();
23162745 } catch (error) {
23172746 console.error("docs could not handle", message.body?.type, String(error));
23192748 }
23202749 }
23212750 },
2322−} satisfies ExportedHandler<Env, G1tEvent>;
2751+} satisfies ExportedHandler<Env, G1tEvent | DocsJob>;
+149−0
1+import assert from "node:assert/strict";
2+import { readFileSync, readdirSync } from "node:fs";
3+import { DatabaseSync } from "node:sqlite";
4+import { test } from "node:test";
5+
6+import { EMBED_PER_HOUR, indexDoc, type IndexDoc } from "./indexer.ts";
7+import type { Embedder, VectorMetadata, VectorStore } from "./vectors.ts";
8+
9+/** D1, as far as the indexer uses it, over node's SQLite with the service's migrations. */
10+function fakeD1(): D1Database {
11+ const db = new DatabaseSync(":memory:");
12+ const dir = new URL("../migrations/", import.meta.url);
13+ for (const file of readdirSync(dir).sort()) db.exec(readFileSync(new URL(file, dir), "utf8"));
14+ const statement = (sql: string, params: unknown[] = []) => ({
15+ sql,
16+ params,
17+ bind: (...values: unknown[]) => statement(sql, values),
18+ first: async () => (db.prepare(sql).get(...(params as never[])) as unknown) ?? null,
19+ all: async () => ({ results: db.prepare(sql).all(...(params as never[])) }),
20+ run: async () => db.prepare(sql).run(...(params as never[])),
21+ });
22+ return {
23+ prepare: (sql: string) => statement(sql),
24+ batch: async (list: ReturnType<typeof statement>[]) => list.map((s) => db.prepare(s.sql).run(...(s.params as never[]))),
25+ } as unknown as D1Database;
26+}
27+
28+function fakes() {
29+ const vectors = new Map<string, { values: number[]; metadata: VectorMetadata }>();
30+ let embedded = 0;
31+ const embedder: Embedder = {
32+ async embed(texts) {
33+ embedded += texts.length;
34+ return texts.map((t) => [t.length, 1]);
35+ },
36+ };
37+ const store: VectorStore = {
38+ async upsert(list) {
39+ for (const v of list) vectors.set(v.id, { values: v.values, metadata: v.metadata });
40+ },
41+ async get(ids) {
42+ return ids.filter((id) => vectors.has(id)).map((id) => ({ id, values: vectors.get(id)!.values }));
43+ },
44+ async delete(ids) {
45+ for (const id of ids) vectors.delete(id);
46+ },
47+ async query() {
48+ return [];
49+ },
50+ };
51+ return { vectors, embedder, store, embeddedCount: () => embedded };
52+}
53+
54+const section = (name: string) => `## ${name}\n\n${Array.from({ length: 60 }, (_, i) => `${name.toLowerCase()}${i}`).join(" ")}.`;
55+const page = (markdown: string, over: Partial<IndexDoc> = {}): IndexDoc => ({ kind: "page", doc_id: "pag_1", workspace_id: "w1", space_id: "s1", title: "Ops", markdown, ...over });
56+
57+test("a page's passages are stored, embedded once, and only changed ones again", async () => {
58+ const db = fakeD1();
59+ const f = fakes();
60+ const first = await indexDoc(db, f.embedder, f.store, page([section("Deploy"), section("Rollback")].join("\n\n")));
61+ assert.deepEqual(first, { chunks: 2, embedded: 2, capped: false });
62+ assert.deepEqual([...f.vectors.keys()].sort(), ["pag_1:0", "pag_1:1"]);
63+ assert.deepEqual(f.vectors.get("pag_1:0")!.metadata, { workspace_id: "w1", space_id: "s1", kind: "page", page_id: "pag_1" });
64+ // Saved again unchanged: nothing embedded.
65+ assert.equal((await indexDoc(db, f.embedder, f.store, page([section("Deploy"), section("Rollback")].join("\n\n")))).embedded, 0);
66+ // One section changed.
67+ assert.equal((await indexDoc(db, f.embedder, f.store, page([section("Deploy"), section("Restore")].join("\n\n")))).embedded, 1);
68+ assert.equal(f.embeddedCount(), 3);
69+ // Word search finds passages by heading and text.
70+ const hit = await db.prepare("SELECT chunk_id FROM doc_chunks_fts WHERE doc_chunks_fts MATCH ?").bind('"restore"').all<{ chunk_id: string }>();
71+ assert.deepEqual(
72+ hit.results.map((r) => r.chunk_id),
73+ ["pag_1:1"],
74+ );
75+});
76+
77+test("a passage that only moved keeps its vector; gone passages leave the index", async () => {
78+ const db = fakeD1();
79+ const f = fakes();
80+ await indexDoc(db, f.embedder, f.store, page([section("Alpha"), section("Beta"), section("Gamma")].join("\n\n")));
81+ const before = f.embeddedCount();
82+ // Alpha removed: Beta and Gamma move up a place, nothing new to embed.
83+ const r = await indexDoc(db, f.embedder, f.store, page([section("Beta"), section("Gamma")].join("\n\n")));
84+ assert.equal(r.embedded, 0);
85+ assert.equal(f.embeddedCount(), before);
86+ assert.deepEqual([...f.vectors.keys()].sort(), ["pag_1:0", "pag_1:1"]);
87+ const rows = await db.prepare("SELECT id, heading, hash = vector_hash AS current FROM doc_chunks ORDER BY seq").all<{ id: string; heading: string; current: number }>();
88+ assert.deepEqual(
89+ rows.results.map((x) => [x.id, x.heading, x.current]),
90+ [
91+ ["pag_1:0", "Beta", 1],
92+ ["pag_1:1", "Gamma", 1],
93+ ],
94+ );
95+});
96+
97+test("moving a page to another space files its vectors there without embedding", async () => {
98+ const db = fakeD1();
99+ const f = fakes();
100+ await indexDoc(db, f.embedder, f.store, page(section("Deploy")));
101+ const r = await indexDoc(db, f.embedder, f.store, page(section("Deploy"), { space_id: "s2" }));
102+ assert.equal(r.embedded, 0);
103+ assert.equal(f.vectors.get("pag_1:0")!.metadata.space_id, "s2");
104+ assert.equal((await db.prepare("SELECT space_id FROM doc_chunks WHERE id = 'pag_1:0'").first<{ space_id: string }>())!.space_id, "s2");
105+});
106+
107+test("embedding stops at the hourly cap and the rest waits, kept for words", async () => {
108+ const db = fakeD1();
109+ const f = fakes();
110+ const at = new Date("2026-10-09T10:15:00Z");
111+ await db.prepare("INSERT INTO doc_embed_usage (workspace_id, hour, chunks, tokens) VALUES ('w1', '2026-10-09T10', ?, 0)").bind(EMBED_PER_HOUR - 1).run();
112+ const r = await indexDoc(db, f.embedder, f.store, page([section("One"), section("Two"), section("Three")].join("\n\n")), at);
113+ assert.deepEqual(r, { chunks: 3, embedded: 1, capped: true });
114+ const waiting = await db.prepare("SELECT COUNT(*) AS n FROM doc_chunks WHERE vector_hash IS NULL OR vector_hash <> hash").first<{ n: number }>();
115+ assert.equal(waiting!.n, 2);
116+ // The next hour, a save picks up the rest.
117+ const later = await indexDoc(db, f.embedder, f.store, page([section("One"), section("Two"), section("Three")].join("\n\n")), new Date("2026-10-09T11:01:00Z"));
118+ assert.deepEqual(later, { chunks: 3, embedded: 2, capped: false });
119+});
120+
121+test("an embedding failure leaves passages for the next save, never throws", async () => {
122+ const db = fakeD1();
123+ const f = fakes();
124+ const broken: Embedder = {
125+ async embed() {
126+ throw new Error("model down");
127+ },
128+ };
129+ const r = await indexDoc(db, broken, f.store, page(section("Deploy")));
130+ assert.equal(r.chunks, 1);
131+ assert.equal(f.vectors.size, 0);
132+ assert.equal((await indexDoc(db, f.embedder, f.store, page(section("Deploy")))).embedded, 1);
133+});
134+
135+test("without an embedder, passages are kept for words only", async () => {
136+ const db = fakeD1();
137+ const r = await indexDoc(db, null, null, page(section("Deploy")));
138+ assert.deepEqual(r, { chunks: 1, embedded: 0, capped: false });
139+ assert.equal((await db.prepare("SELECT COUNT(*) AS n FROM doc_chunks_fts").first<{ n: number }>())!.n, 1);
140+});
141+
142+test("a project's docs file is indexed with its repository", async () => {
143+ const db = fakeD1();
144+ const f = fakes();
145+ await indexDoc(db, f.embedder, f.store, { kind: "repo_file", doc_id: "rf_x", workspace_id: "w1", space_id: "rds_1", title: "Setup", markdown: `# Setup\n\n${section("Install")}`, repo_id: "r1", path: "docs/setup.md" });
146+ assert.deepEqual(f.vectors.get("rf_x:0")!.metadata, { workspace_id: "w1", space_id: "rds_1", kind: "repo_file", repo_file_id: "rf_x", repo_id: "r1" });
147+ const row = await db.prepare("SELECT page_id, repo_file_id, path, heading FROM doc_chunks").first<Record<string, unknown>>();
148+ assert.deepEqual({ ...row }, { page_id: null, repo_file_id: "rf_x", path: "docs/setup.md", heading: "Install" });
149+});
+408−0
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+}
+4−4
4242
4343 type PageRow = { id: string; workspace_id: string; space_id: string; title: string; mentioned: string; markdown: string; slug: string; archived_at: string | null };
4444
45−/** Saves; returns whether a version was recorded, and its id. */
46−export async function save(env: SaveEnv, input: Save, now = new Date()): Promise<{ version_id: string | null }> {
45+/** Saves; returns whether a version was recorded, and its id, and whether the Markdown changed (so the room indexes it again, src/indexer.ts). */
46+export async function save(env: SaveEnv, input: Save, now = new Date()): Promise<{ version_id: string | null; changed: boolean }> {
4747 const at = now.toISOString();
4848 const page = await env.DB.prepare(
4949 "SELECT p.id, p.workspace_id, p.space_id, p.title, p.mentioned, p.markdown, p.archived_at, s.slug AS slug FROM pages p JOIN spaces s ON s.id = p.space_id WHERE p.id = ?",
5050 )
5151 .bind(input.page_id)
5252 .first<PageRow>();
53− if (!page) return { version_id: null };
53+ if (!page) return { version_id: null, changed: false };
5454 const changed = page.markdown !== input.markdown;
5555 const last = input.editors[input.editors.length - 1] ?? null;
5656 const statements: D1PreparedStatement[] = [];
146146 await Promise.all(fresh.map((name) => notify.notify({ username: name }, notification(name)).catch(() => undefined)));
147147 }
148148 }
149− return { version_id: versionId };
149+ return { version_id: versionId, changed };
150150 }
+92−0
1+import assert from "node:assert/strict";
2+import { test } from "node:test";
3+
4+import { MAX_IN_FILTER, MEANING_FLOOR, QueryCache, TOP_K, TOP_K_UNFILTERED, fuseRanks, pickPassages, queryKey, recallLimit, requiredSpaces, vectorQueryPlan, type Candidate } from "./recall.ts";
5+import { ftsAnyQuery } from "./search.ts";
6+
7+const c = (id: string, space: string, score: number, by: Candidate["by"] = "meaning"): Candidate => ({ id, doc_id: id.split(":")[0]!, space_id: space, score, by });
8+
9+test("the index is filtered by the allowed spaces, or by workspace and after when there are too many", () => {
10+ assert.equal(vectorQueryPlan("w1", []), null);
11+ assert.deepEqual(vectorQueryPlan("w1", ["s1", "s2", "s1"]), { topK: TOP_K, filter: { workspace_id: "w1", space_ids: ["s1", "s2"] }, filterAfter: false });
12+ const many = Array.from({ length: MAX_IN_FILTER + 1 }, (_, i) => `s${i}`);
13+ assert.deepEqual(vectorQueryPlan("w1", many), { topK: TOP_K_UNFILTERED, filter: { workspace_id: "w1" }, filterAfter: true });
14+});
15+
16+test("required reading never widens what may be read", () => {
17+ assert.deepEqual(requiredSpaces(["s1", "s2"], ["s2", "s3", "s2"]), ["s2"]);
18+ assert.deepEqual(requiredSpaces(["s1"], null), []);
19+ assert.deepEqual(requiredSpaces(["s1"], "s1"), []);
20+});
21+
22+test("limits: five by default, ten at most", () => {
23+ assert.equal(recallLimit(undefined), 5);
24+ assert.equal(recallLimit(0), 5);
25+ assert.equal(recallLimit(3), 3);
26+ assert.equal(recallLimit(50), 10);
27+ assert.equal(recallLimit("x"), 5);
28+});
29+
30+test("only close enough, only allowed, at most two per page, best first", () => {
31+ const picked = pickPassages(
32+ [
33+ c("p1:0", "s1", 0.81),
34+ c("p1:1", "s1", 0.8),
35+ c("p1:2", "s1", 0.79),
36+ c("p2:0", "s1", 0.7),
37+ c("p3:0", "secret", 0.95),
38+ c("p4:0", "s1", MEANING_FLOOR - 0.01),
39+ ],
40+ { allowed: new Set(["s1"]), limit: 10 },
41+ );
42+ assert.deepEqual(
43+ picked.map((p) => p.id),
44+ ["p1:0", "p1:1", "p2:0"],
45+ );
46+});
47+
48+test("required spaces come first, then the rest", () => {
49+ const picked = pickPassages([c("a:0", "s1", 0.9), c("b:0", "s2", 0.65), c("c:0", "s1", 0.85)], { allowed: new Set(["s1", "s2"]), required: ["s2"], limit: 2 });
50+ assert.deepEqual(
51+ picked.map((p) => p.id),
52+ ["b:0", "a:0"],
53+ );
54+});
55+
56+test("words fill in only when meaning finds too few, and never twice", () => {
57+ const candidates = [c("a:0", "s1", 0.9), c("a:0", "s1", 0.5, "words"), c("b:0", "s1", 0.5, "words"), c("c:0", "s1", 0.5, "words"), c("d:0", "nope", 0.5, "words")];
58+ assert.deepEqual(
59+ pickPassages(candidates, { allowed: new Set(["s1"]), limit: 3 }).map((p) => `${p.id}/${p.by}`),
60+ ["a:0/meaning", "b:0/words", "c:0/words"],
61+ );
62+ assert.deepEqual(
63+ pickPassages(candidates, { allowed: new Set(["s1"]), limit: 1 }).map((p) => p.id),
64+ ["a:0"],
65+ );
66+ // Nothing close and no words: nothing.
67+ assert.deepEqual(pickPassages([c("x:0", "s1", 0.3)], { allowed: new Set(["s1"]), limit: 5 }), []);
68+});
69+
70+test("hybrid search fuses ranks: found both ways first", () => {
71+ assert.deepEqual(fuseRanks(["a", "b", "c"], ["c", "d"]), ["c", "a", "b", "d"]);
72+ assert.deepEqual(fuseRanks([], ["x", "y"]), ["x", "y"]);
73+});
74+
75+test("query embeddings are kept a minute", () => {
76+ const cache = new QueryCache(60_000, 2);
77+ cache.set("a", [1], 0);
78+ assert.deepEqual(cache.get("a", 59_000), [1]);
79+ assert.equal(cache.get("a", 61_000), null);
80+ cache.set("a", [1], 0);
81+ cache.set("b", [2], 1);
82+ cache.set("c", [3], 2);
83+ assert.equal(cache.get("a", 3), null);
84+ assert.deepEqual(cache.get("c", 3), [3]);
85+ assert.equal(queryKey(" How do we Deploy? "), "how do we deploy?");
86+});
87+
88+test("recall's word fallback matches any meaningful word", () => {
89+ assert.equal(ftsAnyQuery("How do we roll back a deploy?"), '"roll" OR "back" OR "deploy"');
90+ assert.equal(ftsAnyQuery("is it a"), null);
91+ assert.equal(ftsAnyQuery('NEAR("x") deploy-key'), '"near" OR "deploy" OR "key"');
92+});
+152−0
1+/**
2+ * Recall's choices, apart from where the passages come from: how the
3+ * vector query is filtered, which matches are close enough, and how
4+ * passages are picked (an agent's required reading first, at most two per
5+ * document, words only to fill). Pure; src/index.ts `recallForAgent` and
6+ * hybrid search run it over Vectorize and D1.
7+ */
8+
9+/**
10+ * The least cosine similarity a passage needs to count as being about the
11+ * query. bge-base-en-v1.5 scores unrelated English around 0.4 to 0.55 and
12+ * a passage on the asked-about topic from about 0.65; 0.6 keeps recall
13+ * quiet when the docs say nothing about it.
14+ */
15+export const MEANING_FLOOR = 0.6;
16+/** The score a passage found only by its words carries: below the floor, so callers can tell. */
17+export const WORDS_SCORE = 0.5;
18+/** Nearest passages asked of the index. */
19+export const TOP_K = 24;
20+/** Asked when the index can't filter by space and recall filters after: more, so enough are left. */
21+export const TOP_K_UNFILTERED = 50;
22+/** The most space ids put in one `$in` filter; past it, filter after the query (a Vectorize filter is at most 2 KB of JSON). */
23+export const MAX_IN_FILTER = 40;
24+/** Passages from one page or file at most. */
25+export const PER_DOC = 2;
26+export const DEFAULT_LIMIT = 5;
27+export const MAX_LIMIT = 10;
28+
29+export function recallLimit(limit: unknown): number {
30+ const n = Math.floor(Number(limit));
31+ if (!Number.isFinite(n) || n < 1) return DEFAULT_LIMIT;
32+ return Math.min(n, MAX_LIMIT);
33+}
34+
35+/**
36+ * How to ask the index: by the allowed spaces when there are few enough
37+ * to name, otherwise by workspace alone, more of them, filtered after.
38+ * Null when nothing may be read.
39+ */
40+export function vectorQueryPlan(workspaceId: string, allowed: string[]): { topK: number; filter: { workspace_id: string; space_ids?: string[] }; filterAfter: boolean } | null {
41+ const ids = [...new Set(allowed)];
42+ if (!ids.length) return null;
43+ if (ids.length <= MAX_IN_FILTER) return { topK: TOP_K, filter: { workspace_id: workspaceId, space_ids: ids }, filterAfter: false };
44+ return { topK: TOP_K_UNFILTERED, filter: { workspace_id: workspaceId }, filterAfter: true };
45+}
46+
47+/**
48+ * The spaces recall looks in first: those asked for (an agent's required
49+ * reading) that may be read. Never wider than what may be read.
50+ */
51+export function requiredSpaces(allowed: string[], asked: unknown): string[] {
52+ if (!Array.isArray(asked)) return [];
53+ const may = new Set(allowed);
54+ return [...new Set(asked.map(String))].filter((id) => may.has(id));
55+}
56+
57+export type Candidate = {
58+ /** The passage's id: `<doc>:<seq>`. */
59+ id: string;
60+ /** Its page's or file's id. */
61+ doc_id: string;
62+ space_id: string;
63+ score: number;
64+ /** Found by meaning (the index) or by its words (full text). */
65+ by: "meaning" | "words";
66+};
67+
68+/**
69+ * The passages to hand over, best first: by meaning above the floor,
70+ * required spaces first, then the rest; then, if that is fewer than
71+ * `limit`, by words, required spaces first. At most PER_DOC from one
72+ * document, each passage once, only from `allowed`.
73+ */
74+export function pickPassages<T extends Candidate>(candidates: T[], options: { allowed: Set<string>; required?: string[]; limit: number; floor?: number }): T[] {
75+ const floor = options.floor ?? MEANING_FLOOR;
76+ const required = new Set(options.required ?? []);
77+ const usable = candidates.filter((c) => options.allowed.has(c.space_id));
78+ const meaning = usable.filter((c) => c.by === "meaning" && c.score >= floor).sort((a, b) => b.score - a.score);
79+ const words = usable.filter((c) => c.by === "words");
80+ const out: T[] = [];
81+ const taken = new Set<string>();
82+ const perDoc = new Map<string, number>();
83+ const take = (list: T[]) => {
84+ for (const c of list) {
85+ if (out.length >= options.limit) return;
86+ if (taken.has(c.id)) continue;
87+ const n = perDoc.get(c.doc_id) ?? 0;
88+ if (n >= PER_DOC) continue;
89+ taken.add(c.id);
90+ perDoc.set(c.doc_id, n + 1);
91+ out.push(c);
92+ }
93+ };
94+ take(meaning.filter((c) => required.has(c.space_id)));
95+ take(meaning.filter((c) => !required.has(c.space_id)));
96+ take(words.filter((c) => required.has(c.space_id)));
97+ take(words.filter((c) => !required.has(c.space_id)));
98+ return out;
99+}
100+
101+/**
102+ * Hybrid search's order for people: each document's place in the word
103+ * results and in the meaning results, fused by reciprocal rank (k = 60),
104+ * so a page both find comes first and either alone still counts.
105+ */
106+export function fuseRanks(words: string[], meaning: string[], k = 60): string[] {
107+ const score = new Map<string, number>();
108+ const add = (ids: string[]) =>
109+ [...new Set(ids)].forEach((id, rank) => {
110+ score.set(id, (score.get(id) ?? 0) + 1 / (k + rank + 1));
111+ });
112+ add(words);
113+ add(meaning);
114+ return [...score.entries()].sort((a, b) => b[1] - a[1]).map(([id]) => id);
115+}
116+
117+/**
118+ * A query's embeddings, kept a minute per isolate: an agent asked the same
119+ * thing again in a session, or a search page reloaded, embeds once.
120+ */
121+export class QueryCache {
122+ private readonly entries = new Map<string, { at: number; vector: number[] }>();
123+ private readonly ttlMs: number;
124+ private readonly max: number;
125+ constructor(ttlMs = 60_000, max = 200) {
126+ this.ttlMs = ttlMs;
127+ this.max = max;
128+ }
129+
130+ get(key: string, now = Date.now()): number[] | null {
131+ const hit = this.entries.get(key);
132+ if (!hit) return null;
133+ if (now - hit.at > this.ttlMs) {
134+ this.entries.delete(key);
135+ return null;
136+ }
137+ return hit.vector;
138+ }
139+
140+ set(key: string, vector: number[], now = Date.now()): void {
141+ if (this.entries.size >= this.max) {
142+ for (const [k, v] of this.entries) if (now - v.at > this.ttlMs) this.entries.delete(k);
143+ while (this.entries.size >= this.max) this.entries.delete(this.entries.keys().next().value!);
144+ }
145+ this.entries.set(key, { at: now, vector });
146+ }
147+}
148+
149+/** A query as the cache keys it: the same words, however spaced or cased. */
150+export function queryKey(query: string): string {
151+ return String(query ?? "").replace(/\s+/g, " ").trim().toLowerCase().slice(0, 2000);
152+}
+18−5
4444 * repos without a viewer (`listFiles`, `rawBlobs`), so callers check the
4545 * repository can be read first, as adding one does.
4646 */
47−export async function indexRepoSpace(env: { DB: D1Database; REPOS: ServiceBinding }, space: RepoSpaceRow, now = new Date()): Promise<{ files: number; read: number }> {
47+export async function indexRepoSpace(
48+ env: { DB: D1Database; REPOS: ServiceBinding },
49+ space: RepoSpaceRow,
50+ now = new Date(),
51+): Promise<{ files: number; read: number; changed: string[]; gone: string[] }> {
4852 const db = env.DB;
4953 const repos = reposClient(env.REPOS);
5054 const listing = await repos.listFiles(space.repo_id, null, 10_000);
7882 }
7983 statements.push(db.prepare("UPDATE repo_spaces SET commit_sha = ?, indexed_at = ? WHERE id = ?").bind(listing.commit, now.toISOString(), space.id));
8084 for (let i = 0; i < statements.length; i += 50) await db.batch(statements.slice(i, i + 50));
81− return { files: wanted.length, read: changed.length };
85+ return { files: wanted.length, read: changed.length, changed: changed.map((f) => f.path), gone };
8286 }
8387
84−/** Every space showing a repository's docs, read again after a push. Never throws for one space's sake. */
85−export async function reindexRepo(env: { DB: D1Database; REPOS: ServiceBinding }, repoId: string): Promise<void> {
88+/**
89+ * Every space showing a repository's docs, read again after a push; then
90+ * `indexed` with what changed in each (the semantic index, src/indexer.ts).
91+ * Never throws for one space's sake.
92+ */
93+export async function reindexRepo(
94+ env: { DB: D1Database; REPOS: ServiceBinding },
95+ repoId: string,
96+ indexed: (spaceId: string, changed: string[], gone: string[]) => Promise<void> = async () => {},
97+): Promise<void> {
8698 const spaces = (await env.DB.prepare("SELECT * FROM repo_spaces WHERE repo_id = ?").bind(repoId).all<RepoSpaceRow>()).results;
8799 for (const space of spaces) {
88100 try {
89− await indexRepoSpace(env, space);
101+ const read = await indexRepoSpace(env, space);
102+ await indexed(space.id, read.changed, read.gone);
90103 } catch (error) {
91104 console.error("docs could not read a project's docs", space.repo, String(error));
92105 }
+26−3
1212 * Storage (the object's own SQLite): the document as a snapshot plus the
1313 * updates since, compacted every so often. A few seconds after a burst of
1414 * edits (an alarm), the room saves the Markdown rendition, the search
15− * index, backlinks and history to D1 (src/persist.ts).
15+ * index, backlinks and history to D1 (src/persist.ts). Half a minute after
16+ * the Markdown first changes, it brings the page's passages in the
17+ * semantic index up to date (src/indexer.ts), on the same alarm, so typing
18+ * never waits on embedding.
1619 *
1720 * The room authorizes nothing about who may open the page: the Worker
1821 * checks the viewer's role before forwarding a socket (src/index.ts,
3235 import { anchorThread, applyEdit, findTarget, rangeIds, rangeMarkdown, restoreFrom, unanchorThread } from "./edits.ts";
3336 import { bodyCitations } from "./citations.ts";
3437 import { citationNodes, documentMarkdown, mentionedIds, outline, type Outline } from "./markdown.ts";
38+import { indexPage, type DocsJob } from "./indexer.ts";
3539 import { save } from "./persist.ts";
3640 import { applyThreadAction, listThreads, setQuote, type ThreadResult } from "./threads.ts";
3741
5256
5357 type Attachment = RoomMember & { clients: number[] };
5458
55−type Env = { DB: D1Database; NOTIFY?: ServiceBinding; EVENTS?: ServiceBinding };
59+type Env = { DB: D1Database; NOTIFY?: ServiceBinding; EVENTS?: ServiceBinding; AI?: Ai; VECTORS?: Vectorize; JOBS?: Queue<DocsJob> };
5660
5761 const FRAGMENT = "document-store";
5862 const MESSAGE_SYNC = 0;
6064 const MESSAGE_QUERY_AWARENESS = 3;
6165 /** Save this long after the last change. */
6266 const SAVE_AFTER_MS = 4_000;
67+/** Index the page's passages this long after its Markdown first changed: at most twice a minute while someone types. */
68+const INDEX_AFTER_MS = 30_000;
6369 /** Compact the stored updates into one snapshot past this many. */
6470 const COMPACT_AT = 300;
6571
150156 }
151157
152158 private scheduleSave(): void {
159+ this.alarmBy(Date.now() + SAVE_AFTER_MS);
160+ }
161+
162+ /** Makes sure the alarm rings by `when`: the save's, or the index's, whichever is first. */
163+ private alarmBy(when: number): void {
153164 void this.ctx.storage.getAlarm().then((at) => {
154− if (at == null) return this.ctx.storage.setAlarm(Date.now() + SAVE_AFTER_MS);
165+ if (at == null || at > when) return this.ctx.storage.setAlarm(when);
155166 });
156167 }
157168
189200 });
190201 this.setMeta("editors", []);
191202 this.setMeta("editor_names", []);
203+ if (result.changed && this.meta<number | null>("index_at", null) == null) {
204+ const when = Date.now() + INDEX_AFTER_MS;
205+ this.setMeta("index_at", when);
206+ this.alarmBy(when);
207+ }
192208 if (result.version_id) {
193209 this.setMeta("last_version_at", Date.now());
194210 this.setMeta("pending_authors", []);
198214
199215 override async alarm(): Promise<void> {
200216 await this.persist();
217+ const due = this.meta<number | null>("index_at", null);
218+ if (due == null) return;
219+ if (Date.now() < due) return this.alarmBy(due);
220+ this.setMeta("index_at", null);
221+ const pageId = this.meta<string | null>("page_id", null);
222+ // Never throws; a failure is retried on the next change.
223+ if (pageId) await indexPage(this.env, pageId);
201224 }
202225
203226 // ── Calls from the Worker ──────────────────────────────────────────────
+24−0
3939 const s = String(value ?? "").trim().replace(/^\/+|\/+$/g, "").toLowerCase();
4040 return /^[a-z0-9][a-z0-9._-]{0,99}\/[a-z0-9._-]{1,100}$/.test(s) ? s : null;
4141 }
42+
43+/** Words too common to say what a question is about. */
44+const STOP_WORDS = new Set(
45+ "a an and are as at be but by can could did do does for from had has have how i if in into is it its me my of on or our should so that the their them then there these they this to us was we were what when where which who why will with would you your about after before any been being both each few more most other over same some such than through too under until very".split(" "),
46+);
47+
48+/**
49+ * A question as a loose FTS5 query: its meaningful words (quoted, longer
50+ * than two letters, at most ten), any of them; bm25 ranks passages with
51+ * more of them first. For recall's word fallback, where a whole sentence
52+ * rarely matches every word. Null when no word is left.
53+ */
54+export function ftsAnyQuery(text: string): string | null {
55+ const words = [
56+ ...new Set(
57+ String(text ?? "")
58+ .toLowerCase()
59+ .split(/[^\p{L}\p{N}_]+/u)
60+ .filter((w) => w.length > 2 && !STOP_WORDS.has(w)),
61+ ),
62+ ].slice(0, 10);
63+ if (!words.length) return null;
64+ return words.map((w) => `"${w.replace(/"/g, "")}"`).join(" OR ");
65+}
+98−0
1+/**
2+ * The semantic index's two adapters: what turns text into vectors
3+ * (`Embedder`) and where vectors are kept and searched (`VectorStore`).
4+ * Docs only ever talks to these, so a self-hosted g1t could put another
5+ * model or vector database behind them. Only Cloudflare's are built
6+ * today: Workers AI (`@cf/baai/bge-base-en-v1.5`, 768 dimensions) and
7+ * Vectorize (the `g1t-docs` index, cosine, metadata indexes on
8+ * `workspace_id` and `space_id`). Without them (no AI or VECTORS
9+ * binding), Docs keeps its passages in D1 and recall matches words only.
10+ */
11+
12+/**
13+ * Workers AI's embedding model, as the index was made with. Pinned:
14+ * vectors from another model mean nothing beside these, so changing it
15+ * means a new index, rebuilt (the backfill, src/indexer.ts).
16+ */
17+export const EMBED_MODEL = "@cf/baai/bge-base-en-v1.5";
18+/** Texts per embedding call. */
19+export const EMBED_BATCH = 50;
20+
21+export type Embedder = {
22+ /** One vector per text, in order. Throws when the model can't answer. */
23+ embed(texts: string[]): Promise<number[][]>;
24+};
25+
26+export type VectorMetadata = {
27+ workspace_id: string;
28+ space_id: string;
29+ kind: "page" | "repo_file";
30+ page_id?: string;
31+ repo_file_id?: string;
32+ repo_id?: string;
33+};
34+
35+export type VectorFilter = {
36+ workspace_id: string;
37+ /** Only these spaces; absent for every space (then the caller filters what comes back). */
38+ space_ids?: string[];
39+};
40+
41+export type VectorMatch = { id: string; score: number };
42+
43+export type VectorStore = {
44+ upsert(vectors: { id: string; values: number[]; metadata: VectorMetadata }[]): Promise<void>;
45+ /** Stored vectors' values by id, for passages that only moved. Missing ids are left out. */
46+ get(ids: string[]): Promise<{ id: string; values: number[] }[]>;
47+ delete(ids: string[]): Promise<void>;
48+ query(vector: number[], options: { topK: number; filter: VectorFilter }): Promise<VectorMatch[]>;
49+};
50+
51+/** Workers AI as the embedder. */
52+export function cloudflareEmbedder(ai: Ai): Embedder {
53+ return {
54+ async embed(texts) {
55+ const out: number[][] = [];
56+ for (let at = 0; at < texts.length; at += EMBED_BATCH) {
57+ const batch = texts.slice(at, at + EMBED_BATCH);
58+ const embedded = (await ai.run(EMBED_MODEL as Parameters<Ai["run"]>[0], { text: batch } as never)) as { data?: number[][] };
59+ const data = embedded.data ?? [];
60+ if (data.length !== batch.length) throw new Error(`the embedding model answered ${data.length} of ${batch.length}`);
61+ out.push(...data);
62+ }
63+ return out;
64+ },
65+ };
66+}
67+
68+/** Vectorize's `getByIds`, `deleteByIds` and `upsert` take at most this many at once (kept well under its limits). */
69+const STORE_BATCH = 20;
70+const UPSERT_BATCH = 100;
71+
72+/** Vectorize as the store. */
73+export function cloudflareVectors(index: Vectorize): VectorStore {
74+ return {
75+ async upsert(vectors) {
76+ for (let at = 0; at < vectors.length; at += UPSERT_BATCH) {
77+ await index.upsert(vectors.slice(at, at + UPSERT_BATCH).map((v) => ({ id: v.id, values: v.values, metadata: v.metadata as unknown as Record<string, VectorizeVectorMetadata> })));
78+ }
79+ },
80+ async get(ids) {
81+ const out: { id: string; values: number[] }[] = [];
82+ for (let at = 0; at < ids.length; at += STORE_BATCH) {
83+ const found = await index.getByIds(ids.slice(at, at + STORE_BATCH));
84+ for (const v of found) if (v.values) out.push({ id: v.id, values: Array.from(v.values as ArrayLike<number>) });
85+ }
86+ return out;
87+ },
88+ async delete(ids) {
89+ for (let at = 0; at < ids.length; at += STORE_BATCH * 5) await index.deleteByIds(ids.slice(at, at + STORE_BATCH * 5));
90+ },
91+ async query(vector, options) {
92+ const filter: Record<string, unknown> = { workspace_id: options.filter.workspace_id };
93+ if (options.filter.space_ids) filter.space_id = { $in: options.filter.space_ids };
94+ const found = await index.query(vector, { topK: options.topK, returnMetadata: "none", returnValues: false, filter: filter as VectorizeVectorMetadataFilter });
95+ return found.matches.map((m) => ({ id: m.id, score: m.score }));
96+ },
97+ };
98+}
+14−0
2626 // Self-hosted, any S3-compatible store instead: DOCS_FILES=s3 and the
2727 // DOCS_S3_* settings (src/files.ts, docs/SELF_HOSTING.md).
2828 "r2_buckets": [{ "binding": "FILES", "bucket_name": "g1t-docs-files" }],
29+ // The semantic index agents recall from (src/indexer.ts, src/vectors.ts):
30+ // Workers AI embeds passages (@cf/baai/bge-base-en-v1.5, 768 dimensions)
31+ // into Vectorize, with metadata indexes on workspace_id and space_id so a
32+ // query only reads spaces its reader may. Create it once:
33+ // npx wrangler vectorize create g1t-docs --dimensions=768 --metric=cosine
34+ // npx wrangler vectorize create-metadata-index g1t-docs --property-name=workspace_id --type=string
35+ // npx wrangler vectorize create-metadata-index g1t-docs --property-name=space_id --type=string
36+ // Without them Docs still keeps passages in D1 and recall matches words.
37+ "ai": { "binding": "AI" },
38+ "vectorize": [{ "binding": "VECTORS", "index_name": "g1t-docs" }],
2939 "services": [
3040 // Workspaces by slug, people's names and avatars, and teams.
3141 { "binding": "IDENTITY", "service": "g1t-identity" },
4656 // service sends it what g1t_contracts::subscribers routes to
4757 // SUBSCRIBER_DOCS. Create it once: npx wrangler queues create g1t-events-docs
4858 "queues": {
59+ // Its own backfill jobs ride the same queue (`docs.index`,
60+ // src/indexer.ts): a workspace's pages and projects' docs indexed in
61+ // batches, and waiting out the hourly embedding cap.
62+ "producers": [{ "binding": "JOBS", "queue": "g1t-events-docs" }],
4963 "consumers": [{ "queue": "g1t-events-docs", "max_batch_size": 20, "max_batch_timeout": 2, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" }]
5064 },
5165 // One room per page: it holds the page's Yjs document and open sockets,