| 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 | |
| 19 | use g1t_contracts::repos::{Branch, Commit, GitAccess, Repo, TreeEntry}; |
| 20 | use g1t_contracts::time::rfc3339; |
| 21 | use g1t_kit::now_ms; |
| 22 | use worker::wasm_bindgen::JsValue; |
| 23 | use worker::{Env, Result}; |
| 24 | |
| 25 | use crate::registry::{Registry, RepoRow, store_key}; |
| 26 | use crate::store::{GitRepo, GitStore, Scope}; |
| 27 | use 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. |
| 31 | pub const DEFAULT_RETENTION_DAYS: u64 = 7; |
| 32 | /// How many working copies one sweep removes. |
| 33 | const RETIRES_PER_SWEEP: u32 = 25; |
| 34 | |
| 35 | pub 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. |
| 43 | pub 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. |
| 48 | pub 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)] |
| 54 | pub 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 | |
| 61 | pub 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. |
| 70 | pub 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. |
| 76 | pub struct Viewed<R> { |
| 77 | inner: R, |
| 78 | /// The copy's branch, and the commit it pointed to. |
| 79 | alias: Option<(String, String)>, |
| 80 | } |
| 81 | |
| 82 | impl<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. |
| 101 | fn aliases(git_ref: &str, branch: &str) -> bool { |
| 102 | git_ref == branch || git_ref == "HEAD" || git_ref.strip_prefix("refs/heads/") == Some(branch) |
| 103 | } |
| 104 | |
| 105 | impl<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 | |
| 149 | impl 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 | |
| 219 | impl<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 | // 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 | } |
| 286 | 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 | }; |
| 291 | // 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), |
| 298 | 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)] |
| 378 | mod 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 | } |