| 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 { foliosCite, recordFolioChanges } from "./folios/staleness.ts"; |
| 36 | import type { FolioRoom } from "./folios/room.ts"; |
| 37 | import type { PageRoom } from "./room.ts"; |
| 38 | import { publishDocEvent } from "./events.ts"; |
| 39 | import { pageSlug } from "./slugs.ts"; |
| 40 | |
| 41 | export 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. */ |
| 54 | export 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. */ |
| 66 | export function systemReader(namespace: string): User { |
| 67 | return { id: "g1t", username: "g1t", kind: "system", verified: true, workspaces: [{ slug: namespace.toLowerCase(), role: "owner" }] }; |
| 68 | } |
| 69 | |
| 70 | const ZERO = /^0+$/; |
| 71 | |
| 72 | async 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. */ |
| 79 | async 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. */ |
| 89 | async 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. */ |
| 97 | async 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 | |
| 112 | type CitedRow = { page_id: string; path: string }; |
| 113 | type 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 | */ |
| 119 | export 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. */ |
| 157 | async 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. */ |
| 251 | async 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). */ |
| 275 | export 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 | } |