g1t/services/repos/src/backups.rs

992 lines40,474 bytesCodeBlame
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
30use std::collections::{BTreeMap, BTreeSet};
31
32use g1t_blobstore::{BlobStore, Config, Part, Store};
33use g1t_contracts::backups::{
34 BackupClaim, BackupComplete, BackupFail, BackupJobArgs, BackupKind, BackupPart, BackupSpec, ClaimBackupsArgs, PART_BYTES,
35};
36use g1t_contracts::repos::RepoPath;
37use g1t_contracts::time::rfc3339;
38use g1t_contracts::{FailureCode, Outcome, new_id};
39use serde::{Deserialize, Serialize};
40use worker::wasm_bindgen::JsValue;
41use worker::{D1Database, Env, Result};
42
43use crate::PULLS_NAMESPACE;
44use crate::meters;
45use crate::registry::{Registry, store_key};
46use 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.
50pub 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).
60pub const FETCH_METER: &str = "internal.git.backup_fetch";
61
62/// How long a claimed job may run before it is given to another sandbox.
63pub const LEASE_MS: u64 = 3 * 60 * 60 * 1000;
64/// How many times a night a backup is tried.
65pub const MAX_ATTEMPTS: u32 = 3;
66/// How many backups of deleted repositories one night removes.
67const PRUNES_PER_NIGHT: u32 = 50;
68
69/// How backups are paced, from the service's variables.
70#[derive(Clone, Copy, Debug, PartialEq, Eq)]
71pub 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
78impl Default for Settings {
79 fn default() -> Self {
80 Settings { per_night: 200, full_every: 30 }
81 }
82}
83
84impl 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.
97pub 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)]
113pub 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.
135pub 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.
152pub 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.
165const DUE_SQL: &str = "
166SELECT 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
173FROM repos r LEFT JOIN repo_backups b ON b.repo_id = r.id
174WHERE 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))))
180ORDER BY coalesce(b.backed_up_ms, 0), r.created_at, r.id
181LIMIT ?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)]
191pub 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)]
211pub 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
223pub const MANIFEST_VERSION: u32 = 1;
224
225impl 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)]
270pub 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.
281pub 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.
302pub 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
308pub fn manifest_key(repo_id: &str) -> String {
309 format!("backups/{repo_id}/manifest.json")
310}
311
312pub fn bundle_key(repo_id: &str, id: &str, kind: BackupKind) -> String {
313 format!("backups/{repo_id}/{id}-{}.bundle", kind.suffix())
314}
315
316fn 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.
322pub 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.
336pub 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)]
347struct 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
361fn refused<T>(message: &str) -> Outcome<T> {
362 Outcome::fail(FailureCode::NotFound, message)
363}
364
365fn n(value: u64) -> JsValue {
366 JsValue::from_f64(value as f64)
367}
368
369fn 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)]
375pub 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.
382pub 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)]
413struct 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
426impl 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.
450async 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.
483pub 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.
567async 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
580const 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.
585pub 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 = &registry.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.
660pub 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.
676pub async fn complete(registry: &Registry, blobs: &Store, a: &BackupComplete, now: u64) -> Result<Outcome<bool>> {
677 let db = &registry.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.
750pub 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.
764fn 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
771async 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.
781async 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)]
798mod 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}