g1t/services/repos/src/forks.rs

402 lines16,865 bytesCodeBlame
1//! 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 };
287 let head = match self.store.open(&key).await {
288 Ok(git) => git.log(&fork.default_branch, 1).await?.into_iter().next().map(|commit| commit.hash),
289 // Already gone from the store.
290 Err(error) if error.to_string().contains("NOT_FOUND") => None,
291 Err(error) => return Err(error),
292 };
293 if let (Some(parent), Some(head)) = (&parent, &head) {
294 self.keep_head(fork, parent, head).await?;
295 }
296 let now = rfc3339(now_ms());
297 self.registry.set_retired(&fork.id, Some((&now, head.as_deref()))).await?;
298 if let Err(error) = self.store.delete(&key).await {
299 // Back as it was, for the next sweep.
300 self.registry.set_retired(&fork.id, None).await?;
301 return Err(error);
302 }
303 Ok(())
304 }
305
306 /// Points `refs/pull/<pull id>/head` in `parent` at `head`, sending
307 /// only the objects the parent lacks.
308 async fn keep_head(&self, fork: &Repo, parent: &Repo, head: &str) -> Result<()> {
309 let parent_git = self.store.open(&store_key(parent)).await?;
310 let reference = pull_ref(&fork.name);
311 let parent_read = parent_git.access(Scope::Read).await?;
312 let refs = refs::all(&parent_read).await?;
313 let existing = refs.iter().find(|(name, _)| *name == reference).map(|(_, hash)| hash.clone());
314 if existing.as_deref() == Some(head) {
315 return Ok(());
316 }
317 let has_it = !parent_git.log(head, 1).await?.is_empty();
318 let pack = if has_it {
319 land::EMPTY_PACK.to_vec()
320 } else {
321 let base = refs.iter().find(|(name, _)| *name == format!("refs/heads/{}", parent.default_branch)).map(|(_, hash)| hash.clone());
322 let fork_read = self.store.open(&store_key(fork)).await?.access(Scope::Read).await?;
323 land::fetch_pack(&fork_read, head, base.as_deref()).await?
324 };
325 let parent_write = parent_git.access(Scope::Write).await?;
326 let kept = land::push_ref(&parent_write, &reference, existing.as_deref(), head, Some(pack)).await;
327 self.refs_moved(&parent.id).await;
328 kept?.map_err(|reason| worker::Error::RustError(format!("{reference} not kept in {}: {reason}", parent.id)))
329 }
330
331 /// Makes a removed working copy again: a fork of its repository, with
332 /// its branch moved back to the head it had.
333 async fn revive(&self, fork: &Repo, retire_after: Option<String>) -> Result<()> {
334 let head = crate::registry::retired(&fork.id);
335 let parent = match &fork.fork_of {
336 Some(id) => self.registry.by_id(id).await?,
337 None => None,
338 };
339 let (Some(head), Some(parent)) = (head, parent) else {
340 return Err(worker::Error::RustError(format!(
341 "working copy {} cannot be made again: its repository or head is gone",
342 fork.id
343 )));
344 };
345 let key = store_key(fork);
346 let parent_git = self.store.open(&store_key(&parent)).await?;
347 parent_git.fork(&key).await?;
348 let git = self.store.open(&key).await?;
349 let current = git.log(&fork.default_branch, 1).await?.into_iter().next().map(|commit| commit.hash);
350 if current.as_deref() != Some(head.as_str()) {
351 let pack = land::fetch_pack(&parent_git.access(Scope::Read).await?, &head, current.as_deref()).await?;
352 let moved = land::push_ref(
353 &git.access(Scope::Write).await?,
354 &format!("refs/heads/{}", fork.default_branch),
355 current.as_deref(),
356 &head,
357 Some(pack),
358 )
359 .await;
360 self.refs_moved(&fork.id).await;
361 moved?.map_err(|reason| worker::Error::RustError(format!("working copy {} made again, but not at {head}: {reason}", fork.id)))?;
362 }
363 self.registry.revived(&fork.id, retire_after.as_deref()).await?;
364 self.refs_moved(&fork.id).await;
365 Ok(())
366 }
367}
368
369#[cfg(test)]
370mod tests {
371 use super::*;
372
373 #[test]
374 fn a_copy_is_removed_days_after_its_pull_request_settles() {
375 let settled = 1_791_936_000_000; // 2026-10-14T00:00:00Z
376 assert_eq!(retire_after(settled, 7), "2026-10-21T00:00:00.000Z");
377 assert_eq!(retire_after(settled, 0), "2026-10-14T00:00:00.000Z");
378 // The sweep compares these as text, which orders like time.
379 assert!(retire_after(settled, 7) > retire_after(settled, 6));
380 assert_eq!(pull_change("pull.merged"), Some(PullChange::Settled));
381 assert_eq!(pull_change("pull.closed"), Some(PullChange::Settled));
382 assert_eq!(pull_change("pull.reopened"), Some(PullChange::Reopened));
383 assert_eq!(pull_change("pull.updated"), None);
384 assert_eq!(pull_id_of(&serde_json::json!({ "pullId": "pul_1", "number": 3 })).as_deref(), Some("pul_1"));
385 assert_eq!(pull_id_of(&serde_json::json!({ "pull_id": "pul_2" })).as_deref(), Some("pul_2"));
386 assert_eq!(pull_id_of(&serde_json::json!({})), None);
387 assert_eq!(pull_ref("pul_1"), "refs/pull/pul_1/head");
388 }
389
390 #[test]
391 fn a_removed_copy_reads_its_branch_at_the_head_it_had() {
392 let viewed = Viewed::retired((), "main", &"a".repeat(40));
393 assert_eq!(viewed.resolve("main"), "a".repeat(40));
394 assert_eq!(viewed.resolve("refs/heads/main"), "a".repeat(40));
395 assert_eq!(viewed.resolve("HEAD"), "a".repeat(40));
396 // Anything else is read from the repository as it is.
397 assert_eq!(viewed.resolve("dev"), "dev");
398 assert_eq!(viewed.resolve(&"b".repeat(40)), "b".repeat(40));
399 let plain = Viewed::plain(());
400 assert_eq!(plain.resolve("main"), "main");
401 }
402}