diff --git a/changelog.md b/changelog.md index de40648..f1665f9 100644 --- a/changelog.md +++ b/changelog.md @@ -6,11 +6,13 @@ _2026-08-03_ - new: `eat-rocks push ` (lib: `push()`) uploads a rocksdb checkpoint directory as a BackupEngine-compatible backup. incremental: unchanged SSTs already in the bucket are recognized via rocks' session-id naming and their recorded checksums reused — no re-read, no re-upload. works against buckets whose earlier backups were written by rocks' own BackupEngine, continuing the same id lineage and dedup. -- new: retention — `purge --keep N` (keep the N newest backups, rocks `PurgeOldBackups` semantics), `delete `, and `gc` (reclaim debris of crashed pushes and interrupted deletions). all support `--dry-run`. `purge` deliberately refuses `--keep 0`. +- new: retention — `retain N` (keep the N newest backups, rocks `PurgeOldBackups` semantics), `delete `, and `cleanup` (reclaim debris of crashed pushes and interrupted deletions). all support `--dry-run`. `retain 0` is deliberately unrepresentable. -- new: `push --retain N` runs the purge right after a successful push — the one-line cron shape. +- new: `push --retain N` runs the retention pass right after a successful push — the one-line cron shape. -- pushed metas default to `schema_version 2.1` with `size` fields (readable by rocksdb ≥ 6.19); `--schema v1` writes byte-identical-to-rocks v1 metas. +- new: `push --consume-checkpoint` deletes each checkpoint file as soon as it lands and removes the directory at the end. checkpoint entries are hard links, so releasing them progressively keeps peak disk usage down rather than pinning the whole checkpoint until the push finishes. + +- pushed metas are byte-identical to rocks' own (schema v1, no sizes) by default; `--schema v2` opts into `schema_version 2.1` with `size` fields so restores verify them (readable by rocksdb ≥ 6.19). - crate description grew: it's backup *and* restore now. the restore API is unchanged. diff --git a/hacking.md b/hacking.md index 73c371d..cf956c9 100644 --- a/hacking.md +++ b/hacking.md @@ -58,77 +58,70 @@ rocks itself uses a file called `CURRENT` as its entrypoint to the db. when restoring, we write all other files first, then atomically rename the new `CURRENT` into place, so a partial restore won't corrupt things. (just following rocks here) -## pushing backups - -the input is a [checkpoint](https://github.com/facebook/rocksdb/wiki/Checkpoints) -directory. important property (verified in rocks source, -`checkpoint_impl.cc` + `db_filesnapshot.cc`): once `CreateCheckpoint` -returns, every file in it is byte-exact — rocks already truncated the -MANIFEST and any open WAL to their valid lengths while checkpointing, and the -dir itself appears via atomic rename. so "read whole file, upload" is correct. - -classification is strict: `CURRENT`, one `MANIFEST-*`, `OPTIONS-*`, `*.log` -go to `private//`; `*.sst` and `*.blob` go to `shared_checksum/`; -*anything* else is a hard error. that strictness is a safety feature: a live -db dir always contains `IDENTITY`/`LOCK`/`LOG` (a checkpoint never does), so -pointing push at a live database — whose MANIFEST may be mid-write — refuses -instead of uploading garbage. - -### commit protocol - -object PUTs are atomic, which simplifies rocks' tmp-file dance a lot: - -1. allocate id = max(committed metas, `.{id}.tmp` markers, `private/` debris) + 1. - counting debris means a crashed push can never share an id with a new one. -2. PUT `meta/.{id}.tmp` with if-none-match (claims the id; both our list and - rocks' scan ignore the name). unlike rocks — which can only write its tmp - meta after copying everything — we know the full file list up front, so - the marker doubles as an advertisement that protects in-flight uploads - from a concurrent gc. -3. upload files. shared files directly to their final names (no tmp+rename: - a partially-uploaded object is simply not visible). -4. PUT `meta/{id}` (if-none-match). **this is the commit point.** -5. DELETE the marker (best-effort; a stale marker is invisible and gc reaps - it once it's old). - -crash anywhere before step 4 leaves only invisible debris for gc. - -### shared file naming, and why we read sst table properties - -naming shared files the legacy way (`__`) requires -checksumming every sst just to learn its name — a full-db read per push. -rocks' fix (and ours): name by db session id -(`_s_`, read from a few KB of table properties; -`src/sst.rs` is a tiny reader for exactly that), and for files that already -exist remotely *and* are referenced by a committed meta, reuse the checksum -recorded there instead of re-reading (`backup_engine.cc:2676`: "to save I/O -on incremental backups, we copy prior known checksum"). blob files and -pre-session-id ssts fall back to legacy naming (full read), same as rocks. - -any `sst.rs` parse failure just means legacy naming for that file — it is -never load-bearing for correctness. - -the honest tradeoff (rocks documents it at `backup_engine.cc:2682-2695`): -reused checksums won't notice an object that corrupted *in the bucket* since -it was pushed. restore always verifies crc32c on download, so it's caught at -the moment it matters. - -### retention + gc - -`purge --keep N` = rocks' `PurgeOldBackups`: the N highest ids survive. -deletion order matches rocks too: the meta DELETE commits the deletion -(abort if it fails); files are reclaimed by gc, which recomputes rocks' -in-memory refcounts from every committed meta in the bucket. two hard rules: - -- **an unparseable committed meta aborts gc.** we can't know what a backup - we can't read references; deleting anything would be guessing. -- gc treats fresh `.{id}.tmp` markers as references (in-flight push - protection). markers older than `--stale-after` (default 7d) are dead - pushes: reaped along with their debris. - -single active writer per prefix is the operating model (rocks requires the -same). `push --retain N` runs push-then-purge in one process, which is the -no-races cron shape. +## do pushups + +eat-rocks doesn't depend on rocksdb directly, so the app needs to set things up a bit for it to push a backup for a live database. we want rocks' `BackupEngine` behaviour, but we we can't use `BackupEngine` directly because it can't read existing contents in object storage to initialize its state, and because it forces full file copies of the backup contents. + +so we drop a layer down, to rocksdb +[checkpoint](https://github.com/facebook/rocksdb/wiki/Checkpoints), which is what `BackupEngine` uses internally. checkpoint writes a consistent set of SSTs and logs to a directory, and if the target directory is on the same filesystem as the db, it can just create hard-links instead of copying file contents. fast and io-cheap. + +eat-rocks takes a checkpoint target directory and reimplements a `BackupEngine`-compatible writer on top, which pushes directly to object storage instead of a filesystem path. + +push tries to make sure it's only working in a checkpoint folder: `CURRENT`, one `MANIFEST-*`, `OPTIONS-*`, and `*.log` files to `private//`; `*.sst` and `*.blob` to `shared_checksum/`. if any other files are found (like `IDENTITY`/`LOCK`/`LOG` from a rocks db folder), we bail hard. + + +### commitment issues + +object PUTs are atomic, which makes things simpler for us than rocks: + +1. id = max(committed metas, `.{id}.tmp` markers, `private/` debris) + 1. + +2. PUT `meta/.{id}.tmp` with if-none-match, which claims the id, and writes the + list of files we'll be referencing, protecting any existing from cleanup. + +3. upload files directly to their final names (in-progress uploads are not visible in object storage) + +4. PUT `meta/{id}`, with if-none-match again to guard against races between id allocation and a very-fast-completing other backup. this finalizes the backup push. + +5. DELETE the `.tmp` marker. failure here is non-fatal, it'll get cleaned up later once stale. + +any crash before step 4 can leave useless objects behind, but that's generally fine since cleanup will catch them. + + +### naming things + +ssts get named like `_s_.sst`, which can be deterministically generated with mimimal content reading from the file itself. + +eat-rocks also supports the legacy naming format, `__`, which require reading the entire file to get the crc32c first to generate. +which is not ideal so is avoided when possible! + +note: eat-rocks' restore process always verifies crc32c on download. + + +### eats rocks and ~~leaves~~ checkpoints + +the checkpoint consumes minimal extra space initially, but prevents SSTs from being actually deleted by the filesystem while we're hard-linked, if rocks did a compaction and released them. + +`--consume-checkpoint` unlinks files as soon as it doesn't need them (already-existing files first), so that no more files stay pinned, occupying fs space, any longer than actually necessary. +this reduces the peak filesystem margin needed to keep the checkpoint valid during backup. + + +### incremental retention and cleanup + +like rocks `BackupEngine`'s `PurgeOldBackups`, we can keep contents for N-most-recent backups alive and drop everything older and unneeded. + +rocks keeps in-memory refcounts of files; we currently don't and re-read stuff from the bucket each time we try to clean up. + +a few things to be aware of + +- a broken meta prevents cleanup, since we can't know what it might have been referencing, which shouldn't be cleaned up. + +- at the start of a push we write a `.{id}.tmp` marker with references to the files that exist (or are about to exist) for this backup -- cleanup does read these and won't remove those. + +- marker files themselves are normally removed after a backup fully succeeds. if they are left behind for any reason, the do also get cleaned up eventually, see `--stale-after` (default 7d). + +you typically should avoid running multiple active pushers at the same time if it can be avoided. + ## with integrity diff --git a/readme.md b/readme.md index ef1b942..05c6c06 100644 --- a/readme.md +++ b/readme.md @@ -51,10 +51,17 @@ eat-rocks --endpoint https://constellation.t3.storage.dev \ eat-rocks --endpoint https://constellation.t3.storage.dev \ push --retain 5 /data/checkpoints/2026-08-02 +# delete each checkpoint file as it's pushed, then remove the directory. +# checkpoint files are hard links, so this keeps peak disk usage down -- +# handy when free space is tight. (destroys the checkpoint either way: if +# the push fails partway, retry from a fresh one.) +eat-rocks --endpoint https://constellation.t3.storage.dev \ + push --consume-checkpoint /data/checkpoints/2026-08-02 + # manage retention -eat-rocks --endpoint https://constellation.t3.storage.dev purge --keep 5 +eat-rocks --endpoint https://constellation.t3.storage.dev retain 5 eat-rocks --endpoint https://constellation.t3.storage.dev delete 3 -eat-rocks --endpoint https://constellation.t3.storage.dev gc --dry-run +eat-rocks --endpoint https://constellation.t3.storage.dev cleanup --dry-run ``` ## lib @@ -79,9 +86,9 @@ to push a backup, you first have to ask rocks to create a checkpoint. the checkp rocksdb::checkpoint::Checkpoint::new(&db)?.create_checkpoint(&cp_dir)?; eat_rocks::push(store, "", &cp_dir, eat_rocks::PushOptions { retain: std::num::NonZeroUsize::new(5), + consume_checkpoint: true, ..Default::default() }).await?; -std::fs::remove_dir_all(&cp_dir)?; // TODO: clean up checkpoint as the push proceeds ``` @@ -99,8 +106,6 @@ std::fs::remove_dir_all(&cp_dir)?; // TODO: clean up checkpoint as the push proc - cleanup of orphaned file: when replacing an existing rocksdb database, `eat-rocks` will delete old no-longer-used rocksdb-like files from the target directory. rocks' own restore function leaves them in place and they never get cleaned up. -- pushed backup metas default to `schema_version 2.1` with `size` fields, so restores verify sizes (rocks still writes schema v1 without sizes). pass `--schema v1` to match rocks' behaviour. - TODO: keep an instance in memory to avoid reading all metas from bucket on every push (like rocks) diff --git a/src/lib.rs b/src/lib.rs index 75d01a3..2a398df 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,8 +1,7 @@ -//! Backup and restore a rocks database with object storage, in the -//! [rocks backup](https://github.com/facebook/rocksdb/wiki/How-to-backup-RocksDB) -//! format — without depending on rocksdb. +//! Rapidly restore or push a [rocks backup](https://github.com/facebook/rocksdb/wiki/How-to-backup-RocksDB) +//! from/to object storage //! -//! restoring: +//! ## restoring //! //! ```rust,no_run //! use eat_rocks::{public_bucket, restore}; @@ -13,37 +12,24 @@ //! # } //! ``` //! -//! pushing: have the app that owns the db make a -//! [checkpoint](https://github.com/facebook/rocksdb/wiki/Checkpoints) -//! (hard links — cheap), then hand the directory to [`push`]: +//! the `public_bucket` function (`easy` feature, enabled by default) works with +//! s3-compatible stores (like tigris) which are configured for public read access. +//! +//! use any [`ObjectStore`](object_store::ObjectStore) to talk to pretty much any backend (S3, GCS, Azure, local filesystem, ...) +//! +//! ## pushing +//! +//! make a rocksdb [checkpoint](https://github.com/facebook/rocksdb/wiki/Checkpoints) +//! directory, then [`push`]: //! //! ```rust,no_run //! # async fn example(store: std::sync::Arc) -> Result<(), eat_rocks::push::PushError> { -//! use std::num::NonZeroUsize; -//! let outcome = eat_rocks::push( -//! store, -//! "", -//! "/data/checkpoints/2026-08-02", -//! eat_rocks::PushOptions { -//! retain: NonZeroUsize::new(5), // keep the newest 5 backups -//! ..Default::default() -//! }, -//! ) -//! .await?; +//! use eat_rocks::push; +//! // let store: std::sync::Arc = ...?; +//! push(store, "", "/data/checkpoints/2026-08-02", Default::default()).await?; //! # Ok(()) //! # } //! ``` -//! -//! pushes are incremental: unchanged SSTs already present in the bucket are -//! recognized (by rocks' own session-id naming) without re-reading them. -//! the resulting backups are interchangeable with rocks' `BackupEngine` — -//! either implementation can restore what the other wrote, into or out of -//! the same bucket. -//! -//! the `public_bucket` function (`easy` feature, enabled by default) works with -//! s3-compatible stores (like tigris) which are configured for public read access. -//! -//! you can use any [`ObjectStore`](object_store::ObjectStore) to talk to pretty much any backend (S3, GCS, Azure, local filesystem, ...) pub mod meta; pub mod push; @@ -57,7 +43,7 @@ pub use restore::{ DEFAULT_CONCURRENCY, RestoreOptions, TargetMode, fetch_meta, list_backup_ids, restore, }; pub use retention::{ - GcOutcome, RetireOptions, RetireOutcome, delete_backup, garbage_collect, purge_old_backups, + CleanupOutcome, RetentionOptions, RetentionOutcome, cleanup, delete_backup, retain_backups, }; use std::io; diff --git a/src/main.rs b/src/main.rs index 4d45ada..5a92e73 100644 --- a/src/main.rs +++ b/src/main.rs @@ -6,7 +6,7 @@ use clap::{Parser, Subcommand, ValueEnum}; use object_store::{ObjectStore, aws::AmazonS3Builder}; use tracing::warn; -use eat_rocks::retention::{GcOutcome, RetireOutcome}; +use eat_rocks::retention::{CleanupOutcome, RetentionOutcome}; #[derive(Debug, thiserror::Error)] enum CliError { @@ -39,10 +39,10 @@ enum CliError { }, #[error("{operation} failed")] - Retire { + Retention { operation: &'static str, #[source] - source: eat_rocks::retention::RetireError, + source: eat_rocks::retention::RetentionError, }, } @@ -134,16 +134,21 @@ enum Command { }, /// push a rocksdb checkpoint directory as a new backup /// - /// make the checkpoint with rocks' Checkpoint::CreateCheckpoint (on the - /// same filesystem as the db: it hard-links). never point this at a live - /// database directory -- its MANIFEST/WAL may be mid-write. (a live dir - /// is detected and refused by its IDENTITY/LOCK/LOG files.) + /// requires a pre-made checkpoint via `Checkpoint::CreateCheckpoint`. this + /// should ideally be on the same filesystem as the database, so that rocks + /// can do hard links (no extra space, fast) instead of full file copies. Push { - /// retain only the newest N backups after a successful push - /// (same as running `purge --keep N` afterwards) + /// keep only the newest N backups after a successful push + /// (same as running `retain N` afterwards) #[arg(long)] retain: Option, + /// delete files once pushed, then remove the checkpoint directory + /// + /// lets the filesystem reclaim space when possible + #[arg(long)] + consume_checkpoint: bool, + /// informational sequence number to record in the backup meta #[arg(long, default_value_t = 0)] sequence_number: u64, @@ -154,83 +159,75 @@ enum Command { /// max concurrent actions against object storage /// - /// uploads buffer up to ~16MiB each, so peak memory is roughly - /// 16MiB x this + /// uploads buffer up to ~16MiB, so max peak memory roughly this * 16MiB #[arg(long, default_value_t = eat_rocks::DEFAULT_PUSH_CONCURRENCY)] concurrency: usize, /// backup meta schema to write - #[arg(long, value_enum, default_value_t = SchemaArg::V2)] + #[arg(long, value_enum, default_value_t = SchemaArg::V1)] schema: SchemaArg, /// shared file naming scheme #[arg(long, value_enum, default_value_t = NamingArg::Session)] naming: NamingArg, - /// path to the checkpoint directory to push + /// path to the checkpoint directory to push a backup from checkpoint_dir: PathBuf, }, - /// delete one backup (and garbage-collect its files) + /// delete one backup and clean up Delete { - /// print what would be deleted without deleting anything + /// print a summary of actions without deleting anything #[arg(long)] dry_run: bool, - /// backup to delete. use `list` to see available backups. + /// which backup to delete. `list` shows existing backup ids. backup_id: u64, }, - /// keep only the newest N backups, deleting the rest - Purge { - /// how many backups to keep (the highest-numbered ones survive). - /// deleting *every* backup is deliberately not expressible here -- - /// use `delete` per backup id for that. + /// keep only the newest N backups, deleting any older + Retain { + /// print a summary of actions without deleting anything #[arg(long)] - keep: NonZeroUsize, + dry_run: bool, - /// print what would be deleted without deleting anything + /// how many most-recent backups to keep + keep: NonZeroUsize, + }, + /// remove orphaned objects not referenced from any remaining backup + Cleanup { + /// print a summary of actions without deleting anything #[arg(long)] dry_run: bool, - }, - /// remove objects no backup references (debris of crashed pushes, - /// leftovers of interrupted deletions) - Gc { + /// treat in-progress push markers older than this as dead #[arg(long, default_value = "7d")] stale_after: humantime::Duration, - - /// print what would be deleted without deleting anything - #[arg(long)] - dry_run: bool, }, } /// see [`eat_rocks::MetaSchema`] #[derive(Clone, Debug, ValueEnum)] enum SchemaArg { - /// schema_version 2.1 with size fields: restores verify sizes; readable - /// by rocksdb >= 6.19 (2021) - V2, - /// byte-identical to what rocks itself writes (no sizes); readable by - /// any rocksdb version + /// what rocks itself currently writes (no sizes), most compatible V1, + /// schema_version 2.1 with size fields: readable by rocksdb >= 6.19 (2021) + V2, } /// see [`eat_rocks::SharedNaming`] #[derive(Clone, Debug, ValueEnum)] enum NamingArg { - /// name shared files by db session id (rocks' default) -- unchanged - /// files are recognized without re-reading them + /// name shared files by db session id (rocks' default), more efficient Session, - /// name shared files by checksum (rocks' legacy scheme) -- forces a - /// full read of every shared file on every push + /// name shared files by checksum (rocks' old way), requires full reads Legacy, } /// optional value for `--replace`. absent = `Create`-mode restore. #[derive(Clone, Debug, ValueEnum)] enum ReplaceMode { - /// (the default when `--replace` is given without a value) require an - /// existing rocks db in target. errors if the target doesn't look like a db. + /// require there to be a db to replace, errors if doesnt' look like one. + /// + /// default behaviour when `--replace` is passed without a value. Replace, /// proceed whether or not target contains a rocks db. Force, @@ -238,9 +235,9 @@ enum ReplaceMode { impl Cli { fn build_store(&self) -> Result, Box> { - // if `--bucket` is passed, then we get path-style buckets in the URL. - // if not, the caller is responsible for putting the bucket in the endpoint - // url, but we still need to set it, hence the `_` placeholder. + // if `--bucket` is passed: path-style buckets in the URL. + // if not, caller has to put the bucket in the endpoint url, but we + // still need to set it to a placeholder. let bucket = self.bucket.as_deref().unwrap_or("_"); let mut builder = AmazonS3Builder::new() @@ -254,7 +251,7 @@ impl Cli { .with_access_key_id(key_id) .with_secret_access_key(secret), (None, None) => builder.with_skip_signature(true), - _ => unreachable!("clap `requires` ensures both or neither are present"), + _ => unreachable!("clap `requires` ensures both or neither"), }; let store = builder.build().map_err(|source| CliError::StoreInit { @@ -283,8 +280,8 @@ async fn main() { Command::Restore { .. } => cmd_restore(&cli).await, Command::Push { .. } => cmd_push(&cli).await, Command::Delete { .. } => cmd_delete(&cli).await, - Command::Purge { .. } => cmd_purge(&cli).await, - Command::Gc { .. } => cmd_gc(&cli).await, + Command::Retain { .. } => cmd_retain(&cli).await, + Command::Cleanup { .. } => cmd_cleanup(&cli).await, }; if let Err(e) = result { @@ -378,6 +375,7 @@ async fn cmd_restore(cli: &Cli) -> Result<(), Box> { async fn cmd_push(cli: &Cli) -> Result<(), Box> { let Command::Push { retain, + consume_checkpoint, sequence_number, app_metadata, concurrency, @@ -399,6 +397,7 @@ async fn cmd_push(cli: &Cli) -> Result<(), Box> { sequence_number: *sequence_number, app_metadata: app_metadata.as_ref().map(|s| s.clone().into_bytes()), retain: *retain, + consume_checkpoint: *consume_checkpoint, schema: match schema { SchemaArg::V2 => eat_rocks::MetaSchema::V2WithSizes, SchemaArg::V1 => eat_rocks::MetaSchema::V1, @@ -420,40 +419,40 @@ async fn cmd_push(cli: &Cli) -> Result<(), Box> { "backup {} pushed: {} files uploaded ({:.1} MiB), {} reused ({:.1} MiB)", outcome.backup_id, outcome.uploaded_files, - outcome.uploaded_bytes as f64 / 2f64.powf(20.), + outcome.uploaded_bytes as f64 / 2_f64.powf(20.), outcome.reused_files, - outcome.reused_bytes as f64 / 2f64.powf(20.), + outcome.reused_bytes as f64 / 2_f64.powf(20.), ); - for id in &outcome.purged { - println!("purged backup {id}"); + for id in &outcome.retired { + println!("deleted backup {id}"); } Ok(()) } -fn print_gc(gc: &GcOutcome) { - let verb = if gc.dry_run { +fn print_cleanup(cleanup: &CleanupOutcome) { + let verb = if cleanup.dry_run { "would delete" } else { "deleted" }; - for path in &gc.deleted_shared { + for path in &cleanup.deleted_shared { println!("{verb} {path}"); } - for path in &gc.deleted_private { + for path in &cleanup.deleted_private { println!("{verb} {path}"); } - for id in &gc.deleted_markers { + for id in &cleanup.deleted_markers { println!("{verb} stale in-progress marker for backup {id}"); } - if gc.deleted_shared.is_empty() - && gc.deleted_private.is_empty() - && gc.deleted_markers.is_empty() + if cleanup.deleted_shared.is_empty() + && cleanup.deleted_private.is_empty() + && cleanup.deleted_markers.is_empty() { println!("no unreferenced objects"); } } -fn print_retire(outcome: &RetireOutcome) { +fn print_retention(outcome: &RetentionOutcome) { let verb = if outcome.dry_run { "would delete" } else { @@ -465,7 +464,7 @@ fn print_retire(outcome: &RetireOutcome) { if outcome.deleted_backups.is_empty() { println!("no backups to delete"); } - print_gc(&outcome.gc); + print_cleanup(&outcome.cleanup); } async fn cmd_delete(cli: &Cli) -> Result<(), Box> { @@ -477,45 +476,45 @@ async fn cmd_delete(cli: &Cli) -> Result<(), Box> { store, &cli.prefix, *backup_id, - &eat_rocks::RetireOptions { + &eat_rocks::RetentionOptions { dry_run: *dry_run, ..Default::default() }, ) .await - .map_err(|source| CliError::Retire { + .map_err(|source| CliError::Retention { operation: "delete", source, })?; - print_retire(&outcome); + print_retention(&outcome); Ok(()) } -async fn cmd_purge(cli: &Cli) -> Result<(), Box> { - let Command::Purge { keep, dry_run } = &cli.command else { +async fn cmd_retain(cli: &Cli) -> Result<(), Box> { + let Command::Retain { keep, dry_run } = &cli.command else { unreachable!() }; let store = cli.build_store()?; - let outcome = eat_rocks::purge_old_backups( + let outcome = eat_rocks::retain_backups( store, &cli.prefix, *keep, - &eat_rocks::RetireOptions { + &eat_rocks::RetentionOptions { dry_run: *dry_run, ..Default::default() }, ) .await - .map_err(|source| CliError::Retire { - operation: "purge", + .map_err(|source| CliError::Retention { + operation: "retain", source, })?; - print_retire(&outcome); + print_retention(&outcome); Ok(()) } -async fn cmd_gc(cli: &Cli) -> Result<(), Box> { - let Command::Gc { +async fn cmd_cleanup(cli: &Cli) -> Result<(), Box> { + let Command::Cleanup { stale_after, dry_run, } = &cli.command @@ -523,20 +522,20 @@ async fn cmd_gc(cli: &Cli) -> Result<(), Box> { unreachable!() }; let store = cli.build_store()?; - let outcome = eat_rocks::garbage_collect( + let outcome = eat_rocks::cleanup( store, &cli.prefix, - &eat_rocks::RetireOptions { + &eat_rocks::RetentionOptions { dry_run: *dry_run, stale_marker_after: (*stale_after).into(), ..Default::default() }, ) .await - .map_err(|source| CliError::Retire { - operation: "gc", + .map_err(|source| CliError::Retention { + operation: "cleanup", source, })?; - print_gc(&outcome); + print_cleanup(&outcome); Ok(()) } diff --git a/src/meta.rs b/src/meta.rs index a8d1989..4921df0 100644 --- a/src/meta.rs +++ b/src/meta.rs @@ -44,18 +44,18 @@ pub enum SerializeError { InvalidMetadata(String), } -/// Which meta file schema flavour to write. +/// Which meta file schema format to write. /// -/// Both are readable by our own [`BackupMeta::parse`]; rocks reads `V1` -/// everywhere and `V2WithSizes` in versions >= 6.19 (2021). +/// eat-rocks reads both; rocks >= 6.19 can read `V2WithSizes` (2021). #[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] pub enum MetaSchema { - /// exactly what rocks itself writes: no `schema_version` header, and only - /// a `crc32` field per file line - V1, - /// a `schema_version 2.1` header, plus `size` fields so restores can - /// verify file sizes (rocks accepts but never writes these itself) + /// old style, but still what rocksdb writes by default. + /// + /// meat has no `schema_version` header, only a `crc32` field per file line. #[default] + V1, + /// new style, with a `schema_version 2.1` header and `size` fields. + /// readable by rocksdb >= 6.19 (2021). V2WithSizes, } @@ -148,15 +148,14 @@ impl BackupMeta { /// Serialize to the rocks backup meta file format. /// - /// `V1` output is byte-identical to what rocks' own `BackupMeta::StoreToFile` - /// writes for the same contents. every line (including the last) ends in a - /// newline, and `crc32` values are plain decimal — rocks' loader - /// round-trips them through `std::to_string` and rejects any padding. + /// `V1` output is what rocks' own `BackupMeta::StoreToFile` writes. every + /// line (including the last) ends in a newline, `crc32` values are plain + /// decimal. rocks' loader uses `std::to_string` and rejects padding. /// - /// a file's `crc32` field is omitted when [`BackupFile::crc32c`] is `None`. - /// rocks itself never writes such a line (and its v1 parser rejects it) — - /// it's allowed here because in-progress marker files advertise paths - /// before their checksums are known, and only we read those. + /// `crc32` field is omitted when [`BackupFile::crc32c`] is `None`. rocks + /// never writes without (+ its v1 parser rejects it), but we use it bc we + /// write in-progress marker files with paths before crc32 is known. (only + /// read by us) pub fn serialize(&self, schema: MetaSchema) -> Result { if let Some(hex) = &self.metadata && hex.chars().any(|c| c.is_whitespace()) diff --git a/src/push.rs b/src/push.rs index d7c8283..db1481f 100644 --- a/src/push.rs +++ b/src/push.rs @@ -1,13 +1,11 @@ -//! Push a rocksdb checkpoint directory to object storage as a -//! backup-engine-compatible backup. +//! Push a rocks checkpoint to a backup-engine-compatible backup //! //! The caller makes a [checkpoint](https://github.com/facebook/rocksdb/wiki/Checkpoints) -//! (hard links, cheap, atomic) and hands us the directory; we assemble the -//! `meta/` + `private/` + `shared_checksum/` layout that rocks' own -//! `BackupEngine` writes, so either implementation can restore the result. +//! (cheap hard links when on same fs as db) and hands us the directory; we +//! assemble the `meta/`, `private/`, and `shared_checksum/` layout that +//! rocks' `BackupEngine` writes. //! -//! see `hacking.md` for the commit protocol (the `meta/` object is the -//! commit point; a `meta/..tmp` marker advertises in-flight pushes to gc). +//! a `meta/..tmp` file marks the in-flight push, `meta/` commits it. use std::collections::HashMap; use std::num::NonZeroUsize; @@ -25,74 +23,72 @@ use tracing::{info, warn}; use crate::meta::{BackupFile, BackupMeta, MetaSchema}; use crate::restore::crc32c_local_file; -use crate::retention::{self, RetireOptions}; +use crate::retention::{self, RetentionOptions}; use crate::sst; /// default max concurrent object store operations for a push -/// -/// lower than restore's 64: each in-flight upload can buffer -/// [`UPLOAD_PART_SIZE`] bytes (twice, worst case: the accumulation buffer plus -/// in-flight multipart parts), so peak memory is roughly -/// `2 * concurrency * UPLOAD_PART_SIZE` — ~256 MiB at these defaults, and much -/// less in practice since most uploads are skipped or smaller than one part. pub const DEFAULT_PUSH_CONCURRENCY: usize = 16; -/// multipart upload part size (files smaller than this go up in a single PUT) -const UPLOAD_PART_SIZE: usize = 8 * 1024 * 1024; +/// multipart upload part size +const UPLOAD_PART_SIZE: usize = 8 * 2_usize.pow(20); /// rocks caps app metadata at 1 MiB raw (`kMaxAppMetaSize`) -const MAX_APP_METADATA: usize = 1024 * 1024; +const MAX_APP_METADATA: usize = 2_usize.pow(20); /// How `shared_checksum/` files are named. #[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] pub enum SharedNaming { - /// `_s_.sst` — rocks' modern default. the + /// `_s_.sst`, rocks' modern default. the /// session id comes from a few KB of SST table properties, so unchanged /// files can be recognized (and their recorded checksums reused) without /// reading them. falls back per-file to [`SharedNaming::LegacyCrc`] when - /// no session id is readable, exactly like rocks does. blob files always - /// use legacy naming (the blob format carries no session id). + /// no session id is readable, like rocks does. blob files always use legacy + /// naming (the blob format carries no session id). #[default] SessionId, - /// `__.sst` — rocks' legacy scheme. naming a file - /// requires checksumming it, i.e. a full read of every shared file on - /// every push. only interesting for byte-compatibility experiments. + /// `__.sst`, rocks' legacy scheme. naming any file + /// requires checksumming it (full read of every shared file on any push). LegacyCrc, } -/// configure how a backup push happens #[derive(Debug, Clone)] pub struct PushOptions { - /// max concurrent object store operations + /// max concurrent object store operations for push /// /// default: [`DEFAULT_PUSH_CONCURRENCY`] pub concurrency: usize, - /// sequence number recorded in the meta file. informational only — rocks - /// documents it as "approximate, should not be used by clients" and never - /// reads it during restore. callers with the db at hand can pass - /// `db.latest_sequence_number()`. + /// sequence number recorded in the meta file, informational. + /// + /// rocks documents it as "approximate, should not be used by clients", + /// never reads it during restore. /// /// default: 0 pub sequence_number: u64, - /// unix-seconds timestamp recorded in the meta file + /// unix-seconds timestamp to record in the meta file /// /// default: `None` (current system time) pub timestamp: Option, - /// application metadata recorded in the meta file (max 1 MiB, matching - /// rocks' cap) + /// arbitrary metadata to record in the meta file (max 1 MiB) /// /// default: None pub app_metadata: Option>, - /// after a successful push, retain only the newest N backups: the same - /// operation as [`retention::purge_old_backups`] + garbage collection + /// after a successful push, delete old backups, keeping N most recent /// /// default: None (keep everything) pub retain: Option, + /// delete files as soon as they are pushed, then the checkpoint dir itself + /// + /// if the push fails partway, some files may already be gone and the + /// checkpoint will be left incomplete, so a fresh checkpoint is needed to + /// retry. + /// + /// default: false (the checkpoint is left alone) + pub consume_checkpoint: bool, /// meta file schema to write /// - /// default: [`MetaSchema::V2WithSizes`] + /// default: [`MetaSchema::V1`] (matches current rocks) pub schema: MetaSchema, - /// shared file naming scheme + /// shared filename scheme /// /// default: [`SharedNaming::SessionId`] pub naming: SharedNaming, @@ -106,28 +102,28 @@ impl Default for PushOptions { timestamp: None, app_metadata: None, retain: None, + consume_checkpoint: false, schema: MetaSchema::default(), naming: SharedNaming::default(), } } } -/// Metadata about the result of a backup push #[derive(Debug)] pub struct PushOutcome { pub backup_id: u64, pub total_files: usize, pub uploaded_files: usize, - /// shared files recognized from prior backups: checksum reused, zero bytes read + /// shared files recognized from prior backups pub reused_files: usize, pub uploaded_bytes: u64, pub reused_bytes: u64, - /// backups removed by `retain` - pub purged: Vec, + /// backups deleted by `retain` + pub removed: Vec, pub elapsed: std::time::Duration, } -/// Errors that can occur while pushing a backup. +/// Backup push possible errors #[derive(Debug, thiserror::Error)] pub enum PushError { #[error("{}: {source}", path.display())] @@ -138,9 +134,9 @@ pub enum PushError { }, #[error( - "unexpected file in checkpoint dir: {name:?} — a checkpoint contains only \ + "unexpected file in checkpoint dir: {name:?}: a checkpoint contains only \ CURRENT, MANIFEST-*, OPTIONS-*, *.log, *.sst, *.blob. (a live db dir has \ - IDENTITY/LOCK/LOG files: never push one, its MANIFEST may be mid-write)" + IDENTITY/LOCK/LOG files)" )] UnexpectedCheckpointFile { name: String }, @@ -222,23 +218,37 @@ pub enum PushError { #[error("failed to serialize backup meta")] Serialize(#[from] crate::meta::SerializeError), - #[error("backup {backup_id} pushed ok, but the retain={keep} purge failed")] - Purge { + #[error( + "backup {backup_id} pushed ok, but deleting all but the newest \ + {keep} backups afterwards failed" + )] + Retain { backup_id: u64, keep: usize, #[source] - source: Box, + source: Box, + }, + + #[error( + "failed to delete {} after pushing it (consume_checkpoint). the \ + checkpoint is now partly consumed: make a fresh one to retry from", + path.display() + )] + ConsumeFile { + path: PathBuf, + #[source] + source: std::io::Error, }, } /// what kind of checkpoint file, and where it goes in the backup #[derive(Debug, Clone, Copy, PartialEq, Eq)] enum FileKind { - /// `.sst` → `shared_checksum/` + /// `.sst` goes in `shared_checksum/` SharedSst, - /// `.blob` → `shared_checksum/`, always legacy-named + /// `.blob` goes in `shared_checksum/`, legacy-named SharedBlob, - /// CURRENT, MANIFEST-*, OPTIONS-*, *.log → `private//` + /// `CURRENT`, `MANIFEST-*`, `OPTIONS-*`, `*.log` go in `private//` Private, } @@ -379,12 +389,10 @@ fn legacy_name(name: &str, crc: u32, size: u64) -> String { /// a shared file's fate for this push #[derive(Debug)] enum SharedAction { - /// present remotely and referenced by a committed meta: reuse its - /// recorded checksum, read nothing + /// exists already and referenced by a committed meta: reuse don't re-upload Reuse { crc: u32 }, - /// upload (new file, unreferenced orphan at that name, or repair of a - /// referenced-but-missing object). crc known already for legacy-named - /// files, verified against the upload stream. + /// upload: new file, unreferenced orphan at that name, or repair of a + /// referenced-but-missing object Upload { known_crc: Option }, } @@ -400,12 +408,12 @@ struct PlannedFile { /// /// mirrors rocks' `AddBackupFileWorkItem` (`backup_engine.cc:2540-2718`): /// session-id naming avoids reading the file; a name that's both present -/// remotely and referenced by a committed meta is trusted and its checksum -/// reused; an *unreferenced* object at the target name gets overwritten -/// (rocks deletes + re-copies there). the one divergence: for legacy names -/// (which embed the crc we just computed) a present-but-unreferenced object -/// is provably identical, so we skip the re-upload that rocks — with -/// filesystems and their partially-written files in mind — performs. +/// remotely and referenced by a committed meta is reused; an *unreferenced* +/// object at the target name gets overwritten (rocks deletes + re-copies). +/// +/// difference from rock: for legacy names a present-but-unreferenced object is +/// identical since its crc matches, so we skip the re-upload. rocks re-copies +/// in that case. async fn plan_shared( store: Arc, store_prefix: StorePath, @@ -710,7 +718,7 @@ pub async fn push( backup_id = new_id, reused_files, upload_files = planned.len() - reused_files, - upload_mb = format_args!("{:.1}", upload_bytes as f64 / 2f64.powf(20.)), + upload_mb = format_args!("{:.1}", upload_bytes as f64 / 2_f64.powf(20.)), "upload plan ready" ); @@ -739,6 +747,16 @@ pub async fn push( let marker_key = full_key(&store_prefix, &format!("meta/.{new_id}.tmp")); put_new(&store, new_id, &marker_key, meta.serialize(opts.schema)?).await?; + // reused files are already safely in the bucket, so their local links can + // go immediately — no need to wait for the uploads to finish + if opts.consume_checkpoint { + for p in &planned { + if matches!(p.action, SharedAction::Reuse { .. }) { + consume(&p.local.path).await?; + } + } + } + // -- upload -------------------------------------------------------------- let uploads = planned.iter().filter_map(|p| { let known_crc = match p.action { @@ -750,6 +768,7 @@ pub async fn push( let backup_path = p.backup_path.clone(); let local_path = p.local.path.clone(); let size = p.local.size; + let consume_checkpoint = opts.consume_checkpoint; Some(async move { let crc = upload_file(store, key, local_path.clone(), size).await?; if let Some(expected) = known_crc @@ -757,6 +776,11 @@ pub async fn push( { return Err(PushError::ChecksumChanged { path: local_path }); } + // durably uploaded: drop the link now rather than at the end, so + // the filesystem can reclaim space while the push continues + if consume_checkpoint { + consume(&local_path).await?; + } Ok::<_, PushError>((backup_path, crc, size)) }) }); @@ -773,7 +797,7 @@ pub async fn push( uploaded_bytes += size; if uploaded_files.is_multiple_of(25) || uploaded_files == total_uploads { let secs = started.elapsed().as_secs_f64(); - let mb = uploaded_bytes as f64 / 2f64.powf(20.); + let mb = uploaded_bytes as f64 / 2_f64.powf(20.); info!( uploaded = uploaded_files, total = total_uploads, @@ -801,24 +825,37 @@ pub async fn push( if let Err(e) = store.delete(&marker_key).await { // the backup is committed; a stale marker is invisible to restore and - // list, and gc removes it once it goes stale + // list, and cleanup removes it once it goes stale warn!(key = %marker_key, error = %e, "failed to delete in-progress marker"); } + // every file is uploaded and the backup is committed, so all that's left + // of the checkpoint is an empty directory. past the commit point a failure + // here costs nothing but a stray dir, so it's a warning, not an error. + if opts.consume_checkpoint + && let Err(e) = tokio::fs::remove_dir(checkpoint_dir).await + { + warn!( + dir = %checkpoint_dir.display(), + error = %e, + "pushed backup ok, but could not remove the consumed checkpoint directory" + ); + } + // -- optional retention --------------------------------------------------- - let purged = match opts.retain { + let removed = match opts.retain { Some(keep) => { - let outcome = retention::purge_old_backups( + let outcome = retention::retain_backups( raw_store, prefix, keep, - &RetireOptions { + &RetentionOptions { concurrency: opts.concurrency, ..Default::default() }, ) .await - .map_err(|source| PushError::Purge { + .map_err(|source| PushError::Retain { backup_id: new_id, keep: keep.get(), source: Box::new(source), @@ -834,9 +871,9 @@ pub async fn push( total_files, uploaded_files, reused_files, - uploaded_mb = format_args!("{:.1}", uploaded_bytes as f64 / 2f64.powf(20.)), - reused_mb = format_args!("{:.1}", reused_bytes as f64 / 2f64.powf(20.)), - purged = purged.len(), + uploaded_mb = format_args!("{:.1}", uploaded_bytes as f64 / 2_f64.powf(20.)), + reused_mb = format_args!("{:.1}", reused_bytes as f64 / 2_f64.powf(20.)), + removed = removed.len(), elapsed_secs = format_args!("{:.1}", elapsed.as_secs_f64()), "push complete" ); @@ -848,15 +885,25 @@ pub async fn push( reused_files, uploaded_bytes, reused_bytes, - purged, + removed, elapsed, }) } -/// scan_backup_ids returns RetireError; surface as the push List error -fn retire_to_push(e: retention::RetireError) -> PushError { +/// unlink a checkpoint file we're done with (see `consume_checkpoint`) +async fn consume(path: &Path) -> Result<(), PushError> { + tokio::fs::remove_file(path) + .await + .map_err(|source| PushError::ConsumeFile { + path: path.to_path_buf(), + source, + }) +} + +/// scan_backup_ids returns RetentionError; surface as the push List error +fn retire_to_push(e: retention::RetentionError) -> PushError { match e { - retention::RetireError::List { prefix, source } => PushError::List { prefix, source }, + retention::RetentionError::List { prefix, source } => PushError::List { prefix, source }, other => unreachable!("scan_backup_ids only lists: {other}"), } } @@ -916,11 +963,12 @@ mod tests { ] ); - // meta parses, has all crcs + sizes (V2WithSizes default), sorted paths + // meta parses, has every crc, sorted paths. no sizes: the default + // schema is V1, exactly what rocks writes. let meta = crate::fetch_meta(&*store, "", 1).await.unwrap(); assert_eq!(meta.files.len(), 5); assert!(meta.files.iter().all(|f| f.crc32c.is_some())); - assert!(meta.files.iter().all(|f| f.size.is_some())); + assert!(meta.files.iter().all(|f| f.size.is_none())); assert_eq!(meta.files.last().unwrap().path, expected_sst); // and our own restore accepts it end-to-end @@ -962,8 +1010,19 @@ mod tests { assert_eq!(keys.iter().filter(|k| k.starts_with("meta/")).count(), 2); } + async fn pushed_meta_text(store: &InMemory) -> String { + let raw = store + .get(&StorePath::from("meta/1")) + .await + .unwrap() + .bytes() + .await + .unwrap(); + String::from_utf8(raw.to_vec()).unwrap() + } + #[tokio::test] - async fn v1_schema_writes_rocks_identical_lines() { + async fn default_schema_writes_rocks_identical_lines() { let cp = tempfile::tempdir().unwrap(); fake_checkpoint(cp.path()); let store = Arc::new(InMemory::new()); @@ -973,7 +1032,6 @@ mod tests { "", cp.path(), PushOptions { - schema: MetaSchema::V1, timestamp: Some(1234), sequence_number: 42, ..Default::default() @@ -982,19 +1040,44 @@ mod tests { .await .unwrap(); - let raw = store - .get(&StorePath::from("meta/1")) - .await - .unwrap() - .bytes() - .await - .unwrap(); - let text = String::from_utf8(raw.to_vec()).unwrap(); + // v1 by default: no schema header, no size fields -- byte-for-byte + // the shape rocks' own StoreToFile produces + let text = pushed_meta_text(&store).await; assert!(text.starts_with("1234\n42\n5\n"), "got: {text:?}"); + assert!(!text.contains("schema_version")); assert!(!text.contains("size")); assert!(text.ends_with('\n')); } + #[tokio::test] + async fn v2_schema_opt_in_adds_header_and_sizes() { + let cp = tempfile::tempdir().unwrap(); + fake_checkpoint(cp.path()); + let store = Arc::new(InMemory::new()); + + push( + store.clone(), + "", + cp.path(), + PushOptions { + schema: MetaSchema::V2WithSizes, + timestamp: Some(1234), + sequence_number: 42, + ..Default::default() + }, + ) + .await + .unwrap(); + + let text = pushed_meta_text(&store).await; + assert!( + text.starts_with("schema_version 2.1\n1234\n42\n5\n"), + "got: {text:?}" + ); + let meta = crate::fetch_meta(&*store, "", 1).await.unwrap(); + assert!(meta.files.iter().all(|f| f.size.is_some())); + } + #[tokio::test] async fn respects_prefix() { let cp = tempfile::tempdir().unwrap(); @@ -1179,6 +1262,94 @@ mod tests { assert!(!keys.iter().any(|k| k.contains(".tmp")), "{keys:?}"); } + #[tokio::test] + async fn consume_checkpoint_removes_everything() { + let parent = tempfile::tempdir().unwrap(); + let cp = parent.path().join("checkpoint"); + std::fs::create_dir(&cp).unwrap(); + fake_checkpoint(&cp); + let store = Arc::new(InMemory::new()); + + push( + store.clone(), + "", + &cp, + PushOptions { + consume_checkpoint: true, + ..Default::default() + }, + ) + .await + .unwrap(); + + assert!(!cp.exists(), "checkpoint dir should be gone"); + // and the backup is intact and restorable + let target = tempfile::tempdir().unwrap(); + crate::restore(store, "", target.path(), Default::default()) + .await + .unwrap(); + assert_eq!( + std::fs::read(target.path().join("000009.sst")).unwrap(), + b"sst-bytes-here" + ); + } + + #[tokio::test] + async fn consume_checkpoint_drops_reused_files_too() { + // second push of the same checkpoint reuses the sst (no upload), but + // consuming must still release its local link + let store = Arc::new(InMemory::new()); + let first = tempfile::tempdir().unwrap(); + fake_checkpoint(first.path()); + push(store.clone(), "", first.path(), PushOptions::default()) + .await + .unwrap(); + + let parent = tempfile::tempdir().unwrap(); + let cp = parent.path().join("checkpoint"); + std::fs::create_dir(&cp).unwrap(); + fake_checkpoint(&cp); + + let second = push( + store.clone(), + "", + &cp, + PushOptions { + consume_checkpoint: true, + ..Default::default() + }, + ) + .await + .unwrap(); + + assert_eq!(second.reused_files, 1); + assert!(!cp.exists(), "checkpoint dir should be gone"); + } + + #[tokio::test] + async fn checkpoint_untouched_by_default() { + let cp = tempfile::tempdir().unwrap(); + fake_checkpoint(cp.path()); + let store = Arc::new(InMemory::new()); + + push(store, "", cp.path(), PushOptions::default()) + .await + .unwrap(); + + for name in [ + "CURRENT", + "MANIFEST-000005", + "OPTIONS-000007", + "000004.log", + "000009.sst", + ] { + assert!( + cp.path().join(name).exists(), + "{name} should still be there" + ); + } + } + #[tokio::test] async fn oversized_metadata_rejected() { let cp = tempfile::tempdir().unwrap(); diff --git a/src/retention.rs b/src/retention.rs index fe79ba8..533e459 100644 --- a/src/retention.rs +++ b/src/retention.rs @@ -1,18 +1,15 @@ -//! Backup retention: delete, purge to a count, and garbage collection. +//! Backup retention: delete one, delete all except newest N, and gc / cleanup. //! -//! Mirrors rocks' `BackupEngine` semantics (`backup_engine.cc`): -//! `PurgeOldBackups` keeps the N highest-numbered backups; deleting a backup -//! removes its meta file *first* (the commit point of the deletion — rocks -//! aborts if that fails, `:1815`), and shared files are removed only once no -//! remaining meta references them (rocks refcounts in memory `:3336-3402`; -//! we recompute the reference set from the bucket, which is why an -//! unreadable meta hard-aborts gc: deleting anything without knowing a -//! backup's references would be guessing). +//! compared to rocks' `BackupEngine` operations: +//! - [`retain_backups`] is like `PurgeOldBackups`, keeping N most recent +//! - deleting backups removes its meta file first (same behaviour) +//! - shared files only removed once no metas reference them (same) +//! - we recompute the references from the bucket every time, where a rocks +//! `BackupEngine` keeps refcounts in memory (we should probably switch to +//! that behaviour) //! -//! The documented operating model is one active writer (pusher/purger) per -//! backup prefix, same as rocks requires exclusive `BackupEngine` access. -//! In-flight pushes advertise their file list in a `meta/..tmp` marker, -//! which gc treats as references until the marker goes stale. +//! in-flight pushes (`meta/..tmp`) are considered references for cleanup +//! unless they are old enough to be considred stale. use std::collections::HashSet; use std::num::NonZeroUsize; @@ -30,26 +27,27 @@ use crate::meta::BackupMeta; /// default age after which an in-progress marker is considered a dead push pub const DEFAULT_STALE_MARKER_AFTER: Duration = Duration::from_secs(7 * 24 * 60 * 60); -/// configure retention operations #[derive(Debug, Clone)] -pub struct RetireOptions { +pub struct RetentionOptions { /// max concurrent object store operations /// /// default: [`crate::DEFAULT_CONCURRENCY`] pub concurrency: usize, - /// report what would be deleted without deleting anything + /// print a summary of actions without deleting anything /// /// default: false pub dry_run: bool, - /// `meta/..tmp` markers younger than this are treated as in-flight - /// pushes: their advertised files are protected from garbage collection. - /// older markers are dead pushes — gc deletes them and their debris. + /// `meta/..tmp` markers newer than this are treated as in-flight, and + /// their files are protected from deletion. older markers are assumed to be + /// dead pushes, and [`cleanup`] removes them. + /// + /// (slightly different from rocks) /// /// default: [`DEFAULT_STALE_MARKER_AFTER`] (7 days) pub stale_marker_after: Duration, } -impl Default for RetireOptions { +impl Default for RetentionOptions { fn default() -> Self { Self { concurrency: crate::DEFAULT_CONCURRENCY, @@ -59,9 +57,8 @@ impl Default for RetireOptions { } } -/// Errors from retention operations. #[derive(Debug, thiserror::Error)] -pub enum RetireError { +pub enum RetentionError { #[error("failed to list objects under {prefix:?}")] List { prefix: StorePath, @@ -70,7 +67,7 @@ pub enum RetireError { }, #[error( - "cannot read backup meta {backup_id}; refusing to gc without knowing \ + "cannot read backup meta {backup_id}; refusing to clean up without knowing \ what it references" )] MetaFetch { @@ -87,9 +84,10 @@ pub enum RetireError { }, #[error( - "cannot parse fresh in-progress marker for backup {id}; refusing to gc \ - without knowing what it references (delete meta/.{id}.tmp manually if \ - it is known-dead, or retry once it is older than stale_marker_after)" + "cannot parse fresh in-progress marker for backup {id}; refusing to \ + clean up without knowing what it references (delete meta/.{id}.tmp \ + manually if it is known-dead, or retry once it is older than \ + stale_marker_after)" )] MarkerParse { id: u64, @@ -115,9 +113,9 @@ pub enum RetireError { }, } -/// What a garbage collection pass removed (or would remove, if `dry_run`). +/// What a cleanup pass removed (or would remove, if `dry_run`). #[derive(Debug, Default)] -pub struct GcOutcome { +pub struct CleanupOutcome { /// unreferenced `shared_checksum/` / `shared/` objects (backup-relative paths) pub deleted_shared: Vec, /// `private//` objects belonging to no committed backup (backup-relative paths) @@ -127,12 +125,12 @@ pub struct GcOutcome { pub dry_run: bool, } -/// The result of [`delete_backup`] or [`purge_old_backups`]. +/// The result of [`delete_backup`] or [`retain_backups`]. #[derive(Debug)] -pub struct RetireOutcome { - /// backups whose meta files were deleted (their files show up in `gc`) +pub struct RetentionOutcome { + /// backups whose meta files were deleted (their files show up in `cleanup`) pub deleted_backups: Vec, - pub gc: GcOutcome, + pub cleanup: CleanupOutcome, pub dry_run: bool, } @@ -176,14 +174,14 @@ fn canonical_id(name: &str) -> Option { pub(crate) async fn scan_backup_ids( store: &dyn ObjectStore, prefix: &str, -) -> Result { +) -> Result { let meta_prefix = StorePath::from(prefix).join("meta"); let mut committed = Vec::new(); let mut markers = Vec::new(); let mut stream = store.list(Some(&meta_prefix)); while let Some(item) = stream.next().await { - let item = item.map_err(|source| RetireError::List { + let item = item.map_err(|source| RetentionError::List { prefix: meta_prefix.clone(), source, })?; @@ -211,7 +209,7 @@ pub(crate) async fn scan_backup_ids( let listing = store .list_with_delimiter(Some(&private_prefix)) .await - .map_err(|source| RetireError::List { + .map_err(|source| RetentionError::List { prefix: private_prefix.clone(), source, })?; @@ -228,16 +226,15 @@ pub(crate) async fn scan_backup_ids( }) } -/// delete, treating already-gone as success: that IS the goal state, and S3 -/// proper doesn't even distinguish it (DELETE of a missing key returns 204) -async fn delete_object(store: &dyn ObjectStore, key: &StorePath) -> Result<(), RetireError> { +/// delete, treating already-gone as success +async fn delete_object(store: &dyn ObjectStore, key: &StorePath) -> Result<(), RetentionError> { match store.delete(key).await { Ok(()) => Ok(()), Err(object_store::Error::NotFound { .. }) => { debug!(key = %key, "already deleted"); Ok(()) } - Err(source) => Err(RetireError::Delete { + Err(source) => Err(RetentionError::Delete { key: key.clone(), source, }), @@ -249,8 +246,8 @@ async fn list_dir( store: &dyn ObjectStore, prefix: &str, dir: &str, -) -> Result, RetireError> { - // NB: extend, not join — `dir` can span components ("private/3"), and +) -> Result, RetentionError> { + // NB: extend, not join: `dir` can span components ("private/3"), and // join percent-encodes an embedded slash as part of a single segment let mut full_prefix = StorePath::from(prefix); full_prefix.extend(&StorePath::from(dir)); @@ -258,7 +255,7 @@ async fn list_dir( let mut out = Vec::new(); let mut stream = store.list(Some(&full_prefix)); while let Some(item) = stream.next().await { - let item = item.map_err(|source| RetireError::List { + let item = item.map_err(|source| RetentionError::List { prefix: full_prefix.clone(), source, })?; @@ -278,33 +275,33 @@ async fn fetch_marker( store: &dyn ObjectStore, prefix: &str, id: u64, -) -> Result { +) -> Result { let key = StorePath::from(prefix) .join("meta") .join(format!(".{id}.tmp")); let data = store .get(&key) .await - .map_err(|source| RetireError::MarkerFetch { id, source })? + .map_err(|source| RetentionError::MarkerFetch { id, source })? .bytes() .await - .map_err(|source| RetireError::MarkerFetch { id, source })?; + .map_err(|source| RetentionError::MarkerFetch { id, source })?; let text = String::from_utf8_lossy(&data); - BackupMeta::parse(&text).map_err(|source| RetireError::MarkerParse { id, source }) + BackupMeta::parse(&text).map_err(|source| RetentionError::MarkerParse { id, source }) } -/// Remove objects no committed backup (or fresh in-flight push) references. +/// Remove unfreferenced objects /// /// See the module docs for the safety rules. `treat_deleted` supports -/// dry runs of [`delete_backup`]/[`purge_old_backups`]: those ids are -/// simulated as already deleted so the gc report matches what a real run +/// dry runs of [`delete_backup`]/[`retain_backups`]: those ids are +/// simulated as already deleted so the cleanup report matches what a real run /// would additionally remove. -async fn gc_inner( +async fn cleanup_inner( store: &dyn ObjectStore, prefix: &str, - opts: &RetireOptions, + opts: &RetentionOptions, treat_deleted: &HashSet, -) -> Result { +) -> Result { let ids = scan_backup_ids(store, prefix).await?; let committed: Vec = ids .committed @@ -318,7 +315,7 @@ async fn gc_inner( for &id in &committed { let meta = crate::fetch_meta(store, prefix, id) .await - .map_err(|source| RetireError::MetaFetch { + .map_err(|source| RetentionError::MetaFetch { backup_id: id, source: Box::new(source), })?; @@ -347,12 +344,12 @@ async fn gc_inner( } } // dry-run simulation may target an id that also has a marker; the - // simulated deletion wins (matches a real purge deleting both) + // simulated deletion wins (matches a real retain deleting both) for id in treat_deleted { protected_ids.remove(id); } - let mut outcome = GcOutcome { + let mut outcome = CleanupOutcome { dry_run: opts.dry_run, ..Default::default() }; @@ -399,43 +396,43 @@ async fn gc_inner( deleted_private = outcome.deleted_private.len(), deleted_markers = outcome.deleted_markers.len(), dry_run = opts.dry_run, - "garbage collection complete" + "cleanup complete" ); Ok(outcome) } /// Remove objects no committed backup (or fresh in-flight push) references: /// leftovers of crashed pushes, and shared files orphaned by deletions. -pub async fn garbage_collect( +pub async fn cleanup( store: Arc, prefix: &str, - opts: &RetireOptions, -) -> Result { + opts: &RetentionOptions, +) -> Result { let store: Arc = Arc::new(LimitStore::new(store, opts.concurrency)); - gc_inner(&*store, prefix, opts, &HashSet::new()).await + cleanup_inner(&*store, prefix, opts, &HashSet::new()).await } -/// Delete one backup by id, then garbage-collect. +/// Delete one backup by id, then clean up. /// -/// The `meta/` delete is the commit point (rocks deletes meta first and +/// Deleting `meta/` is the commit point (rocks deletes meta first and /// aborts on failure too, `backup_engine.cc:1815`); the backup's private /// files and any shared files it exclusively referenced are then removed by -/// the gc pass, and show up in [`RetireOutcome::gc`]. +/// the cleanup pass, and show up in [`RetentionOutcome::cleanup`]. pub async fn delete_backup( store: Arc, prefix: &str, id: u64, - opts: &RetireOptions, -) -> Result { + opts: &RetentionOptions, +) -> Result { let store: Arc = Arc::new(LimitStore::new(store, opts.concurrency)); let meta_key = StorePath::from(prefix).join("meta").join(id.to_string()); match store.head(&meta_key).await { Ok(_) => {} Err(object_store::Error::NotFound { .. }) => { - return Err(RetireError::NoSuchBackup { id }); + return Err(RetentionError::NoSuchBackup { id }); } Err(source) => { - return Err(RetireError::Head { + return Err(RetentionError::Head { key: meta_key, source, }); @@ -450,26 +447,25 @@ pub async fn delete_backup( info!(backup_id = id, "deleted backup meta"); } - let gc = gc_inner(&*store, prefix, opts, &simulated).await?; - Ok(RetireOutcome { + let cleanup = cleanup_inner(&*store, prefix, opts, &simulated).await?; + Ok(RetentionOutcome { deleted_backups: vec![id], - gc, + cleanup, dry_run: opts.dry_run, }) } /// Keep only the newest `keep` backups (the highest ids), deleting the rest -/// oldest-first — the semantics of rocks' `PurgeOldBackups` -/// (`backup_engine.cc:1753-1782`), followed by one gc pass. +/// oldest-first like rocks' `PurgeOldBackups` (`backup_engine.cc:1753-1782`), +/// followed by a cleanup pass. /// -/// `keep` can't be zero: "delete every backup" must be spelled out with -/// per-id [`delete_backup`] calls instead of being one typo away. -pub async fn purge_old_backups( +/// use [`delete_backup`] on each id to delete _all_ backups +pub async fn retain_backups( store: Arc, prefix: &str, keep: NonZeroUsize, - opts: &RetireOptions, -) -> Result { + opts: &RetentionOptions, +) -> Result { let store: Arc = Arc::new(LimitStore::new(store, opts.concurrency)); let ids = scan_backup_ids(&*store, prefix).await?; @@ -485,7 +481,7 @@ pub async fn purge_old_backups( } else { // oldest first, meta only: each delete atomically retires that // backup, and an error part-way leaves a consistent (just larger - // than asked) set. files are reclaimed by the gc pass below. + // than asked) set. files are reclaimed by the cleanup pass below. for &id in &doomed { let key = StorePath::from(prefix).join("meta").join(id.to_string()); delete_object(&*store, &key).await?; @@ -493,16 +489,16 @@ pub async fn purge_old_backups( } } - let gc = gc_inner(&*store, prefix, opts, &simulated).await?; + let cleanup = cleanup_inner(&*store, prefix, opts, &simulated).await?; info!( deleted = doomed.len(), kept = ids.committed.len() - doomed.len(), dry_run = opts.dry_run, - "purge complete" + "retention complete" ); - Ok(RetireOutcome { + Ok(RetentionOutcome { deleted_backups: doomed, - gc, + cleanup, dry_run: opts.dry_run, }) } @@ -575,15 +571,15 @@ mod tests { } #[tokio::test] - async fn purge_keeps_n_highest() { + async fn retain_keeps_n_highest() { let store = Arc::new(InMemory::new()); seed_three_backups(&store).await; - let outcome = purge_old_backups( + let outcome = retain_backups( store.clone(), "", NonZeroUsize::new(1).unwrap(), - &RetireOptions::default(), + &RetentionOptions::default(), ) .await .unwrap(); @@ -591,7 +587,7 @@ mod tests { assert_eq!(outcome.deleted_backups, vec![1, 2]); // A is exclusive to 1+2 -> gone; B still referenced by 3 -> kept assert_eq!( - outcome.gc.deleted_shared, + outcome.cleanup.deleted_shared, vec!["shared_checksum/000001_1_2.sst".to_string()] ); assert_eq!( @@ -606,16 +602,16 @@ mod tests { } #[tokio::test] - async fn purge_noop_when_under_keep() { + async fn retain_noop_when_under_keep() { let store = Arc::new(InMemory::new()); seed_three_backups(&store).await; let before = keys(&store).await; - let outcome = purge_old_backups( + let outcome = retain_backups( store.clone(), "", NonZeroUsize::new(5).unwrap(), - &RetireOptions::default(), + &RetentionOptions::default(), ) .await .unwrap(); @@ -625,16 +621,16 @@ mod tests { } #[tokio::test] - async fn purge_dry_run_deletes_nothing_but_reports_everything() { + async fn retain_dry_run_deletes_nothing_but_reports_everything() { let store = Arc::new(InMemory::new()); seed_three_backups(&store).await; let before = keys(&store).await; - let dry = purge_old_backups( + let dry = retain_backups( store.clone(), "", NonZeroUsize::new(1).unwrap(), - &RetireOptions { + &RetentionOptions { dry_run: true, ..Default::default() }, @@ -642,19 +638,19 @@ mod tests { .await .unwrap(); assert_eq!(keys(&store).await, before, "dry run must not delete"); - assert!(dry.dry_run && dry.gc.dry_run); + assert!(dry.dry_run && dry.cleanup.dry_run); - let real = purge_old_backups( + let real = retain_backups( store.clone(), "", NonZeroUsize::new(1).unwrap(), - &RetireOptions::default(), + &RetentionOptions::default(), ) .await .unwrap(); assert_eq!(dry.deleted_backups, real.deleted_backups); - assert_eq!(dry.gc.deleted_shared, real.gc.deleted_shared); - assert_eq!(dry.gc.deleted_private, real.gc.deleted_private); + assert_eq!(dry.cleanup.deleted_shared, real.cleanup.deleted_shared); + assert_eq!(dry.cleanup.deleted_private, real.cleanup.deleted_private); } #[tokio::test] @@ -662,14 +658,14 @@ mod tests { let store = Arc::new(InMemory::new()); seed_three_backups(&store).await; - let outcome = delete_backup(store.clone(), "", 2, &RetireOptions::default()) + let outcome = delete_backup(store.clone(), "", 2, &RetentionOptions::default()) .await .unwrap(); assert_eq!(outcome.deleted_backups, vec![2]); // everything b2 referenced is still referenced by b1 or b3 - assert!(outcome.gc.deleted_shared.is_empty()); + assert!(outcome.cleanup.deleted_shared.is_empty()); assert_eq!( - outcome.gc.deleted_private, + outcome.cleanup.deleted_private, vec!["private/2/CURRENT".to_string()] ); @@ -681,8 +677,11 @@ mod tests { async fn delete_nonexistent_backup_errors() { let store = Arc::new(InMemory::new()); seed_three_backups(&store).await; - let result = delete_backup(store, "", 99, &RetireOptions::default()).await; - assert!(matches!(result, Err(RetireError::NoSuchBackup { id: 99 }))); + let result = delete_backup(store, "", 99, &RetentionOptions::default()).await; + assert!(matches!( + result, + Err(RetentionError::NoSuchBackup { id: 99 }) + )); } #[tokio::test] @@ -693,10 +692,10 @@ mod tests { put(&store, "private/9/CURRENT", b"debris").await; put(&store, "meta/.9.tmp", b"1000\n100\n0\n").await; - let outcome = garbage_collect( + let outcome = cleanup( store.clone(), "", - &RetireOptions { + &RetentionOptions { // InMemory stamps insertion time; zero = everything is stale stale_marker_after: Duration::ZERO, ..Default::default() @@ -737,7 +736,7 @@ mod tests { ) .await; - let outcome = garbage_collect(store.clone(), "", &RetireOptions::default()) + let outcome = cleanup(store.clone(), "", &RetentionOptions::default()) .await .unwrap(); @@ -756,10 +755,10 @@ mod tests { put(&store, "meta/4", b"not a meta file at all").await; put(&store, "shared_checksum/000099_9_5.sst", b"orphan").await; - let result = garbage_collect(store.clone(), "", &RetireOptions::default()).await; + let result = cleanup(store.clone(), "", &RetentionOptions::default()).await; assert!(matches!( result, - Err(RetireError::MetaFetch { backup_id: 4, .. }) + Err(RetentionError::MetaFetch { backup_id: 4, .. }) )); // and nothing was deleted assert!( @@ -780,7 +779,7 @@ mod tests { // outside the prefix: must be untouched put(&store, "other/shared_checksum/000099_9_5.sst", b"keep").await; - let outcome = garbage_collect(store.clone(), "pfx", &RetireOptions::default()) + let outcome = cleanup(store.clone(), "pfx", &RetentionOptions::default()) .await .unwrap(); assert_eq!( diff --git a/src/sst.rs b/src/sst.rs index e4e99ae..474d0dc 100644 --- a/src/sst.rs +++ b/src/sst.rs @@ -1,13 +1,12 @@ -//! Minimal read-only extraction of the db session id from a rocksdb SST file. +//! Minimal parser to get the session id from a rocksdb SST file. //! -//! Backups name shared SSTs `_s_.sst` (rocks' -//! default `kUseDbSessionId | kFlagIncludeFileSize` naming), which lets an -//! incremental push identify an already-uploaded file from a few KB of table -//! properties instead of checksumming the whole thing. +//! The session id is used in rocks' BackupEngine filenames, so we need to read +//! it to check if a remote copy of an SST already exists. The filenem looks +//! `_s_.sst`. //! -//! This is *never* load-bearing for correctness: any failure here just makes -//! the caller fall back to legacy `__` naming (a full -//! read), exactly like rocks treats SSTs without a session id. +//! Any failure falls back to the legacy naming format, which embeds the file +//! crc32c in the name so requires reading the *whole* sst to compute: +//! `__`. //! //! Format reference: `rocksdb/table/format.cc` (footer layout, versions 2-7), //! `rocksdb/table/meta_blocks.cc` (`"rocksdb.properties"` metaindex entry), @@ -41,10 +40,6 @@ const SESSION_ID_KEY: &[u8] = b"rocksdb.creating.session.identity"; /// meta blocks are a few KB; anything huge means we misparsed a handle const MAX_META_BLOCK_SIZE: u64 = 64 * 1024 * 1024; -/// Why a session id could not be extracted. -/// -/// Only interesting for logs and tests — callers should treat any of these as -/// "use legacy naming for this file". #[derive(Debug, thiserror::Error)] pub(crate) enum SstError { #[error("io error reading sst")] @@ -73,7 +68,7 @@ pub(crate) enum SstError { /// Best-effort read of an SST's db session id. /// -/// `None` means "name this file the legacy way" — the reason (old file, +/// `None` means "name this file the legacy way". the reason (old file, /// unsupported format, corruption, io error) is only logged, since the file /// is about to get a full read on the fallback path anyway, which will /// surface any real io problem. diff --git a/tests/e2e.rs b/tests/e2e.rs index 1a09235..f759d63 100644 --- a/tests/e2e.rs +++ b/tests/e2e.rs @@ -651,6 +651,84 @@ async fn push_migration_continuity() { ); } +/// consuming a checkpoint releases its hard links as the push proceeds +#[tokio::test] +async fn push_consume_checkpoint_releases_hard_links() { + // the point of consume_checkpoint: checkpoint entries are hard links to + // live db files, and dropping them lets the fs reclaim blocks for files + // the db no longer references. verify with real link counts. + let db_dir = tempfile::tempdir().unwrap(); + let cp_parent = tempfile::tempdir().unwrap(); + + let db = rocksdb::DB::open_default(db_dir.path()).unwrap(); + for i in 0..500u32 { + db.put(format!("k{i:05}").as_bytes(), format!("v{i}").as_bytes()) + .unwrap(); + } + let cp = make_checkpoint(&db, cp_parent.path()); + + // a checkpointed sst is hard-linked: two links to the same inode + let sst = std::fs::read_dir(&cp) + .unwrap() + .map(|e| e.unwrap().path()) + .find(|p| p.extension().is_some_and(|e| e == "sst")) + .expect("checkpoint should contain an sst"); + let db_sst = db_dir.path().join(sst.file_name().unwrap()); + { + use std::os::unix::fs::MetadataExt; + assert_eq!( + std::fs::metadata(&sst).unwrap().nlink(), + 2, + "checkpoint sst should be a hard link to the db's copy" + ); + assert_eq!( + std::fs::metadata(&sst).unwrap().ino(), + std::fs::metadata(&db_sst).unwrap().ino(), + "same inode" + ); + } + + let store = Arc::new(InMemory::new()); + eat_rocks::push( + store.clone(), + "", + &cp, + eat_rocks::PushOptions { + consume_checkpoint: true, + ..Default::default() + }, + ) + .await + .unwrap(); + + // checkpoint fully gone, db untouched, and the inode is back down to the + // db's single link (so the blocks free the moment rocks drops that file) + assert!(!cp.exists(), "consumed checkpoint dir should be removed"); + { + use std::os::unix::fs::MetadataExt; + assert_eq!( + std::fs::metadata(&db_sst).unwrap().nlink(), + 1, + "db sst should be down to one link after consumption" + ); + } + drop(db); + + // the backup itself is complete: rocks restores it + let backup_dir = tempfile::tempdir().unwrap(); + store_to_dir(&store, backup_dir.path()).await; + let restored = tempfile::tempdir().unwrap(); + rocksdb_restore(backup_dir.path(), restored.path(), restored.path(), 1).unwrap(); + let rdb = rocksdb::DB::open_for_read_only(&rocksdb::Options::default(), restored.path(), false) + .unwrap(); + for i in 0..500u32 { + assert_eq!( + rdb.get(format!("k{i:05}").as_bytes()).unwrap().unwrap(), + format!("v{i}").as_bytes() + ); + } +} + /// two eat-rocks pushes of a growing db: second is incremental, and both /// restore via rocks AND via our own restore #[tokio::test] @@ -734,7 +812,7 @@ async fn push_purge_survivor_restores() { } drop(db); - let outcome = eat_rocks::purge_old_backups( + let outcome = eat_rocks::retain_backups( store.clone(), "", std::num::NonZeroUsize::new(1).unwrap(), @@ -899,10 +977,10 @@ async fn push_crash_leaves_no_visible_backup() { // gc with markers instantly stale: debris of push 1 is reclaimed, // everything backup 2 references survives - let gc = eat_rocks::garbage_collect( + let gc = eat_rocks::cleanup( inner.clone() as Arc, "", - &eat_rocks::RetireOptions { + &eat_rocks::RetentionOptions { stale_marker_after: std::time::Duration::ZERO, ..Default::default() },