flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/context/src/index.ts

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