g1t/services/context/src/index.ts

1,261 lines58,620 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

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