Skip to content
1,450 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, interesting, type FileFacts } from "./extract";
73import { composeRunContext, type ContextNote, type ProjectContext } from "./runcontext";
74import { evaluate } from "./scorecards";
75import { allowedKinds, countVisible, indexFilter, memoryReadable, merge, projectReadable, 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 };
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 } | 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 const look = async (dir: string) => {
480 const found = await repos.tree(repo, actor, ref, at(dir)).catch(() => null);
481 return found?.ok ? found.value.entries : [];
482 };
483 for (const dir of ["docs", "doc", "runbooks"]) if (dirs.has(dir)) blobs(await look(dir), dir);
484 for (const dir of [".g1t", ".github"]) {
485 if (!dirs.has(dir)) continue;
486 const inside = await look(dir);
487 blobs(inside, dir);
488 if (inside.some((entry) => entry.name === "workflows" && entry.kind === "tree")) blobs(await look(`${dir}/workflows`), `${dir}/workflows`);
489 }
490 return { files, siblings, head: top.value.head?.hash ?? null };
491 }
492
493 /**
494 * Builds one project's place in the catalog from its default branch (or
495 * `commit`), reading only files whose blobs changed unless `force`.
496 * Indexes what changed and sends what its docs say to memory.
497 */
498 async scan(project: Project, cache: WorkspaceCache, commit: string | null, force: boolean): Promise<ScanStats> {
499 const stats: ScanStats = { entities: 0, candidates: 0, kept: 0, indexed: 0 };
500 if (project.source.kind !== "hosted") return stats;
501 const { repo, repoId, rootDir, defaultBranch } = project.source;
502 const ref = commit ?? defaultBranch;
503 const found = await this.candidates(project, cache.actor, ref);
504 if (!found) return stats;
505 const workspace = project.workspace;
506
507 // What each file says: stored for unchanged blobs, read for the rest.
508 const stored = new Map(
509 (
510 await this.db.prepare("SELECT path, hash, facts FROM files WHERE project_id = ?").bind(project.id).all<{ path: string; hash: string; facts: string }>()
511 ).results.map((row) => [row.path, row]),
512 );
513 const files: FileRecord[] = [];
514 const changed: string[] = [];
515 const writes: D1PreparedStatement[] = [];
516 let reads = 0;
517 let workflowReads = 0;
518 const repos = reposClient(this.env.REPOS);
519 for (const file of found.files) {
520 const before = stored.get(file.path);
521 const workflow = file.path.includes("/workflows/");
522 const fresh = before && before.hash === file.hash && !force;
523 const canRead = reads < MAX_READS && (!workflow || workflowReads < MAX_WORKFLOW_READS);
524 if (fresh || !canRead) {
525 if (before) files.push({ path: file.path, facts: JSON.parse(before.facts) as FileFacts });
526 continue;
527 }
528 reads++;
529 if (workflow) workflowReads++;
530 const blob = await repos.blob(repo, cache.actor, found.head ?? ref, [rootDir, file.path].filter(Boolean).join("/")).catch(() => null);
531 if (!blob?.ok || blob.value.text == null) continue;
532 const facts = extract(file.path, blob.value.text, { project: project.name, siblings: found.siblings });
533 files.push({ path: file.path, facts });
534 changed.push(file.path);
535 writes.push(
536 this.db
537 .prepare(
538 `INSERT INTO files (project_id, path, hash, facts, read_at) VALUES (?, ?, ?, ?, ?)
539 ON CONFLICT (project_id, path) DO UPDATE SET hash = excluded.hash, facts = excluded.facts, read_at = excluded.read_at`,
540 )
541 .bind(project.id, file.path, file.hash, JSON.stringify(facts), now()),
542 );
543 }
544 const present = new Set(found.files.map((file) => file.path));
545 for (const path of stored.keys()) {
546 if (!present.has(path)) writes.push(this.db.prepare("DELETE FROM files WHERE project_id = ? AND path = ?").bind(project.id, path));
547 }
548
549 // What g1t knows about it besides its files.
550 const [members, deploys, integrations, graph, log] = await Promise.all([
551 cache.memberNames(),
552 cache.deployments(),
553 cache.integrations(),
554 projectsClient(this.env.PROJECTS).graph(project.id).catch(() => ({ dependsOn: [], usedBy: [] })),
555 repos.log(repo, cache.actor, found.head ?? ref, 100).catch(() => null),
556 ]);
557 const deploy = deploys.find((d) => d.slug === project.slug) ?? null;
558 const around: Surroundings = {
559 owners: authorsOf(log?.ok ? log.value : [], members),
560 dependsOn: graph.dependsOn.map((dep) => ({ slug: dep.slug, as: dep.as })),
561 deploy: deploy
562 ? {
563 enabled: deploy.enabled,
564 production: deploy.production ? { url: deploy.production.url, commit: deploy.production.commit, deployedAt: deploy.production.deployedAt } : null,
565 previews: deploy.previews,
566 latest: deploy.latest ? { status: deploy.latest.status, kind: deploy.latest.kind, error: deploy.latest.error, createdAt: deploy.latest.createdAt } : null,
567 }
568 : null,
569 integrations,
570 };
571 const input: ProjectInput = {
572 id: project.id,
573 workspace,
574 slug: project.slug,
575 name: project.name,
576 description: project.description,
577 private: project.private,
578 repoId,
579 repo,
580 rootDir,
581 defaultBranch,
582 };
583 const built = assemble(input, files, around);
584 const drafts: EntityDraft[] = [...built.entities, ...integrationEntities(integrations)];
585
586 // Entities: upserted; the project's own that it no longer has, removed.
587 const previous = new Map(
588 (
589 await this.db
590 .prepare("SELECT id, name, summary FROM entities WHERE workspace = ? AND (project_id = ? OR project_id IS NULL)")
591 .bind(workspace, project.id)
592 .all<{ id: string; name: string; summary: string | null }>()
593 ).results.map((row) => [row.id, row]),
594 );
595 const ids: string[] = [];
596 const toIndex: { id: string; text: string; meta: IndexMeta }[] = [];
597 const at = now();
598 for (const draft of drafts) {
599 const id = await entityId(workspace, draft.kind, draft.key);
600 ids.push(id);
601 const shared = SHARED.has(draft.kind);
602 const source = draft.kind === "app" || draft.kind === "environment" ? "deployments" : draft.kind === "integration" ? "integrations" : "scan";
603 writes.push(
604 this.db
605 .prepare(
606 `INSERT INTO entities (id, workspace, kind, key, name, summary, project_id, project, repo_id, private, data, source, ref, updated_at)
607 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
608 ON CONFLICT (workspace, kind, key) DO UPDATE SET name = excluded.name, summary = excluded.summary,
609 project_id = excluded.project_id, project = excluded.project, repo_id = excluded.repo_id, private = excluded.private,
610 data = excluded.data, source = excluded.source, ref = excluded.ref, updated_at = excluded.updated_at`,
611 )
612 .bind(
613 id,
614 workspace,
615 draft.kind,
616 draft.key,
617 draft.name,
618 draft.summary,
619 shared ? null : project.id,
620 shared ? null : project.slug,
621 shared ? null : repoId,
622 shared ? 0 : project.private ? 1 : 0,
623 JSON.stringify(draft.data),
624 source,
625 draft.ref,
626 at,
627 ),
628 );
629 const before = previous.get(id);
630 if (force || !before || before.name !== draft.name || before.summary !== draft.summary) {
631 toIndex.push({
632 id,
633 text: `${draft.kind} ${draft.name}. ${draft.summary ?? ""}`,
634 meta: {
635 workspace,
636 kind: draft.kind,
637 project: shared ? "" : project.slug,
638 private: shared ? false : project.private,
639 title: draft.name,
640 snippet: snippet(draft.summary),
641 url: draft.ref ?? "",
642 source: "catalog",
643 by: "",
644 at,
645 },
646 });
647 }
648 }
649 const gone = (
650 await this.db
651 .prepare("SELECT id FROM entities WHERE project_id = ? AND id NOT IN (SELECT value FROM json_each(?))")
652 .bind(project.id, JSON.stringify(ids))
653 .all<{ id: string }>()
654 ).results.map((row) => row.id);
655 if (gone.length) {
656 writes.push(this.db.prepare("DELETE FROM entities WHERE id IN (SELECT value FROM json_each(?))").bind(JSON.stringify(gone)));
657 writes.push(this.db.prepare("DELETE FROM items WHERE entity_id IN (SELECT value FROM json_each(?))").bind(JSON.stringify(gone)));
658 }
659
660 // Relations: the project's own, replaced whole.
661 writes.push(this.db.prepare("DELETE FROM relations WHERE project_id = ?").bind(project.id));
662 for (const relation of built.relations) {
663 const [from, to] = await Promise.all([entityId(workspace, relation.from.kind, relation.from.key), entityId(workspace, relation.to.kind, relation.to.key)]);
664 writes.push(
665 this.db
666 .prepare("INSERT OR REPLACE INTO relations (workspace, from_id, kind, to_id, project_id, updated_at) VALUES (?, ?, ?, ?, ?, ?)")
667 .bind(workspace, from, relation.kind, to, project.id, at),
668 );
669 }
670
671 // Docs that changed: their pieces, for search.
672 const repoPath = `/${repo.namespace}/${repo.name}`;
673 const chunkIds: string[] = [];
674 for (const file of files) {
675 const doc = file.facts.doc;
676 if (!doc || (!changed.includes(file.path) && !force)) continue;
677 const docId = await entityId(workspace, "doc", `${project.slug}:${doc.path}`);
678 writes.push(this.db.prepare("DELETE FROM items WHERE entity_id = ?").bind(docId));
679 const url = `${repoPath}/blob/${defaultBranch}/${[rootDir, doc.path].filter(Boolean).join("/")}`;
680 doc.chunks.forEach((text, i) => {
681 const id = `${docId}:${i}`;
682 chunkIds.push(id);
683 writes.push(
684 this.db
685 .prepare(
686 `INSERT OR REPLACE INTO items (id, workspace, kind, entity_id, project_id, project, repo_id, private, title, text, url, by, updated_at)
687 VALUES (?, ?, 'doc', ?, ?, ?, ?, ?, ?, ?, ?, NULL, ?)`,
688 )
689 .bind(id, workspace, docId, project.id, project.slug, repoId, project.private ? 1 : 0, `${doc.title} (${doc.path})`, text, url, at),
690 );
691 toIndex.push({
692 id,
693 text: `${doc.title}\n${text}`,
694 meta: { workspace, kind: "doc", project: project.slug, private: project.private, title: `${doc.title} (${doc.path})`, snippet: snippet(text), url, source: doc.path, by: "", at },
695 });
696 });
697 }
698 writes.push(
699 this.db
700 .prepare(
701 `INSERT INTO scans (project_id, workspace, repo_id, "commit", tests, scanned_at) VALUES (?, ?, ?, ?, ?, ?)
702 ON CONFLICT (project_id) DO UPDATE SET workspace = excluded.workspace, "commit" = excluded."commit", tests = excluded.tests, scanned_at = excluded.scanned_at`,
703 )
704 .bind(project.id, workspace, repoId, found.head, built.tests ? 1 : 0, at),
705 );
706 for (let i = 0; i < writes.length; i += 50) await this.db.batch(writes.slice(i, i + 50));
707 await this.unindex(gone);
708 stats.entities = drafts.length;
709 stats.indexed = await this.index(workspace, toIndex);
710
711 // What changed files say worth remembering, as memory candidates.
712 const items: CaptureItem[] = built.hints
713 .filter((hint) => force || changed.includes(hint.path))
714 .map((hint) => ({
715 scope: "project",
716 repoId,
717 kind: hint.kind,
718 text: hint.text,
719 confidence: hint.confidence,
720 source: "doc",
721 reference: `doc:${repoId}:${[rootDir, hint.path].filter(Boolean).join("/")}`,
722 evidence: hint.evidence,
723 }));
724 for (let i = 0; i < items.length; i += 50) {
725 const captured = await memoryReviewClient(this.env.WORK)
726 .captureMemories(workspace, items.slice(i, i + 50), "g1t")
727 .catch(() => null);
728 stats.candidates += captured?.added ?? 0;
729 stats.kept += captured?.kept ?? 0;
730 }
731 return stats;
732 }
733
734 // ---- Issues, pull requests and memory in search ------------------------
735
736 private async repoOf(repoId: string): Promise<{ namespace: string; name: string } | null> {
737 const response = await this.env.REPOS.fetch("https://repos/rpc/path_by_id", {
738 method: "POST",
739 headers: { "content-type": "application/json" },
740 body: JSON.stringify({ id: repoId }),
741 });
742 return response.ok ? ((await response.json()) as { namespace: string; name: string } | null) : null;
743 }
744
745 /** Stores and indexes issues and pull requests as search rows. */
746 private async putItems(
747 workspace: string,
748 project: Project | null,
749 repoId: string,
750 rows: { kind: "issue" | "pull"; number: number; title: string; body: string | null; by: string | null; updatedAt: string; url: string; private: boolean }[],
751 ): Promise<number> {
752 const writes: D1PreparedStatement[] = [];
753 const toIndex: { id: string; text: string; meta: IndexMeta }[] = [];
754 for (const row of rows) {
755 const id = `${row.kind}:${repoId}#${row.number}`;
756 const title = `#${row.number} ${row.title}`;
757 const text = (row.body ?? "").slice(0, ITEM_CHARS);
758 writes.push(
759 this.db
760 .prepare(
761 `INSERT OR REPLACE INTO items (id, workspace, kind, entity_id, project_id, project, repo_id, private, title, text, url, by, updated_at)
762 VALUES (?, ?, ?, NULL, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
763 )
764 .bind(id, workspace, row.kind, project?.id ?? null, project?.slug ?? null, repoId, row.private ? 1 : 0, title, text, row.url, row.by, row.updatedAt),
765 );
766 toIndex.push({
767 id,
768 text: `${row.title}\n${text}`,
769 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 },
770 });
771 }
772 for (let i = 0; i < writes.length; i += 50) await this.db.batch(writes.slice(i, i + 50));
773 return this.index(workspace, toIndex);
774 }
775
776 private async indexIssueOrPull(kind: "issue" | "pull", repoId: string, number: number): Promise<void> {
777 const path = await this.repoOf(repoId);
778 if (!path) return;
779 const workspace = path.namespace.toLowerCase();
780 const actor = await this.workspaceActor(workspace);
781 if (!actor) return;
782 const [repo, projects] = await Promise.all([reposClient(this.env.REPOS).get(path, actor), projectsClient(this.env.PROJECTS).byRepo(repoId)]);
783 if (!repo.ok || repo.value.forkOf) return;
784 const work = workClient(this.env.WORK);
785 const base = `/${path.namespace}/${path.name}`;
786 if (kind === "issue") {
787 const found = await work.getIssue(path, number, actor);
788 if (!found.ok) return;
789 const issue = found.value.issue;
790 await this.putItems(workspace, projects[0] ?? null, repoId, [
791 { kind, number, title: issue.title, body: issue.body, by: issue.author.username, updatedAt: issue.updatedAt, url: `${base}/issues/${number}`, private: repo.value.isPrivate },
792 ]);
793 } else {
794 const found = await work.getPull(path, number, actor);
795 if (!found.ok) return;
796 const pull = found.value.pull as { title: string; body: string | null; agent: string; updatedAt?: string; author?: { username: string } };
797 await this.putItems(workspace, projects[0] ?? null, repoId, [
798 { 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 },
799 ]);
800 }
801 }
802
803 /** A memory in search while it is kept, and out of it otherwise. */
804 private async indexMemories(workspace: string, memories: Memory[]): Promise<number> {
805 const kept = memories.filter((memory) => (memory.status ?? "kept") === "kept");
806 await this.unindex(memories.filter((memory) => !kept.includes(memory)).map((memory) => `memory:${memory.id}`));
807 const projects = await projectsClient(this.env.PROJECTS)
808 .list(workspace, (await this.workspaceActor(workspace)) ?? null)
809 .then((found) => (found.ok ? found.value : []))
810 .catch(() => [] as Project[]);
811 const slugOf = (memory: Memory) =>
812 memory.scope === "project" && memory.repo
813 ? (projects.find(
814 (project) =>
815 project.source.kind === "hosted" &&
816 project.primary &&
817 project.source.repo.namespace.toLowerCase() === memory.repo!.namespace.toLowerCase() &&
818 project.source.repo.name.toLowerCase() === memory.repo!.name.toLowerCase(),
819 )?.slug ?? "")
820 : "";
821 return this.index(
822 workspace,
823 kept.map((memory) => ({
824 id: `memory:${memory.id}`,
825 text: `${memory.kind}: ${memory.text}`,
826 meta: {
827 workspace,
828 kind: "memory",
829 project: slugOf(memory),
830 // Memory is for members, whatever its project.
831 private: true,
832 title: `${memory.kind[0].toUpperCase()}${memory.kind.slice(1)}${memory.scope === "workspace" ? " (workspace)" : ""}`,
833 snippet: snippet(memory.text),
834 url: memory.repo ? `/${memory.repo.namespace}/${memory.repo.name}/memory` : `/${workspace}/-/memory`,
835 source: memorySource(memory),
836 by: memory.createdBy,
837 at: memory.updatedAt,
838 },
839 })),
840 );
841 }
842
843 // ---- Events --------------------------------------------------------------
844
845 async onEvent(event: G1tEvent): Promise<void> {
846 switch (event.type) {
847 case "git.push": {
848 if (!event.data.defaultBranch) return;
849 const projects = (await projectsClient(this.env.PROJECTS).byRepo(event.data.repoId)).slice(0, MAX_PROJECTS_PER_PUSH);
850 if (projects.length === 0) return;
851 const actor = await this.workspaceActor(projects[0].workspace);
852 if (!actor) return;
853 const cache = new WorkspaceCache(this.env, projects[0].workspace, actor);
854 for (const project of projects) await this.scan(project, cache, event.data.after, false);
855 return;
856 }
857 case "issue.opened":
858 case "issue.updated":
859 case "issue.closed":
860 return this.indexIssueOrPull("issue", event.data.repoId, event.data.number);
861 case "pull.ready":
862 case "pull.merged":
863 return this.indexIssueOrPull("pull", event.data.repoId, event.data.number);
864 case "memory.changed": {
865 const { memoryId, workspace, status } = event.data;
866 if (status !== "kept") return this.unindex([`memory:${memoryId}`]);
867 const memories = await memoryReviewClient(this.env.WORK).memoriesById(workspace, [memoryId]);
868 await this.indexMemories(workspace, memories);
869 return;
870 }
871 case "repo.transferred":
872 case "repo.renamed": {
873 // What the hub knew about the repository under its old path goes,
874 // index included (its rows carry the path: a transfer's old
875 // workspace, a rename's old name); the workspace it is in now is
876 // built again, which reads it where and as it is now. Nothing here
877 // is the only copy.
878 const move = repoMove(event)!;
879 const current = await currentMovedPath(this.env.REPOS, move);
880 const stale = staleMovedPaths(move, current);
881 if (stale.length === 0) return;
882 const workspace = current.split("/")[0]!.toLowerCase();
883 await this.forgetRepo(move.repoId, [...new Set(stale.map((path) => path.split("/")[0]!.toLowerCase()))]);
884 const actor = await this.workspaceActor(workspace);
885 if (actor) await this.backfill({ actor, workspace });
886 return;
887 }
888 case "repo.deleted":
889 case "repo.purged": {
890 // Deleted, it is hidden: what the hub knew of it goes, index
891 // included, so no search or agent finds it. A restore builds it
892 // again; a purge finds nothing left.
893 await this.forgetRepo(event.data.repoId, [event.data.namespace.toLowerCase()]);
894 return;
895 }
896 case "repo.restored": {
897 const workspace = event.data.namespace.toLowerCase();
898 const actor = await this.workspaceActor(workspace);
899 if (actor) await this.backfill({ actor, workspace });
900 return;
901 }
902 case "workspace.deleted": {
903 await this.forgetWorkspace(event.data.slug);
904 return;
905 }
906 case "workspace.renamed": {
907 const current = await currentWorkspaceSlug(this.env.IDENTITY, event.data);
908 for (const old of staleSlugs(event.data, current)) {
909 await this.db.batch(
910 ["entities", "relations", "items", "scans", "usage"].map((table) =>
911 this.db.prepare(`UPDATE OR IGNORE ${table} SET workspace = ? WHERE workspace = ?`).bind(current, old),
912 ),
913 );
914 await this.db.prepare("DELETE FROM backfills WHERE workspace = ?").bind(old).run();
915 }
916 // The index's rows carry the slug: build them again under the new one.
917 const actor = await this.workspaceActor(current);
918 if (actor) await this.backfill({ actor, workspace: current });
919 return;
920 }
921 default:
922 return;
923 }
924 }
925
926 // ---- Backfill ------------------------------------------------------------
927
928 private async backfillRow(workspace: string): Promise<Backfill | null> {
929 const row = await this.db.prepare("SELECT * FROM backfills WHERE workspace = ?").bind(workspace).first<Record<string, unknown>>();
930 if (!row) return null;
931 return {
932 workspace,
933 status: row.status as Backfill["status"],
934 by: String(row.by),
935 projects: Number(row.projects),
936 done: Number(row.done),
937 entities: Number(row.entities),
938 candidates: Number(row.candidates),
939 kept: Number(row.kept),
940 indexed: Number(row.indexed),
941 error: (row.error as string | null) ?? null,
942 startedAt: String(row.started_at),
943 finishedAt: (row.finished_at as string | null) ?? null,
944 };
945 }
946
947 async backfill(a: { actor: User; workspace: string }): Promise<Result<Backfill>> {
948 const workspace = a.workspace.toLowerCase();
949 if (!isMember(a.actor, workspace)) return fail("forbidden", "Only members can rebuild a workspace's context.");
950 if (await platformPaused(this.env.BILLING, "indexing")) return fail("paused", INDEXING_PAUSED);
951 const running = await this.backfillRow(workspace);
952 if (running?.status === "running" && Date.now() - Date.parse(running.startedAt) < BACKFILL_STALE_MS) return ok(running);
953 const actor = (await this.workspaceActor(workspace)) ?? a.actor;
954 const listed = await projectsClient(this.env.PROJECTS).list(workspace, actor);
955 if (!listed.ok) return listed;
956 const projects = listed.value.filter((project) => project.source.kind === "hosted").slice(0, BACKFILL_PROJECTS);
957 const at = now();
958 await this.db
959 .prepare(
960 `INSERT OR REPLACE INTO backfills (workspace, status, by, projects, done, entities, candidates, kept, indexed, error, started_at, finished_at)
961 VALUES (?, ?, ?, ?, 0, 0, 0, 0, 0, NULL, ?, ?)`,
962 )
963 .bind(workspace, projects.length ? "running" : "done", a.actor.username, projects.length, at, projects.length ? null : at)
964 .run();
965 const jobs: { body: Job }[] = [
966 ...projects.map((project) => ({ body: { type: "backfill_project", workspace, slug: project.slug } as Job })),
967 { body: { type: "backfill_memory", workspace } },
968 ];
969 for (let i = 0; i < jobs.length; i += 100) await this.env.JOBS.sendBatch(jobs.slice(i, i + 100));
970 return ok((await this.backfillRow(workspace))!);
971 }
972
973 async runJob(job: Job): Promise<void> {
974 const actor = await this.workspaceActor(job.workspace);
975 if (!actor) return;
976 // Indexing paused across g1t: the job is counted done with the pause
977 // as its note, so the backfill finishes and can be run again later.
978 if (await platformPaused(this.env.BILLING, "indexing")) {
979 if (job.type === "backfill_project") {
980 await this.db
981 .prepare(
982 `UPDATE backfills SET done = done + 1, error = ?,
983 status = CASE WHEN done + 1 >= projects THEN 'done' ELSE status END,
984 finished_at = CASE WHEN done + 1 >= projects THEN ? ELSE finished_at END
985 WHERE workspace = ?`,
986 )
987 .bind(INDEXING_PAUSED, now(), job.workspace)
988 .run();
989 }
990 return;
991 }
992 if (job.type === "backfill_memory") {
993 const memories = await memoryReviewClient(this.env.WORK).searchMemories(job.workspace, null, { limit: 100 });
994 const indexed = await this.indexMemories(job.workspace, memories);
995 await this.db.prepare("UPDATE backfills SET indexed = indexed + ? WHERE workspace = ?").bind(indexed, job.workspace).run();
996 return;
997 }
998 const stats: ScanStats = { entities: 0, candidates: 0, kept: 0, indexed: 0 };
999 let error: string | null = null;
1000 try {
1001 const found = await projectsClient(this.env.PROJECTS).get(job.workspace, job.slug, actor);
1002 if (found.ok && found.value.source.kind === "hosted") {
1003 const project = found.value;
1004 const source = project.source as Extract<Project["source"], { kind: "hosted" }>;
1005 const cache = new WorkspaceCache(this.env, job.workspace, actor);
1006 Object.assign(stats, await this.scan(project, cache, null, true));
1007 if (project.primary) {
1008 // Decisions from merged pull requests, and people's corrections in their reviews.
1009 const seeded = await memoryReviewClient(this.env.WORK).seedFromPulls(source.repoId, BACKFILL_PULLS).catch(() => null);
1010 stats.candidates += seeded?.added ?? 0;
1011 stats.kept += seeded?.kept ?? 0;
1012 const work = workClient(this.env.WORK);
1013 const base = `/${source.repo.namespace}/${source.repo.name}`;
1014 const [issues, pulls] = await Promise.all([
1015 work.listIssues(source.repo, actor).catch(() => null),
1016 work.listPulls(source.repo, actor, "closed").catch(() => null),
1017 ]);
1018 stats.indexed += await this.putItems(job.workspace, project, source.repoId, [
1019 ...(issues?.ok ? issues.value.slice(0, BACKFILL_ITEMS) : []).map((issue) => ({
1020 kind: "issue" as const,
1021 number: issue.number,
1022 title: issue.title,
1023 body: issue.body,
1024 by: issue.author.username,
1025 updatedAt: issue.updatedAt,
1026 url: `${base}/issues/${issue.number}`,
1027 private: project.private,
1028 })),
1029 ...(pulls?.ok ? pulls.value.filter((pull) => pull.status === "merged").slice(0, BACKFILL_ITEMS) : []).map((pull) => ({
1030 kind: "pull" as const,
1031 number: pull.number,
1032 title: pull.title,
1033 body: pull.body,
1034 by: (pull as { author?: { username: string } }).author?.username ?? pull.agent,
1035 updatedAt: pull.mergedAt ?? now(),
1036 url: `${base}/pull/${pull.number}`,
1037 private: project.private,
1038 })),
1039 ]);
1040 }
1041 }
1042 } catch (caught) {
1043 error = String(caught).slice(0, 300);
1044 console.error("backfill of", job.workspace, job.slug, "failed", caught);
1045 }
1046 await this.db
1047 .prepare(
1048 `UPDATE backfills SET done = done + 1, entities = entities + ?, candidates = candidates + ?, kept = kept + ?, indexed = indexed + ?,
1049 error = COALESCE(?, error),
1050 status = CASE WHEN done + 1 >= projects THEN 'done' ELSE status END,
1051 finished_at = CASE WHEN done + 1 >= projects THEN ? ELSE finished_at END
1052 WHERE workspace = ?`,
1053 )
1054 .bind(stats.entities, stats.candidates, stats.kept, stats.indexed, error, now(), job.workspace)
1055 .run();
1056 }
1057
1058 // ---- Reading -------------------------------------------------------------
1059
1060 /**
1061 * Who is reading, and what of the workspace they may see. Someone who can
1062 * read every repository in it (an owner, a member while its base
1063 * permission is Read or more, its own token) sees everything; anyone else
1064 * sees the projects the projects service lists for them, which are those
1065 * whose repositories they can read.
1066 */
1067 private async reader(workspace: string, viewer: Viewer): Promise<Reader> {
1068 const member = isMember(viewer, workspace);
1069 const full = member && !!viewer && granted(viewer, { id: "", namespace: workspace, isPrivate: true }) != null;
1070 if (full) return { workspace, member, full, visible: new Set() };
1071 const listed = await projectsClient(this.env.PROJECTS).list(workspace, viewer).catch(() => null);
1072 const projects = listed?.ok ? listed.value : [];
1073 return {
1074 workspace,
1075 member,
1076 full,
1077 visible: new Set(projects.map((p) => p.slug)),
1078 privateVisible: projects.some((p) => p.private),
1079 repos: new Set(
1080 projects.flatMap((p) => (p.source.kind === "hosted" ? [`${p.source.repo.namespace}/${p.source.repo.name}`.toLowerCase()] : [])),
1081 ),
1082 };
1083 }
1084
1085 private visibleRow(row: { private: number; project: string | null }, reader: Reader): boolean {
1086 return readable({ workspace: reader.workspace, kind: "entity", project: row.project ?? "", private: !!row.private }, reader);
1087 }
1088
1089 async status(a: { workspace: string; viewer: Viewer }): Promise<Result<ContextStatus>> {
1090 const workspace = a.workspace.toLowerCase();
1091 if (!isMember(a.viewer, workspace)) return fail("not_found", "There is no such workspace.");
1092 let backfill = await this.backfillRow(workspace);
1093 // The hub is never empty for long: a workspace's first look starts it.
1094 if (!backfill) {
1095 const actor = await this.workspaceActor(workspace);
1096 if (actor) {
1097 const started = await this.backfill({ actor: { ...actor, username: "g1t" }, workspace });
1098 backfill = started.ok ? started.value : null;
1099 }
1100 }
1101 // Counted over what the viewer may read: a member whose base permission
1102 // is None counts only the projects they were given.
1103 const reader = await this.reader(workspace, a.viewer);
1104 const [counted, usage] = await Promise.all([
1105 this.db
1106 .prepare("SELECT kind, project, private, COUNT(*) AS n FROM entities WHERE workspace = ? GROUP BY kind, project, private")
1107 .bind(workspace)
1108 .all<{ kind: EntityKind; project: string | null; private: number; n: number }>(),
1109 this.db.prepare("SELECT tokens, cost_micros FROM usage WHERE workspace = ? AND month = ?").bind(workspace, month()).first<{ tokens: number; cost_micros: number }>(),
1110 ]);
1111 return ok({
1112 backfill,
1113 counts: countVisible(counted.results, reader),
1114 usage: { month: month(), tokens: usage?.tokens ?? 0, costMicros: usage?.cost_micros ?? 0 },
1115 // Embeddings are compute: a paid plan or the trial (semanticOpen).
1116 semantic: !!(this.env.AI && this.env.VECTORS) && (await this.semanticOpen(workspace)),
1117 });
1118 }
1119
1120 async catalog(a: { workspace: string; viewer: Viewer; kind?: EntityKind | null; project?: string | null }): Promise<Result<Catalog>> {
1121 const workspace = a.workspace.toLowerCase();
1122 const reader = await this.reader(workspace, a.viewer);
1123 let sql = "SELECT * FROM entities WHERE workspace = ?";
1124 const params: unknown[] = [workspace];
1125 if (a.kind) {
1126 sql += " AND kind = ?";
1127 params.push(a.kind);
1128 }
1129 if (a.project) {
1130 sql += " AND (project = ? OR project IS NULL)";
1131 params.push(a.project.toLowerCase());
1132 }
1133 sql += " ORDER BY kind, name COLLATE NOCASE LIMIT 2000";
1134 const rows = (await this.db.prepare(sql).bind(...params).all<EntityRow>()).results.filter((row) => this.visibleRow(row, reader));
1135 const ids = new Set(rows.map((row) => row.id));
1136 const relations = (
1137 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 }>()
1138 ).results
1139 .filter((row) => ids.has(row.from_id) && ids.has(row.to_id))
1140 .map((row) => ({ from: row.from_id, kind: row.kind, to: row.to_id }));
1141 const built = await this.db.prepare("SELECT MAX(scanned_at) AS at FROM scans WHERE workspace = ?").bind(workspace).first<{ at: string | null }>();
1142 return ok({ entities: rows.map(toEntity), relations, builtAt: built?.at ?? null });
1143 }
1144
1145 async entity(a: { workspace: string; viewer: Viewer; kind: EntityKind; id: string }): Promise<Result<EntityDetail>> {
1146 const workspace = a.workspace.toLowerCase();
1147 if (!ENTITY_KINDS.includes(a.kind)) return fail("invalid", `kind is one of ${ENTITY_KINDS.join(", ")}.`);
1148 const reader = await this.reader(workspace, a.viewer);
1149 const row =
1150 (await this.db.prepare("SELECT * FROM entities WHERE workspace = ? AND id = ?").bind(workspace, a.id).first<EntityRow>()) ??
1151 (await this.db.prepare("SELECT * FROM entities WHERE workspace = ? AND kind = ? AND key = ? COLLATE NOCASE").bind(workspace, a.kind, a.id).first<EntityRow>());
1152 if (!row || row.kind !== a.kind || !this.visibleRow(row, reader)) return fail("not_found", `There is no such ${a.kind} in ${workspace}'s catalog.`);
1153 const [out, into] = await Promise.all([
1154 this.db
1155 .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")
1156 .bind(row.id)
1157 .all<EntityRow & { rel: RelationKind }>(),
1158 this.db
1159 .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")
1160 .bind(row.id)
1161 .all<EntityRow & { rel: RelationKind }>(),
1162 ]);
1163 const relations = [
1164 ...out.results.filter((r) => this.visibleRow(r, reader)).map((r) => ({ kind: r.rel, direction: "out" as const, entity: toEntity(r) })),
1165 ...into.results.filter((r) => this.visibleRow(r, reader)).map((r) => ({ kind: r.rel, direction: "in" as const, entity: toEntity(r) })),
1166 ];
1167 return ok({ entity: toEntity(row), relations });
1168 }
1169
1170 async search(a: { workspace: string; viewer: Viewer; query: string; project?: string | null; kinds?: SearchKind[] | null; limit?: number | null }): Promise<Result<SearchResult>> {
1171 const workspace = a.workspace.toLowerCase();
1172 const query = (a.query ?? "").trim().slice(0, 500);
1173 if (!query) return fail("invalid", "Say what to search for.");
1174 const reader = await this.reader(workspace, a.viewer);
1175 if (!reader.full && !reader.member && reader.visible.size === 0) return ok({ query, hits: [], mode: "text" });
1176 const limit = Math.min(Math.max(a.limit ?? 20, 1), 50);
1177 const kinds = allowedKinds(reader, a.kinds);
1178 const wants = (kind: SearchKind) => !kinds || kinds.includes(kind);
1179 const project = a.project?.toLowerCase() || null;
1180
1181 // Semantic: the index, filtered to what this reader may see.
1182 let semantic: SearchHit[] = [];
1183 let mode: SearchResult["mode"] = "text";
1184 const vector = await this.embedQuery(workspace, query);
1185 if (vector && this.env.VECTORS) {
1186 try {
1187 const found = await this.env.VECTORS.query(vector, { topK: 20, returnMetadata: "all", filter: indexFilter(reader, { project, kinds }) as VectorizeVectorMetadataFilter });
1188 mode = "semantic";
1189 semantic = found.matches
1190 .map((match) => ({ id: match.id, score: match.score, meta: match.metadata as unknown as IndexMeta }))
1191 .filter((match) => match.meta && readable(match.meta, reader))
1192 .map((match) => ({
1193 kind: match.meta.kind as SearchKind,
1194 id: match.id.startsWith("memory:") ? match.id.slice(7) : match.id,
1195 title: match.meta.title,
1196 snippet: match.meta.snippet,
1197 project: match.meta.project || null,
1198 url: match.meta.url || null,
1199 score: match.score,
1200 source: match.meta.source,
1201 by: match.meta.by || null,
1202 updatedAt: match.meta.at || null,
1203 }));
1204 // Memory changes after it is indexed: show only what is still kept.
1205 const memoryIds = semantic.filter((hit) => hit.kind === "memory").map((hit) => hit.id);
1206 if (memoryIds.length) {
1207 const kept = new Set(
1208 (await memoryReviewClient(this.env.WORK).memoriesById(workspace, memoryIds))
1209 .filter((memory) => (memory.status ?? "kept") === "kept" && memoryReadable(memory.repo, reader))
1210 .map((memory) => memory.id),
1211 );
1212 semantic = semantic.filter((hit) => hit.kind !== "memory" || kept.has(hit.id));
1213 }
1214 } catch (error) {
1215 console.error("semantic search failed; matching words instead", error);
1216 }
1217 }
1218
1219 // Text: every word, in the catalog, docs, issues and pull requests, and memory.
1220 const words = query.toLowerCase().split(/\s+/).filter(Boolean).slice(0, 6);
1221 const like = (columns: string) => words.map(() => `(${columns}) LIKE ?`).join(" AND ");
1222 const patterns = words.map((word) => `%${word.replace(/[%_]/g, "")}%`);
1223 const text: SearchHit[] = [];
1224 const entityKinds = ENTITY_KINDS.filter(wants);
1225 if (entityKinds.length) {
1226 const rows = await this.db
1227 .prepare(
1228 `SELECT * FROM entities WHERE workspace = ? AND kind IN (SELECT value FROM json_each(?)) ${project ? "AND (project = ? OR project IS NULL)" : ""}
1229 AND ${like("lower(name || ' ' || COALESCE(summary, ''))")} LIMIT 30`,
1230 )
1231 .bind(workspace, JSON.stringify(entityKinds), ...(project ? [project] : []), ...patterns)
1232 .all<EntityRow>();
1233 for (const row of rows.results.filter((r) => this.visibleRow(r, reader))) {
1234 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 });
1235 }
1236 }
1237 const itemKinds = (["doc", "issue", "pull"] as const).filter(wants);
1238 if (itemKinds.length) {
1239 const rows = await this.db
1240 .prepare(
1241 `SELECT * FROM items WHERE workspace = ? AND kind IN (SELECT value FROM json_each(?)) ${project ? "AND project = ?" : ""}
1242 AND ${like("lower(title || ' ' || text)")} ORDER BY updated_at DESC LIMIT 30`,
1243 )
1244 .bind(workspace, JSON.stringify(itemKinds), ...(project ? [project] : []), ...patterns)
1245 .all<ItemRow>();
1246 for (const row of rows.results.filter((r) => this.visibleRow(r, reader))) {
1247 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 });
1248 }
1249 }
1250 if (reader.member && wants("memory")) {
1251 const memories = await memoryReviewClient(this.env.WORK).searchMemories(workspace, query, { limit: 10 }).catch(() => [] as Memory[]);
1252 for (const memory of memories.filter((m) => memoryReadable(m.repo, reader))) {
1253 text.push({
1254 kind: "memory",
1255 id: memory.id,
1256 title: `${memory.kind[0].toUpperCase()}${memory.kind.slice(1)}${memory.scope === "workspace" ? " (workspace)" : ""}`,
1257 snippet: snippet(memory.text),
1258 project: null,
1259 url: memory.repo ? `/${memory.repo.namespace}/${memory.repo.name}/memory` : `/${workspace}/-/memory`,
1260 score: memory.pinned ? 0.35 : 0.25,
1261 source: memorySource(memory),
1262 by: memory.createdBy,
1263 updatedAt: memory.updatedAt,
1264 });
1265 }
1266 }
1267 return ok({ query, hits: merge(semantic, text, limit), mode });
1268 }
1269
1270 async scorecards(a: { workspace: string; viewer: Viewer }): Promise<Result<Scorecard[]>> {
1271 const workspace = a.workspace.toLowerCase();
1272 if (!isMember(a.viewer, workspace)) return fail("forbidden", "Scorecards are for members of the workspace.");
1273 const listed = await projectsClient(this.env.PROJECTS).list(workspace, a.viewer);
1274 if (!listed.ok) return listed;
1275 const [entities, scans, deploys, security] = await Promise.all([
1276 this.db
1277 .prepare("SELECT * FROM entities WHERE workspace = ? AND kind IN ('project', 'doc')")
1278 .bind(workspace)
1279 .all<EntityRow>(),
1280 this.db.prepare("SELECT project_id, tests FROM scans WHERE workspace = ?").bind(workspace).all<{ project_id: string; tests: number }>(),
1281 deploymentsClient(this.env.DEPLOYMENTS)
1282 .overview(workspace, a.viewer)
1283 .then((found) => (found.ok ? found.value : []))
1284 .catch(() => [] as ProjectDeploys[]),
1285 this.env.SECURITY
1286 ? securityClient(this.env.SECURITY)
1287 .workspace(workspace, a.viewer)
1288 .then((found) => (found.ok ? found.value : null))
1289 .catch(() => null)
1290 : Promise.resolve(null),
1291 ]);
1292 const tests = new Map(scans.results.map((row) => [row.project_id, !!row.tests]));
1293 const cards: Scorecard[] = [];
1294 for (const project of listed.value) {
1295 if (project.source.kind !== "hosted") continue;
1296 const own = entities.results.filter((row) => row.project_id === project.id);
1297 const entry = own.find((row) => row.kind === "project");
1298 const data = entry ? toEntity(entry).data : {};
1299 const deploy = deploys.find((d) => d.slug === project.slug);
1300 const repoId = project.source.repoId;
1301 const findings = security?.find((repo) => repo.repoId === repoId);
1302 const rules = evaluate({
1303 name: project.name,
1304 owners: Array.isArray(data.owners) ? (data.owners as string[]) : [],
1305 docs: own.filter((row) => row.kind === "doc").map((row) => String(toEntity(row).data.path ?? "")),
1306 tests: tests.get(project.id) ?? false,
1307 testCommand: Array.isArray(data.testCommands) && data.testCommands.length ? String(data.testCommands[0]) : null,
1308 deploy: deploy
1309 ? {
1310 enabled: deploy.enabled,
1311 production: deploy.production ? { url: deploy.production.url } : null,
1312 latest: deploy.latest ? { kind: deploy.latest.kind, status: deploy.latest.status, error: deploy.latest.error } : null,
1313 }
1314 : null,
1315 secretFindings: security ? (findings?.secrets ?? 0) : null,
1316 });
1317 const applies = rules.filter((rule) => rule.status !== "na");
1318 cards.push({
1319 project: project.slug,
1320 name: project.name,
1321 repo: project.source.repo,
1322 passed: applies.filter((rule) => rule.status === "pass").length,
1323 total: applies.length,
1324 rules,
1325 });
1326 }
1327 return ok(cards);
1328 }
1329
1330 /**
1331 * The Context section for an agent starting work, holding only what the
1332 * person it acts for (`requester`) may read: an outside collaborator's run
1333 * is told the project's memory, never the workspace's, and only the
1334 * projects around it they can read. No requester is the workspace's own
1335 * step. Never fails a run: on any trouble, nothing.
1336 */
1337 async runContext(a: { repoId: string; task?: string; budget?: number; requester?: Viewer }): Promise<RunContext> {
1338 try {
1339 const projects = (await projectsClient(this.env.PROJECTS).byRepo(a.repoId)).slice(0, 2);
1340 if (projects.length === 0) return { text: null, sources: [] };
1341 const workspace = projects[0].workspace;
1342 const actor = await this.workspaceActor(workspace);
1343 if (!actor) return { text: null, sources: [] };
1344 const reader = a.requester ? await this.reader(workspace, a.requester) : null;
1345 const home = projects.find((project) => project.source.kind === "hosted");
1346 const repoKey = home?.source.kind === "hosted" ? `${home.source.repo.namespace}/${home.source.repo.name}` : "";
1347 const cache = new WorkspaceCache(this.env, workspace, actor);
1348 const deploys = await cache.deployments();
1349 const live = new Map(deploys.map((d) => [d.slug, d]));
1350 const contexts: ProjectContext[] = [];
1351 for (const project of projects) {
1352 if (project.source.kind !== "hosted") continue;
1353 const rows = (await this.db.prepare("SELECT * FROM entities WHERE workspace = ? AND project_id = ?").bind(workspace, project.id).all<EntityRow>()).results.map(toEntity);
1354 const entry = rows.find((row) => row.kind === "project");
1355 const graph = await projectsClient(this.env.PROJECTS).graph(project.id).catch(() => ({ dependsOn: [], usedBy: [] }));
1356 const deploy = live.get(project.slug);
1357 contexts.push({
1358 slug: project.slug,
1359 name: project.name,
1360 repo: `${project.source.repo.namespace}/${project.source.repo.name}`,
1361 rootDir: project.source.rootDir,
1362 languages: Array.isArray(entry?.data.languages) ? (entry!.data.languages as string[]) : [],
1363 packages: rows.filter((row) => row.kind === "package").map((row) => row.name).slice(0, 5),
1364 testCommands: Array.isArray(entry?.data.testCommands) ? (entry!.data.testCommands as string[]) : [],
1365 owners: Array.isArray(entry?.data.owners) ? (entry!.data.owners as string[]) : [],
1366 dependsOn: graph.dependsOn
1367 .filter((dep) => projectReadable(dep.slug, reader))
1368 .map((dep) => ({ slug: dep.slug, as: dep.as, url: live.get(dep.slug)?.production?.url ?? null })),
1369 usedBy: graph.usedBy.filter((dep) => projectReadable(dep.slug, reader)).map((dep) => ({ slug: dep.slug })),
1370 environments: deploy?.enabled
1371 ? [{ name: "Production", url: deploy.production?.url ?? null, status: deploy.latest?.kind === "production" ? deploy.latest.status : deploy.production ? "ready" : null }]
1372 : [],
1373 docs: rows.filter((row) => row.kind === "doc").map((row) => String(row.data.path ?? row.name)),
1374 });
1375 }
1376 // Workspace memory (no project) only for a run its members may be told it.
1377 const slugs = new Set([...(reader && !reader.member ? [] : [""]), ...projects.map((project) => project.slug)]);
1378 const review = memoryReviewClient(this.env.WORK);
1379 // The memories closest to the task, less the pinned ones every run already has.
1380 let memories: ContextNote[] = [];
1381 const task = (a.task ?? "").trim();
1382 const vector = task ? await this.embedQuery(workspace, task.slice(0, 1500)) : null;
1383 if (vector && this.env.VECTORS) {
1384 const found = await this.env.VECTORS.query(vector, { topK: 20, returnMetadata: "all", filter: { workspace, kind: "memory" } }).catch(() => null);
1385 const ids = (found?.matches ?? [])
1386 .filter((match) => slugs.has(String((match.metadata as { project?: string } | undefined)?.project ?? "")) && match.score >= 0.5)
1387 .map((match) => match.id.slice("memory:".length));
1388 if (ids.length) {
1389 const byId = new Map((await review.memoriesById(workspace, ids)).map((memory) => [memory.id, memory]));
1390 memories = ids
1391 .map((id) => byId.get(id))
1392 .filter((memory): memory is Memory => !!memory && (memory.status ?? "kept") === "kept" && !memory.pinned && runMemoryReadable(memory, reader, repoKey))
1393 .slice(0, 6)
1394 .map((memory) => ({ id: memory.id, kind: memory.kind, text: memory.text, source: memorySource(memory) }));
1395 }
1396 }
1397 const shown = new Set(memories.map((note) => note.id));
1398 const decisions = (await review.searchMemories(workspace, null, { repoIds: [a.repoId], limit: 100 }).catch(() => [] as Memory[]))
1399 .filter((memory) => memory.kind === "decision" && !memory.pinned && !shown.has(memory.id) && runMemoryReadable(memory, reader, repoKey))
1400 .sort((x, y) => y.updatedAt.localeCompare(x.updatedAt))
1401 .slice(0, 3)
1402 .map((memory) => ({ id: memory.id, kind: memory.kind, text: memory.text, source: memorySource(memory) }));
1403 return composeRunContext({ projects: contexts, memories, decisions, budget: Math.min(Math.max(a.budget ?? 4000, 500), 12_000) });
1404 } catch (error) {
1405 console.error("no run context for", a.repoId, error);
1406 return { text: null, sources: [] };
1407 }
1408 }
1409}
1410
1411export default {
1412 async fetch(request: Request, env: Env): Promise<Response> {
1413 const match = new URL(request.url).pathname.match(/^\/rpc\/([a-z_]+)$/);
1414 if (request.method !== "POST" || !match) return new Response("Not found\n", { status: 404 });
1415 const service = new Context(env);
1416 const args = (await request.json().catch(() => ({}))) as any;
1417 switch (match[1]) {
1418 case "catalog":
1419 return Response.json(await service.catalog(args));
1420 case "entity":
1421 return Response.json(await service.entity(args));
1422 case "search":
1423 return Response.json(await service.search(args));
1424 case "scorecards":
1425 return Response.json(await service.scorecards(args));
1426 case "backfill":
1427 return Response.json(await service.backfill(args));
1428 case "status":
1429 return Response.json(await service.status(args));
1430 case "run_context":
1431 return Response.json(await service.runContext(args));
1432 default:
1433 return new Response("Unknown method\n", { status: 404 });
1434 }
1435 },
1436
1437 async queue(batch: MessageBatch<G1tEvent | Job>, env: Env): Promise<void> {
1438 const service = new Context(env);
1439 for (const message of batch.messages) {
1440 try {
1441 if (batch.queue === "g1t-context-jobs") await service.runJob(message.body as Job);
1442 else await service.onEvent(message.body as G1tEvent);
1443 message.ack();
1444 } catch (error) {
1445 console.error("context could not handle", (message.body as { type?: string }).type, error);
1446 message.retry();
1447 }
1448 }
1449 },
1450} satisfies ExportedHandler<Env, G1tEvent | Job>;