g1t/services/repos/src/forks.rs

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