g1t/services/repos/src/forks.rs

406 lines17,006 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.

Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily1//! Pull requests' working copies after their pull request is done.
2//!
3//! Every pull request made from an issue works in its own copy of the
4//! repository in the git store (`pulls--<pull id>`, `fork_for_pull`).
5//! Cloudflare does not say whether a fork shares objects with its source,
6//! and storage is billed and capped per account (1 TB), so copies are not
7//! kept forever: [`retention_days`] after the pull request merges or
8//! closes (`pull.merged`, `pull.closed`), the hourly sweep removes the
9//! copy's git data. Reopened within the window (`pull.reopened`), it stays.
10//!
11//! Nothing is lost. Before the copy goes, its head is kept in the
12//! repository it came from as `refs/pull/<pull id>/head` (only the
13//! objects the repository lacks travel), and the row stays, with that
14//! commit. Reads of the copy (the pull request's changes, its files) are
15//! answered from the repository ([`Viewed`]). Anything that writes to it
16//! (a push, git access, catching up, landing) makes it again first, from
17//! that ref ([`Repos::live`]).
18
19use g1t_contracts::repos::{Branch, Commit, GitAccess, Repo, TreeEntry};
20use g1t_contracts::time::rfc3339;
21use g1t_kit::now_ms;
22use worker::wasm_bindgen::JsValue;
23use worker::{Env, Result};
24
25use crate::registry::{Registry, RepoRow, store_key};
26use crate::store::{GitRepo, GitStore, Scope};
27use crate::{PULLS_NAMESPACE, Repos, land, refs};
28
29/// Days a working copy is kept after its pull request merges or closes,
30/// unless `FORK_RETENTION_DAYS` says otherwise.
31pub const DEFAULT_RETENTION_DAYS: u64 = 7;
32/// How many working copies one sweep removes.
33const RETIRES_PER_SWEEP: u32 = 25;
34
35pub fn retention_days(env: &Env) -> u64 {
36 env.var("FORK_RETENTION_DAYS")
37 .ok()
38 .and_then(|value| value.to_string().parse().ok())
39 .unwrap_or(DEFAULT_RETENTION_DAYS)
40}
41
42/// When a working copy whose pull request settled at `now` is removed.
43pub fn retire_after(now: u64, days: u64) -> String {
44 rfc3339(now + days * 24 * 3600 * 1000)
45}
46
47/// Where a pull request's head is kept in its repository.
48pub fn pull_ref(pull_id: &str) -> String {
49 format!("refs/pull/{pull_id}/head")
50}
51
52/// What a pull request event means for its working copy.
53#[derive(Debug, PartialEq, Eq)]
54pub enum PullChange {
55 /// Merged or closed: remove it after the window.
56 Settled,
57 /// Open again: keep it, or make it again.
58 Reopened,
59}
60
61pub fn pull_change(kind: &str) -> Option<PullChange> {
62 match kind {
63 "pull.merged" | "pull.closed" => Some(PullChange::Settled),
64 "pull.reopened" => Some(PullChange::Reopened),
65 _ => None,
66 }
67}
68
69/// The pull request an event is about.
70pub fn pull_id_of(data: &serde_json::Value) -> Option<String> {
71 data.get("pullId").or_else(|| data.get("pull_id")).and_then(|id| id.as_str()).map(str::to_owned)
72}
73
74/// A repository as reads see it: a working copy that was removed reads
75/// as the repository it came from, with its branch at the head it had.
76pub struct Viewed<R> {
77 inner: R,
78 /// The copy's branch, and the commit it pointed to.
79 alias: Option<(String, String)>,
80}
81
82impl<R> Viewed<R> {
83 pub fn plain(inner: R) -> Self {
84 Viewed { inner, alias: None }
85 }
86
87 pub fn retired(inner: R, branch: &str, head: &str) -> Self {
88 Viewed { inner, alias: Some((branch.to_owned(), head.to_owned())) }
89 }
90
91 /// `git_ref`, or the head a removed copy's branch pointed to.
92 fn resolve<'a>(&'a self, git_ref: &'a str) -> &'a str {
93 match &self.alias {
94 Some((branch, head)) if aliases(git_ref, branch) => head,
95 _ => git_ref,
96 }
97 }
98}
99
100/// Whether `git_ref` names `branch`, or HEAD.
101fn aliases(git_ref: &str, branch: &str) -> bool {
102 git_ref == branch || git_ref == "HEAD" || git_ref.strip_prefix("refs/heads/") == Some(branch)
103}
104
105impl<R: GitRepo> GitRepo for Viewed<R> {
106 async fn access(&self, scope: Scope) -> Result<GitAccess> {
107 self.inner.access(scope).await
108 }
109
110 async fn branches(&self) -> Result<Vec<Branch>> {
111 match &self.alias {
112 Some((branch, head)) => Ok(vec![Branch { name: branch.clone(), hash: head.clone() }]),
113 None => self.inner.branches().await,
114 }
115 }
116
117 async fn log(&self, git_ref: &str, limit: u32) -> Result<Vec<Commit>> {
118 self.inner.log(self.resolve(git_ref), limit).await
119 }
120
121 async fn parents(&self, commit_hash: &str) -> Result<Option<Vec<String>>> {
122 self.inner.parents(commit_hash).await
123 }
124
125 async fn read_tree(&self, tree_hash: &str) -> Result<Option<Vec<TreeEntry>>> {
126 self.inner.read_tree(tree_hash).await
127 }
128
129 async fn read_blob(&self, blob_hash: &str) -> Result<Option<Vec<u8>>> {
130 self.inner.read_blob(blob_hash).await
131 }
132
133 async fn read_file(&self, git_ref: &str, path: &str) -> Result<Option<Vec<u8>>> {
134 self.inner.read_file(self.resolve(git_ref), path).await
135 }
136
137 async fn fork(&self, target_key: &str) -> Result<()> {
138 if self.alias.is_some() {
139 return Err(worker::Error::RustError("a removed working copy is not forked".into()));
140 }
141 self.inner.fork(target_key).await
142 }
143
144 fn at_refs_version(&mut self, version: Option<u64>) {
145 self.inner.at_refs_version(version);
146 }
147}
148
149impl Registry {
150 /// The working copy of a pull request, if it has one.
151 pub async fn pull_fork(&self, pull_id: &str) -> Result<Option<Repo>> {
152 Ok(self
153 .db
154 .prepare("SELECT * FROM repos WHERE namespace = ? AND name = ? AND fork_of IS NOT NULL AND deleted_at IS NULL")
155 .bind(&[PULLS_NAMESPACE.into(), pull_id.to_lowercase().into()])?
156 .first::<RepoRow>(None)
157 .await?
158 .map(Repo::from))
159 }
160
161 /// Sets, or with `None` clears, when a working copy is removed.
162 pub async fn set_retire_after(&self, id: &str, after: Option<&str>) -> Result<()> {
163 self.db
164 .prepare("UPDATE repos SET retire_after = ? WHERE id = ? AND retired_at IS NULL")
165 .bind(&[after.map_or(JsValue::NULL, JsValue::from), id.into()])?
166 .run()
167 .await?;
168 Ok(())
169 }
170
171 /// Working copies whose time has come, oldest first.
172 pub async fn retiring(&self, now: &str, limit: u32) -> Result<Vec<Repo>> {
173 Ok(self
174 .db
175 .prepare(
176 "SELECT * FROM repos
177 WHERE retire_after IS NOT NULL AND retire_after <= ? AND retired_at IS NULL
178 AND fork_of IS NOT NULL AND deleted_at IS NULL
179 ORDER BY retire_after LIMIT ?",
180 )
181 .bind(&[now.into(), limit.into()])?
182 .all()
183 .await?
184 .results::<RepoRow>()?
185 .into_iter()
186 .map(Repo::from)
187 .collect())
188 }
189
190 /// Records that a working copy's git data is gone, and the head it had;
191 /// `None` puts it back as it was, for a removal that failed.
192 pub async fn set_retired(&self, id: &str, retired: Option<(&str, Option<&str>)>) -> Result<()> {
193 let (at, head) = match retired {
194 Some((at, head)) => (JsValue::from(at), head.map_or(JsValue::NULL, JsValue::from)),
195 None => (JsValue::NULL, JsValue::NULL),
196 };
197 self.db
198 .prepare("UPDATE repos SET retired_at = ?, retired_head = ? WHERE id = ?")
199 .bind(&[at, head, id.into()])?
200 .run()
201 .await?;
202 crate::registry::note_retired(id, retired.and_then(|(_, head)| head));
203 Ok(())
204 }
205
206 /// Records that a working copy was made again; `retire_after`, when to
207 /// remove it next.
208 pub async fn revived(&self, id: &str, retire_after: Option<&str>) -> Result<()> {
209 self.db
210 .prepare("UPDATE repos SET retired_at = NULL, retired_head = NULL, retire_after = ? WHERE id = ?")
211 .bind(&[retire_after.map_or(JsValue::NULL, JsValue::from), id.into()])?
212 .run()
213 .await?;
214 crate::registry::note_retired(id, None);
215 Ok(())
216 }
217}
218
219impl<S: GitStore> Repos<S> {
220 /// The repository's store as reads see it: a removed working copy
221 /// reads from the repository it came from. Reads by branch are kept
222 /// under its refs version (store.rs).
223 pub(crate) async fn read_git(&self, repo: &Repo) -> Result<Viewed<S::Repo>> {
224 let now = now_ms();
225 if let (Some(head), Some(parent_id)) = (crate::registry::retired(&repo.id), &repo.fork_of)
226 && let Some(parent) = self.registry.by_id(parent_id).await?
227 {
228 let mut git = self.store.open(&store_key(&parent)).await?;
229 git.at_refs_version(crate::refs_cache::usable(crate::registry::refs_state(&parent.id), now));
230 return Ok(Viewed::retired(git, &repo.default_branch, &head));
231 }
232 let mut git = self.store.open(&store_key(repo)).await?;
233 git.at_refs_version(crate::refs_cache::usable(crate::registry::refs_state(&repo.id), now));
234 Ok(Viewed::plain(git))
235 }
236
237 /// Makes a removed working copy again before anything writes to it.
238 /// Its pull request is still settled, so it is removed again after the
239 /// window unless the pull request reopens.
240 pub(crate) async fn live(&self, repo: &Repo) -> Result<()> {
241 if crate::registry::retired(&repo.id).is_none() {
242 return Ok(());
243 }
244 self.revive(repo, Some(retire_after(now_ms(), self.fork_days))).await
245 }
246
247 /// `pull.merged` or `pull.closed`: the working copy goes in `days`.
248 pub(crate) async fn pull_settled(&self, pull_id: &str) -> Result<()> {
249 if let Some(fork) = self.registry.pull_fork(pull_id).await? {
250 self.registry.set_retire_after(&fork.id, Some(&retire_after(now_ms(), self.fork_days))).await?;
251 }
252 Ok(())
253 }
254
255 /// `pull.reopened`: the working copy stays, made again if it had gone.
256 pub(crate) async fn pull_reopened(&self, pull_id: &str) -> Result<()> {
257 let Some(fork) = self.registry.pull_fork(pull_id).await? else {
258 return Ok(());
259 };
260 if crate::registry::retired(&fork.id).is_some() {
261 return self.revive(&fork, None).await;
262 }
263 self.registry.set_retire_after(&fork.id, None).await
264 }
265
266 /// The sweep: removes the working copies whose time has come.
267 pub(crate) async fn retire_due(&self) -> Result<u32> {
268 let due = self.registry.retiring(&rfc3339(now_ms()), RETIRES_PER_SWEEP).await?;
269 let mut retired = 0;
270 for fork in due {
271 match self.retire(&fork).await {
272 Ok(()) => retired += 1,
273 Err(error) => worker::console_error!("working copy {} not removed: {error}", fork.id),
274 }
275 }
276 Ok(retired)
277 }
278
279 /// Keeps a working copy's head in its repository, then removes its git
280 /// data. A failure before the removal leaves everything as it was.
281 async fn retire(&self, fork: &Repo) -> Result<()> {
282 let key = store_key(fork);
283 let parent = match &fork.fork_of {
284 Some(id) => self.registry.by_id(id).await?,
285 None => None,
286 };
Fewer Artifacts reads: the store is asked for a handle only when needed, objects are kept in the isolate, and issue events read no workflows nobody listens for287 // Already gone from the store: it says so when first asked.
288 let read = match self.store.open(&key).await {
289 Ok(git) => git.log(&fork.default_branch, 1).await,
290 Err(error) => Err(error),
291 };
292 let head = match read {
293 Ok(commits) => commits.into_iter().next().map(|commit| commit.hash),
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily294 Err(error) if error.to_string().contains("NOT_FOUND") => None,
295 Err(error) => return Err(error),
296 };
297 if let (Some(parent), Some(head)) = (&parent, &head) {
298 self.keep_head(fork, parent, head).await?;
299 }
300 let now = rfc3339(now_ms());
301 self.registry.set_retired(&fork.id, Some((&now, head.as_deref()))).await?;
302 if let Err(error) = self.store.delete(&key).await {
303 // Back as it was, for the next sweep.
304 self.registry.set_retired(&fork.id, None).await?;
305 return Err(error);
306 }
307 Ok(())
308 }
309
310 /// Points `refs/pull/<pull id>/head` in `parent` at `head`, sending
311 /// only the objects the parent lacks.
312 async fn keep_head(&self, fork: &Repo, parent: &Repo, head: &str) -> Result<()> {
313 let parent_git = self.store.open(&store_key(parent)).await?;
314 let reference = pull_ref(&fork.name);
315 let parent_read = parent_git.access(Scope::Read).await?;
316 let refs = refs::all(&parent_read).await?;
317 let existing = refs.iter().find(|(name, _)| *name == reference).map(|(_, hash)| hash.clone());
318 if existing.as_deref() == Some(head) {
319 return Ok(());
320 }
321 let has_it = !parent_git.log(head, 1).await?.is_empty();
322 let pack = if has_it {
323 land::EMPTY_PACK.to_vec()
324 } else {
325 let base = refs.iter().find(|(name, _)| *name == format!("refs/heads/{}", parent.default_branch)).map(|(_, hash)| hash.clone());
326 let fork_read = self.store.open(&store_key(fork)).await?.access(Scope::Read).await?;
327 land::fetch_pack(&fork_read, head, base.as_deref()).await?
328 };
329 let parent_write = parent_git.access(Scope::Write).await?;
330 let kept = land::push_ref(&parent_write, &reference, existing.as_deref(), head, Some(pack)).await;
331 self.refs_moved(&parent.id).await;
332 kept?.map_err(|reason| worker::Error::RustError(format!("{reference} not kept in {}: {reason}", parent.id)))
333 }
334
335 /// Makes a removed working copy again: a fork of its repository, with
336 /// its branch moved back to the head it had.
337 async fn revive(&self, fork: &Repo, retire_after: Option<String>) -> Result<()> {
338 let head = crate::registry::retired(&fork.id);
339 let parent = match &fork.fork_of {
340 Some(id) => self.registry.by_id(id).await?,
341 None => None,
342 };
343 let (Some(head), Some(parent)) = (head, parent) else {
344 return Err(worker::Error::RustError(format!(
345 "working copy {} cannot be made again: its repository or head is gone",
346 fork.id
347 )));
348 };
349 let key = store_key(fork);
350 let parent_git = self.store.open(&store_key(&parent)).await?;
351 parent_git.fork(&key).await?;
352 let git = self.store.open(&key).await?;
353 let current = git.log(&fork.default_branch, 1).await?.into_iter().next().map(|commit| commit.hash);
354 if current.as_deref() != Some(head.as_str()) {
355 let pack = land::fetch_pack(&parent_git.access(Scope::Read).await?, &head, current.as_deref()).await?;
356 let moved = land::push_ref(
357 &git.access(Scope::Write).await?,
358 &format!("refs/heads/{}", fork.default_branch),
359 current.as_deref(),
360 &head,
361 Some(pack),
362 )
363 .await;
364 self.refs_moved(&fork.id).await;
365 moved?.map_err(|reason| worker::Error::RustError(format!("working copy {} made again, but not at {head}: {reason}", fork.id)))?;
366 }
367 self.registry.revived(&fork.id, retire_after.as_deref()).await?;
368 self.refs_moved(&fork.id).await;
369 Ok(())
370 }
371}
372
373#[cfg(test)]
374mod tests {
375 use super::*;
376
377 #[test]
378 fn a_copy_is_removed_days_after_its_pull_request_settles() {
379 let settled = 1_791_936_000_000; // 2026-10-14T00:00:00Z
380 assert_eq!(retire_after(settled, 7), "2026-10-21T00:00:00.000Z");
381 assert_eq!(retire_after(settled, 0), "2026-10-14T00:00:00.000Z");
382 // The sweep compares these as text, which orders like time.
383 assert!(retire_after(settled, 7) > retire_after(settled, 6));
384 assert_eq!(pull_change("pull.merged"), Some(PullChange::Settled));
385 assert_eq!(pull_change("pull.closed"), Some(PullChange::Settled));
386 assert_eq!(pull_change("pull.reopened"), Some(PullChange::Reopened));
387 assert_eq!(pull_change("pull.updated"), None);
388 assert_eq!(pull_id_of(&serde_json::json!({ "pullId": "pul_1", "number": 3 })).as_deref(), Some("pul_1"));
389 assert_eq!(pull_id_of(&serde_json::json!({ "pull_id": "pul_2" })).as_deref(), Some("pul_2"));
390 assert_eq!(pull_id_of(&serde_json::json!({})), None);
391 assert_eq!(pull_ref("pul_1"), "refs/pull/pul_1/head");
392 }
393
394 #[test]
395 fn a_removed_copy_reads_its_branch_at_the_head_it_had() {
396 let viewed = Viewed::retired((), "main", &"a".repeat(40));
397 assert_eq!(viewed.resolve("main"), "a".repeat(40));
398 assert_eq!(viewed.resolve("refs/heads/main"), "a".repeat(40));
399 assert_eq!(viewed.resolve("HEAD"), "a".repeat(40));
400 // Anything else is read from the repository as it is.
401 assert_eq!(viewed.resolve("dev"), "dev");
402 assert_eq!(viewed.resolve(&"b".repeat(40)), "b".repeat(40));
403 let plain = Viewed::plain(());
404 assert_eq!(plain.resolve("main"), "main");
405 }
406}