| 1 | + | //! Nightly backups of every repository, outside the git store |
| 2 | + | //! (docs/ARTIFACTS.md, R11; the flow is in `g1t_contracts::backups`). |
| 3 | + | //! |
| 4 | + | //! Each repository whose refs moved since its last backup gets a |
| 5 | + | //! `git bundle`: a full one first, then incremental ones whose |
| 6 | + | //! prerequisites are the commits the one before ended at, and a full one |
| 7 | + | //! again after [`Settings::full_every`] incremental ones, so a restore |
| 8 | + | //! never reads a long chain. Bundles and a manifest that lists the chain |
| 9 | + | //! are kept in object storage through the `BlobStore` port: the BACKUPS R2 |
| 10 | + | //! bucket hosted, any S3-compatible store (MinIO in the compose file) |
| 11 | + | //! self-hosted, as BACKUP_STORE says. |
| 12 | + | //! |
| 13 | + | //! ```text |
| 14 | + | //! backups/<repo id>/manifest.json |
| 15 | + | //! backups/<repo id>/<20261006T025300Z>-full.bundle |
| 16 | + | //! backups/<repo id>/<20261007T025300Z>-incr.bundle |
| 17 | + | //! ``` |
| 18 | + | //! |
| 19 | + | //! Restoring is fetching each bundle of `chain` in order into an empty |
| 20 | + | //! repository, then setting every ref to what the last entry says |
| 21 | + | //! (scripts/ops/backup-restore-drill.mjs does it and compares). |
| 22 | + | //! |
| 23 | + | //! `repo_backups` keeps, per repository, the last backup (the refs version |
| 24 | + | //! it was cut at, its tips, when) and the job in hand, if any: |
| 25 | + | //! `idle` → `queued` (the nightly cron) → `running` (claimed by the |
| 26 | + | //! runner's sweep) → `idle` again, done or failed. A job that has been |
| 27 | + | //! running longer than [`LEASE_MS`] is queued again; one that failed |
| 28 | + | //! [`MAX_ATTEMPTS`] times waits for the next night. |
| 29 | + | |
| 30 | + | use std::collections::{BTreeMap, BTreeSet}; |
| 31 | + | |
| 32 | + | use g1t_blobstore::{BlobStore, Config, Part, Store}; |
| 33 | + | use g1t_contracts::backups::{ |
| 34 | + | BackupClaim, BackupComplete, BackupFail, BackupJobArgs, BackupKind, BackupPart, BackupSpec, ClaimBackupsArgs, PART_BYTES, |
| 35 | + | }; |
| 36 | + | use g1t_contracts::repos::RepoPath; |
| 37 | + | use g1t_contracts::time::rfc3339; |
| 38 | + | use g1t_contracts::{FailureCode, Outcome, new_id}; |
| 39 | + | use serde::{Deserialize, Serialize}; |
| 40 | + | use worker::wasm_bindgen::JsValue; |
| 41 | + | use worker::{D1Database, Env, Result}; |
| 42 | + | |
| 43 | + | use crate::PULLS_NAMESPACE; |
| 44 | + | use crate::meters; |
| 45 | + | use crate::registry::{Registry, store_key}; |
| 46 | + | use crate::store::{GitStore, Scope}; |
| 47 | + | |
| 48 | + | /// Where backups are kept: the BACKUPS bucket, or, when BACKUP_STORE is |
| 49 | + | /// `s3`, the bucket BACKUP_S3_BUCKET names on the installation's S3 store. |
| 50 | + | pub const STORAGE: Config = Config { |
| 51 | + | kind: "BACKUP_STORE", |
| 52 | + | binding: "BACKUPS", |
| 53 | + | r2_signer: None, |
| 54 | + | s3_bucket: "BACKUP_S3_BUCKET", |
| 55 | + | s3_public_endpoint: None, |
| 56 | + | }; |
| 57 | + | |
| 58 | + | /// The meter a backup's clone is counted under: an operation for g1t's own |
| 59 | + | /// bill, never for the workspace's (migrations/0013). |
| 60 | + | pub const FETCH_METER: &str = "internal.git.backup_fetch"; |
| 61 | + | |
| 62 | + | /// How long a claimed job may run before it is given to another sandbox. |
| 63 | + | pub const LEASE_MS: u64 = 3 * 60 * 60 * 1000; |
| 64 | + | /// How many times a night a backup is tried. |
| 65 | + | pub const MAX_ATTEMPTS: u32 = 3; |
| 66 | + | /// How many backups of deleted repositories one night removes. |
| 67 | + | const PRUNES_PER_NIGHT: u32 = 50; |
| 68 | + | |
| 69 | + | /// How backups are paced, from the service's variables. |
| 70 | + | #[derive(Clone, Copy, Debug, PartialEq, Eq)] |
| 71 | + | pub struct Settings { |
| 72 | + | /// BACKUPS_PER_NIGHT: how many repositories one night queues. |
| 73 | + | pub per_night: u32, |
| 74 | + | /// BACKUP_FULL_EVERY: incremental bundles before the next full one. |
| 75 | + | pub full_every: u32, |
| 76 | + | } |
| 77 | + | |
| 78 | + | impl Default for Settings { |
| 79 | + | fn default() -> Self { |
| 80 | + | Settings { per_night: 200, full_every: 30 } |
| 81 | + | } |
| 82 | + | } |
| 83 | + | |
| 84 | + | impl Settings { |
| 85 | + | pub fn from_env(env: &Env) -> Settings { |
| 86 | + | let number = |name: &str| env.var(name).ok().and_then(|value| value.to_string().parse::<u32>().ok()); |
| 87 | + | let defaults = Settings::default(); |
| 88 | + | Settings { |
| 89 | + | per_night: number("BACKUPS_PER_NIGHT").unwrap_or(defaults.per_night), |
| 90 | + | full_every: number("BACKUP_FULL_EVERY").unwrap_or(defaults.full_every).max(1), |
| 91 | + | } |
| 92 | + | } |
| 93 | + | } |
| 94 | + | |
| 95 | + | /// The storage backups go to, or None when this installation has none |
| 96 | + | /// (no BACKUPS binding and BACKUP_STORE is not `s3`): backups are then off. |
| 97 | + | pub fn storage(env: &Env) -> Option<Store> { |
| 98 | + | match Store::from_env(env, &STORAGE) { |
| 99 | + | Ok(store) => Some(store), |
| 100 | + | Err(error) => { |
| 101 | + | worker::console_log!("repos: backups are off: {error}"); |
| 102 | + | None |
| 103 | + | } |
| 104 | + | } |
| 105 | + | } |
| 106 | + | |
| 107 | + | // --------------------------------------------------------------------- |
| 108 | + | // Which repositories are due |
| 109 | + | // --------------------------------------------------------------------- |
| 110 | + | |
| 111 | + | /// A repository and its last backup, as the nightly query reads them. |
| 112 | + | #[derive(Clone, Debug, Default, PartialEq, Eq, Deserialize)] |
| 113 | + | pub struct Candidate { |
| 114 | + | pub repo_id: String, |
| 115 | + | pub namespace: String, |
| 116 | + | pub created_at: String, |
| 117 | + | pub refs_version: u64, |
| 118 | + | /// Until when a credential that can push was out of g1t's hands |
| 119 | + | /// (`git_access`); a push with it does not move `refs_version`. |
| 120 | + | pub refs_open_until: u64, |
| 121 | + | pub deleted: bool, |
| 122 | + | pub retired: bool, |
| 123 | + | /// None: never backed up, and no row. |
| 124 | + | pub status: Option<String>, |
| 125 | + | pub backed_version: Option<u64>, |
| 126 | + | pub backed_up_ms: u64, |
| 127 | + | } |
| 128 | + | |
| 129 | + | /// Whether a repository needs a backup tonight: it is live, it is not a |
| 130 | + | /// pull request's working copy (whose work lands in its repository, and |
| 131 | + | /// whose head is kept there once it goes), nothing is already queued or |
| 132 | + | /// running for it, and its refs moved since the last backup: its |
| 133 | + | /// `refs_version` went past the one backed up, or a credential that could |
| 134 | + | /// push was handed out after the last backup started. |
| 135 | + | pub fn is_due(c: &Candidate) -> bool { |
| 136 | + | if c.deleted || c.retired || c.namespace == PULLS_NAMESPACE { |
| 137 | + | return false; |
| 138 | + | } |
| 139 | + | match c.status.as_deref() { |
| 140 | + | None => true, |
| 141 | + | Some("idle") => match c.backed_version { |
| 142 | + | None => true, |
| 143 | + | Some(version) => version < c.refs_version || c.refs_open_until > c.backed_up_ms, |
| 144 | + | }, |
| 145 | + | _ => false, |
| 146 | + | } |
| 147 | + | } |
| 148 | + | |
| 149 | + | /// The repositories to queue: those due, the longest since their last |
| 150 | + | /// backup first (never backed up first of all, oldest repository first), |
| 151 | + | /// `limit` at most. |
| 152 | + | pub fn pick_due(candidates: &[Candidate], limit: usize) -> Vec<String> { |
| 153 | + | let mut due: Vec<&Candidate> = candidates.iter().filter(|c| is_due(c)).collect(); |
| 154 | + | due.sort_by(|a, b| { |
| 155 | + | a.backed_up_ms |
| 156 | + | .cmp(&b.backed_up_ms) |
| 157 | + | .then_with(|| a.created_at.cmp(&b.created_at)) |
| 158 | + | .then_with(|| a.repo_id.cmp(&b.repo_id)) |
| 159 | + | }); |
| 160 | + | due.into_iter().take(limit).map(|c| c.repo_id.clone()).collect() |
| 161 | + | } |
| 162 | + | |
| 163 | + | /// The same choice in SQL, so a night reads only what it queues. `?1`: |
| 164 | + | /// the working copies' namespace, `?2`: how many. |
| 165 | + | const DUE_SQL: &str = " |
| 166 | + | SELECT r.id AS repo_id, r.namespace, r.created_at, |
| 167 | + | coalesce(r.refs_version, 0) AS refs_version, |
| 168 | + | coalesce(r.refs_open_until, 0) AS refs_open_until, |
| 169 | + | (r.deleted_at IS NOT NULL) AS deleted, |
| 170 | + | (r.retired_at IS NOT NULL) AS retired, |
| 171 | + | b.status, b.refs_version AS backed_version, |
| 172 | + | coalesce(b.backed_up_ms, 0) AS backed_up_ms |
| 173 | + | FROM repos r LEFT JOIN repo_backups b ON b.repo_id = r.id |
| 174 | + | WHERE r.deleted_at IS NULL AND r.retired_at IS NULL AND r.namespace != ?1 |
| 175 | + | AND (b.repo_id IS NULL |
| 176 | + | OR (b.status = 'idle' |
| 177 | + | AND (b.refs_version IS NULL |
| 178 | + | OR b.refs_version < coalesce(r.refs_version, 0) |
| 179 | + | OR coalesce(r.refs_open_until, 0) > coalesce(b.backed_up_ms, 0)))) |
| 180 | + | ORDER BY coalesce(b.backed_up_ms, 0), r.created_at, r.id |
| 181 | + | LIMIT ?2"; |
| 182 | + | |
| 183 | + | // --------------------------------------------------------------------- |
| 184 | + | // The chain and its manifest |
| 185 | + | // --------------------------------------------------------------------- |
| 186 | + | |
| 187 | + | /// One backup in a chain: a bundle, or, when nothing new was there to |
| 188 | + | /// bundle (a branch deleted, a ref pointed at a commit already kept), only |
| 189 | + | /// the refs it ended with. |
| 190 | + | #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] |
| 191 | + | pub struct Entry { |
| 192 | + | /// When it was cut, `20261006T025300Z`: also its name in the chain. |
| 193 | + | pub id: String, |
| 194 | + | pub kind: BackupKind, |
| 195 | + | /// The bundle's key in storage; None when only the refs moved. |
| 196 | + | pub key: Option<String>, |
| 197 | + | pub created_at: String, |
| 198 | + | /// The repository's `refs_version` when the clone began. |
| 199 | + | pub refs_version: u64, |
| 200 | + | /// Every ref, by name, once this backup is applied. |
| 201 | + | pub refs: BTreeMap<String, String>, |
| 202 | + | /// The commits the bundle leaves out: the previous entry's tips. |
| 203 | + | pub prerequisites: Vec<String>, |
| 204 | + | pub size: u64, |
| 205 | + | pub sha256: Option<String>, |
| 206 | + | } |
| 207 | + | |
| 208 | + | /// What restoring a repository reads first: its chain, oldest first, and |
| 209 | + | /// the chain before it, kept until the next full backup replaces it. |
| 210 | + | #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] |
| 211 | + | pub struct Manifest { |
| 212 | + | pub version: u32, |
| 213 | + | pub repo_id: String, |
| 214 | + | /// The repository's name in the git store, and its path, when last cut. |
| 215 | + | pub store_key: String, |
| 216 | + | pub path: Option<RepoPath>, |
| 217 | + | pub updated_at: String, |
| 218 | + | pub chain: Vec<Entry>, |
| 219 | + | #[serde(default)] |
| 220 | + | pub previous: Vec<Entry>, |
| 221 | + | } |
| 222 | + | |
| 223 | + | pub const MANIFEST_VERSION: u32 = 1; |
| 224 | + | |
| 225 | + | impl Manifest { |
| 226 | + | pub fn new(repo_id: &str, store_key: &str) -> Manifest { |
| 227 | + | Manifest { |
| 228 | + | version: MANIFEST_VERSION, |
| 229 | + | repo_id: repo_id.to_owned(), |
| 230 | + | store_key: store_key.to_owned(), |
| 231 | + | path: None, |
| 232 | + | updated_at: String::new(), |
| 233 | + | chain: Vec::new(), |
| 234 | + | previous: Vec::new(), |
| 235 | + | } |
| 236 | + | } |
| 237 | + | |
| 238 | + | pub fn last(&self) -> Option<&Entry> { |
| 239 | + | self.chain.last() |
| 240 | + | } |
| 241 | + | |
| 242 | + | /// Adds `entry` to the chain. A full one starts a new chain, and the |
| 243 | + | /// current one becomes `previous`. Returns the bundles no longer kept. |
| 244 | + | pub fn add(&mut self, entry: Entry) -> Vec<String> { |
| 245 | + | if entry.kind == BackupKind::Full { |
| 246 | + | let dropped = std::mem::take(&mut self.previous); |
| 247 | + | self.previous = std::mem::replace(&mut self.chain, vec![entry]); |
| 248 | + | return dropped.into_iter().filter_map(|entry| entry.key).collect(); |
| 249 | + | } |
| 250 | + | self.chain.push(entry); |
| 251 | + | Vec::new() |
| 252 | + | } |
| 253 | + | |
| 254 | + | /// Every bundle it lists. |
| 255 | + | pub fn keys(&self) -> Vec<String> { |
| 256 | + | self.chain.iter().chain(&self.previous).filter_map(|entry| entry.key.clone()).collect() |
| 257 | + | } |
| 258 | + | |
| 259 | + | pub fn to_bytes(&self) -> Vec<u8> { |
| 260 | + | serde_json::to_vec_pretty(self).unwrap_or_default() |
| 261 | + | } |
| 262 | + | |
| 263 | + | pub fn from_bytes(bytes: &[u8]) -> Option<Manifest> { |
| 264 | + | serde_json::from_slice::<Manifest>(bytes).ok().filter(|manifest| manifest.version == MANIFEST_VERSION) |
| 265 | + | } |
| 266 | + | } |
| 267 | + | |
| 268 | + | /// What the next backup of a repository cuts. |
| 269 | + | #[derive(Clone, Debug, PartialEq, Eq)] |
| 270 | + | pub struct Plan { |
| 271 | + | pub kind: BackupKind, |
| 272 | + | pub prerequisites: Vec<String>, |
| 273 | + | pub previous_refs: BTreeMap<String, String>, |
| 274 | + | } |
| 275 | + | |
| 276 | + | /// The next backup, from the manifest (what a restore reads) and the |
| 277 | + | /// entry the repository's row says was last (`last_entry`). Full when |
| 278 | + | /// there is no chain, the two disagree, the last entry had no refs, or the |
| 279 | + | /// chain already holds `full_every` incremental backups. Otherwise |
| 280 | + | /// incremental, leaving out every commit the last entry's refs reach. |
| 281 | + | pub fn plan(manifest: Option<&Manifest>, last_entry: Option<&str>, full_every: u32) -> Plan { |
| 282 | + | let full = Plan { kind: BackupKind::Full, prerequisites: Vec::new(), previous_refs: BTreeMap::new() }; |
| 283 | + | let Some(manifest) = manifest else { return full }; |
| 284 | + | let Some(last) = manifest.last() else { return full }; |
| 285 | + | let previous_refs = last.refs.clone(); |
| 286 | + | if last_entry != Some(last.id.as_str()) || last.refs.is_empty() { |
| 287 | + | return full; |
| 288 | + | } |
| 289 | + | let incrementals = manifest.chain.len().saturating_sub(1); |
| 290 | + | if incrementals >= full_every as usize { |
| 291 | + | return Plan { previous_refs, ..full }; |
| 292 | + | } |
| 293 | + | let prerequisites: BTreeSet<&String> = last.refs.values().collect(); |
| 294 | + | Plan { |
| 295 | + | kind: BackupKind::Incremental, |
| 296 | + | prerequisites: prerequisites.into_iter().cloned().collect(), |
| 297 | + | previous_refs, |
| 298 | + | } |
| 299 | + | } |
| 300 | + | |
| 301 | + | /// `20261006T025300Z`, from milliseconds since the epoch. |
| 302 | + | pub fn stamp(now_ms: u64) -> String { |
| 303 | + | let text = rfc3339(now_ms); |
| 304 | + | let whole = text.split('.').next().unwrap_or(&text).trim_end_matches('Z'); |
| 305 | + | format!("{}Z", whole.replace(['-', ':'], "")) |
| 306 | + | } |
| 307 | + | |
| 308 | + | pub fn manifest_key(repo_id: &str) -> String { |
| 309 | + | format!("backups/{repo_id}/manifest.json") |
| 310 | + | } |
| 311 | + | |
| 312 | + | pub fn bundle_key(repo_id: &str, id: &str, kind: BackupKind) -> String { |
| 313 | + | format!("backups/{repo_id}/{id}-{}.bundle", kind.suffix()) |
| 314 | + | } |
| 315 | + | |
| 316 | + | fn is_hash(text: &str) -> bool { |
| 317 | + | text.len() == 40 && text.bytes().all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b)) |
| 318 | + | } |
| 319 | + | |
| 320 | + | /// Whether what a sandbox says a bundle holds can be: ref names git would |
| 321 | + | /// make (`HEAD` or under `refs/`) pointing at full commit hashes. |
| 322 | + | pub fn valid_refs(refs: &BTreeMap<String, String>) -> bool { |
| 323 | + | refs.iter().all(|(name, hash)| { |
| 324 | + | let named = name == "HEAD" |
| 325 | + | || (name.starts_with("refs/") |
| 326 | + | && name.len() <= 1024 |
| 327 | + | && !name.contains("..") |
| 328 | + | && !name.ends_with('/') |
| 329 | + | && !name.bytes().any(|b| b <= b' ' || b"~^:?*[\\".contains(&b))); |
| 330 | + | named && is_hash(hash) |
| 331 | + | }) |
| 332 | + | } |
| 333 | + | |
| 334 | + | /// Whether the parts a sandbox says it sent are the ones a bundle of |
| 335 | + | /// `size` bytes makes, numbered from 1 with none missing. |
| 336 | + | pub fn parts_fit(parts: &[BackupPart], size: u64, part_bytes: u64) -> bool { |
| 337 | + | let wanted = size.div_ceil(part_bytes.max(1)); |
| 338 | + | parts.len() as u64 == wanted && parts.iter().enumerate().all(|(index, part)| part.number as usize == index + 1) |
| 339 | + | } |
| 340 | + | |
| 341 | + | // --------------------------------------------------------------------- |
| 342 | + | // The record in D1, and the job in hand |
| 343 | + | // --------------------------------------------------------------------- |
| 344 | + | |
| 345 | + | /// A repository's row in `repo_backups`, as a job reads it. |
| 346 | + | #[derive(Clone, Debug, Default, Deserialize)] |
| 347 | + | struct JobRow { |
| 348 | + | repo_id: String, |
| 349 | + | token_hash: Option<String>, |
| 350 | + | attempts: Option<f64>, |
| 351 | + | target_version: Option<f64>, |
| 352 | + | store_key: Option<String>, |
| 353 | + | upload_key: Option<String>, |
| 354 | + | upload_id: Option<String>, |
| 355 | + | upload_kind: Option<String>, |
| 356 | + | upload_entry: Option<String>, |
| 357 | + | prerequisites: Option<String>, |
| 358 | + | last_entry: Option<String>, |
| 359 | + | } |
| 360 | + | |
| 361 | + | fn refused<T>(message: &str) -> Outcome<T> { |
| 362 | + | Outcome::fail(FailureCode::NotFound, message) |
| 363 | + | } |
| 364 | + | |
| 365 | + | fn n(value: u64) -> JsValue { |
| 366 | + | JsValue::from_f64(value as f64) |
| 367 | + | } |
| 368 | + | |
| 369 | + | fn text(value: Option<&str>) -> JsValue { |
| 370 | + | value.map_or(JsValue::NULL, JsValue::from) |
| 371 | + | } |
| 372 | + | |
| 373 | + | /// What a nightly run did. |
| 374 | + | #[derive(Debug, Default)] |
| 375 | + | pub struct Night { |
| 376 | + | pub queued: u32, |
| 377 | + | pub pruned: u32, |
| 378 | + | } |
| 379 | + | |
| 380 | + | /// Queues tonight's backups: the repositories due, `settings.per_night` |
| 381 | + | /// at most. Also removes the backups of repositories that were purged. |
| 382 | + | pub async fn nightly(db: &D1Database, blobs: &Store, settings: Settings, now: u64) -> Result<Night> { |
| 383 | + | let candidates = db |
| 384 | + | .prepare(DUE_SQL) |
| 385 | + | .bind(&[PULLS_NAMESPACE.into(), n(settings.per_night as u64)])? |
| 386 | + | .all() |
| 387 | + | .await? |
| 388 | + | .results::<CandidateRow>()? |
| 389 | + | .into_iter() |
| 390 | + | .map(Candidate::from) |
| 391 | + | .collect::<Vec<_>>(); |
| 392 | + | let due = pick_due(&candidates, settings.per_night as usize); |
| 393 | + | let mut statements = Vec::new(); |
| 394 | + | for repo_id in &due { |
| 395 | + | statements.push( |
| 396 | + | db.prepare( |
| 397 | + | "INSERT INTO repo_backups (repo_id, status, queued_ms, attempts) VALUES (?1, 'queued', ?2, 0) |
| 398 | + | ON CONFLICT (repo_id) DO UPDATE SET status = 'queued', queued_ms = ?2, attempts = 0, last_error = NULL |
| 399 | + | WHERE repo_backups.status = 'idle'", |
| 400 | + | ) |
| 401 | + | .bind(&[repo_id.as_str().into(), n(now)])?, |
| 402 | + | ); |
| 403 | + | } |
| 404 | + | if !statements.is_empty() { |
| 405 | + | db.batch(statements).await?; |
| 406 | + | } |
| 407 | + | let pruned = prune(db, blobs).await?; |
| 408 | + | Ok(Night { queued: due.len() as u32, pruned }) |
| 409 | + | } |
| 410 | + | |
| 411 | + | /// The row the due query reads; D1 gives numbers as floats. |
| 412 | + | #[derive(Deserialize)] |
| 413 | + | struct CandidateRow { |
| 414 | + | repo_id: String, |
| 415 | + | namespace: String, |
| 416 | + | created_at: Option<serde_json::Value>, |
| 417 | + | refs_version: Option<f64>, |
| 418 | + | refs_open_until: Option<f64>, |
| 419 | + | deleted: Option<f64>, |
| 420 | + | retired: Option<f64>, |
| 421 | + | status: Option<String>, |
| 422 | + | backed_version: Option<f64>, |
| 423 | + | backed_up_ms: Option<f64>, |
| 424 | + | } |
| 425 | + | |
| 426 | + | impl From<CandidateRow> for Candidate { |
| 427 | + | fn from(row: CandidateRow) -> Candidate { |
| 428 | + | Candidate { |
| 429 | + | repo_id: row.repo_id, |
| 430 | + | namespace: row.namespace, |
| 431 | + | created_at: match row.created_at { |
| 432 | + | Some(serde_json::Value::String(text)) => text, |
| 433 | + | Some(other) => other.to_string(), |
| 434 | + | None => String::new(), |
| 435 | + | }, |
| 436 | + | refs_version: row.refs_version.unwrap_or(0.0) as u64, |
| 437 | + | refs_open_until: row.refs_open_until.unwrap_or(0.0) as u64, |
| 438 | + | deleted: row.deleted.unwrap_or(0.0) != 0.0, |
| 439 | + | retired: row.retired.unwrap_or(0.0) != 0.0, |
| 440 | + | status: row.status, |
| 441 | + | backed_version: row.backed_version.map(|v| v as u64), |
| 442 | + | backed_up_ms: row.backed_up_ms.unwrap_or(0.0) as u64, |
| 443 | + | } |
| 444 | + | } |
| 445 | + | } |
| 446 | + | |
| 447 | + | /// Removes the backups of repositories that no longer exist (purged, so |
| 448 | + | /// their data is gone for good): every bundle the manifest lists, the |
| 449 | + | /// manifest, and the row. |
| 450 | + | async fn prune(db: &D1Database, blobs: &Store) -> Result<u32> { |
| 451 | + | #[derive(Deserialize)] |
| 452 | + | struct Gone { |
| 453 | + | repo_id: String, |
| 454 | + | } |
| 455 | + | let gone = db |
| 456 | + | .prepare( |
| 457 | + | "SELECT b.repo_id FROM repo_backups b LEFT JOIN repos r ON r.id = b.repo_id |
| 458 | + | WHERE r.id IS NULL LIMIT ?1", |
| 459 | + | ) |
| 460 | + | .bind(&[n(PRUNES_PER_NIGHT as u64)])? |
| 461 | + | .all() |
| 462 | + | .await? |
| 463 | + | .results::<Gone>()?; |
| 464 | + | for row in &gone { |
| 465 | + | let key = manifest_key(&row.repo_id); |
| 466 | + | if let Some(manifest) = blobs.read(&key).await?.as_deref().and_then(Manifest::from_bytes) { |
| 467 | + | for bundle in manifest.keys() { |
| 468 | + | blobs.delete(&bundle).await?; |
| 469 | + | } |
| 470 | + | } |
| 471 | + | blobs.delete(&key).await?; |
| 472 | + | db.prepare("DELETE FROM repo_backups WHERE repo_id = ?1") |
| 473 | + | .bind(&[row.repo_id.as_str().into()])? |
| 474 | + | .run() |
| 475 | + | .await?; |
| 476 | + | } |
| 477 | + | Ok(gone.len() as u32) |
| 478 | + | } |
| 479 | + | |
| 480 | + | /// For the runner's sweep: up to `limit` queued backups, each marked |
| 481 | + | /// running with a token of its own, so long as no more than `max_running` |
| 482 | + | /// are then running. Jobs past their lease go back in the queue first. |
| 483 | + | pub async fn claim(db: &D1Database, blobs: Option<&Store>, a: &ClaimBackupsArgs, now: u64) -> Result<Vec<BackupClaim>> { |
| 484 | + | // Past their lease: the sandbox died without saying so. |
| 485 | + | let stale = db |
| 486 | + | .prepare("SELECT * FROM repo_backups WHERE status = 'running' AND claimed_ms < ?1") |
| 487 | + | .bind(&[n(now.saturating_sub(LEASE_MS))])? |
| 488 | + | .all() |
| 489 | + | .await? |
| 490 | + | .results::<JobRow>()?; |
| 491 | + | for row in &stale { |
| 492 | + | give_up_upload(blobs, row).await; |
| 493 | + | settle_failure(db, row, "The backup ran past its time and was started again.").await?; |
| 494 | + | } |
| 495 | + | if blobs.is_none() { |
| 496 | + | return Ok(Vec::new()); |
| 497 | + | } |
| 498 | + | |
| 499 | + | #[derive(Deserialize)] |
| 500 | + | struct Count { |
| 501 | + | running: f64, |
| 502 | + | } |
| 503 | + | let running = db |
| 504 | + | .prepare("SELECT count(*) AS running FROM repo_backups WHERE status = 'running'") |
| 505 | + | .first::<Count>(None) |
| 506 | + | .await? |
| 507 | + | .map_or(0, |row| row.running as u32); |
| 508 | + | let room = a.max_running.saturating_sub(running).min(a.limit); |
| 509 | + | if room == 0 { |
| 510 | + | return Ok(Vec::new()); |
| 511 | + | } |
| 512 | + | |
| 513 | + | #[derive(Deserialize)] |
| 514 | + | struct Queued { |
| 515 | + | repo_id: String, |
| 516 | + | namespace: Option<String>, |
| 517 | + | name: Option<String>, |
| 518 | + | deleted: Option<f64>, |
| 519 | + | } |
| 520 | + | let queued = db |
| 521 | + | .prepare( |
| 522 | + | "SELECT b.repo_id, r.namespace, r.name, (r.id IS NULL OR r.deleted_at IS NOT NULL) AS deleted |
| 523 | + | FROM repo_backups b LEFT JOIN repos r ON r.id = b.repo_id |
| 524 | + | WHERE b.status = 'queued' ORDER BY b.queued_ms, b.repo_id LIMIT ?1", |
| 525 | + | ) |
| 526 | + | .bind(&[n(room as u64)])? |
| 527 | + | .all() |
| 528 | + | .await? |
| 529 | + | .results::<Queued>()?; |
| 530 | + | let mut claims = Vec::new(); |
| 531 | + | for row in queued { |
| 532 | + | let (Some(namespace), Some(name)) = (row.namespace, row.name) else { continue }; |
| 533 | + | if row.deleted.unwrap_or(0.0) != 0.0 { |
| 534 | + | // Deleted since it was queued: nothing to back up until it is |
| 535 | + | // restored, if it is. |
| 536 | + | db.prepare("UPDATE repo_backups SET status = 'idle' WHERE repo_id = ?1 AND status = 'queued'") |
| 537 | + | .bind(&[row.repo_id.as_str().into()])? |
| 538 | + | .run() |
| 539 | + | .await?; |
| 540 | + | continue; |
| 541 | + | } |
| 542 | + | let job_id = new_id("bkp", now); |
| 543 | + | let token = g1t_secrets::random_hex(32); |
| 544 | + | let claimed = db |
| 545 | + | .prepare( |
| 546 | + | "UPDATE repo_backups SET status = 'running', job_id = ?2, token_hash = ?3, claimed_ms = ?4, |
| 547 | + | attempts = attempts + 1, target_version = NULL, store_key = NULL, upload_key = NULL, |
| 548 | + | upload_id = NULL, upload_kind = NULL, upload_entry = NULL, prerequisites = NULL |
| 549 | + | WHERE repo_id = ?1 AND status = 'queued' RETURNING repo_id", |
| 550 | + | ) |
| 551 | + | .bind(&[ |
| 552 | + | row.repo_id.as_str().into(), |
| 553 | + | job_id.as_str().into(), |
| 554 | + | g1t_secrets::sha256_hex(&token).into(), |
| 555 | + | n(now), |
| 556 | + | ])? |
| 557 | + | .first::<serde_json::Value>(None) |
| 558 | + | .await?; |
| 559 | + | if claimed.is_some() { |
| 560 | + | claims.push(BackupClaim { job_id, token, repo_id: row.repo_id, path: RepoPath { namespace, name } }); |
| 561 | + | } |
| 562 | + | } |
| 563 | + | Ok(claims) |
| 564 | + | } |
| 565 | + | |
| 566 | + | /// The running job `a` names, if its token is the one it was given. |
| 567 | + | async fn job(db: &D1Database, a: &BackupJobArgs) -> Result<Option<JobRow>> { |
| 568 | + | let row = db |
| 569 | + | .prepare("SELECT * FROM repo_backups WHERE job_id = ?1 AND status = 'running'") |
| 570 | + | .bind(&[a.job_id.as_str().into()])? |
| 571 | + | .first::<JobRow>(None) |
| 572 | + | .await?; |
| 573 | + | Ok(row.filter(|row| { |
| 574 | + | row.token_hash |
| 575 | + | .as_deref() |
| 576 | + | .is_some_and(|hash| g1t_secrets::same(hash, &g1t_secrets::sha256_hex(&a.token))) |
| 577 | + | })) |
| 578 | + | } |
| 579 | + | |
| 580 | + | const NO_JOB: &str = "No such backup job, or it is not running."; |
| 581 | + | |
| 582 | + | /// The job, for its sandbox: what to cut, and a read-only credential for |
| 583 | + | /// the repository that lasts minutes. Asked again, the same bundle with a |
| 584 | + | /// new credential. |
| 585 | + | pub async fn spec<S: GitStore>( |
| 586 | + | registry: &Registry, |
| 587 | + | blobs: &Store, |
| 588 | + | store: &S, |
| 589 | + | a: &BackupJobArgs, |
| 590 | + | full_every: u32, |
| 591 | + | now: u64, |
| 592 | + | ) -> Result<Outcome<BackupSpec>> { |
| 593 | + | let db = ®istry.db; |
| 594 | + | let Some(row) = job(db, a).await? else { return Ok(refused(NO_JOB)) }; |
| 595 | + | #[derive(Deserialize)] |
| 596 | + | struct Live { |
| 597 | + | refs_version: Option<f64>, |
| 598 | + | } |
| 599 | + | let Some(repo) = registry.by_id(&row.repo_id).await? else { |
| 600 | + | settle_failure(db, &row, "The repository was deleted.").await?; |
| 601 | + | return Ok(refused("The repository was deleted.")); |
| 602 | + | }; |
| 603 | + | let key = store_key(&repo); |
| 604 | + | let access = store.handout(&key, Scope::Read).await?; |
| 605 | + | // Asked before: the same bundle, so parts already sent still fit. |
| 606 | + | if let (Some(kind), Some(_)) = (row.upload_kind.as_deref(), row.upload_id.as_deref()) { |
| 607 | + | let manifest = blobs.read(&manifest_key(&row.repo_id)).await?.as_deref().and_then(Manifest::from_bytes); |
| 608 | + | let prerequisites: Vec<String> = row.prerequisites.as_deref().and_then(|p| serde_json::from_str(p).ok()).unwrap_or_default(); |
| 609 | + | let previous_refs = manifest.as_ref().and_then(|m| m.last()).map(|e| e.refs.clone()).unwrap_or_default(); |
| 610 | + | return Ok(Outcome::Ok(BackupSpec { |
| 611 | + | kind: if kind == "full" { BackupKind::Full } else { BackupKind::Incremental }, |
| 612 | + | remote: access.remote, |
| 613 | + | git_token: access.token, |
| 614 | + | prerequisites, |
| 615 | + | previous_refs, |
| 616 | + | part_bytes: PART_BYTES, |
| 617 | + | })); |
| 618 | + | } |
| 619 | + | let version = db |
| 620 | + | .prepare("SELECT refs_version FROM repos WHERE id = ?1") |
| 621 | + | .bind(&[row.repo_id.as_str().into()])? |
| 622 | + | .first::<Live>(None) |
| 623 | + | .await? |
| 624 | + | .and_then(|live| live.refs_version) |
| 625 | + | .unwrap_or(0.0) as u64; |
| 626 | + | let manifest = blobs.read(&manifest_key(&row.repo_id)).await?.as_deref().and_then(Manifest::from_bytes); |
| 627 | + | let next = plan(manifest.as_ref(), row.last_entry.as_deref(), full_every); |
| 628 | + | let entry = stamp(now); |
| 629 | + | let upload_key = bundle_key(&row.repo_id, &entry, next.kind); |
| 630 | + | let upload_id = blobs.create_multipart(&upload_key).await?; |
| 631 | + | db.prepare( |
| 632 | + | "UPDATE repo_backups SET target_version = ?2, store_key = ?3, upload_key = ?4, upload_id = ?5, |
| 633 | + | upload_kind = ?6, upload_entry = ?7, prerequisites = ?8, backed_from_ms = ?9 |
| 634 | + | WHERE job_id = ?1", |
| 635 | + | ) |
| 636 | + | .bind(&[ |
| 637 | + | a.job_id.as_str().into(), |
| 638 | + | n(version), |
| 639 | + | key.as_str().into(), |
| 640 | + | upload_key.as_str().into(), |
| 641 | + | upload_id.as_str().into(), |
| 642 | + | next.kind.suffix().into(), |
| 643 | + | entry.as_str().into(), |
| 644 | + | serde_json::to_string(&next.prerequisites)?.into(), |
| 645 | + | n(now), |
| 646 | + | ])? |
| 647 | + | .run() |
| 648 | + | .await?; |
| 649 | + | Ok(Outcome::Ok(BackupSpec { |
| 650 | + | kind: next.kind, |
| 651 | + | remote: access.remote, |
| 652 | + | git_token: access.token, |
| 653 | + | prerequisites: next.prerequisites, |
| 654 | + | previous_refs: next.previous_refs, |
| 655 | + | part_bytes: PART_BYTES, |
| 656 | + | })) |
| 657 | + | } |
| 658 | + | |
| 659 | + | /// One part of the job's bundle, kept. |
| 660 | + | pub async fn part(db: &D1Database, blobs: &Store, a: &BackupJobArgs, number: u16, bytes: Vec<u8>) -> Result<Outcome<BackupPart>> { |
| 661 | + | let Some(row) = job(db, a).await? else { return Ok(refused(NO_JOB)) }; |
| 662 | + | let (Some(key), Some(upload)) = (row.upload_key.as_deref(), row.upload_id.as_deref()) else { |
| 663 | + | return Ok(Outcome::fail(FailureCode::Conflict, "Ask for the job's spec first.")); |
| 664 | + | }; |
| 665 | + | if number == 0 || bytes.len() as u64 > PART_BYTES { |
| 666 | + | return Ok(Outcome::fail(FailureCode::Invalid, "Parts are numbered from 1, and hold 32 MiB at most.")); |
| 667 | + | } |
| 668 | + | let Part { number, etag } = blobs.upload_part(key, upload, number, bytes).await?; |
| 669 | + | Ok(Outcome::Ok(BackupPart { number, etag })) |
| 670 | + | } |
| 671 | + | |
| 672 | + | /// The bundle is cut and sent: the upload is completed, the chain gains |
| 673 | + | /// an entry (unless nothing changed), bundles of the chain before last |
| 674 | + | /// are removed, and the repository is backed up as of the refs version it |
| 675 | + | /// had when the clone began. |
| 676 | + | pub async fn complete(registry: &Registry, blobs: &Store, a: &BackupComplete, now: u64) -> Result<Outcome<bool>> { |
| 677 | + | let db = ®istry.db; |
| 678 | + | let args = BackupJobArgs { job_id: a.job_id.clone(), token: a.token.clone() }; |
| 679 | + | let Some(row) = job(db, &args).await? else { return Ok(refused(NO_JOB)) }; |
| 680 | + | let (Some(upload_key), Some(upload_id), Some(entry_id), Some(store_key)) = |
| 681 | + | (row.upload_key.clone(), row.upload_id.clone(), row.upload_entry.clone(), row.store_key.clone()) |
| 682 | + | else { |
| 683 | + | return Ok(Outcome::fail(FailureCode::Conflict, "Ask for the job's spec first.")); |
| 684 | + | }; |
| 685 | + | if !valid_refs(&a.refs) { |
| 686 | + | return Ok(Outcome::fail(FailureCode::Invalid, "A ref name or commit is not one git makes.")); |
| 687 | + | } |
| 688 | + | if !parts_fit(&a.parts, a.size, PART_BYTES) { |
| 689 | + | return Ok(Outcome::fail(FailureCode::Invalid, "The parts do not add up to the bundle's size.")); |
| 690 | + | } |
| 691 | + | meter_fetch(&store_key, a.fetched_bytes); |
| 692 | + | let kind = if row.upload_kind.as_deref() == Some("full") { BackupKind::Full } else { BackupKind::Incremental }; |
| 693 | + | let mut manifest = blobs |
| 694 | + | .read(&manifest_key(&row.repo_id)) |
| 695 | + | .await? |
| 696 | + | .as_deref() |
| 697 | + | .and_then(Manifest::from_bytes) |
| 698 | + | .unwrap_or_else(|| Manifest::new(&row.repo_id, &store_key)); |
| 699 | + | let unchanged = a.size == 0 && kind == BackupKind::Incremental && manifest.last().is_some_and(|last| last.refs == a.refs); |
| 700 | + | let version = row.target_version.unwrap_or(0.0) as u64; |
| 701 | + | let mut last_entry = row.last_entry.clone(); |
| 702 | + | if a.size == 0 { |
| 703 | + | blobs.abort_multipart(&upload_key, &upload_id).await?; |
| 704 | + | } else { |
| 705 | + | let parts: Vec<Part> = a.parts.iter().map(|p| Part { number: p.number, etag: p.etag.clone() }).collect(); |
| 706 | + | blobs.complete_multipart(&upload_key, &upload_id, &parts).await?; |
| 707 | + | } |
| 708 | + | if !unchanged { |
| 709 | + | let prerequisites: Vec<String> = row.prerequisites.as_deref().and_then(|p| serde_json::from_str(p).ok()).unwrap_or_default(); |
| 710 | + | let dropped = manifest.add(Entry { |
| 711 | + | id: entry_id.clone(), |
| 712 | + | kind, |
| 713 | + | key: (a.size > 0).then(|| upload_key.clone()), |
| 714 | + | created_at: rfc3339(now), |
| 715 | + | refs_version: version, |
| 716 | + | refs: a.refs.clone(), |
| 717 | + | prerequisites, |
| 718 | + | size: a.size, |
| 719 | + | sha256: a.sha256.clone(), |
| 720 | + | }); |
| 721 | + | manifest.store_key = store_key.clone(); |
| 722 | + | manifest.updated_at = rfc3339(now); |
| 723 | + | manifest.path = registry |
| 724 | + | .by_id(&row.repo_id) |
| 725 | + | .await? |
| 726 | + | .map(|repo| RepoPath { namespace: repo.namespace, name: repo.name }); |
| 727 | + | blobs.put(&manifest_key(&row.repo_id), manifest.to_bytes()).await?; |
| 728 | + | for key in dropped { |
| 729 | + | if let Err(error) = blobs.delete(&key).await { |
| 730 | + | worker::console_error!("repos: backup {key} not removed: {error}"); |
| 731 | + | } |
| 732 | + | } |
| 733 | + | last_entry = Some(entry_id); |
| 734 | + | } |
| 735 | + | db.prepare( |
| 736 | + | "UPDATE repo_backups SET status = 'idle', job_id = NULL, token_hash = NULL, claimed_ms = NULL, |
| 737 | + | attempts = 0, last_error = NULL, refs_version = ?2, backed_up_ms = backed_from_ms, |
| 738 | + | tips = ?3, last_entry = ?4, upload_key = NULL, upload_id = NULL, upload_kind = NULL, |
| 739 | + | upload_entry = NULL, prerequisites = NULL, target_version = NULL |
| 740 | + | WHERE job_id = ?1", |
| 741 | + | ) |
| 742 | + | .bind(&[a.job_id.as_str().into(), n(version), serde_json::to_string(&a.refs)?.into(), text(last_entry.as_deref())])? |
| 743 | + | .run() |
| 744 | + | .await?; |
| 745 | + | Ok(Outcome::Ok(true)) |
| 746 | + | } |
| 747 | + | |
| 748 | + | /// The sandbox could not do the job: tried again later tonight, up to |
| 749 | + | /// [`MAX_ATTEMPTS`] times, and tomorrow night after that. |
| 750 | + | pub async fn fail(db: &D1Database, blobs: &Store, a: &BackupFail) -> Result<Outcome<bool>> { |
| 751 | + | let args = BackupJobArgs { job_id: a.job_id.clone(), token: a.token.clone() }; |
| 752 | + | let Some(row) = job(db, &args).await? else { return Ok(refused(NO_JOB)) }; |
| 753 | + | if let Some(key) = row.store_key.as_deref() { |
| 754 | + | meter_fetch(key, a.fetched_bytes); |
| 755 | + | } |
| 756 | + | give_up_upload(Some(blobs), &row).await; |
| 757 | + | let said: String = a.error.chars().take(500).collect(); |
| 758 | + | settle_failure(db, &row, &said).await?; |
| 759 | + | Ok(Outcome::Ok(true)) |
| 760 | + | } |
| 761 | + | |
| 762 | + | /// A clone counts once, with what it read, whenever it got as far as |
| 763 | + | /// reading anything. |
| 764 | + | fn meter_fetch(store_key: &str, fetched_bytes: u64) { |
| 765 | + | if fetched_bytes > 0 { |
| 766 | + | meters::record("internal.git.info_refs", store_key, 0, 0); |
| 767 | + | meters::record(FETCH_METER, store_key, 0, fetched_bytes); |
| 768 | + | } |
| 769 | + | } |
| 770 | + | |
| 771 | + | async fn give_up_upload(blobs: Option<&Store>, row: &JobRow) { |
| 772 | + | if let (Some(blobs), Some(key), Some(upload)) = (blobs, row.upload_key.as_deref(), row.upload_id.as_deref()) |
| 773 | + | && let Err(error) = blobs.abort_multipart(key, upload).await |
| 774 | + | { |
| 775 | + | worker::console_error!("repos: backup upload {key} not given up: {error}"); |
| 776 | + | } |
| 777 | + | } |
| 778 | + | |
| 779 | + | /// Back in the queue for another try, or idle until tomorrow night once |
| 780 | + | /// it has been tried enough. |
| 781 | + | async fn settle_failure(db: &D1Database, row: &JobRow, error: &str) -> Result<()> { |
| 782 | + | let attempts = row.attempts.unwrap_or(0.0) as u32; |
| 783 | + | let next = if attempts >= MAX_ATTEMPTS { "idle" } else { "queued" }; |
| 784 | + | worker::console_error!("repos: backup of {} failed ({attempts} of {MAX_ATTEMPTS}): {error}", row.repo_id); |
| 785 | + | db.prepare( |
| 786 | + | "UPDATE repo_backups SET status = ?2, last_error = ?3, job_id = NULL, token_hash = NULL, |
| 787 | + | claimed_ms = NULL, upload_key = NULL, upload_id = NULL, upload_kind = NULL, |
| 788 | + | upload_entry = NULL, prerequisites = NULL, target_version = NULL |
| 789 | + | WHERE repo_id = ?1", |
| 790 | + | ) |
| 791 | + | .bind(&[row.repo_id.as_str().into(), next.into(), error.into()])? |
| 792 | + | .run() |
| 793 | + | .await?; |
| 794 | + | Ok(()) |
| 795 | + | } |
| 796 | + | |
| 797 | + | #[cfg(test)] |
| 798 | + | mod tests { |
| 799 | + | use super::*; |
| 800 | + | |
| 801 | + | const A: &str = "c71546fcd893ef8b0f57388b65e620d759705dda"; |
| 802 | + | const B: &str = "4807077b296e6edbf410d55e72749d3e1170c291"; |
| 803 | + | const C: &str = "0000000000000000000000000000000000000abc"; |
| 804 | + | |
| 805 | + | fn candidate(id: &str) -> Candidate { |
| 806 | + | Candidate { |
| 807 | + | repo_id: id.into(), |
| 808 | + | namespace: "acme".into(), |
| 809 | + | created_at: "2026-01-01T00:00:00.000Z".into(), |
| 810 | + | refs_version: 3, |
| 811 | + | ..Candidate::default() |
| 812 | + | } |
| 813 | + | } |
| 814 | + | |
| 815 | + | fn backed(id: &str, version: u64, at: u64) -> Candidate { |
| 816 | + | Candidate { status: Some("idle".into()), backed_version: Some(version), backed_up_ms: at, ..candidate(id) } |
| 817 | + | } |
| 818 | + | |
| 819 | + | #[test] |
| 820 | + | fn a_repository_never_backed_up_is_due() { |
| 821 | + | assert!(is_due(&candidate("r1"))); |
| 822 | + | // Even one whose refs never moved: rows from before refs_version. |
| 823 | + | assert!(is_due(&Candidate { refs_version: 0, ..candidate("r1") })); |
| 824 | + | } |
| 825 | + | |
| 826 | + | #[test] |
| 827 | + | fn a_repository_is_due_only_once_its_refs_moved() { |
| 828 | + | assert!(!is_due(&backed("r1", 3, 1_000))); |
| 829 | + | assert!(is_due(&backed("r1", 2, 1_000))); |
| 830 | + | // A credential that could push went out after the last backup began. |
| 831 | + | assert!(is_due(&Candidate { refs_open_until: 2_000, ..backed("r1", 3, 1_000) })); |
| 832 | + | assert!(!is_due(&Candidate { refs_open_until: 900, ..backed("r1", 3, 1_000) })); |
| 833 | + | // Tried and failed: no version yet. |
| 834 | + | assert!(is_due(&Candidate { backed_version: None, ..backed("r1", 0, 0) })); |
| 835 | + | } |
| 836 | + | |
| 837 | + | #[test] |
| 838 | + | fn deleted_retired_working_copies_and_jobs_in_hand_are_not_due() { |
| 839 | + | assert!(!is_due(&Candidate { deleted: true, ..candidate("r1") })); |
| 840 | + | assert!(!is_due(&Candidate { retired: true, ..candidate("r1") })); |
| 841 | + | assert!(!is_due(&Candidate { namespace: PULLS_NAMESPACE.into(), ..candidate("r1") })); |
| 842 | + | assert!(!is_due(&Candidate { status: Some("queued".into()), ..backed("r1", 1, 0) })); |
| 843 | + | assert!(!is_due(&Candidate { status: Some("running".into()), ..backed("r1", 1, 0) })); |
| 844 | + | } |
| 845 | + | |
| 846 | + | #[test] |
| 847 | + | fn the_longest_waiting_go_first_and_no_more_than_the_limit() { |
| 848 | + | let candidates = vec![ |
| 849 | + | backed("recent", 1, 5_000), |
| 850 | + | backed("older", 1, 1_000), |
| 851 | + | backed("current", 3, 0), |
| 852 | + | Candidate { created_at: "2026-02-01T00:00:00.000Z".into(), ..candidate("new-b") }, |
| 853 | + | Candidate { created_at: "2025-02-01T00:00:00.000Z".into(), ..candidate("new-a") }, |
| 854 | + | ]; |
| 855 | + | assert_eq!(pick_due(&candidates, 10), ["new-a", "new-b", "older", "recent"]); |
| 856 | + | assert_eq!(pick_due(&candidates, 2), ["new-a", "new-b"]); |
| 857 | + | assert!(pick_due(&candidates, 0).is_empty()); |
| 858 | + | } |
| 859 | + | |
| 860 | + | fn refs(pairs: &[(&str, &str)]) -> BTreeMap<String, String> { |
| 861 | + | pairs.iter().map(|(name, hash)| ((*name).to_owned(), (*hash).to_owned())).collect() |
| 862 | + | } |
| 863 | + | |
| 864 | + | fn entry(id: &str, kind: BackupKind, tips: &[(&str, &str)]) -> Entry { |
| 865 | + | Entry { |
| 866 | + | id: id.into(), |
| 867 | + | kind, |
| 868 | + | key: Some(format!("backups/r1/{id}-{}.bundle", kind.suffix())), |
| 869 | + | created_at: String::new(), |
| 870 | + | refs_version: 1, |
| 871 | + | refs: refs(tips), |
| 872 | + | prerequisites: Vec::new(), |
| 873 | + | size: 10, |
| 874 | + | sha256: None, |
| 875 | + | } |
| 876 | + | } |
| 877 | + | |
| 878 | + | fn manifest(incrementals: usize) -> Manifest { |
| 879 | + | let mut manifest = Manifest::new("r1", "acme--rocket"); |
| 880 | + | manifest.add(entry("e0", BackupKind::Full, &[("HEAD", A), ("refs/heads/main", A)])); |
| 881 | + | for index in 0..incrementals { |
| 882 | + | manifest.add(entry(&format!("e{}", index + 1), BackupKind::Incremental, &[("HEAD", A), ("refs/heads/main", A), ("refs/tags/v1", B)])); |
| 883 | + | } |
| 884 | + | manifest |
| 885 | + | } |
| 886 | + | |
| 887 | + | #[test] |
| 888 | + | fn the_first_backup_is_full() { |
| 889 | + | let next = plan(None, None, 30); |
| 890 | + | assert_eq!(next.kind, BackupKind::Full); |
| 891 | + | assert!(next.prerequisites.is_empty() && next.previous_refs.is_empty()); |
| 892 | + | assert_eq!(plan(Some(&Manifest::new("r1", "k")), None, 30).kind, BackupKind::Full); |
| 893 | + | } |
| 894 | + | |
| 895 | + | #[test] |
| 896 | + | fn the_next_is_incremental_from_the_last_tips_each_once() { |
| 897 | + | let manifest = manifest(1); |
| 898 | + | let next = plan(Some(&manifest), Some("e1"), 30); |
| 899 | + | assert_eq!(next.kind, BackupKind::Incremental); |
| 900 | + | assert_eq!(next.prerequisites, [B, A]); |
| 901 | + | assert_eq!(next.previous_refs, manifest.last().unwrap().refs); |
| 902 | + | } |
| 903 | + | |
| 904 | + | #[test] |
| 905 | + | fn a_full_bundle_is_cut_again_after_enough_incrementals() { |
| 906 | + | assert_eq!(plan(Some(&manifest(29)), Some("e29"), 30).kind, BackupKind::Incremental); |
| 907 | + | let next = plan(Some(&manifest(30)), Some("e30"), 30); |
| 908 | + | assert_eq!(next.kind, BackupKind::Full); |
| 909 | + | assert!(next.prerequisites.is_empty()); |
| 910 | + | // What it last held is still said, so an unchanged clone is seen. |
| 911 | + | assert!(!next.previous_refs.is_empty()); |
| 912 | + | assert_eq!(plan(Some(&manifest(1)), Some("e1"), 1).kind, BackupKind::Full); |
| 913 | + | } |
| 914 | + | |
| 915 | + | #[test] |
| 916 | + | fn a_chain_the_row_does_not_know_or_an_empty_one_starts_again() { |
| 917 | + | assert_eq!(plan(Some(&manifest(2)), Some("e1"), 30).kind, BackupKind::Full); |
| 918 | + | assert_eq!(plan(Some(&manifest(2)), None, 30).kind, BackupKind::Full); |
| 919 | + | let mut empty = Manifest::new("r1", "k"); |
| 920 | + | empty.add(entry("e0", BackupKind::Full, &[])); |
| 921 | + | assert_eq!(plan(Some(&empty), Some("e0"), 30).kind, BackupKind::Full); |
| 922 | + | } |
| 923 | + | |
| 924 | + | #[test] |
| 925 | + | fn a_full_backup_starts_a_chain_and_the_one_before_last_goes() { |
| 926 | + | let mut manifest = manifest(2); |
| 927 | + | assert!(manifest.previous.is_empty()); |
| 928 | + | let dropped = manifest.add(entry("f1", BackupKind::Full, &[("refs/heads/main", C)])); |
| 929 | + | assert!(dropped.is_empty(), "the chain before is kept until the next full one"); |
| 930 | + | assert_eq!(manifest.chain.len(), 1); |
| 931 | + | assert_eq!(manifest.previous.len(), 3); |
| 932 | + | manifest.add(entry("f1a", BackupKind::Incremental, &[("refs/heads/main", C)])); |
| 933 | + | let dropped = manifest.add(entry("f2", BackupKind::Full, &[("refs/heads/main", C)])); |
| 934 | + | assert_eq!(dropped, ["backups/r1/e0-full.bundle", "backups/r1/e1-incr.bundle", "backups/r1/e2-incr.bundle"]); |
| 935 | + | assert_eq!(manifest.previous.iter().map(|e| e.id.as_str()).collect::<Vec<_>>(), ["f1", "f1a"]); |
| 936 | + | assert_eq!(manifest.keys().len(), 3); |
| 937 | + | } |
| 938 | + | |
| 939 | + | #[test] |
| 940 | + | fn a_manifest_round_trips_and_an_unknown_version_is_not_read() { |
| 941 | + | let mut manifest = manifest(2); |
| 942 | + | manifest.path = Some(RepoPath { namespace: "acme".into(), name: "rocket".into() }); |
| 943 | + | manifest.add(Entry { key: None, size: 0, ..entry("e3", BackupKind::Incremental, &[("refs/heads/main", A)]) }); |
| 944 | + | let bytes = manifest.to_bytes(); |
| 945 | + | let text = String::from_utf8(bytes.clone()).unwrap(); |
| 946 | + | // What the restore drill reads: snake_case, the kind in words. |
| 947 | + | assert!(text.contains("\"repo_id\": \"r1\"") && text.contains("\"kind\": \"incremental\"") && text.contains("\"key\": null")); |
| 948 | + | assert_eq!(Manifest::from_bytes(&bytes), Some(manifest.clone())); |
| 949 | + | let mut later = manifest; |
| 950 | + | later.version = 2; |
| 951 | + | assert_eq!(Manifest::from_bytes(&later.to_bytes()), None); |
| 952 | + | assert_eq!(Manifest::from_bytes(b"not json"), None); |
| 953 | + | } |
| 954 | + | |
| 955 | + | #[test] |
| 956 | + | fn a_backup_clone_is_g1ts_operation_not_the_workspaces() { |
| 957 | + | let mapping = meters::Mapping::defaults(); |
| 958 | + | assert_eq!(mapping.cost(FETCH_METER), 1.0); |
| 959 | + | assert_eq!(mapping.billable(FETCH_METER), 0.0); |
| 960 | + | assert_eq!(mapping.billable("internal.git.fetch"), 1.0); |
| 961 | + | } |
| 962 | + | |
| 963 | + | #[test] |
| 964 | + | fn keys_and_stamps() { |
| 965 | + | assert_eq!(stamp(1_369_353_600_123), "20130524T000000Z"); |
| 966 | + | assert_eq!(manifest_key("r1"), "backups/r1/manifest.json"); |
| 967 | + | assert_eq!(bundle_key("r1", "20130524T000000Z", BackupKind::Incremental), "backups/r1/20130524T000000Z-incr.bundle"); |
| 968 | + | assert_eq!(bundle_key("r1", "20130524T000000Z", BackupKind::Full), "backups/r1/20130524T000000Z-full.bundle"); |
| 969 | + | } |
| 970 | + | |
| 971 | + | #[test] |
| 972 | + | fn only_refs_git_makes_are_taken() { |
| 973 | + | assert!(valid_refs(&refs(&[("HEAD", A), ("refs/heads/main", A), ("refs/pull/pr_1/head", B)]))); |
| 974 | + | assert!(valid_refs(&BTreeMap::new())); |
| 975 | + | assert!(!valid_refs(&refs(&[("main", A)]))); |
| 976 | + | assert!(!valid_refs(&refs(&[("refs/heads/../x", A)]))); |
| 977 | + | assert!(!valid_refs(&refs(&[("refs/heads/a b", A)]))); |
| 978 | + | assert!(!valid_refs(&refs(&[("refs/heads/main", "abc")]))); |
| 979 | + | assert!(!valid_refs(&refs(&[("refs/heads/main", &A.to_uppercase())]))); |
| 980 | + | } |
| 981 | + | |
| 982 | + | #[test] |
| 983 | + | fn parts_must_cover_the_bundle_in_order() { |
| 984 | + | let part = |number| BackupPart { number, etag: "e".into() }; |
| 985 | + | assert!(parts_fit(&[], 0, 10)); |
| 986 | + | assert!(parts_fit(&[part(1)], 10, 10)); |
| 987 | + | assert!(parts_fit(&[part(1), part(2)], 11, 10)); |
| 988 | + | assert!(!parts_fit(&[part(1)], 11, 10)); |
| 989 | + | assert!(!parts_fit(&[part(2), part(1)], 11, 10)); |
| 990 | + | assert!(!parts_fit(&[part(1)], 0, 10)); |
| 991 | + | } |
| 992 | + | } |