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