g1t/services/repos/src/forks.rs

402 lines16,865 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 };
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}