g1t/services/repos/src/backups.rs

998 lines40,863 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Merge branch 'worktree-agent-ac5b181a013e54348'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);
Repositories shard across git store namespaces, move between them, and can keep to the EU; a namespace can be served read-only from the self-hosted git store, rebuilt from the nightly backups (#20, #25)604 // Served from the fallback store (fallback.rs), which holds what the
605 // backups hold: backing that up would only record an older state.
606 if store.on_fallback(&key) {
607 settle_failure(db, &row, "The git store is on its fallback; backups wait until it is back.").await?;
608 return Ok(refused("The git store is on its fallback; backups wait until it is back."));
609 }
Merge branch 'worktree-agent-ac5b181a013e54348'610 let access = store.handout(&key, Scope::Read).await?;
611 // Asked before: the same bundle, so parts already sent still fit.
612 if let (Some(kind), Some(_)) = (row.upload_kind.as_deref(), row.upload_id.as_deref()) {
613 let manifest = blobs.read(&manifest_key(&row.repo_id)).await?.as_deref().and_then(Manifest::from_bytes);
614 let prerequisites: Vec<String> = row.prerequisites.as_deref().and_then(|p| serde_json::from_str(p).ok()).unwrap_or_default();
615 let previous_refs = manifest.as_ref().and_then(|m| m.last()).map(|e| e.refs.clone()).unwrap_or_default();
616 return Ok(Outcome::Ok(BackupSpec {
617 kind: if kind == "full" { BackupKind::Full } else { BackupKind::Incremental },
618 remote: access.remote,
619 git_token: access.token,
620 prerequisites,
621 previous_refs,
622 part_bytes: PART_BYTES,
623 }));
624 }
625 let version = db
626 .prepare("SELECT refs_version FROM repos WHERE id = ?1")
627 .bind(&[row.repo_id.as_str().into()])?
628 .first::<Live>(None)
629 .await?
630 .and_then(|live| live.refs_version)
631 .unwrap_or(0.0) as u64;
632 let manifest = blobs.read(&manifest_key(&row.repo_id)).await?.as_deref().and_then(Manifest::from_bytes);
633 let next = plan(manifest.as_ref(), row.last_entry.as_deref(), full_every);
634 let entry = stamp(now);
635 let upload_key = bundle_key(&row.repo_id, &entry, next.kind);
636 let upload_id = blobs.create_multipart(&upload_key).await?;
637 db.prepare(
638 "UPDATE repo_backups SET target_version = ?2, store_key = ?3, upload_key = ?4, upload_id = ?5,
639 upload_kind = ?6, upload_entry = ?7, prerequisites = ?8, backed_from_ms = ?9
640 WHERE job_id = ?1",
641 )
642 .bind(&[
643 a.job_id.as_str().into(),
644 n(version),
645 key.as_str().into(),
646 upload_key.as_str().into(),
647 upload_id.as_str().into(),
648 next.kind.suffix().into(),
649 entry.as_str().into(),
650 serde_json::to_string(&next.prerequisites)?.into(),
651 n(now),
652 ])?
653 .run()
654 .await?;
655 Ok(Outcome::Ok(BackupSpec {
656 kind: next.kind,
657 remote: access.remote,
658 git_token: access.token,
659 prerequisites: next.prerequisites,
660 previous_refs: next.previous_refs,
661 part_bytes: PART_BYTES,
662 }))
663}
664
665/// One part of the job's bundle, kept.
666pub async fn part(db: &D1Database, blobs: &Store, a: &BackupJobArgs, number: u16, bytes: Vec<u8>) -> Result<Outcome<BackupPart>> {
667 let Some(row) = job(db, a).await? else { return Ok(refused(NO_JOB)) };
668 let (Some(key), Some(upload)) = (row.upload_key.as_deref(), row.upload_id.as_deref()) else {
669 return Ok(Outcome::fail(FailureCode::Conflict, "Ask for the job's spec first."));
670 };
671 if number == 0 || bytes.len() as u64 > PART_BYTES {
672 return Ok(Outcome::fail(FailureCode::Invalid, "Parts are numbered from 1, and hold 32 MiB at most."));
673 }
674 let Part { number, etag } = blobs.upload_part(key, upload, number, bytes).await?;
675 Ok(Outcome::Ok(BackupPart { number, etag }))
676}
677
678/// The bundle is cut and sent: the upload is completed, the chain gains
679/// an entry (unless nothing changed), bundles of the chain before last
680/// are removed, and the repository is backed up as of the refs version it
681/// had when the clone began.
682pub async fn complete(registry: &Registry, blobs: &Store, a: &BackupComplete, now: u64) -> Result<Outcome<bool>> {
683 let db = &registry.db;
684 let args = BackupJobArgs { job_id: a.job_id.clone(), token: a.token.clone() };
685 let Some(row) = job(db, &args).await? else { return Ok(refused(NO_JOB)) };
686 let (Some(upload_key), Some(upload_id), Some(entry_id), Some(store_key)) =
687 (row.upload_key.clone(), row.upload_id.clone(), row.upload_entry.clone(), row.store_key.clone())
688 else {
689 return Ok(Outcome::fail(FailureCode::Conflict, "Ask for the job's spec first."));
690 };
691 if !valid_refs(&a.refs) {
692 return Ok(Outcome::fail(FailureCode::Invalid, "A ref name or commit is not one git makes."));
693 }
694 if !parts_fit(&a.parts, a.size, PART_BYTES) {
695 return Ok(Outcome::fail(FailureCode::Invalid, "The parts do not add up to the bundle's size."));
696 }
697 meter_fetch(&store_key, a.fetched_bytes);
698 let kind = if row.upload_kind.as_deref() == Some("full") { BackupKind::Full } else { BackupKind::Incremental };
699 let mut manifest = blobs
700 .read(&manifest_key(&row.repo_id))
701 .await?
702 .as_deref()
703 .and_then(Manifest::from_bytes)
704 .unwrap_or_else(|| Manifest::new(&row.repo_id, &store_key));
705 let unchanged = a.size == 0 && kind == BackupKind::Incremental && manifest.last().is_some_and(|last| last.refs == a.refs);
706 let version = row.target_version.unwrap_or(0.0) as u64;
707 let mut last_entry = row.last_entry.clone();
708 if a.size == 0 {
709 blobs.abort_multipart(&upload_key, &upload_id).await?;
710 } else {
711 let parts: Vec<Part> = a.parts.iter().map(|p| Part { number: p.number, etag: p.etag.clone() }).collect();
712 blobs.complete_multipart(&upload_key, &upload_id, &parts).await?;
713 }
714 if !unchanged {
715 let prerequisites: Vec<String> = row.prerequisites.as_deref().and_then(|p| serde_json::from_str(p).ok()).unwrap_or_default();
716 let dropped = manifest.add(Entry {
717 id: entry_id.clone(),
718 kind,
719 key: (a.size > 0).then(|| upload_key.clone()),
720 created_at: rfc3339(now),
721 refs_version: version,
722 refs: a.refs.clone(),
723 prerequisites,
724 size: a.size,
725 sha256: a.sha256.clone(),
726 });
727 manifest.store_key = store_key.clone();
728 manifest.updated_at = rfc3339(now);
729 manifest.path = registry
730 .by_id(&row.repo_id)
731 .await?
732 .map(|repo| RepoPath { namespace: repo.namespace, name: repo.name });
733 blobs.put(&manifest_key(&row.repo_id), manifest.to_bytes()).await?;
734 for key in dropped {
735 if let Err(error) = blobs.delete(&key).await {
736 worker::console_error!("repos: backup {key} not removed: {error}");
737 }
738 }
739 last_entry = Some(entry_id);
740 }
741 db.prepare(
742 "UPDATE repo_backups SET status = 'idle', job_id = NULL, token_hash = NULL, claimed_ms = NULL,
743 attempts = 0, last_error = NULL, refs_version = ?2, backed_up_ms = backed_from_ms,
744 tips = ?3, last_entry = ?4, upload_key = NULL, upload_id = NULL, upload_kind = NULL,
745 upload_entry = NULL, prerequisites = NULL, target_version = NULL
746 WHERE job_id = ?1",
747 )
748 .bind(&[a.job_id.as_str().into(), n(version), serde_json::to_string(&a.refs)?.into(), text(last_entry.as_deref())])?
749 .run()
750 .await?;
751 Ok(Outcome::Ok(true))
752}
753
754/// The sandbox could not do the job: tried again later tonight, up to
755/// [`MAX_ATTEMPTS`] times, and tomorrow night after that.
756pub async fn fail(db: &D1Database, blobs: &Store, a: &BackupFail) -> Result<Outcome<bool>> {
757 let args = BackupJobArgs { job_id: a.job_id.clone(), token: a.token.clone() };
758 let Some(row) = job(db, &args).await? else { return Ok(refused(NO_JOB)) };
759 if let Some(key) = row.store_key.as_deref() {
760 meter_fetch(key, a.fetched_bytes);
761 }
762 give_up_upload(Some(blobs), &row).await;
763 let said: String = a.error.chars().take(500).collect();
764 settle_failure(db, &row, &said).await?;
765 Ok(Outcome::Ok(true))
766}
767
768/// A clone counts once, with what it read, whenever it got as far as
769/// reading anything.
770fn meter_fetch(store_key: &str, fetched_bytes: u64) {
771 if fetched_bytes > 0 {
772 meters::record("internal.git.info_refs", store_key, 0, 0);
773 meters::record(FETCH_METER, store_key, 0, fetched_bytes);
774 }
775}
776
777async fn give_up_upload(blobs: Option<&Store>, row: &JobRow) {
778 if let (Some(blobs), Some(key), Some(upload)) = (blobs, row.upload_key.as_deref(), row.upload_id.as_deref())
779 && let Err(error) = blobs.abort_multipart(key, upload).await
780 {
781 worker::console_error!("repos: backup upload {key} not given up: {error}");
782 }
783}
784
785/// Back in the queue for another try, or idle until tomorrow night once
786/// it has been tried enough.
787async fn settle_failure(db: &D1Database, row: &JobRow, error: &str) -> Result<()> {
788 let attempts = row.attempts.unwrap_or(0.0) as u32;
789 let next = if attempts >= MAX_ATTEMPTS { "idle" } else { "queued" };
790 worker::console_error!("repos: backup of {} failed ({attempts} of {MAX_ATTEMPTS}): {error}", row.repo_id);
791 db.prepare(
792 "UPDATE repo_backups SET status = ?2, last_error = ?3, job_id = NULL, token_hash = NULL,
793 claimed_ms = NULL, upload_key = NULL, upload_id = NULL, upload_kind = NULL,
794 upload_entry = NULL, prerequisites = NULL, target_version = NULL
795 WHERE repo_id = ?1",
796 )
797 .bind(&[row.repo_id.as_str().into(), next.into(), error.into()])?
798 .run()
799 .await?;
800 Ok(())
801}
802
803#[cfg(test)]
804mod tests {
805 use super::*;
806
807 const A: &str = "c71546fcd893ef8b0f57388b65e620d759705dda";
808 const B: &str = "4807077b296e6edbf410d55e72749d3e1170c291";
809 const C: &str = "0000000000000000000000000000000000000abc";
810
811 fn candidate(id: &str) -> Candidate {
812 Candidate {
813 repo_id: id.into(),
814 namespace: "acme".into(),
815 created_at: "2026-01-01T00:00:00.000Z".into(),
816 refs_version: 3,
817 ..Candidate::default()
818 }
819 }
820
821 fn backed(id: &str, version: u64, at: u64) -> Candidate {
822 Candidate { status: Some("idle".into()), backed_version: Some(version), backed_up_ms: at, ..candidate(id) }
823 }
824
825 #[test]
826 fn a_repository_never_backed_up_is_due() {
827 assert!(is_due(&candidate("r1")));
828 // Even one whose refs never moved: rows from before refs_version.
829 assert!(is_due(&Candidate { refs_version: 0, ..candidate("r1") }));
830 }
831
832 #[test]
833 fn a_repository_is_due_only_once_its_refs_moved() {
834 assert!(!is_due(&backed("r1", 3, 1_000)));
835 assert!(is_due(&backed("r1", 2, 1_000)));
836 // A credential that could push went out after the last backup began.
837 assert!(is_due(&Candidate { refs_open_until: 2_000, ..backed("r1", 3, 1_000) }));
838 assert!(!is_due(&Candidate { refs_open_until: 900, ..backed("r1", 3, 1_000) }));
839 // Tried and failed: no version yet.
840 assert!(is_due(&Candidate { backed_version: None, ..backed("r1", 0, 0) }));
841 }
842
843 #[test]
844 fn deleted_retired_working_copies_and_jobs_in_hand_are_not_due() {
845 assert!(!is_due(&Candidate { deleted: true, ..candidate("r1") }));
846 assert!(!is_due(&Candidate { retired: true, ..candidate("r1") }));
847 assert!(!is_due(&Candidate { namespace: PULLS_NAMESPACE.into(), ..candidate("r1") }));
848 assert!(!is_due(&Candidate { status: Some("queued".into()), ..backed("r1", 1, 0) }));
849 assert!(!is_due(&Candidate { status: Some("running".into()), ..backed("r1", 1, 0) }));
850 }
851
852 #[test]
853 fn the_longest_waiting_go_first_and_no_more_than_the_limit() {
854 let candidates = vec![
855 backed("recent", 1, 5_000),
856 backed("older", 1, 1_000),
857 backed("current", 3, 0),
858 Candidate { created_at: "2026-02-01T00:00:00.000Z".into(), ..candidate("new-b") },
859 Candidate { created_at: "2025-02-01T00:00:00.000Z".into(), ..candidate("new-a") },
860 ];
861 assert_eq!(pick_due(&candidates, 10), ["new-a", "new-b", "older", "recent"]);
862 assert_eq!(pick_due(&candidates, 2), ["new-a", "new-b"]);
863 assert!(pick_due(&candidates, 0).is_empty());
864 }
865
866 fn refs(pairs: &[(&str, &str)]) -> BTreeMap<String, String> {
867 pairs.iter().map(|(name, hash)| ((*name).to_owned(), (*hash).to_owned())).collect()
868 }
869
870 fn entry(id: &str, kind: BackupKind, tips: &[(&str, &str)]) -> Entry {
871 Entry {
872 id: id.into(),
873 kind,
874 key: Some(format!("backups/r1/{id}-{}.bundle", kind.suffix())),
875 created_at: String::new(),
876 refs_version: 1,
877 refs: refs(tips),
878 prerequisites: Vec::new(),
879 size: 10,
880 sha256: None,
881 }
882 }
883
884 fn manifest(incrementals: usize) -> Manifest {
885 let mut manifest = Manifest::new("r1", "acme--rocket");
886 manifest.add(entry("e0", BackupKind::Full, &[("HEAD", A), ("refs/heads/main", A)]));
887 for index in 0..incrementals {
888 manifest.add(entry(&format!("e{}", index + 1), BackupKind::Incremental, &[("HEAD", A), ("refs/heads/main", A), ("refs/tags/v1", B)]));
889 }
890 manifest
891 }
892
893 #[test]
894 fn the_first_backup_is_full() {
895 let next = plan(None, None, 30);
896 assert_eq!(next.kind, BackupKind::Full);
897 assert!(next.prerequisites.is_empty() && next.previous_refs.is_empty());
898 assert_eq!(plan(Some(&Manifest::new("r1", "k")), None, 30).kind, BackupKind::Full);
899 }
900
901 #[test]
902 fn the_next_is_incremental_from_the_last_tips_each_once() {
903 let manifest = manifest(1);
904 let next = plan(Some(&manifest), Some("e1"), 30);
905 assert_eq!(next.kind, BackupKind::Incremental);
906 assert_eq!(next.prerequisites, [B, A]);
907 assert_eq!(next.previous_refs, manifest.last().unwrap().refs);
908 }
909
910 #[test]
911 fn a_full_bundle_is_cut_again_after_enough_incrementals() {
912 assert_eq!(plan(Some(&manifest(29)), Some("e29"), 30).kind, BackupKind::Incremental);
913 let next = plan(Some(&manifest(30)), Some("e30"), 30);
914 assert_eq!(next.kind, BackupKind::Full);
915 assert!(next.prerequisites.is_empty());
916 // What it last held is still said, so an unchanged clone is seen.
917 assert!(!next.previous_refs.is_empty());
918 assert_eq!(plan(Some(&manifest(1)), Some("e1"), 1).kind, BackupKind::Full);
919 }
920
921 #[test]
922 fn a_chain_the_row_does_not_know_or_an_empty_one_starts_again() {
923 assert_eq!(plan(Some(&manifest(2)), Some("e1"), 30).kind, BackupKind::Full);
924 assert_eq!(plan(Some(&manifest(2)), None, 30).kind, BackupKind::Full);
925 let mut empty = Manifest::new("r1", "k");
926 empty.add(entry("e0", BackupKind::Full, &[]));
927 assert_eq!(plan(Some(&empty), Some("e0"), 30).kind, BackupKind::Full);
928 }
929
930 #[test]
931 fn a_full_backup_starts_a_chain_and_the_one_before_last_goes() {
932 let mut manifest = manifest(2);
933 assert!(manifest.previous.is_empty());
934 let dropped = manifest.add(entry("f1", BackupKind::Full, &[("refs/heads/main", C)]));
935 assert!(dropped.is_empty(), "the chain before is kept until the next full one");
936 assert_eq!(manifest.chain.len(), 1);
937 assert_eq!(manifest.previous.len(), 3);
938 manifest.add(entry("f1a", BackupKind::Incremental, &[("refs/heads/main", C)]));
939 let dropped = manifest.add(entry("f2", BackupKind::Full, &[("refs/heads/main", C)]));
940 assert_eq!(dropped, ["backups/r1/e0-full.bundle", "backups/r1/e1-incr.bundle", "backups/r1/e2-incr.bundle"]);
941 assert_eq!(manifest.previous.iter().map(|e| e.id.as_str()).collect::<Vec<_>>(), ["f1", "f1a"]);
942 assert_eq!(manifest.keys().len(), 3);
943 }
944
945 #[test]
946 fn a_manifest_round_trips_and_an_unknown_version_is_not_read() {
947 let mut manifest = manifest(2);
948 manifest.path = Some(RepoPath { namespace: "acme".into(), name: "rocket".into() });
949 manifest.add(Entry { key: None, size: 0, ..entry("e3", BackupKind::Incremental, &[("refs/heads/main", A)]) });
950 let bytes = manifest.to_bytes();
951 let text = String::from_utf8(bytes.clone()).unwrap();
952 // What the restore drill reads: snake_case, the kind in words.
953 assert!(text.contains("\"repo_id\": \"r1\"") && text.contains("\"kind\": \"incremental\"") && text.contains("\"key\": null"));
954 assert_eq!(Manifest::from_bytes(&bytes), Some(manifest.clone()));
955 let mut later = manifest;
956 later.version = 2;
957 assert_eq!(Manifest::from_bytes(&later.to_bytes()), None);
958 assert_eq!(Manifest::from_bytes(b"not json"), None);
959 }
960
961 #[test]
962 fn a_backup_clone_is_g1ts_operation_not_the_workspaces() {
963 let mapping = meters::Mapping::defaults();
964 assert_eq!(mapping.cost(FETCH_METER), 1.0);
965 assert_eq!(mapping.billable(FETCH_METER), 0.0);
966 assert_eq!(mapping.billable("internal.git.fetch"), 1.0);
967 }
968
969 #[test]
970 fn keys_and_stamps() {
971 assert_eq!(stamp(1_369_353_600_123), "20130524T000000Z");
972 assert_eq!(manifest_key("r1"), "backups/r1/manifest.json");
973 assert_eq!(bundle_key("r1", "20130524T000000Z", BackupKind::Incremental), "backups/r1/20130524T000000Z-incr.bundle");
974 assert_eq!(bundle_key("r1", "20130524T000000Z", BackupKind::Full), "backups/r1/20130524T000000Z-full.bundle");
975 }
976
977 #[test]
978 fn only_refs_git_makes_are_taken() {
979 assert!(valid_refs(&refs(&[("HEAD", A), ("refs/heads/main", A), ("refs/pull/pr_1/head", B)])));
980 assert!(valid_refs(&BTreeMap::new()));
981 assert!(!valid_refs(&refs(&[("main", A)])));
982 assert!(!valid_refs(&refs(&[("refs/heads/../x", A)])));
983 assert!(!valid_refs(&refs(&[("refs/heads/a b", A)])));
984 assert!(!valid_refs(&refs(&[("refs/heads/main", "abc")])));
985 assert!(!valid_refs(&refs(&[("refs/heads/main", &A.to_uppercase())])));
986 }
987
988 #[test]
989 fn parts_must_cover_the_bundle_in_order() {
990 let part = |number| BackupPart { number, etag: "e".into() };
991 assert!(parts_fit(&[], 0, 10));
992 assert!(parts_fit(&[part(1)], 10, 10));
993 assert!(parts_fit(&[part(1), part(2)], 11, 10));
994 assert!(!parts_fit(&[part(1)], 11, 10));
995 assert!(!parts_fit(&[part(2), part(1)], 11, 10));
996 assert!(!parts_fit(&[part(1)], 0, 10));
997 }
998}

This file's history is long; its oldest lines are credited to the oldest commit read.