oxedyne/fe2o3/fe2o3_o3db_sync/examples/o3db_sweep.rs
9.3 KiB, 1 run
created by r1870400018:38129, which is this file's identity for as long as the history lasts, whatever it is later renamed to
download · who wrote it · its history
| 1 | //! In-place orphan sweep for an encrypted Ozone store. |
| 2 | //! |
| 3 | //! Reclaims the orphaned chunk-data records a churned store leaks -- chunk records left behind when |
| 4 | //! a chunked value was overwritten at a new geometry, or, on a pre-fix build, on any overwrite -- |
| 5 | //! by tombstoning each orphaned chunk key so the running collector reclaims its bytes. Unlike the |
| 6 | //! migration, this runs against a LIVE store with garbage collection ON and needs no second store |
| 7 | //! and no downtime. |
| 8 | //! |
| 9 | //! ```ignore |
| 10 | //! cargo run -p oxedyne_fe2o3_o3db_sync --example o3db_sweep -- \ |
| 11 | //! --source ./o3db --key ./keys/db/at_rest [--scan-secs 600] [--skew-secs 5] |
| 12 | //! ``` |
| 13 | //! |
| 14 | //! The scheme parameterisation here MUST match the store's own: this binary is wired for the |
| 15 | //! daimond gateway's store -- 16-byte `u128` user ids, AES-256-GCM at rest, CRC-32 checksums. A |
| 16 | //! store written under a different parameterisation will not read back and the sweep will fail |
| 17 | //! loudly rather than silently retire the wrong keys. |
| 18 | //! |
| 19 | //! # Safety |
| 20 | //! - The store is opened with garbage collection ON: the sweep relies on the collector to reclaim |
| 21 | //! what it tombstones. It issues no write to any value; it only tombstones chunk keys it has |
| 22 | //! proven orphaned. |
| 23 | //! - The key is read from the path given at runtime. It is NEVER hardcoded and there is NO default: |
| 24 | //! a missing or wrong-sized key aborts the run. |
| 25 | //! - `--skew-secs` guards concurrent writers: only chunk records stamped more than this many seconds |
| 26 | //! before the sweep started are retired. The default is deliberately generous. Pass `0` only for |
| 27 | //! a store nothing is writing. |
| 28 | //! |
| 29 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 30 | //! Anthropic Claude |
| 31 | |
| 32 | use oxedyne_fe2o3_core::prelude::*; |
| 33 | use oxedyne_fe2o3_crypto::enc::EncryptionScheme; |
| 34 | use oxedyne_fe2o3_hash::{ |
| 35 | csum::ChecksumScheme, |
| 36 | hash::HashScheme, |
| 37 | }; |
| 38 | use oxedyne_fe2o3_jdat::id::IdDat; |
| 39 | use oxedyne_fe2o3_o3db_sync::{ |
| 40 | base::constant, |
| 41 | comm::response::Wait, |
| 42 | data::core::RestSchemesInput, |
| 43 | db::O3db, |
| 44 | sweep, |
| 45 | }; |
| 46 | |
| 47 | use std::{ |
| 48 | fs, |
| 49 | path::{ |
| 50 | Path, |
| 51 | PathBuf, |
| 52 | }, |
| 53 | process, |
| 54 | thread, |
| 55 | time::Duration, |
| 56 | }; |
| 57 | |
| 58 | |
| 59 | // The gateway's own parameterisation. Change these only to match a store |
| 60 | // written by a differently configured application. |
| 61 | type Uid = IdDat<16, u128>; |
| 62 | type Db = O3db<16, Uid, EncryptionScheme, HashScheme, HashScheme, ChecksumScheme>; |
| 63 | |
| 64 | const DB_KEY_LEN: usize = 32; |
| 65 | |
| 66 | // How long to let the running collector settle after the sweep returns, before measuring the |
| 67 | // reclaimed footprint. Reclamation is asynchronous; this is the operator tool's own settle. |
| 68 | const GC_SETTLE_SECS: u64 = 5; |
| 69 | |
| 70 | |
| 71 | fn main() { |
| 72 | if let Err(e) = run() { |
| 73 | eprintln!("o3db_sweep FAILED: {}", e); |
| 74 | process::exit(1); |
| 75 | } |
| 76 | } |
| 77 | |
| 78 | fn run() -> Outcome<()> { |
| 79 | log_set_level!("info"); |
| 80 | |
| 81 | let args = res!(Args::parse()); |
| 82 | |
| 83 | // Read the at-rest key from disk at runtime. No default, ever. |
| 84 | let key = res!(load_key(&args.key_path)); |
| 85 | |
| 86 | // Open the store with garbage collection ON: the sweep needs the collector running to reclaim |
| 87 | // what it tombstones. |
| 88 | info!("Opening store at {:?} (gc on)...", args.source); |
| 89 | let db = res!(open_store(&args.source, &key, true, "sweep")); |
| 90 | |
| 91 | let scan_wait = Wait { |
| 92 | max_wait: Duration::from_secs(args.scan_secs), |
| 93 | check_interval: constant::CHECK_INTERVAL, |
| 94 | }; |
| 95 | info!( |
| 96 | "Sweeping orphaned chunk records (scan deadline {}s, epoch skew {}s)...", |
| 97 | args.scan_secs, args.skew_secs); |
| 98 | let report = res!(sweep::sweep_orphans( |
| 99 | db.api(), |
| 100 | Uid::default(), |
| 101 | None, |
| 102 | scan_wait, |
| 103 | Duration::from_secs(args.skew_secs), |
| 104 | )); |
| 105 | |
| 106 | if report.orphans_retired != report.orphans_found { |
| 107 | res!(db.shutdown()); |
| 108 | return Err(err!( |
| 109 | "Sweep retired {} of {} orphans found; a tombstone write did not land. Investigate \ |
| 110 | before relying on the reclaim.", report.orphans_retired, report.orphans_found; |
| 111 | Data)); |
| 112 | } |
| 113 | |
| 114 | // The sweep returns as soon as every tombstone is acknowledged; reclamation is asynchronous, so |
| 115 | // the reclaimed footprint is the caller's to measure. Let the collector -- still running -- settle, |
| 116 | // then measure the root before shutting down. |
| 117 | info!("Letting the collector settle for {}s before measuring the reclaimed footprint...", |
| 118 | GC_SETTLE_SECS); |
| 119 | thread::sleep(Duration::from_secs(GC_SETTLE_SECS)); |
| 120 | let bytes_after = res!(zone_data_bytes(&args.source)); |
| 121 | |
| 122 | res!(db.shutdown()); |
| 123 | thread::sleep(Duration::from_secs(1)); |
| 124 | |
| 125 | println!("\n{}", report.summary()); |
| 126 | |
| 127 | let pct = if report.bytes_before > 0 { |
| 128 | (bytes_after as f64 / report.bytes_before as f64) * 100.0 |
| 129 | } else { |
| 130 | 100.0 |
| 131 | }; |
| 132 | println!( |
| 133 | "\nDONE: {} orphaned chunk records retired, {} skipped as too recent to be sure.\n\ |
| 134 | Data bytes {} -> {} ({:.2}% of the pre-sweep footprint) after a {}s settle.", |
| 135 | report.orphans_retired, report.skipped_recent, |
| 136 | report.bytes_before, bytes_after, pct, GC_SETTLE_SECS); |
| 137 | |
| 138 | Ok(()) |
| 139 | } |
| 140 | |
| 141 | /// Total bytes of the store's data files (`.dat`) under a root, walked recursively; index files and |
| 142 | /// anything else are not counted. The reclaim the sweep drives shows up as data-file bytes. |
| 143 | fn zone_data_bytes(root: &Path) -> Outcome<u64> { |
| 144 | let mut total = 0u64; |
| 145 | let mut stack = vec![root.to_path_buf()]; |
| 146 | while let Some(dir) = stack.pop() { |
| 147 | let rd = match fs::read_dir(&dir) { |
| 148 | Ok(rd) => rd, |
| 149 | Err(_) => continue, |
| 150 | }; |
| 151 | for entry in rd { |
| 152 | let entry = res!(entry); |
| 153 | let path = entry.path(); |
| 154 | if path.is_dir() { |
| 155 | stack.push(path); |
| 156 | } else if path.extension().map(|e| e == constant::DATA_FILE_EXT).unwrap_or(false) { |
| 157 | let meta = res!(entry.metadata()); |
| 158 | total += meta.len(); |
| 159 | } |
| 160 | } |
| 161 | } |
| 162 | Ok(total) |
| 163 | } |
| 164 | |
| 165 | /// Opens and starts a store. |
| 166 | fn open_store( |
| 167 | root: &Path, |
| 168 | key: &[u8; DB_KEY_LEN], |
| 169 | gc_on: bool, |
| 170 | label: &str, |
| 171 | ) |
| 172 | -> Outcome<Db> |
| 173 | { |
| 174 | let aes_gcm = res!(EncryptionScheme::new_aes_256_gcm_with_key(&key[..])); |
| 175 | let crc32 = ChecksumScheme::new_crc32(); |
| 176 | let schms_input = RestSchemesInput::new( |
| 177 | Some(aes_gcm), |
| 178 | None::<HashScheme>, |
| 179 | None::<HashScheme>, |
| 180 | Some(crc32), |
| 181 | ); |
| 182 | let mut db: Db = res!(O3db::new(root, None, schms_input, Uid::default())); |
| 183 | res!(db.start(label.to_string())); |
| 184 | res!(ok!(db.updated_api()).activate_gc(gc_on)); |
| 185 | thread::sleep(Duration::from_millis(500)); |
| 186 | let (_, msgs) = res!(db.api().ping_bots(constant::USER_REQUEST_WAIT)); |
| 187 | info!("{}: {} bots responded.", label, msgs.len()); |
| 188 | Ok(db) |
| 189 | } |
| 190 | |
| 191 | /// Reads the 32-byte at-rest key, failing loudly if it is absent or the wrong size. There is |
| 192 | /// deliberately no fallback: a wrong key cannot decrypt a byte. |
| 193 | fn load_key(path: &Path) -> Outcome<[u8; DB_KEY_LEN]> { |
| 194 | if !path.exists() { |
| 195 | return Err(err!( |
| 196 | "No database key at {:?}. Supply the store's at-rest key with --key; \ |
| 197 | without it the store cannot be read.", path; |
| 198 | Missing, Key, Input)); |
| 199 | } |
| 200 | let bytes = res!(fs::read(path)); |
| 201 | if bytes.len() != DB_KEY_LEN { |
| 202 | return Err(err!( |
| 203 | "The database key at {:?} is {} bytes, but must be {}.", |
| 204 | path, bytes.len(), DB_KEY_LEN; |
| 205 | Invalid, Input, Key)); |
| 206 | } |
| 207 | let mut key = [0u8; DB_KEY_LEN]; |
| 208 | key.copy_from_slice(&bytes); |
| 209 | Ok(key) |
| 210 | } |
| 211 | |
| 212 | struct Args { |
| 213 | source: PathBuf, |
| 214 | key_path: PathBuf, |
| 215 | scan_secs: u64, |
| 216 | skew_secs: u64, |
| 217 | } |
| 218 | |
| 219 | impl Args { |
| 220 | fn parse() -> Outcome<Self> { |
| 221 | let mut source: Option<PathBuf> = None; |
| 222 | let mut key_path: Option<PathBuf> = None; |
| 223 | let mut scan_secs: u64 = 600; |
| 224 | let mut skew_secs: u64 = 5; |
| 225 | |
| 226 | let mut it = std::env::args().skip(1); |
| 227 | while let Some(arg) = it.next() { |
| 228 | match arg.as_str() { |
| 229 | "--source" => source = Some(PathBuf::from(res!(next(&mut it, "--source")))), |
| 230 | "--key" => key_path = Some(PathBuf::from(res!(next(&mut it, "--key")))), |
| 231 | "--scan-secs" => { |
| 232 | let v = res!(next(&mut it, "--scan-secs")); |
| 233 | scan_secs = res!(v.parse::<u64>().map_err(|_| err!( |
| 234 | "--scan-secs must be a whole number of seconds, got {:?}.", v; |
| 235 | Invalid, Input))); |
| 236 | }, |
| 237 | "--skew-secs" => { |
| 238 | let v = res!(next(&mut it, "--skew-secs")); |
| 239 | skew_secs = res!(v.parse::<u64>().map_err(|_| err!( |
| 240 | "--skew-secs must be a whole number of seconds, got {:?}.", v; |
| 241 | Invalid, Input))); |
| 242 | }, |
| 243 | other => return Err(err!( |
| 244 | "Unrecognised argument {:?}. Usage: --source DIR --key PATH \ |
| 245 | [--scan-secs N] [--skew-secs N].", other; |
| 246 | Invalid, Input)), |
| 247 | } |
| 248 | } |
| 249 | Ok(Self { |
| 250 | source: res!(source.ok_or_else(|| err!("Missing --source DIR."; Missing, Input))), |
| 251 | key_path: res!(key_path.ok_or_else(|| err!("Missing --key PATH."; Missing, Input))), |
| 252 | scan_secs, |
| 253 | skew_secs, |
| 254 | }) |
| 255 | } |
| 256 | } |
| 257 | |
| 258 | fn next(it: &mut impl Iterator<Item = String>, flag: &str) -> Outcome<String> { |
| 259 | it.next().ok_or_else(|| err!( |
| 260 | "Argument {} needs a value.", flag; Missing, Input)) |
| 261 | } |