Skip to content
1,465 linesCodeBlameRaw
1/**
2 * The context service: the context hub's catalog, search and scorecards.
3 *
4 * - **Catalog.** What a workspace builds and runs, as entities (projects,
5 * apps, APIs, packages, languages, owners, environments, integrations,
6 * docs) and relations between them. Built by itself: on every push to a
7 * project's default branch it reads the project's manifests and docs
8 * whose blobs changed (at most `MAX_READS` files a push) and joins them
9 * with what projects, deployments, identity and integrations know.
10 * Rebuilding is idempotent: entity ids are hashes of what they are.
11 * - **Search.** Docs, catalog entities, issues, pull requests and kept
12 * memory, embedded with Workers AI into a Vectorize index whose rows
13 * carry their workspace, project and whether they are private, so a
14 * query reads only what its reader may. Matching words in D1 answers
15 * when the index cannot, and fills in after it.
16 * - **Backfill.** Once per workspace (and again on request): the catalog
17 * for every project, memory candidates from docs, manifests and the last
18 * merged pull requests, and the search index filled. One queued job per
19 * project, capped and metered.
20 * - **Run context.** The Context section every g1t agent run starts with.
21 *
22 * Reached through service bindings: `POST /rpc/<method>`.
23 */
24
25import {
26 type Backfill,
27 type CaptureItem,
28 type Catalog,
29 type ContextStatus,
30 type Entity,
31 type EntityDetail,
32 type EntityKind,
33 ENTITY_KINDS,
34 type G1tEvent,
35 type Memory,
36 type Project,
37 type ProjectDeploys,
38 type RelationKind,
39 type RunContext,
40 type Scorecard,
41 type SearchHit,
42 type SearchKind,
43 type SearchResult,
44 type ServiceBinding,
45 type User,
46 type Viewer,
47 ComputeGate,
48 billingClient,
49 embeddingEstimateMicros,
50 localRefusal,
51 platformPaused,
52 currentMovedPath,
53 currentWorkspaceSlug,
54 repoMove,
55 staleMovedPaths,
56 deploymentsClient,
57 fail,
58 granted,
59 identityClient,
60 integrationsClient,
61 memoryReviewClient,
62 ok,
63 projectsClient,
64 reposClient,
65 securityClient,
66 staleSlugs,
67 workClient,
68 type Result,
69} from "@g1t/contracts";
70
71import { assemble, authorsOf, integrationEntities, type EntityDraft, type FileRecord, type ProjectInput, type Surroundings } from "./assemble";
72import { EXTRACT_VERSION, extract, interesting, type FileFacts } from "./extract";
73import { composeRunContext, type ContextNote, type ProjectContext } from "./runcontext";
74import { evaluate } from "./scorecards";
75import { allowedKinds, countVisible, indexFilter, memoryReadable, merge, readable, runMemoryReadable, type IndexMeta, type Reader } from "./visibility";
76
77type Job =
78 | { type: "backfill_project"; workspace: string; slug: string }
79 | { type: "backfill_memory"; workspace: string };
80
81type Env = {
82 DB: D1Database;
83 AI?: Ai;
84 VECTORS?: Vectorize;
85 JOBS: Queue<Job>;
86 REPOS: ServiceBinding;
87 PROJECTS: ServiceBinding;
88 WORK: ServiceBinding;
89 IDENTITY: ServiceBinding;
90 DEPLOYMENTS: ServiceBinding;
91 INTEGRATIONS: ServiceBinding;
92 SECURITY?: ServiceBinding;
93 /** Told the month's embedding cost so far, which it charges once the month is over. */
94 BILLING?: ServiceBinding;
95};
96
97/**
98 * Workers AI's embedding model: 768 dimensions, as the index was made with.
99 * Pinned, not a default staff choose in sudo: vectors from another model
100 * mean nothing beside these, so changing it means a new index, rebuilt.
101 * The catalogue (billing's `gateway_models`) lists it with its price.
102 */
103const EMBED_MODEL = "@cf/baai/bge-base-en-v1.5";
104/** What it costs, in millionths of a dollar per token ($0.067 per million). */
105const MICROS_PER_TOKEN = 0.067;
106/** Past this many tokens in a month, a workspace's new text goes unindexed; text search still finds it. */
107const MONTHLY_TOKENS = 20_000_000;
108/** The most text embedded for one row; the model reads 512 tokens. */
109const EMBED_CHARS = 2000;
110const EMBED_BATCH = 50;
111/** The most files read from a repository for one project's rebuild. */
112const MAX_READS = 15;
113/** Of them, workflows. */
114const MAX_WORKFLOW_READS = 3;
115/** The most projects one push rebuilds. */
116const MAX_PROJECTS_PER_PUSH = 5;
117/** A backfill's caps. */
118const BACKFILL_PROJECTS = 50;
119const BACKFILL_PULLS = 20;
120const BACKFILL_ITEMS = 30;
121/** A backfill still marked running after this long is taken to have died. */
122const BACKFILL_STALE_MS = 30 * 60 * 1000;
123const ITEM_CHARS = 4000;
124const SNIPPET_CHARS = 280;
125/** Kinds every project shares rather than owns. */
126const SHARED: Set<EntityKind> = new Set(["owner", "language", "integration"]);
127
128/** One compute gate per isolate, so entitlements are kept between calls. */
129/** What a backfill says while indexing is paused across g1t (billing's `platform_pause`). */
130const INDEXING_PAUSED = "g1t has paused indexing across the platform for now. Rebuild the context again later.";
131
132let computeGate: ComputeGate | null = null;
133function gateFor(billing: ServiceBinding): ComputeGate {
134 computeGate ??= new ComputeGate(billing);
135 return computeGate;
136}
137
138const now = () => new Date().toISOString();
139const month = () => now().slice(0, 7);
140
141function isMember(viewer: Viewer, workspace: string): boolean {
142 return !!viewer?.workspaces?.some((m) => m.slug === workspace.toLowerCase());
143}
144
145async function sha(text: string): Promise<string> {
146 const digest = await crypto.subtle.digest("SHA-256", new TextEncoder().encode(text));
147 return [...new Uint8Array(digest)].map((b) => b.toString(16).padStart(2, "0")).join("");
148}
149
150/** An entity's id: the same for the same thing, however often it is rebuilt. */
151export async function entityId(workspace: string, kind: string, key: string): Promise<string> {
152 return `ent_${(await sha(`${workspace}\u0000${kind}\u0000${key.toLowerCase()}`)).slice(0, 24)}`;
153}
154
155function snippet(text: string | null | undefined): string {
156 const line = (text ?? "").replace(/\s+/g, " ").trim();
157 return line.length <= SNIPPET_CHARS ? line : `${line.slice(0, SNIPPET_CHARS - 1)}…`;
158}
159
160type EntityRow = {
161 id: string;
162 workspace: string;
163 kind: EntityKind;
164 key: string;
165 name: string;
166 summary: string | null;
167 project_id: string | null;
168 project: string | null;
169 repo_id: string | null;
170 private: number;
171 data: string;
172 source: string;
173 ref: string | null;
174 updated_at: string;
175};
176
177type ItemRow = {
178 id: string;
179 workspace: string;
180 kind: "doc" | "issue" | "pull";
181 entity_id: string | null;
182 project: string | null;
183 private: number;
184 title: string;
185 text: string;
186 url: string | null;
187 by: string | null;
188 updated_at: string;
189};
190
191function toEntity(row: EntityRow): Entity {
192 let data: Record<string, unknown> = {};
193 try {
194 data = JSON.parse(row.data);
195 } catch {}
196 return {
197 id: row.id,
198 workspace: row.workspace,
199 kind: row.kind,
200 key: row.key,
201 name: row.name,
202 summary: row.summary,
203 project: row.project,
204 private: !!row.private,
205 data,
206 source: row.source,
207 ref: row.ref,
208 updatedAt: row.updated_at,
209 };
210}
211
212/** Where a memory came from, in a few words. */
213function memorySource(memory: Memory): string {
214 const ref = memory.source.reference ?? "";
215 const number = memory.source.number;
216 switch (memory.source.kind) {
217 case "doc":
218 return ref.split(":").slice(2).join(":") || "a doc";
219 case "review":
220 return number ? `review on #${number}` : "a review";
221 case "pr":
222 return number ? `#${number}` : "a merged pull request";
223 case "run":
224 return number ? `agent run on #${number}` : "an agent run";
225 default:
226 return `added by ${memory.createdBy}`;
227 }
228}
229
230/** What one workspace's rebuilds share: asked of other services once. */
231class WorkspaceCache {
232 private members?: Promise<string[]>;
233 private deploys?: Promise<ProjectDeploys[]>;
234 private connections?: Promise<Surroundings["integrations"]>;
235 constructor(
236 private readonly env: Env,
237 readonly workspace: string,
238 readonly actor: User,
239 ) {}
240
241 memberNames(): Promise<string[]> {
242 this.members ??= identityClient(this.env.IDENTITY)
243 .listMembers(this.workspace, this.actor)
244 .then((found) => (found.ok ? found.value.map((member) => member.username) : []))
245 .catch(() => []);
246 return this.members;
247 }
248
249 deployments(): Promise<ProjectDeploys[]> {
250 this.deploys ??= deploymentsClient(this.env.DEPLOYMENTS)
251 .overview(this.workspace, this.actor)
252 .then((found) => (found.ok ? found.value : []))
253 .catch(() => []);
254 return this.deploys;
255 }
256
257 integrations(): Promise<Surroundings["integrations"]> {
258 this.connections ??= integrationsClient(this.env.INTEGRATIONS)
259 .list(this.workspace, this.actor)
260 .then((found) =>
261 found.ok
262 ? found.value
263 .filter((connection) => connection.kind !== "models")
264 .map((connection) => ({
265 id: connection.id,
266 provider: connection.provider,
267 name: connection.name,
268 kind: connection.kind,
269 repo: connection.config.repo ?? null,
270 }))
271 : [],
272 )
273 .catch(() => []);
274 return this.connections;
275 }
276}
277
278type ScanStats = { entities: number; candidates: number; kept: number; indexed: number; pruned?: number };
279
280class Context {
281 constructor(private readonly env: Env) {}
282
283 private get db() {
284 return this.env.DB;
285 }
286
287 private async workspaceActor(slug: string): Promise<User | null> {
288 const workspace = await identityClient(this.env.IDENTITY).getWorkspace(slug);
289 if (!workspace) return null;
290 return {
291 id: workspace.id,
292 username: workspace.slug,
293 kind: "workspace",
294 verified: true,
295 workspaces: [{ slug: workspace.slug, role: "member" }],
296 };
297 }
298
299 // ---- Search index --------------------------------------------------------
300
301 /**
302 * Whether semantic search is open to a workspace: embeddings are compute,
303 * so only on a paid plan or the trial (`localRefusal`). Text search
304 * answers for everyone else. Closed when billing cannot say, which
305 * leaves text search; open where there is no billing service at all.
306 */
307 private async semanticOpen(workspace: string): Promise<boolean> {
308 if (!this.env.BILLING) return true;
309 const ent = await gateFor(this.env.BILLING).entitlements(workspace);
310 return ent != null && localRefusal(ent, "embedding", false) == null;
311 }
312
313 /**
314 * Embeds and stores rows in the search index, within the workspace's
315 * monthly allowance, reserving what it costs with billing first and
316 * settling what it did. Never throws.
317 */
318 private async index(workspace: string, rows: { id: string; text: string; meta: IndexMeta }[]): Promise<number> {
319 const { AI, VECTORS } = this.env;
320 if (!AI || !VECTORS || rows.length === 0) return 0;
321 if (!(await this.semanticOpen(workspace))) return 0;
322 let reservation: string | null = null;
323 let spentTokens = 0;
324 if (this.env.BILLING) {
325 const estimate = rows.reduce((sum, row) => sum + Math.ceil(Math.min(row.text.length, EMBED_CHARS) / 4), 0);
326 const admitted = await gateFor(this.env.BILLING).admit({
327 workspace,
328 repo: { namespace: workspace, name: rows[0].meta.project ?? "" },
329 public: rows.every((row) => !row.meta.private),
330 kind: "embedding",
331 estimateMicros: embeddingEstimateMicros(estimate),
332 });
333 if (!admitted.ok) {
334 console.log("embeddings not started for", workspace, admitted.code, admitted.message);
335 return 0;
336 }
337 reservation = admitted.reservation?.id ?? null;
338 }
339 try {
340 const used = await this.db
341 .prepare("SELECT tokens FROM usage WHERE workspace = ? AND month = ?")
342 .bind(workspace, month())
343 .first<{ tokens: number }>();
344 if ((used?.tokens ?? 0) >= MONTHLY_TOKENS) return 0;
345 let stored = 0;
346 for (let at = 0; at < rows.length; at += EMBED_BATCH) {
347 const batch = rows.slice(at, at + EMBED_BATCH);
348 const texts = batch.map((row) => row.text.slice(0, EMBED_CHARS));
349 const embedded = (await AI.run(EMBED_MODEL, { text: texts })) as { data?: number[][] };
350 const vectors = (embedded.data ?? []).map((values, i) => ({
351 id: batch[i].id,
352 values,
353 metadata: batch[i].meta as unknown as Record<string, VectorizeVectorMetadata>,
354 }));
355 if (vectors.length) await VECTORS.upsert(vectors);
356 stored += vectors.length;
357 const tokens = (rows: { text: string }[]) => rows.reduce((sum, row) => sum + Math.ceil(row.text.length / 4), 0);
358 const all = texts.map((text) => ({ text }));
359 spentTokens += tokens(all);
360 await this.meter(workspace, tokens(all), tokens(all.filter((_, i) => batch[i].meta.private)));
361 }
362 return stored;
363 } catch (error) {
364 console.error("could not index", rows.length, "rows for", workspace, error);
365 return 0;
366 } finally {
367 if (reservation && this.env.BILLING) {
368 await gateFor(this.env.BILLING).settle(reservation, embeddingEstimateMicros(spentTokens));
369 }
370 }
371 }
372
373 /** Drops a repository's rows kept under any of `workspaces`, and their index entries. */
374 private async forgetRepo(repoId: string, workspaces: string[]): Promise<void> {
375 const marks = workspaces.map(() => "?").join(", ");
376 const items = await this.db
377 .prepare(`SELECT id FROM items WHERE repo_id = ? AND workspace IN (${marks})`)
378 .bind(repoId, ...workspaces)
379 .all<{ id: string }>();
380 for (let i = 0; i < items.results.length; i += 100) {
381 await this.unindex(items.results.slice(i, i + 100).map((row) => row.id));
382 }
383 const projects = `SELECT project_id FROM entities WHERE repo_id = ? AND workspace IN (${marks}) AND project_id IS NOT NULL`;
384 await this.db.batch([
385 this.db.prepare(`DELETE FROM relations WHERE project_id IN (${projects})`).bind(repoId, ...workspaces),
386 this.db.prepare(`DELETE FROM files WHERE project_id IN (${projects})`).bind(repoId, ...workspaces),
387 this.db.prepare(`DELETE FROM items WHERE repo_id = ? AND workspace IN (${marks})`).bind(repoId, ...workspaces),
388 this.db.prepare(`DELETE FROM scans WHERE repo_id = ? AND workspace IN (${marks})`).bind(repoId, ...workspaces),
389 this.db.prepare(`DELETE FROM entities WHERE repo_id = ? AND workspace IN (${marks})`).bind(repoId, ...workspaces),
390 ]);
391 }
392
393 /** A deleted workspace's hub: everything but its usage, which billing has already read. */
394 private async forgetWorkspace(slug: string): Promise<void> {
395 const items = await this.db.prepare("SELECT id FROM items WHERE workspace = ?").bind(slug).all<{ id: string }>();
396 for (let i = 0; i < items.results.length; i += 100) {
397 await this.unindex(items.results.slice(i, i + 100).map((row) => row.id));
398 }
399 await this.db.batch(
400 ["items", "relations", "entities", "scans", "backfills"].map((table) =>
401 this.db.prepare(`DELETE FROM ${table} WHERE workspace = ?`).bind(slug),
402 ),
403 );
404 }
405
406 private async unindex(ids: string[]): Promise<void> {
407 if (!this.env.VECTORS || ids.length === 0) return;
408 try {
409 await this.env.VECTORS.deleteByIds(ids);
410 } catch (error) {
411 console.error("could not take", ids.length, "rows out of the index", error);
412 }
413 }
414
415 /**
416 * Records what embedding cost. Of it, `billableTokens` are private
417 * text's: public repositories' text and searches are never charged.
418 * Billing is told the month's billable total, which it charges at cost
419 * plus its margin once the month is over (its `embedding_tokens` meter).
420 * A failure to tell billing only delays it: the next call sends the
421 * whole month again.
422 */
423 private async meter(workspace: string, tokens: number, billableTokens = 0): Promise<void> {
424 const total = await this.db
425 .prepare(
426 `INSERT INTO usage (workspace, month, tokens, cost_micros, billable_micros) VALUES (?1, ?2, ?3, ?4, ?5)
427 ON CONFLICT (workspace, month) DO UPDATE SET
428 tokens = tokens + ?3, cost_micros = cost_micros + ?4, billable_micros = billable_micros + ?5
429 RETURNING billable_micros`,
430 )
431 .bind(workspace, month(), tokens, Math.ceil(tokens * MICROS_PER_TOKEN), billableTokens * MICROS_PER_TOKEN)
432 .first<{ billable_micros: number }>();
433 if (this.env.BILLING && total && billableTokens > 0) {
434 await billingClient(this.env.BILLING)
435 .notePending(workspace, "context", Math.ceil(total.billable_micros))
436 .catch((error) => console.error("could not tell billing what embedding cost", workspace, error));
437 }
438 }
439
440 /**
441 * A query's embedding, for semantic search. Too small to reserve one by
442 * one (a few hundred tokens, a fraction of a cent): it needs the plan or
443 * the trial, as indexing does, and is metered with the month's usage.
444 */
445 private async embedQuery(workspace: string, query: string): Promise<number[] | null> {
446 if (!this.env.AI || !this.env.VECTORS) return null;
447 if (!(await this.semanticOpen(workspace))) return null;
448 try {
449 const embedded = (await this.env.AI.run(EMBED_MODEL, { text: [query.slice(0, EMBED_CHARS)] })) as { data?: number[][] };
450 await this.meter(workspace, Math.ceil(query.length / 4));
451 return embedded.data?.[0] ?? null;
452 } catch (error) {
453 console.error("could not embed a query", error);
454 return null;
455 }
456 }
457
458 // ---- Building the catalog --------------------------------------------------
459
460 /** The files of a project worth reading, with their blobs, found in a few tree reads. */
461 private async candidates(project: Project, actor: User, ref: string): Promise<{ files: { path: string; hash: string }[]; siblings: string[]; head: string | null; partial: boolean } | null> {
462 if (project.source.kind !== "hosted") return null;
463 const repos = reposClient(this.env.REPOS);
464 const { repo, rootDir: root } = project.source;
465 const at = (path: string) => [root, path].filter(Boolean).join("/");
466 const top = await repos.tree(repo, actor, ref, root);
467 if (!top.ok) return null;
468 const files: { path: string; hash: string }[] = [];
469 const siblings = top.value.entries.map((entry) => entry.name);
470 const blobs = (entries: { name: string; hash: string; kind: string }[], dir: string) => {
471 for (const entry of entries) {
472 if (entry.kind !== "blob" && entry.kind !== "exec") continue;
473 const path = dir ? `${dir}/${entry.name}` : entry.name;
474 if (interesting(path)) files.push({ path, hash: entry.hash });
475 }
476 };
477 blobs(top.value.entries, "");
478 const dirs = new Set(top.value.entries.filter((entry) => entry.kind === "tree").map((entry) => entry.name));
479 // A folder that could not be listed leaves the list partial.
480 let partial = false;
481 const look = async (dir: string) => {
482 const found = await repos.tree(repo, actor, ref, at(dir)).catch(() => null);
483 if (!found?.ok) partial = true;
484 return found?.ok ? found.value.entries : [];
485 };
486 for (const dir of ["docs", "doc", "runbooks"]) if (dirs.has(dir)) blobs(await look(dir), dir);
487 for (const dir of [".g1t", ".github"]) {
488 if (!dirs.has(dir)) continue;
489 const inside = await look(dir);
490 blobs(inside, dir);
491 if (inside.some((entry) => entry.name === "workflows" && entry.kind === "tree")) blobs(await look(`${dir}/workflows`), `${dir}/workflows`);
492 }
493 return { files, siblings, head: top.value.head?.hash ?? null, partial };
494 }
495
496 /**
497 * Builds one project's place in the catalog from its default branch (or
498 * `commit`), reading only files whose blobs changed unless `force`.
499 * Indexes what changed and sends what its docs say to memory.
500 */
501 async scan(project: Project, cache: WorkspaceCache, commit: string | null, force: boolean): Promise<ScanStats> {
502 const stats: ScanStats = { entities: 0, candidates: 0, kept: 0, indexed: 0 };
503 if (project.source.kind !== "hosted") return stats;
504 const { repo, repoId, rootDir, defaultBranch } = project.source;
505 const ref = commit ?? defaultBranch;
506 const found = await this.candidates(project, cache.actor, ref);
507 if (!found) return stats;
508 const workspace = project.workspace;
509
510 // What each file says: stored for unchanged blobs, read for the rest.
511 const stored = new Map(
512 (
513 await this.db.prepare("SELECT path, hash, facts FROM files WHERE project_id = ?").bind(project.id).all<{ path: string; hash: string; facts: string }>()
514 ).results.map((row) => [row.path, row]),
515 );
516 const files: FileRecord[] = [];
517 const changed: string[] = [];
518 const writes: D1PreparedStatement[] = [];
519 let reads = 0;
520 let workflowReads = 0;
521 // Whether every file is known as this version of extract reads it: only
522 // then are the doc candidates it no longer suggests let go.
523 let complete = !found.partial;
524 const repos = reposClient(this.env.REPOS);
525 for (const file of found.files) {
526 const before = stored.get(file.path);
527 const kept = before ? (JSON.parse(before.facts) as FileFacts) : null;
528 const workflow = file.path.includes("/workflows/");
529 // Facts an older extract read are read again, as a changed file is.
530 const fresh = before && before.hash === file.hash && kept?.version === EXTRACT_VERSION && !force;
531 const canRead = reads < MAX_READS && (!workflow || workflowReads < MAX_WORKFLOW_READS);
532 if (fresh || !canRead) {
533 if (!fresh) complete = false;
534 if (kept) files.push({ path: file.path, facts: kept });
535 continue;
536 }
537 reads++;
538 if (workflow) workflowReads++;
539 const blob = await repos.blob(repo, cache.actor, found.head ?? ref, [rootDir, file.path].filter(Boolean).join("/")).catch(() => null);
540 if (!blob?.ok || blob.value.text == null) {
541 complete = false;
542 continue;
543 }
544 const facts = extract(file.path, blob.value.text, { project: project.name, siblings: found.siblings });
545 files.push({ path: file.path, facts });
546 changed.push(file.path);
547 writes.push(
548 this.db
549 .prepare(
550 `INSERT INTO files (project_id, path, hash, facts, read_at) VALUES (?, ?, ?, ?, ?)
551 ON CONFLICT (project_id, path) DO UPDATE SET hash = excluded.hash, facts = excluded.facts, read_at = excluded.read_at`,
552 )
553 .bind(project.id, file.path, file.hash, JSON.stringify(facts), now()),
554 );
555 }
556 const present = new Set(found.files.map((file) => file.path));
557 for (const path of stored.keys()) {
558 if (!present.has(path)) writes.push(this.db.prepare("DELETE FROM files WHERE project_id = ? AND path = ?").bind(project.id, path));
559 }
560
561 // What g1t knows about it besides its files.
562 const [members, deploys, integrations, log] = await Promise.all([
563 cache.memberNames(),
564 cache.deployments(),
565 cache.integrations(),
566 repos.log(repo, cache.actor, found.head ?? ref, 100).catch(() => null),
567 ]);
568 const deploy = deploys.find((d) => d.slug === project.slug) ?? null;
569 const around: Surroundings = {
570 owners: authorsOf(log?.ok ? log.value : [], members),
571 deploy: deploy
572 ? {
573 enabled: deploy.enabled,
574 production: deploy.production ? { url: deploy.production.url, commit: deploy.production.commit, deployedAt: deploy.production.deployedAt } : null,
575 previews: deploy.previews,
576 latest: deploy.latest ? { status: deploy.latest.status, kind: deploy.latest.kind, error: deploy.latest.error, createdAt: deploy.latest.createdAt } : null,
577 }
578 : null,
579 integrations,
580 };
581 const input: ProjectInput = {
582 id: project.id,
583 workspace,
584 slug: project.slug,
585 name: project.name,
586 description: project.description,
587 private: project.private,
588 repoId,
589 repo,
590 rootDir,
591 defaultBranch,
592 };
593 const built = assemble(input, files, around);
594 const drafts: EntityDraft[] = [...built.entities, ...integrationEntities(integrations)];
595
596 // Entities: upserted; the project's own that it no longer has, removed.
597 const previous = new Map(
598 (
599 await this.db
600 .prepare("SELECT id, name, summary FROM entities WHERE workspace = ? AND (project_id = ? OR project_id IS NULL)")
601 .bind(workspace, project.id)
602 .all<{ id: string; name: string; summary: string | null }>()
603 ).results.map((row) => [row.id, row]),
604 );
605 const ids: string[] = [];
606 const toIndex: { id: string; text: string; meta: IndexMeta }[] = [];
607 const at = now();
608 for (const draft of drafts) {
609 const id = await entityId(workspace, draft.kind, draft.key);
610 ids.push(id);
611 const shared = SHARED.has(draft.kind);
612 const source = draft.kind === "app" || draft.kind === "environment" ? "deployments" : draft.kind === "integration" ? "integrations" : "scan";
613 writes.push(
614 this.db
615 .prepare(
616 `INSERT INTO entities (id, workspace, kind, key, name, summary, project_id, project, repo_id, private, data, source, ref, updated_at)
617 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
618 ON CONFLICT (workspace, kind, key) DO UPDATE SET name = excluded.name, summary = excluded.summary,
619 project_id = excluded.project_id, project = excluded.project, repo_id = excluded.repo_id, private = excluded.private,
620 data = excluded.data, source = excluded.source, ref = excluded.ref, updated_at = excluded.updated_at`,
621 )
622 .bind(
623 id,
624 workspace,
625 draft.kind,
626 draft.key,
627 draft.name,
628 draft.summary,
629 shared ? null : project.id,
630 shared ? null : project.slug,
631 shared ? null : repoId,
632 shared ? 0 : project.private ? 1 : 0,
633 JSON.stringify(draft.data),
634 source,
635 draft.ref,
636 at,
637 ),
638 );
639 const before = previous.get(id);
640 if (force || !before || before.name !== draft.name || before.summary !== draft.summary) {
641 toIndex.push({
642 id,
643 text: `${draft.kind} ${draft.name}. ${draft.summary ?? ""}`,
644 meta: {
645 workspace,
646 kind: draft.kind,
647 project: shared ? "" : project.slug,
648 private: shared ? false : project.private,
649 title: draft.name,
650 snippet: snippet(draft.summary),
651 url: draft.ref ?? "",
652 source: "catalog",
653 by: "",
654 at,
655 },
656 });
657 }
658 }
659 const gone = (
660 await this.db
661 .prepare("SELECT id FROM entities WHERE project_id = ? AND id NOT IN (SELECT value FROM json_each(?))")
662 .bind(project.id, JSON.stringify(ids))
663 .all<{ id: string }>()
664 ).results.map((row) => row.id);
665 if (gone.length) {
666 writes.push(this.db.prepare("DELETE FROM entities WHERE id IN (SELECT value FROM json_each(?))").bind(JSON.stringify(gone)));
667 writes.push(this.db.prepare("DELETE FROM items WHERE entity_id IN (SELECT value FROM json_each(?))").bind(JSON.stringify(gone)));
668 }
669
670 // Relations: the project's own, replaced whole.
671 writes.push(this.db.prepare("DELETE FROM relations WHERE project_id = ?").bind(project.id));
672 for (const relation of built.relations) {
673 const [from, to] = await Promise.all([entityId(workspace, relation.from.kind, relation.from.key), entityId(workspace, relation.to.kind, relation.to.key)]);
674 writes.push(
675 this.db
676 .prepare("INSERT OR REPLACE INTO relations (workspace, from_id, kind, to_id, project_id, updated_at) VALUES (?, ?, ?, ?, ?, ?)")
677 .bind(workspace, from, relation.kind, to, project.id, at),
678 );
679 }
680
681 // Docs that changed: their pieces, for search.
682 const repoPath = `/${repo.namespace}/${repo.name}`;
683 const chunkIds: string[] = [];
684 for (const file of files) {
685 const doc = file.facts.doc;
686 if (!doc || (!changed.includes(file.path) && !force)) continue;
687 const docId = await entityId(workspace, "doc", `${project.slug}:${doc.path}`);
688 writes.push(this.db.prepare("DELETE FROM items WHERE entity_id = ?").bind(docId));
689 const url = `${repoPath}/blob/${defaultBranch}/${[rootDir, doc.path].filter(Boolean).join("/")}`;
690 doc.chunks.forEach((text, i) => {
691 const id = `${docId}:${i}`;
692 chunkIds.push(id);
693 writes.push(
694 this.db
695 .prepare(
696 `INSERT OR REPLACE INTO items (id, workspace, kind, entity_id, project_id, project, repo_id, private, title, text, url, by, updated_at)
697 VALUES (?, ?, 'doc', ?, ?, ?, ?, ?, ?, ?, ?, NULL, ?)`,
698 )
699 .bind(id, workspace, docId, project.id, project.slug, repoId, project.private ? 1 : 0, `${doc.title} (${doc.path})`, text, url, at),
700 );
701 toIndex.push({
702 id,
703 text: `${doc.title}\n${text}`,
704 meta: { workspace, kind: "doc", project: project.slug, private: project.private, title: `${doc.title} (${doc.path})`, snippet: snippet(text), url, source: doc.path, by: "", at },
705 });
706 });
707 }
708 writes.push(
709 this.db
710 .prepare(
711 `INSERT INTO scans (project_id, workspace, repo_id, "commit", tests, scanned_at) VALUES (?, ?, ?, ?, ?, ?)
712 ON CONFLICT (project_id) DO UPDATE SET workspace = excluded.workspace, "commit" = excluded."commit", tests = excluded.tests, scanned_at = excluded.scanned_at`,
713 )
714 .bind(project.id, workspace, repoId, found.head, built.tests ? 1 : 0, at),
715 );
716 for (let i = 0; i < writes.length; i += 50) await this.db.batch(writes.slice(i, i + 50));
717 await this.unindex(gone);
718 stats.entities = drafts.length;
719 stats.indexed = await this.index(workspace, toIndex);
720
721 // What changed files say worth remembering, as memory candidates.
722 const items: CaptureItem[] = built.hints
723 .filter((hint) => force || changed.includes(hint.path))
724 .map((hint) => ({
725 scope: "project",
726 repoId,
727 kind: hint.kind,
728 text: hint.text,
729 confidence: hint.confidence,
730 source: "doc",
731 reference: `doc:${repoId}:${[rootDir, hint.path].filter(Boolean).join("/")}`,
732 evidence: hint.evidence,
733 }));
734 for (let i = 0; i < items.length; i += 50) {
735 const captured = await memoryReviewClient(this.env.WORK)
736 .captureMemories(workspace, items.slice(i, i + 50), "g1t")
737 .catch(() => null);
738 stats.candidates += captured?.added ?? 0;
739 stats.kept += captured?.kept ?? 0;
740 }
741 // Candidates from this project's docs that are still waiting but that
742 // its docs, as read now, no longer suggest: a line since changed, or one
743 // the rules for what is worth remembering leave out. Kept and dismissed
744 // memory is never touched.
745 if (complete) {
746 const pruned = await memoryReviewClient(this.env.WORK)
747 .pruneDocCandidates(workspace, repoId, built.hints.map((hint) => hint.text))
748 .catch(() => null);
749 stats.pruned = pruned?.removed ?? 0;
750 }
751 return stats;
752 }
753
754 // ---- Issues, pull requests and memory in search ------------------------
755
756 private async repoOf(repoId: string): Promise<{ namespace: string; name: string } | null> {
757 const response = await this.env.REPOS.fetch("https://repos/rpc/path_by_id", {
758 method: "POST",
759 headers: { "content-type": "application/json" },
760 body: JSON.stringify({ id: repoId }),
761 });
762 return response.ok ? ((await response.json()) as { namespace: string; name: string } | null) : null;
763 }
764
765 /** Stores and indexes issues and pull requests as search rows. */
766 private async putItems(
767 workspace: string,
768 project: Project | null,
769 repoId: string,
770 rows: { kind: "issue" | "pull"; number: number; title: string; body: string | null; by: string | null; updatedAt: string; url: string; private: boolean }[],
771 ): Promise<number> {
772 const writes: D1PreparedStatement[] = [];
773 const toIndex: { id: string; text: string; meta: IndexMeta }[] = [];
774 for (const row of rows) {
775 const id = `${row.kind}:${repoId}#${row.number}`;
776 const title = `#${row.number} ${row.title}`;
777 const text = (row.body ?? "").slice(0, ITEM_CHARS);
778 writes.push(
779 this.db
780 .prepare(
781 `INSERT OR REPLACE INTO items (id, workspace, kind, entity_id, project_id, project, repo_id, private, title, text, url, by, updated_at)
782 VALUES (?, ?, ?, NULL, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
783 )
784 .bind(id, workspace, row.kind, project?.id ?? null, project?.slug ?? null, repoId, row.private ? 1 : 0, title, text, row.url, row.by, row.updatedAt),
785 );
786 toIndex.push({
787 id,
788 text: `${row.title}\n${text}`,
789 meta: { workspace, kind: row.kind, project: project?.slug ?? "", private: row.private, title, snippet: snippet(text || row.title), url: row.url, source: row.kind === "issue" ? "issue" : "pull request", by: row.by ?? "", at: row.updatedAt },
790 });
791 }
792 for (let i = 0; i < writes.length; i += 50) await this.db.batch(writes.slice(i, i + 50));
793 return this.index(workspace, toIndex);
794 }
795
796 private async indexIssueOrPull(kind: "issue" | "pull", repoId: string, number: number): Promise<void> {
797 const path = await this.repoOf(repoId);
798 if (!path) return;
799 const workspace = path.namespace.toLowerCase();
800 const actor = await this.workspaceActor(workspace);
801 if (!actor) return;
802 const [repo, projects] = await Promise.all([reposClient(this.env.REPOS).get(path, actor), projectsClient(this.env.PROJECTS).byRepo(repoId)]);
803 if (!repo.ok || repo.value.forkOf) return;
804 const work = workClient(this.env.WORK);
805 const base = `/${path.namespace}/${path.name}`;
806 if (kind === "issue") {
807 const found = await work.getIssue(path, number, actor);
808 if (!found.ok) return;
809 const issue = found.value.issue;
810 await this.putItems(workspace, projects[0] ?? null, repoId, [
811 { kind, number, title: issue.title, body: issue.body, by: issue.author.username, updatedAt: issue.updatedAt, url: `${base}/issues/${number}`, private: repo.value.isPrivate },
812 ]);
813 } else {
814 const found = await work.getPull(path, number, actor);
815 if (!found.ok) return;
816 const pull = found.value.pull as { title: string; body: string | null; agent: string; updatedAt?: string; author?: { username: string } };
817 await this.putItems(workspace, projects[0] ?? null, repoId, [
818 { kind, number, title: pull.title, body: pull.body, by: pull.author?.username ?? pull.agent, updatedAt: pull.updatedAt ?? now(), url: `${base}/pull/${number}`, private: repo.value.isPrivate },
819 ]);
820 }
821 }
822
823 /** A memory in search while it is kept, and out of it otherwise. */
824 private async indexMemories(workspace: string, memories: Memory[]): Promise<number> {
825 const kept = memories.filter((memory) => (memory.status ?? "kept") === "kept");
826 await this.unindex(memories.filter((memory) => !kept.includes(memory)).map((memory) => `memory:${memory.id}`));
827 const projects = await projectsClient(this.env.PROJECTS)
828 .list(workspace, (await this.workspaceActor(workspace)) ?? null)
829 .then((found) => (found.ok ? found.value : []))
830 .catch(() => [] as Project[]);
831 const slugOf = (memory: Memory) =>
832 memory.scope === "project" && memory.repo
833 ? (projects.find(
834 (project) =>
835 project.source.kind === "hosted" &&
836 project.primary &&
837 project.source.repo.namespace.toLowerCase() === memory.repo!.namespace.toLowerCase() &&
838 project.source.repo.name.toLowerCase() === memory.repo!.name.toLowerCase(),
839 )?.slug ?? "")
840 : "";
841 return this.index(
842 workspace,
843 kept.map((memory) => ({
844 id: `memory:${memory.id}`,
845 text: `${memory.kind}: ${memory.text}`,
846 meta: {
847 workspace,
848 kind: "memory",
849 project: slugOf(memory),
850 // Memory is for members, whatever its project.
851 private: true,
852 title: `${memory.kind[0].toUpperCase()}${memory.kind.slice(1)}${memory.scope === "workspace" ? " (workspace)" : ""}`,
853 snippet: snippet(memory.text),
854 url: memory.repo ? `/${memory.repo.namespace}/${memory.repo.name}/memory` : `/${workspace}/-/memory`,
855 source: memorySource(memory),
856 by: memory.createdBy,
857 at: memory.updatedAt,
858 },
859 })),
860 );
861 }
862
863 // ---- Events --------------------------------------------------------------
864
865 async onEvent(event: G1tEvent): Promise<void> {
866 switch (event.type) {
867 case "git.push": {
868 if (!event.data.defaultBranch) return;
869 const projects = (await projectsClient(this.env.PROJECTS).byRepo(event.data.repoId)).slice(0, MAX_PROJECTS_PER_PUSH);
870 if (projects.length === 0) return;
871 const actor = await this.workspaceActor(projects[0].workspace);
872 if (!actor) return;
873 const cache = new WorkspaceCache(this.env, projects[0].workspace, actor);
874 for (const project of projects) await this.scan(project, cache, event.data.after, false);
875 return;
876 }
877 case "issue.opened":
878 case "issue.updated":
879 case "issue.closed":
880 return this.indexIssueOrPull("issue", event.data.repoId, event.data.number);
881 case "pull.ready":
882 case "pull.merged":
883 return this.indexIssueOrPull("pull", event.data.repoId, event.data.number);
884 case "memory.changed": {
885 const { memoryId, workspace, status } = event.data;
886 if (status !== "kept") return this.unindex([`memory:${memoryId}`]);
887 const memories = await memoryReviewClient(this.env.WORK).memoriesById(workspace, [memoryId]);
888 await this.indexMemories(workspace, memories);
889 return;
890 }
891 case "repo.transferred":
892 case "repo.renamed": {
893 // What the hub knew about the repository under its old path goes,
894 // index included (its rows carry the path: a transfer's old
895 // workspace, a rename's old name); the workspace it is in now is
896 // built again, which reads it where and as it is now. Nothing here
897 // is the only copy.
898 const move = repoMove(event)!;
899 const current = await currentMovedPath(this.env.REPOS, move);
900 const stale = staleMovedPaths(move, current);
901 if (stale.length === 0) return;
902 const workspace = current.split("/")[0]!.toLowerCase();
903 await this.forgetRepo(move.repoId, [...new Set(stale.map((path) => path.split("/")[0]!.toLowerCase()))]);
904 const actor = await this.workspaceActor(workspace);
905 if (actor) await this.backfill({ actor, workspace });
906 return;
907 }
908 case "repo.deleted":
909 case "repo.purged": {
910 // Deleted, it is hidden: what the hub knew of it goes, index
911 // included, so no search or agent finds it. A restore builds it
912 // again; a purge finds nothing left.
913 await this.forgetRepo(event.data.repoId, [event.data.namespace.toLowerCase()]);
914 return;
915 }
916 case "repo.restored": {
917 const workspace = event.data.namespace.toLowerCase();
918 const actor = await this.workspaceActor(workspace);
919 if (actor) await this.backfill({ actor, workspace });
920 return;
921 }
922 case "workspace.deleted": {
923 await this.forgetWorkspace(event.data.slug);
924 return;
925 }
926 case "workspace.renamed": {
927 const current = await currentWorkspaceSlug(this.env.IDENTITY, event.data);
928 for (const old of staleSlugs(event.data, current)) {
929 await this.db.batch(
930 ["entities", "relations", "items", "scans", "usage"].map((table) =>
931 this.db.prepare(`UPDATE OR IGNORE ${table} SET workspace = ? WHERE workspace = ?`).bind(current, old),
932 ),
933 );
934 await this.db.prepare("DELETE FROM backfills WHERE workspace = ?").bind(old).run();
935 }
936 // The index's rows carry the slug: build them again under the new one.
937 const actor = await this.workspaceActor(current);
938 if (actor) await this.backfill({ actor, workspace: current });
939 return;
940 }
941 default:
942 return;
943 }
944 }
945
946 // ---- Backfill ------------------------------------------------------------
947
948 private async backfillRow(workspace: string): Promise<Backfill | null> {
949 const row = await this.db.prepare("SELECT * FROM backfills WHERE workspace = ?").bind(workspace).first<Record<string, unknown>>();
950 if (!row) return null;
951 return {
952 workspace,
953 status: row.status as Backfill["status"],
954 by: String(row.by),
955 projects: Number(row.projects),
956 done: Number(row.done),
957 entities: Number(row.entities),
958 candidates: Number(row.candidates),
959 kept: Number(row.kept),
960 indexed: Number(row.indexed),
961 error: (row.error as string | null) ?? null,
962 startedAt: String(row.started_at),
963 finishedAt: (row.finished_at as string | null) ?? null,
964 };
965 }
966
967 async backfill(a: { actor: User; workspace: string }): Promise<Result<Backfill>> {
968 const workspace = a.workspace.toLowerCase();
969 if (!isMember(a.actor, workspace)) return fail("forbidden", "Only members can rebuild a workspace's context.");
970 if (await platformPaused(this.env.BILLING, "indexing")) return fail("paused", INDEXING_PAUSED);
971 const running = await this.backfillRow(workspace);
972 if (running?.status === "running" && Date.now() - Date.parse(running.startedAt) < BACKFILL_STALE_MS) return ok(running);
973 const actor = (await this.workspaceActor(workspace)) ?? a.actor;
974 const listed = await projectsClient(this.env.PROJECTS).list(workspace, actor);
975 if (!listed.ok) return listed;
976 const projects = listed.value.filter((project) => project.source.kind === "hosted").slice(0, BACKFILL_PROJECTS);
977 const at = now();
978 await this.db
979 .prepare(
980 `INSERT OR REPLACE INTO backfills (workspace, status, by, projects, done, entities, candidates, kept, indexed, error, started_at, finished_at)
981 VALUES (?, ?, ?, ?, 0, 0, 0, 0, 0, NULL, ?, ?)`,
982 )
983 .bind(workspace, projects.length ? "running" : "done", a.actor.username, projects.length, at, projects.length ? null : at)
984 .run();
985 const jobs: { body: Job }[] = [
986 ...projects.map((project) => ({ body: { type: "backfill_project", workspace, slug: project.slug } as Job })),
987 { body: { type: "backfill_memory", workspace } },
988 ];
989 for (let i = 0; i < jobs.length; i += 100) await this.env.JOBS.sendBatch(jobs.slice(i, i + 100));
990 return ok((await this.backfillRow(workspace))!);
991 }
992
993 async runJob(job: Job): Promise<void> {
994 const actor = await this.workspaceActor(job.workspace);
995 if (!actor) return;
996 // Indexing paused across g1t: the job is counted done with the pause
997 // as its note, so the backfill finishes and can be run again later.
998 if (await platformPaused(this.env.BILLING, "indexing")) {
999 if (job.type === "backfill_project") {
1000 await this.db
1001 .prepare(
1002 `UPDATE backfills SET done = done + 1, error = ?,
1003 status = CASE WHEN done + 1 >= projects THEN 'done' ELSE status END,
1004 finished_at = CASE WHEN done + 1 >= projects THEN ? ELSE finished_at END
1005 WHERE workspace = ?`,
1006 )
1007 .bind(INDEXING_PAUSED, now(), job.workspace)
1008 .run();
1009 }
1010 return;
1011 }
1012 if (job.type === "backfill_memory") {
1013 const memories = await memoryReviewClient(this.env.WORK).searchMemories(job.workspace, null, { limit: 100 });
1014 const indexed = await this.indexMemories(job.workspace, memories);
1015 await this.db.prepare("UPDATE backfills SET indexed = indexed + ? WHERE workspace = ?").bind(indexed, job.workspace).run();
1016 return;
1017 }
1018 const stats: ScanStats = { entities: 0, candidates: 0, kept: 0, indexed: 0 };
1019 let error: string | null = null;
1020 try {
1021 const found = await projectsClient(this.env.PROJECTS).get(job.workspace, job.slug, actor);
1022 if (found.ok && found.value.source.kind === "hosted") {
1023 const project = found.value;
1024 const source = project.source as Extract<Project["source"], { kind: "hosted" }>;
1025 const cache = new WorkspaceCache(this.env, job.workspace, actor);
1026 Object.assign(stats, await this.scan(project, cache, null, true));
1027 if (project.primary) {
1028 // Decisions from merged pull requests, and people's corrections in their reviews.
1029 const seeded = await memoryReviewClient(this.env.WORK).seedFromPulls(source.repoId, BACKFILL_PULLS).catch(() => null);
1030 stats.candidates += seeded?.added ?? 0;
1031 stats.kept += seeded?.kept ?? 0;
1032 const work = workClient(this.env.WORK);
1033 const base = `/${source.repo.namespace}/${source.repo.name}`;
1034 const [issues, pulls] = await Promise.all([
1035 work.listIssues(source.repo, actor).catch(() => null),
1036 work.listPulls(source.repo, actor, "closed").catch(() => null),
1037 ]);
1038 stats.indexed += await this.putItems(job.workspace, project, source.repoId, [
1039 ...(issues?.ok ? issues.value.slice(0, BACKFILL_ITEMS) : []).map((issue) => ({
1040 kind: "issue" as const,
1041 number: issue.number,
1042 title: issue.title,
1043 body: issue.body,
1044 by: issue.author.username,
1045 updatedAt: issue.updatedAt,
1046 url: `${base}/issues/${issue.number}`,
1047 private: project.private,
1048 })),
1049 ...(pulls?.ok ? pulls.value.filter((pull) => pull.status === "merged").slice(0, BACKFILL_ITEMS) : []).map((pull) => ({
1050 kind: "pull" as const,
1051 number: pull.number,
1052 title: pull.title,
1053 body: pull.body,
1054 by: (pull as { author?: { username: string } }).author?.username ?? pull.agent,
1055 updatedAt: pull.mergedAt ?? now(),
1056 url: `${base}/pull/${pull.number}`,
1057 private: project.private,
1058 })),
1059 ]);
1060 }
1061 }
1062 } catch (caught) {
1063 error = String(caught).slice(0, 300);
1064 console.error("backfill of", job.workspace, job.slug, "failed", caught);
1065 }
1066 await this.db
1067 .prepare(
1068 `UPDATE backfills SET done = done + 1, entities = entities + ?, candidates = candidates + ?, kept = kept + ?, indexed = indexed + ?,
1069 error = COALESCE(?, error),
1070 status = CASE WHEN done + 1 >= projects THEN 'done' ELSE status END,
1071 finished_at = CASE WHEN done + 1 >= projects THEN ? ELSE finished_at END
1072 WHERE workspace = ?`,
1073 )
1074 .bind(stats.entities, stats.candidates, stats.kept, stats.indexed, error, now(), job.workspace)
1075 .run();
1076 }
1077
1078 // ---- Reading -------------------------------------------------------------
1079
1080 /**
1081 * Who is reading, and what of the workspace they may see. Someone who can
1082 * read every repository in it (an owner, a member while its base
1083 * permission is Read or more, its own token) sees everything; anyone else
1084 * sees the projects the projects service lists for them, which are those
1085 * whose repositories they can read.
1086 */
1087 private async reader(workspace: string, viewer: Viewer): Promise<Reader> {
1088 const member = isMember(viewer, workspace);
1089 const full = member && !!viewer && granted(viewer, { id: "", namespace: workspace, isPrivate: true }) != null;
1090 if (full) return { workspace, member, full, visible: new Set() };
1091 const listed = await projectsClient(this.env.PROJECTS).list(workspace, viewer).catch(() => null);
1092 const projects = listed?.ok ? listed.value : [];
1093 return {
1094 workspace,
1095 member,
1096 full,
1097 visible: new Set(projects.map((p) => p.slug)),
1098 privateVisible: projects.some((p) => p.private),
1099 repos: new Set(
1100 projects.flatMap((p) => (p.source.kind === "hosted" ? [`${p.source.repo.namespace}/${p.source.repo.name}`.toLowerCase()] : [])),
1101 ),
1102 };
1103 }
1104
1105 private visibleRow(row: { private: number; project: string | null }, reader: Reader): boolean {
1106 return readable({ workspace: reader.workspace, kind: "entity", project: row.project ?? "", private: !!row.private }, reader);
1107 }
1108
1109 async status(a: { workspace: string; viewer: Viewer }): Promise<Result<ContextStatus>> {
1110 const workspace = a.workspace.toLowerCase();
1111 if (!isMember(a.viewer, workspace)) return fail("not_found", "There is no such workspace.");
1112 let backfill = await this.backfillRow(workspace);
1113 // The hub is never empty for long: a workspace's first look starts it.
1114 if (!backfill) {
1115 const actor = await this.workspaceActor(workspace);
1116 if (actor) {
1117 const started = await this.backfill({ actor: { ...actor, username: "g1t" }, workspace });
1118 backfill = started.ok ? started.value : null;
1119 }
1120 }
1121 // Counted over what the viewer may read: a member whose base permission
1122 // is None counts only the projects they were given.
1123 const reader = await this.reader(workspace, a.viewer);
1124 const [counted, usage] = await Promise.all([
1125 this.db
1126 .prepare("SELECT kind, project, private, COUNT(*) AS n FROM entities WHERE workspace = ? GROUP BY kind, project, private")
1127 .bind(workspace)
1128 .all<{ kind: EntityKind; project: string | null; private: number; n: number }>(),
1129 this.db.prepare("SELECT tokens, cost_micros FROM usage WHERE workspace = ? AND month = ?").bind(workspace, month()).first<{ tokens: number; cost_micros: number }>(),
1130 ]);
1131 return ok({
1132 backfill,
1133 counts: countVisible(counted.results, reader),
1134 usage: { month: month(), tokens: usage?.tokens ?? 0, costMicros: usage?.cost_micros ?? 0 },
1135 // Embeddings are compute: a paid plan or the trial (semanticOpen).
1136 semantic: !!(this.env.AI && this.env.VECTORS) && (await this.semanticOpen(workspace)),
1137 });
1138 }
1139
1140 async catalog(a: { workspace: string; viewer: Viewer; kind?: EntityKind | null; project?: string | null }): Promise<Result<Catalog>> {
1141 const workspace = a.workspace.toLowerCase();
1142 const reader = await this.reader(workspace, a.viewer);
1143 let sql = "SELECT * FROM entities WHERE workspace = ?";
1144 const params: unknown[] = [workspace];
1145 if (a.kind) {
1146 sql += " AND kind = ?";
1147 params.push(a.kind);
1148 }
1149 if (a.project) {
1150 sql += " AND (project = ? OR project IS NULL)";
1151 params.push(a.project.toLowerCase());
1152 }
1153 sql += " ORDER BY kind, name COLLATE NOCASE LIMIT 2000";
1154 const rows = (await this.db.prepare(sql).bind(...params).all<EntityRow>()).results.filter((row) => this.visibleRow(row, reader));
1155 const ids = new Set(rows.map((row) => row.id));
1156 const relations = (
1157 await this.db.prepare("SELECT from_id, kind, to_id FROM relations WHERE workspace = ? LIMIT 10000").bind(workspace).all<{ from_id: string; kind: RelationKind; to_id: string }>()
1158 ).results
1159 .filter((row) => ids.has(row.from_id) && ids.has(row.to_id))
1160 .map((row) => ({ from: row.from_id, kind: row.kind, to: row.to_id }));
1161 const built = await this.db.prepare("SELECT MAX(scanned_at) AS at FROM scans WHERE workspace = ?").bind(workspace).first<{ at: string | null }>();
1162 return ok({ entities: rows.map(toEntity), relations, builtAt: built?.at ?? null });
1163 }
1164
1165 async entity(a: { workspace: string; viewer: Viewer; kind: EntityKind; id: string }): Promise<Result<EntityDetail>> {
1166 const workspace = a.workspace.toLowerCase();
1167 if (!ENTITY_KINDS.includes(a.kind)) return fail("invalid", `kind is one of ${ENTITY_KINDS.join(", ")}.`);
1168 const reader = await this.reader(workspace, a.viewer);
1169 const row =
1170 (await this.db.prepare("SELECT * FROM entities WHERE workspace = ? AND id = ?").bind(workspace, a.id).first<EntityRow>()) ??
1171 (await this.db.prepare("SELECT * FROM entities WHERE workspace = ? AND kind = ? AND key = ? COLLATE NOCASE").bind(workspace, a.kind, a.id).first<EntityRow>());
1172 if (!row || row.kind !== a.kind || !this.visibleRow(row, reader)) return fail("not_found", `There is no such ${a.kind} in ${workspace}'s catalog.`);
1173 const [out, into] = await Promise.all([
1174 this.db
1175 .prepare("SELECT r.kind AS rel, e.* FROM relations r JOIN entities e ON e.id = r.to_id WHERE r.from_id = ? LIMIT 200")
1176 .bind(row.id)
1177 .all<EntityRow & { rel: RelationKind }>(),
1178 this.db
1179 .prepare("SELECT r.kind AS rel, e.* FROM relations r JOIN entities e ON e.id = r.from_id WHERE r.to_id = ? LIMIT 200")
1180 .bind(row.id)
1181 .all<EntityRow & { rel: RelationKind }>(),
1182 ]);
1183 const relations = [
1184 ...out.results.filter((r) => this.visibleRow(r, reader)).map((r) => ({ kind: r.rel, direction: "out" as const, entity: toEntity(r) })),
1185 ...into.results.filter((r) => this.visibleRow(r, reader)).map((r) => ({ kind: r.rel, direction: "in" as const, entity: toEntity(r) })),
1186 ];
1187 return ok({ entity: toEntity(row), relations });
1188 }
1189
1190 async search(a: { workspace: string; viewer: Viewer; query: string; project?: string | null; kinds?: SearchKind[] | null; limit?: number | null }): Promise<Result<SearchResult>> {
1191 const workspace = a.workspace.toLowerCase();
1192 const query = (a.query ?? "").trim().slice(0, 500);
1193 if (!query) return fail("invalid", "Say what to search for.");
1194 const reader = await this.reader(workspace, a.viewer);
1195 if (!reader.full && !reader.member && reader.visible.size === 0) return ok({ query, hits: [], mode: "text" });
1196 const limit = Math.min(Math.max(a.limit ?? 20, 1), 50);
1197 const kinds = allowedKinds(reader, a.kinds);
1198 const wants = (kind: SearchKind) => !kinds || kinds.includes(kind);
1199 const project = a.project?.toLowerCase() || null;
1200
1201 // Semantic: the index, filtered to what this reader may see.
1202 let semantic: SearchHit[] = [];
1203 let mode: SearchResult["mode"] = "text";
1204 const vector = await this.embedQuery(workspace, query);
1205 if (vector && this.env.VECTORS) {
1206 try {
1207 const found = await this.env.VECTORS.query(vector, { topK: 20, returnMetadata: "all", filter: indexFilter(reader, { project, kinds }) as VectorizeVectorMetadataFilter });
1208 mode = "semantic";
1209 semantic = found.matches
1210 .map((match) => ({ id: match.id, score: match.score, meta: match.metadata as unknown as IndexMeta }))
1211 .filter((match) => match.meta && readable(match.meta, reader))
1212 .map((match) => ({
1213 kind: match.meta.kind as SearchKind,
1214 id: match.id.startsWith("memory:") ? match.id.slice(7) : match.id,
1215 title: match.meta.title,
1216 snippet: match.meta.snippet,
1217 project: match.meta.project || null,
1218 url: match.meta.url || null,
1219 score: match.score,
1220 source: match.meta.source,
1221 by: match.meta.by || null,
1222 updatedAt: match.meta.at || null,
1223 }));
1224 // Memory changes after it is indexed: show only what is still kept.
1225 const memoryIds = semantic.filter((hit) => hit.kind === "memory").map((hit) => hit.id);
1226 if (memoryIds.length) {
1227 const kept = new Set(
1228 (await memoryReviewClient(this.env.WORK).memoriesById(workspace, memoryIds))
1229 .filter((memory) => (memory.status ?? "kept") === "kept" && memoryReadable(memory.repo, reader))
1230 .map((memory) => memory.id),
1231 );
1232 semantic = semantic.filter((hit) => hit.kind !== "memory" || kept.has(hit.id));
1233 }
1234 } catch (error) {
1235 console.error("semantic search failed; matching words instead", error);
1236 }
1237 }
1238
1239 // Text: every word, in the catalog, docs, issues and pull requests, and memory.
1240 const words = query.toLowerCase().split(/\s+/).filter(Boolean).slice(0, 6);
1241 const like = (columns: string) => words.map(() => `(${columns}) LIKE ?`).join(" AND ");
1242 const patterns = words.map((word) => `%${word.replace(/[%_]/g, "")}%`);
1243 const text: SearchHit[] = [];
1244 const entityKinds = ENTITY_KINDS.filter(wants);
1245 if (entityKinds.length) {
1246 const rows = await this.db
1247 .prepare(
1248 `SELECT * FROM entities WHERE workspace = ? AND kind IN (SELECT value FROM json_each(?)) ${project ? "AND (project = ? OR project IS NULL)" : ""}
1249 AND ${like("lower(name || ' ' || COALESCE(summary, ''))")} LIMIT 30`,
1250 )
1251 .bind(workspace, JSON.stringify(entityKinds), ...(project ? [project] : []), ...patterns)
1252 .all<EntityRow>();
1253 for (const row of rows.results.filter((r) => this.visibleRow(r, reader))) {
1254 text.push({ kind: row.kind, id: row.id, title: row.name, snippet: snippet(row.summary), project: row.project, url: row.ref, score: 0.3, source: "catalog", by: null, updatedAt: row.updated_at });
1255 }
1256 }
1257 const itemKinds = (["doc", "issue", "pull"] as const).filter(wants);
1258 if (itemKinds.length) {
1259 const rows = await this.db
1260 .prepare(
1261 `SELECT * FROM items WHERE workspace = ? AND kind IN (SELECT value FROM json_each(?)) ${project ? "AND project = ?" : ""}
1262 AND ${like("lower(title || ' ' || text)")} ORDER BY updated_at DESC LIMIT 30`,
1263 )
1264 .bind(workspace, JSON.stringify(itemKinds), ...(project ? [project] : []), ...patterns)
1265 .all<ItemRow>();
1266 for (const row of rows.results.filter((r) => this.visibleRow(r, reader))) {
1267 text.push({ kind: row.kind, id: row.id, title: row.title, snippet: snippet(row.text), project: row.project, url: row.url, score: 0.2, source: row.kind === "doc" ? "doc" : row.kind, by: row.by, updatedAt: row.updated_at });
1268 }
1269 }
1270 if (reader.member && wants("memory")) {
1271 const memories = await memoryReviewClient(this.env.WORK).searchMemories(workspace, query, { limit: 10 }).catch(() => [] as Memory[]);
1272 for (const memory of memories.filter((m) => memoryReadable(m.repo, reader))) {
1273 text.push({
1274 kind: "memory",
1275 id: memory.id,
1276 title: `${memory.kind[0].toUpperCase()}${memory.kind.slice(1)}${memory.scope === "workspace" ? " (workspace)" : ""}`,
1277 snippet: snippet(memory.text),
1278 project: null,
1279 url: memory.repo ? `/${memory.repo.namespace}/${memory.repo.name}/memory` : `/${workspace}/-/memory`,
1280 score: memory.pinned ? 0.35 : 0.25,
1281 source: memorySource(memory),
1282 by: memory.createdBy,
1283 updatedAt: memory.updatedAt,
1284 });
1285 }
1286 }
1287 return ok({ query, hits: merge(semantic, text, limit), mode });
1288 }
1289
1290 async scorecards(a: { workspace: string; viewer: Viewer }): Promise<Result<Scorecard[]>> {
1291 const workspace = a.workspace.toLowerCase();
1292 if (!isMember(a.viewer, workspace)) return fail("forbidden", "Scorecards are for members of the workspace.");
1293 const listed = await projectsClient(this.env.PROJECTS).list(workspace, a.viewer);
1294 if (!listed.ok) return listed;
1295 const [entities, scans, deploys, security] = await Promise.all([
1296 this.db
1297 .prepare("SELECT * FROM entities WHERE workspace = ? AND kind IN ('project', 'doc')")
1298 .bind(workspace)
1299 .all<EntityRow>(),
1300 this.db.prepare("SELECT project_id, tests FROM scans WHERE workspace = ?").bind(workspace).all<{ project_id: string; tests: number }>(),
1301 deploymentsClient(this.env.DEPLOYMENTS)
1302 .overview(workspace, a.viewer)
1303 .then((found) => (found.ok ? found.value : []))
1304 .catch(() => [] as ProjectDeploys[]),
1305 this.env.SECURITY
1306 ? securityClient(this.env.SECURITY)
1307 .workspace(workspace, a.viewer)
1308 .then((found) => (found.ok ? found.value : null))
1309 .catch(() => null)
1310 : Promise.resolve(null),
1311 ]);
1312 const tests = new Map(scans.results.map((row) => [row.project_id, !!row.tests]));
1313 const cards: Scorecard[] = [];
1314 for (const project of listed.value) {
1315 if (project.source.kind !== "hosted") continue;
1316 const own = entities.results.filter((row) => row.project_id === project.id);
1317 const entry = own.find((row) => row.kind === "project");
1318 const data = entry ? toEntity(entry).data : {};
1319 const deploy = deploys.find((d) => d.slug === project.slug);
1320 const repoId = project.source.repoId;
1321 const findings = security?.find((repo) => repo.repoId === repoId);
1322 const rules = evaluate({
1323 name: project.name,
1324 owners: Array.isArray(data.owners) ? (data.owners as string[]) : [],
1325 docs: own.filter((row) => row.kind === "doc").map((row) => String(toEntity(row).data.path ?? "")),
1326 tests: tests.get(project.id) ?? false,
1327 testCommand: Array.isArray(data.testCommands) && data.testCommands.length ? String(data.testCommands[0]) : null,
1328 deploy: deploy
1329 ? {
1330 enabled: deploy.enabled,
1331 production: deploy.production ? { url: deploy.production.url } : null,
1332 latest: deploy.latest ? { kind: deploy.latest.kind, status: deploy.latest.status, error: deploy.latest.error } : null,
1333 }
1334 : null,
1335 secretFindings: security ? (findings?.secrets ?? 0) : null,
1336 });
1337 const applies = rules.filter((rule) => rule.status !== "na");
1338 cards.push({
1339 project: project.slug,
1340 name: project.name,
1341 repo: project.source.repo,
1342 passed: applies.filter((rule) => rule.status === "pass").length,
1343 total: applies.length,
1344 rules,
1345 });
1346 }
1347 return ok(cards);
1348 }
1349
1350 /**
1351 * The Context section for an agent starting work, holding only what the
1352 * person it acts for (`requester`) may read: an outside collaborator's run
1353 * is told the project's memory, never the workspace's, and only the
1354 * projects around it they can read. No requester is the workspace's own
1355 * step. Never fails a run: on any trouble, nothing.
1356 */
1357 async runContext(a: { repoId: string; task?: string; budget?: number; requester?: Viewer }): Promise<RunContext> {
1358 try {
1359 const projects = (await projectsClient(this.env.PROJECTS).byRepo(a.repoId)).slice(0, 2);
1360 if (projects.length === 0) return { text: null, sources: [] };
1361 const workspace = projects[0].workspace;
1362 const actor = await this.workspaceActor(workspace);
1363 if (!actor) return { text: null, sources: [] };
1364 const reader = a.requester ? await this.reader(workspace, a.requester) : null;
1365 const home = projects.find((project) => project.source.kind === "hosted");
1366 const repoKey = home?.source.kind === "hosted" ? `${home.source.repo.namespace}/${home.source.repo.name}` : "";
1367 const cache = new WorkspaceCache(this.env, workspace, actor);
1368 const deploys = await cache.deployments();
1369 const live = new Map(deploys.map((d) => [d.slug, d]));
1370 const contexts: ProjectContext[] = [];
1371 for (const project of projects) {
1372 if (project.source.kind !== "hosted") continue;
1373 const rows = (await this.db.prepare("SELECT * FROM entities WHERE workspace = ? AND project_id = ?").bind(workspace, project.id).all<EntityRow>()).results.map(toEntity);
1374 const entry = rows.find((row) => row.kind === "project");
1375 const deploy = live.get(project.slug);
1376 contexts.push({
1377 slug: project.slug,
1378 name: project.name,
1379 repo: `${project.source.repo.namespace}/${project.source.repo.name}`,
1380 rootDir: project.source.rootDir,
1381 languages: Array.isArray(entry?.data.languages) ? (entry!.data.languages as string[]) : [],
1382 packages: rows.filter((row) => row.kind === "package").map((row) => row.name).slice(0, 5),
1383 testCommands: Array.isArray(entry?.data.testCommands) ? (entry!.data.testCommands as string[]) : [],
1384 owners: Array.isArray(entry?.data.owners) ? (entry!.data.owners as string[]) : [],
1385 environments: deploy?.enabled
1386 ? [{ name: "Production", url: deploy.production?.url ?? null, status: deploy.latest?.kind === "production" ? deploy.latest.status : deploy.production ? "ready" : null }]
1387 : [],
1388 docs: rows.filter((row) => row.kind === "doc").map((row) => String(row.data.path ?? row.name)),
1389 });
1390 }
1391 // Workspace memory (no project) only for a run its members may be told it.
1392 const slugs = new Set([...(reader && !reader.member ? [] : [""]), ...projects.map((project) => project.slug)]);
1393 const review = memoryReviewClient(this.env.WORK);
1394 // The memories closest to the task, less the pinned ones every run already has.
1395 let memories: ContextNote[] = [];
1396 const task = (a.task ?? "").trim();
1397 const vector = task ? await this.embedQuery(workspace, task.slice(0, 1500)) : null;
1398 if (vector && this.env.VECTORS) {
1399 const found = await this.env.VECTORS.query(vector, { topK: 20, returnMetadata: "all", filter: { workspace, kind: "memory" } }).catch(() => null);
1400 const ids = (found?.matches ?? [])
1401 .filter((match) => slugs.has(String((match.metadata as { project?: string } | undefined)?.project ?? "")) && match.score >= 0.5)
1402 .map((match) => match.id.slice("memory:".length));
1403 if (ids.length) {
1404 const byId = new Map((await review.memoriesById(workspace, ids)).map((memory) => [memory.id, memory]));
1405 memories = ids
1406 .map((id) => byId.get(id))
1407 .filter((memory): memory is Memory => !!memory && (memory.status ?? "kept") === "kept" && !memory.pinned && runMemoryReadable(memory, reader, repoKey))
1408 .slice(0, 6)
1409 .map((memory) => ({ id: memory.id, kind: memory.kind, text: memory.text, source: memorySource(memory) }));
1410 }
1411 }
1412 const shown = new Set(memories.map((note) => note.id));
1413 const decisions = (await review.searchMemories(workspace, null, { repoIds: [a.repoId], limit: 100 }).catch(() => [] as Memory[]))
1414 .filter((memory) => memory.kind === "decision" && !memory.pinned && !shown.has(memory.id) && runMemoryReadable(memory, reader, repoKey))
1415 .sort((x, y) => y.updatedAt.localeCompare(x.updatedAt))
1416 .slice(0, 3)
1417 .map((memory) => ({ id: memory.id, kind: memory.kind, text: memory.text, source: memorySource(memory) }));
1418 return composeRunContext({ projects: contexts, memories, decisions, budget: Math.min(Math.max(a.budget ?? 4000, 500), 12_000) });
1419 } catch (error) {
1420 console.error("no run context for", a.repoId, error);
1421 return { text: null, sources: [] };
1422 }
1423 }
1424}
1425
1426export default {
1427 async fetch(request: Request, env: Env): Promise<Response> {
1428 const match = new URL(request.url).pathname.match(/^\/rpc\/([a-z_]+)$/);
1429 if (request.method !== "POST" || !match) return new Response("Not found\n", { status: 404 });
1430 const service = new Context(env);
1431 const args = (await request.json().catch(() => ({}))) as any;
1432 switch (match[1]) {
1433 case "catalog":
1434 return Response.json(await service.catalog(args));
1435 case "entity":
1436 return Response.json(await service.entity(args));
1437 case "search":
1438 return Response.json(await service.search(args));
1439 case "scorecards":
1440 return Response.json(await service.scorecards(args));
1441 case "backfill":
1442 return Response.json(await service.backfill(args));
1443 case "status":
1444 return Response.json(await service.status(args));
1445 case "run_context":
1446 return Response.json(await service.runContext(args));
1447 default:
1448 return new Response("Unknown method\n", { status: 404 });
1449 }
1450 },
1451
1452 async queue(batch: MessageBatch<G1tEvent | Job>, env: Env): Promise<void> {
1453 const service = new Context(env);
1454 for (const message of batch.messages) {
1455 try {
1456 if (batch.queue === "g1t-context-jobs") await service.runJob(message.body as Job);
1457 else await service.onEvent(message.body as G1tEvent);
1458 message.ack();
1459 } catch (error) {
1460 console.error("context could not handle", (message.body as { type?: string }).type, error);
1461 message.retry();
1462 }
1463 }
1464 },
1465} satisfies ExportedHandler<Env, G1tEvent | Job>;