oxedyne/fe2o3/fe2o3_o3db_sync/examples/o3db_migrate.rs
8.6 KiB, 1 run
created by r1870400018:38110, 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 | //! Live-set migration / compaction for an encrypted Ozone store. |
| 2 | //! |
| 3 | //! Copies exactly the live key set of a source store into a fresh target store, |
| 4 | //! dropping orphaned chunk-data records (the leak that grows a churned store |
| 5 | //! without bound). Both stores are opened under the SAME at-rest key and the |
| 6 | //! SAME configuration, so the target holds the same values, freshly encrypted, |
| 7 | //! with matching chunk geometry. |
| 8 | //! |
| 9 | //! ```ignore |
| 10 | //! cargo run -p oxedyne_fe2o3_o3db_sync --example o3db_migrate -- \ |
| 11 | //! --source ./o3db --target ./o3db.new --key ./keys/db/at_rest [--scan-secs 600] |
| 12 | //! ``` |
| 13 | //! |
| 14 | //! The scheme parameterisation here MUST match the store's own: this binary is |
| 15 | //! wired for the daimond gateway's store -- 16-byte `u128` user ids, AES-256-GCM |
| 16 | //! at rest, CRC-32 checksums. A store written under a different parameterisation |
| 17 | //! will not read back and the migration will fail loudly rather than silently |
| 18 | //! copy garbage. |
| 19 | //! |
| 20 | //! # Safety |
| 21 | //! - The source is opened with garbage collection OFF and is only ever read; |
| 22 | //! every write goes to the target. Stop the process that owns the store first, |
| 23 | //! so the store is quiescent. |
| 24 | //! - The key is read from the path given at runtime. It is NEVER hardcoded and |
| 25 | //! there is NO default: a missing or wrong-sized key aborts the run. |
| 26 | //! - The run refuses to report success unless the copy verified: the target's |
| 27 | //! live-key count equals the source's, and every value read back |
| 28 | //! byte-identical. |
| 29 | //! - The operator still owns the shrink gate: this prints the source and target |
| 30 | //! byte totals; a target that did not shrink as expected must NOT be swapped |
| 31 | //! in. See the runbook. |
| 32 | //! |
| 33 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 34 | //! Anthropic Claude |
| 35 | |
| 36 | use oxedyne_fe2o3_core::prelude::*; |
| 37 | use oxedyne_fe2o3_crypto::enc::EncryptionScheme; |
| 38 | use oxedyne_fe2o3_hash::{ |
| 39 | csum::ChecksumScheme, |
| 40 | hash::HashScheme, |
| 41 | }; |
| 42 | use oxedyne_fe2o3_jdat::id::IdDat; |
| 43 | use oxedyne_fe2o3_o3db_sync::{ |
| 44 | base::constant, |
| 45 | comm::response::Wait, |
| 46 | data::core::RestSchemesInput, |
| 47 | db::O3db, |
| 48 | migrate, |
| 49 | }; |
| 50 | |
| 51 | use std::{ |
| 52 | fs, |
| 53 | path::{ |
| 54 | Path, |
| 55 | PathBuf, |
| 56 | }, |
| 57 | process, |
| 58 | thread, |
| 59 | time::Duration, |
| 60 | }; |
| 61 | |
| 62 | |
| 63 | // The gateway's own parameterisation. Change these only to match a store |
| 64 | // written by a differently configured application. |
| 65 | type Uid = IdDat<16, u128>; |
| 66 | type Db = O3db<16, Uid, EncryptionScheme, HashScheme, HashScheme, ChecksumScheme>; |
| 67 | |
| 68 | const DB_KEY_LEN: usize = 32; |
| 69 | |
| 70 | |
| 71 | fn main() { |
| 72 | if let Err(e) = run() { |
| 73 | eprintln!("o3db_migrate 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 | // 1. Open the source read-only: garbage collection off, no writes issued. |
| 87 | info!("Opening source store at {:?} (gc off, read-only)...", args.source); |
| 88 | let src = res!(open_store(&args.source, None, &key, false, "migrate-src")); |
| 89 | |
| 90 | // 2. Create the target with the SOURCE'S configuration, so chunk geometry |
| 91 | // and encryption parameters match exactly. |
| 92 | let src_cfg = src.cfg().clone(); |
| 93 | info!("Creating fresh target store at {:?} with the source's configuration...", args.target); |
| 94 | let tgt = res!(open_store(&args.target, Some(src_cfg), &key, false, "migrate-tgt")); |
| 95 | |
| 96 | // 3. Migrate and verify. |
| 97 | let scan_wait = Wait { |
| 98 | max_wait: Duration::from_secs(args.scan_secs), |
| 99 | check_interval: constant::CHECK_INTERVAL, |
| 100 | }; |
| 101 | info!("Migrating live key set (scan deadline {}s)...", args.scan_secs); |
| 102 | let report = res!(migrate::migrate_live_set( |
| 103 | src.api(), |
| 104 | tgt.api(), |
| 105 | Uid::default(), |
| 106 | None, |
| 107 | scan_wait, |
| 108 | )); |
| 109 | |
| 110 | // 4. Shut both down cleanly so the target's bytes are settled on disk. |
| 111 | res!(tgt.shutdown()); |
| 112 | res!(src.shutdown()); |
| 113 | thread::sleep(Duration::from_secs(1)); |
| 114 | |
| 115 | println!("\n{}", report.summary()); |
| 116 | |
| 117 | if !report.verified_ok { |
| 118 | return Err(err!( |
| 119 | "Migration did NOT verify; the target must not be swapped in."; |
| 120 | Data)); |
| 121 | } |
| 122 | |
| 123 | // The shrink gate is the operator's to enforce, but flag the unexpected. |
| 124 | if report.source_bytes > 0 && report.target_bytes * 2 >= report.source_bytes { |
| 125 | println!( |
| 126 | "\nWARNING: the target did NOT shrink much (target is {} of source bytes). \ |
| 127 | If a large reclaim was expected, DO NOT swap; investigate first.", |
| 128 | fmt!("{:.1}%", (report.target_bytes as f64 / report.source_bytes as f64) * 100.0)); |
| 129 | } |
| 130 | |
| 131 | println!( |
| 132 | "\nPASS: {} live keys copied and verified byte-identical; {} tombstones dropped.\n\ |
| 133 | Source: {} bytes / {} files. Target: {} bytes / {} files.\n\ |
| 134 | The source was opened read-only; verify its data/index files are unchanged before swapping.", |
| 135 | report.copied, report.tombstones, |
| 136 | report.source_bytes, report.source_files, |
| 137 | report.target_bytes, report.target_files); |
| 138 | |
| 139 | Ok(()) |
| 140 | } |
| 141 | |
| 142 | /// Opens and starts a store, leaving garbage collection off unless asked. |
| 143 | fn open_store( |
| 144 | root: &Path, |
| 145 | cfg_opt: Option<OzoneConfigT>, |
| 146 | key: &[u8; DB_KEY_LEN], |
| 147 | gc_on: bool, |
| 148 | label: &str, |
| 149 | ) |
| 150 | -> Outcome<Db> |
| 151 | { |
| 152 | let aes_gcm = res!(EncryptionScheme::new_aes_256_gcm_with_key(&key[..])); |
| 153 | let crc32 = ChecksumScheme::new_crc32(); |
| 154 | let schms_input = RestSchemesInput::new( |
| 155 | Some(aes_gcm), |
| 156 | None::<HashScheme>, |
| 157 | None::<HashScheme>, |
| 158 | Some(crc32), |
| 159 | ); |
| 160 | let mut db: Db = res!(O3db::new(root, cfg_opt, schms_input, Uid::default())); |
| 161 | res!(db.start(label.to_string())); |
| 162 | // Default is off; state it explicitly so a read never collects. |
| 163 | res!(ok!(db.updated_api()).activate_gc(gc_on)); |
| 164 | thread::sleep(Duration::from_millis(500)); |
| 165 | let (_, msgs) = res!(db.api().ping_bots(constant::USER_REQUEST_WAIT)); |
| 166 | info!("{}: {} bots responded.", label, msgs.len()); |
| 167 | Ok(db) |
| 168 | } |
| 169 | |
| 170 | /// Reads the 32-byte at-rest key, failing loudly if it is absent or the wrong |
| 171 | /// size. There is deliberately no fallback: a wrong key cannot decrypt a byte. |
| 172 | fn load_key(path: &Path) -> Outcome<[u8; DB_KEY_LEN]> { |
| 173 | if !path.exists() { |
| 174 | return Err(err!( |
| 175 | "No database key at {:?}. Supply the store's at-rest key with --key; \ |
| 176 | without it the store cannot be read.", path; |
| 177 | Missing, Key, Input)); |
| 178 | } |
| 179 | let bytes = res!(fs::read(path)); |
| 180 | if bytes.len() != DB_KEY_LEN { |
| 181 | return Err(err!( |
| 182 | "The database key at {:?} is {} bytes, but must be {}.", |
| 183 | path, bytes.len(), DB_KEY_LEN; |
| 184 | Invalid, Input, Key)); |
| 185 | } |
| 186 | let mut key = [0u8; DB_KEY_LEN]; |
| 187 | key.copy_from_slice(&bytes); |
| 188 | Ok(key) |
| 189 | } |
| 190 | |
| 191 | // Alias so the signature above does not need the full config path spelled out. |
| 192 | type OzoneConfigT = oxedyne_fe2o3_o3db_sync::base::cfg::OzoneConfig; |
| 193 | |
| 194 | struct Args { |
| 195 | source: PathBuf, |
| 196 | target: PathBuf, |
| 197 | key_path: PathBuf, |
| 198 | scan_secs: u64, |
| 199 | } |
| 200 | |
| 201 | impl Args { |
| 202 | fn parse() -> Outcome<Self> { |
| 203 | let mut source: Option<PathBuf> = None; |
| 204 | let mut target: Option<PathBuf> = None; |
| 205 | let mut key_path: Option<PathBuf> = None; |
| 206 | let mut scan_secs: u64 = 600; |
| 207 | |
| 208 | let mut it = std::env::args().skip(1); |
| 209 | while let Some(arg) = it.next() { |
| 210 | match arg.as_str() { |
| 211 | "--source" => source = Some(PathBuf::from(res!(next(&mut it, "--source")))), |
| 212 | "--target" => target = Some(PathBuf::from(res!(next(&mut it, "--target")))), |
| 213 | "--key" => key_path = Some(PathBuf::from(res!(next(&mut it, "--key")))), |
| 214 | "--scan-secs" => { |
| 215 | let v = res!(next(&mut it, "--scan-secs")); |
| 216 | scan_secs = res!(v.parse::<u64>().map_err(|_| err!( |
| 217 | "--scan-secs must be a whole number of seconds, got {:?}.", v; |
| 218 | Invalid, Input))); |
| 219 | }, |
| 220 | other => return Err(err!( |
| 221 | "Unrecognised argument {:?}. Usage: --source DIR --target DIR \ |
| 222 | --key PATH [--scan-secs N].", other; |
| 223 | Invalid, Input)), |
| 224 | } |
| 225 | } |
| 226 | Ok(Self { |
| 227 | source: res!(source.ok_or_else(|| err!("Missing --source DIR."; Missing, Input))), |
| 228 | target: res!(target.ok_or_else(|| err!("Missing --target DIR."; Missing, Input))), |
| 229 | key_path: res!(key_path.ok_or_else(|| err!("Missing --key PATH."; Missing, Input))), |
| 230 | scan_secs, |
| 231 | }) |
| 232 | } |
| 233 | } |
| 234 | |
| 235 | fn next(it: &mut impl Iterator<Item = String>, flag: &str) -> Outcome<String> { |
| 236 | it.next().ok_or_else(|| err!( |
| 237 | "Argument {} needs a value.", flag; Missing, Input)) |
| 238 | } |