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