Skip to content
306 linesCodeBlameRaw
1/**
2 * Pages that cite code a change touched become possibly out of date
3 * (docs/WORKSPACE.md, "Agents and docs"). The events service sends this
4 * service `git.push` and `pull.merged` (`SUBSCRIBER_DOCS`,
5 * crates/contracts subscribers.rs), and the lifecycle events every
6 * subscriber hears.
7 *
8 * - **What changed** is asked as g1t itself (a system viewer in the
9 * repository's workspace): a merged pull request's files from the work
10 * service, or a push's from comparing it with where the branch was.
11 * Only pushes to the default branch count.
12 * - **Cheap when nothing cites the repository:** the citations table is
13 * asked first, and nothing else is.
14 * - **Once per page and commit:** a merge is told twice (the push and the
15 * pull request); both land on one row, which the pull request names.
16 * Owners hear of it once, and `doc.page.stale` is published once.
17 * - **Privacy:** what is stored is the change itself; what a reader sees
18 * of it (the page, its banner, the inbox) is filtered by whether they can
19 * read the repository.
20 */
21import {
22 currentMovedPath,
23 identityClient,
24 notifyClient,
25 repoMove,
26 reposClient,
27 staleMovedPaths,
28 workClient,
29 type G1tEvent,
30 type ServiceBinding,
31 type User,
32} from "@g1t/contracts";
33
34import { touchedPaths } from "./citations.ts";
35import { foliosCite, recordFolioChanges } from "./folios/staleness.ts";
36import type { FolioRoom } from "./folios/room.ts";
37import type { PageRoom } from "./room.ts";
38import { publishDocEvent } from "./events.ts";
39import { pageSlug } from "./slugs.ts";
40
41export type StaleEnv = {
42 DB: D1Database;
43 IDENTITY: ServiceBinding;
44 REPOS?: ServiceBinding;
45 WORK?: ServiceBinding;
46 NOTIFY?: ServiceBinding;
47 EVENTS?: ServiceBinding;
48 PAGES?: DurableObjectNamespace<PageRoom>;
49 /** Folios' rooms (Artifacts mode): folios cite code too (src/folios/staleness.ts). */
50 FOLIOS?: DurableObjectNamespace<FolioRoom>;
51};
52
53/** What a change touched, as this module records it. */
54export type Change = {
55 repo: string;
56 repo_id: string;
57 commit: string;
58 pull: { number: number; title: string | null } | null;
59 /** Every path the change touched. */
60 changed: string[];
61 /** Who made it (an account id), for the events it causes. */
62 actor: string | null;
63};
64
65/** g1t itself, reading in a repository's workspace: how staleness asks what changed. */
66export function systemReader(namespace: string): User {
67 return { id: "g1t", username: "g1t", kind: "system", verified: true, workspaces: [{ slug: namespace.toLowerCase(), role: "owner" }] };
68}
69
70const ZERO = /^0+$/;
71
72async function pathById(repos: ServiceBinding, id: string): Promise<{ namespace: string; name: string } | null> {
73 const response = await repos.fetch("https://service/rpc/path_by_id", { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ id }) });
74 if (!response.ok) return null;
75 return (await response.json()) as { namespace: string; name: string } | null;
76}
77
78/** Whether any live page cites the repository, or any workspace shows its docs. */
79async function interested(db: D1Database, repo: string, repoId: string): Promise<{ cited: boolean; spaces: boolean }> {
80 const [cited, spaces] = await db.batch<{ yes: number }>([
81 db.prepare("SELECT 1 AS yes FROM citations c JOIN pages p ON p.id = c.page_id WHERE c.repo = ? AND p.archived_at IS NULL LIMIT 1").bind(repo),
82 db.prepare("SELECT 1 AS yes FROM repo_spaces WHERE repo_id = ? LIMIT 1").bind(repoId),
83 ]);
84 const folios = cited?.results.length ? false : await foliosCite(db, repo).catch(() => false);
85 return { cited: !!cited?.results.length || folios, spaces: !!spaces?.results.length };
86}
87
88/** What a push to the default branch changed: the files between where it was and where it is. */
89async function pushChange(env: StaleEnv, path: { namespace: string; name: string }, repoId: string, before: string | undefined, after: string): Promise<string[] | null> {
90 if (!env.REPOS || !before || ZERO.test(before)) return null;
91 const compared = await reposClient(env.REPOS).compare(repoId, systemReader(path.namespace), before, after);
92 if (!compared.ok) return null;
93 return compared.value.files.map((f) => f.path);
94}
95
96/** What a merged pull request changed, and its title; null when it didn't merge into the default branch. */
97async function pullChange(env: StaleEnv, path: { namespace: string; name: string }, repoId: string, number: number): Promise<{ title: string; changed: string[] } | null> {
98 if (!env.WORK || !env.REPOS) return null;
99 const reader = systemReader(path.namespace);
100 const [found, repo] = await Promise.all([workClient(env.WORK).getPull(path, number, reader), reposClient(env.REPOS).getById(repoId, reader)]);
101 if (!found.ok || !repo.ok) return null;
102 const pull = found.value.pull;
103 if (pull.base && pull.base !== repo.value.defaultBranch) return null;
104 let changed = (pull.files ?? []).map((f) => f.path);
105 if (!changed.length && pull.mergeBase && pull.headCommit) {
106 const compared = await reposClient(env.REPOS).compare(repoId, reader, pull.mergeBase, pull.headCommit);
107 if (compared.ok) changed = compared.value.files.map((f) => f.path);
108 }
109 return { title: pull.title, changed };
110}
111
112type CitedRow = { page_id: string; path: string };
113type PageRow = { id: string; workspace_id: string; space_id: string; title: string; space_slug: string };
114
115/**
116 * Records a change against every page whose citations it touches.
117 * Returns the pages newly made stale by it.
118 */
119export async function record(env: StaleEnv, change: Change, now = new Date()): Promise<string[]> {
120 const db = env.DB;
121 const cited = (
122 await db
123 .prepare("SELECT c.page_id, c.path FROM citations c JOIN pages p ON p.id = c.page_id WHERE c.repo = ? AND p.archived_at IS NULL")
124 .bind(change.repo)
125 .all<CitedRow>()
126 ).results;
127 const byPage = new Map<string, CitedRow[]>();
128 for (const row of cited) byPage.set(row.page_id, [...(byPage.get(row.page_id) ?? []), row]);
129 const fresh: string[] = [];
130 const hits = new Map<string, string[]>();
131 for (const [pageId, rows] of byPage) {
132 const touched = touchedPaths(rows, change.changed);
133 if (!touched.length) continue;
134 hits.set(pageId, touched);
135 const at = now.toISOString();
136 const [inserted] = await db.batch<{ page_id: string }>([
137 db
138 .prepare(
139 "INSERT INTO page_changes (page_id, repo, repo_id, commit_sha, pull_number, pull_title, paths, detected_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT (page_id, repo, commit_sha) DO NOTHING RETURNING page_id",
140 )
141 .bind(pageId, change.repo, change.repo_id, change.commit, change.pull?.number ?? null, change.pull?.title ?? null, JSON.stringify(touched), at),
142 ...(change.pull
143 ? [
144 db
145 .prepare("UPDATE page_changes SET pull_number = ?, pull_title = ? WHERE page_id = ? AND repo = ? AND commit_sha = ?")
146 .bind(change.pull.number, change.pull.title, pageId, change.repo, change.commit),
147 ]
148 : []),
149 ]);
150 if (inserted?.results.length) fresh.push(pageId);
151 }
152 if (fresh.length) await tellOf(env, change, fresh, hits);
153 return fresh;
154}
155
156/** Owners hear of pages newly stale, rooms ask their readers to look again, and `doc.page.stale` goes out. */
157async function tellOf(env: StaleEnv, change: Change, pageIds: string[], hits: Map<string, string[]>): Promise<void> {
158 const db = env.DB;
159 const marks = pageIds.map(() => "?").join(",");
160 const [pages, owners] = await Promise.all([
161 db
162 .prepare(`SELECT p.id, p.workspace_id, p.space_id, p.title, s.slug AS space_slug FROM pages p JOIN spaces s ON s.id = p.space_id WHERE p.id IN (${marks})`)
163 .bind(...pageIds)
164 .all<PageRow>(),
165 db.prepare(`SELECT page_id, principal FROM page_owners WHERE page_id IN (${marks})`).bind(...pageIds).all<{ page_id: string; principal: string }>(),
166 ]);
167 const workspaceIds = [...new Set(pages.results.map((p) => p.workspace_id))];
168 const slugs = await identityClient(env.IDENTITY)
169 .usernames(workspaceIds)
170 .catch(() => ({}) as Record<string, string>);
171 // Which owners can read the repository: they are told which change it was.
172 const ownerIds = [...new Set(owners.results.filter((o) => o.principal.startsWith("user:")).map((o) => o.principal.slice(5)))];
173 const people = ownerIds.length ? await identityClient(env.IDENTITY).usersForAudience(ownerIds).catch(() => [] as User[]) : [];
174 const [namespace, name] = change.repo.split("/") as [string, string];
175 const canRead = new Set<string>();
176 if (env.REPOS) {
177 await Promise.all(
178 people.map(async (person) => {
179 const found = await reposClient(env.REPOS!)
180 .get({ namespace, name }, person)
181 .catch(() => null);
182 if (found?.ok) canRead.add(person.id);
183 }),
184 );
185 }
186 const what = change.pull ? `${change.repo}#${change.pull.number}` : `${change.repo}@${change.commit.slice(0, 7)}`;
187 const notify = env.NOTIFY ? notifyClient(env.NOTIFY) : null;
188 for (const page of pages.results) {
189 const slug = slugs[page.workspace_id];
190 if (!slug) continue;
191 const href = `/${slug}/-/docs/${page.space_slug}/${pageSlug(page.title, page.id)}`;
192 const paths = hits.get(page.id) ?? [];
193 const keys = owners.results.filter((o) => o.page_id === page.id).map((o) => o.principal);
194 const work: Promise<unknown>[] = [];
195 if (notify) {
196 for (const key of keys.filter((k) => k.startsWith("user:"))) {
197 const id = key.slice(5);
198 const known = canRead.has(id);
199 work.push(
200 notify
201 .notify(
202 { user_id: id },
203 {
204 id: `doc-stale:${page.id}:${change.commit}:${id}`,
205 kind: "inbox",
206 workspace: slug,
207 title: `${page.title || "Untitled"} may be out of date`,
208 body: known ? `${what} changed ${paths.slice(0, 3).join(", ")}${paths.length > 3 ? ` and ${paths.length - 3} more` : ""}` : "A change to code this page cites was merged.",
209 href,
210 actor: { kind: "system", id: "g1t", name: "g1t", avatar: null, avatar_seed: null },
211 created_at: new Date().toISOString(),
212 },
213 )
214 .catch(() => undefined),
215 );
216 }
217 }
218 if (env.PAGES) {
219 work.push(
220 env.PAGES.get(env.PAGES.idFromName(page.id))
221 .notice({ type: "page.staleness" })
222 .catch(() => undefined),
223 );
224 }
225 work.push(
226 publishDocEvent(
227 env.EVENTS,
228 "doc.page.stale",
229 {
230 workspace: slug,
231 workspaceId: page.workspace_id,
232 pageId: page.id,
233 spaceId: page.space_id,
234 title: page.title,
235 path: href,
236 repoId: change.repo_id,
237 repo: change.repo,
238 commit: change.commit,
239 pull: change.pull?.number ?? null,
240 paths,
241 owners: keys,
242 },
243 change.actor ? `user:${change.actor}` : null,
244 ),
245 );
246 await Promise.all(work);
247 }
248}
249
250/** A repository moved: rows kept under its old path follow it. */
251async function followMove(env: StaleEnv, event: G1tEvent): Promise<void> {
252 const move = repoMove(event);
253 if (!move || !env.REPOS) return;
254 const current = (await currentMovedPath(env.REPOS, move)).toLowerCase();
255 const stale = staleMovedPaths(move, current).map((p) => p.toLowerCase());
256 if (!stale.length) return;
257 const db = env.DB;
258 const statements: D1PreparedStatement[] = [];
259 for (const old of stale) {
260 statements.push(
261 db.prepare("UPDATE OR IGNORE citations SET repo = ? WHERE repo = ?").bind(current, old),
262 db.prepare("UPDATE OR IGNORE page_changes SET repo = ? WHERE repo = ?").bind(current, old),
263 db.prepare("UPDATE OR IGNORE page_projects SET repo = ? WHERE repo = ?").bind(current, old),
264 db.prepare("UPDATE OR IGNORE folio_citations SET repo = ? WHERE repo = ?").bind(current, old),
265 db.prepare("UPDATE OR IGNORE folio_changes SET repo = ? WHERE repo = ?").bind(current, old),
266 db.prepare("UPDATE OR IGNORE folio_projects SET repo = ? WHERE repo = ?").bind(current, old),
267 db.prepare("UPDATE OR IGNORE space_projects SET repo = ? WHERE repo = ?").bind(current, old),
268 db.prepare("UPDATE repo_spaces SET repo = ? WHERE repo_id = ?").bind(current, move.repoId),
269 );
270 }
271 await db.batch(statements);
272}
273
274/** One event, as this service acts on it. `reindex` reads a project's docs again (src/repo-spaces.ts). */
275export async function onEvent(env: StaleEnv, event: G1tEvent, reindex: (repoId: string, commit: string) => Promise<void>): Promise<void> {
276 if (event.type === "repo.renamed" || event.type === "repo.transferred") return followMove(env, event);
277 if (event.type === "repo.purged") {
278 await env.DB.prepare("DELETE FROM repo_spaces WHERE repo_id = ?").bind(event.data.repoId).run();
279 return;
280 }
281 if (event.type !== "git.push" && event.type !== "pull.merged") return;
282 if (!env.REPOS) return;
283 const repoId = event.repoId ?? (event.data as { repoId?: string }).repoId ?? null;
284 if (!repoId) return;
285 if (event.type === "git.push" && !event.data.defaultBranch) return;
286 const path = await pathById(env.REPOS, repoId);
287 if (!path) return;
288 const repo = `${path.namespace}/${path.name}`.toLowerCase();
289 const wants = await interested(env.DB, repo, repoId);
290 if (event.type === "git.push") {
291 if (wants.spaces) await reindex(repoId, event.data.after);
292 if (!wants.cited) return;
293 const changed = await pushChange(env, path, repoId, event.data.before, event.data.after);
294 if (!changed?.length) return;
295 const change: Change = { repo, repo_id: repoId, commit: event.data.after, pull: null, changed, actor: event.actor };
296 await record(env, change);
297 await recordFolioChanges(env, change);
298 return;
299 }
300 if (!wants.cited) return;
301 const pulled = await pullChange(env, path, repoId, event.data.number);
302 if (!pulled?.changed.length) return;
303 const change: Change = { repo, repo_id: repoId, commit: event.data.commit, pull: { number: event.data.number, title: pulled.title }, changed: pulled.changed, actor: event.actor };
304 await record(env, change);
305 await recordFolioChanges(env, change);
306}