oxedyne/fe2o3/fe2o3_o3db_sync/tests/migrate.rs
24.5 KiB, 6 runs
created by r1870400018:38114, 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 integration test. |
| 2 | //! |
| 3 | //! Proves the `migrate` module on synthetic, ENCRYPTED stores that reproduce |
| 4 | //! the exact conditions of the daimond gateway store: |
| 5 | //! |
| 6 | //! - **Orphaned chunks are dropped.** A chunked value overwritten many times |
| 7 | //! under the pre-fix random-ticket scheme leaves its old chunk-data keys as |
| 8 | //! live index entries with nothing referencing them. The source bloats; the |
| 9 | //! target holds only the current value's chunks, so it is dramatically |
| 10 | //! smaller. The current live value reads back byte-identical from the target. |
| 11 | //! - **Unreferenced `Complete` keys survive (the sbody trap).** A small |
| 12 | //! `sync:`-like record holds only hashes; several large `sbody:`-like records |
| 13 | //! are separate `Complete` keys that nothing in the value graph points at. The |
| 14 | //! migration must copy every one of them -- a reachability copy would drop |
| 15 | //! them. This is the data-loss case, tested explicitly. |
| 16 | //! - **Several keyspaces, several keys each, all carried.** The per-prefix |
| 17 | //! report accounts for every keyspace. |
| 18 | //! - **Tombstones are dropped.** A deleted key is emitted by a scan but reads as |
| 19 | //! absent; the target simply lacks it, and reads it as absent too. |
| 20 | //! - **The source is not mutated in any data or index file by a read-only, |
| 21 | //! gc-off open.** The one file a plain open rewrites is `config.jdat`; every |
| 22 | //! `.dat`/`.ind` file is byte-identical after the migration. This is the |
| 23 | //! backup-safety property the rename-based runbook depends on. |
| 24 | //! - **An unchunked-only store migrates and compacts too.** |
| 25 | //! |
| 26 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 27 | //! Anthropic Claude |
| 28 | |
| 29 | use oxedyne_fe2o3_core::prelude::*; |
| 30 | use oxedyne_fe2o3_crypto::enc::EncryptionScheme; |
| 31 | use oxedyne_fe2o3_hash::{ |
| 32 | csum::ChecksumScheme, |
| 33 | hash::HashScheme, |
| 34 | }; |
| 35 | use oxedyne_fe2o3_jdat::{ |
| 36 | prelude::*, |
| 37 | chunk::PartKey, |
| 38 | }; |
| 39 | use oxedyne_fe2o3_o3db_sync::{ |
| 40 | base::{ |
| 41 | cfg::OzoneConfig, |
| 42 | constant, |
| 43 | }, |
| 44 | comm::response::Wait, |
| 45 | data::core::RestSchemesInput, |
| 46 | file::{ |
| 47 | core::{ |
| 48 | FileAccess, |
| 49 | FileType, |
| 50 | }, |
| 51 | stored::{ |
| 52 | StoredIndex, |
| 53 | StoredKey, |
| 54 | }, |
| 55 | zdir::ZoneDir, |
| 56 | }, |
| 57 | migrate, |
| 58 | test::setup, |
| 59 | }; |
| 60 | |
| 61 | use std::{ |
| 62 | collections::BTreeMap, |
| 63 | fs, |
| 64 | io::BufReader, |
| 65 | path::{ |
| 66 | Path, |
| 67 | PathBuf, |
| 68 | }, |
| 69 | thread, |
| 70 | time::Duration, |
| 71 | }; |
| 72 | |
| 73 | |
| 74 | type Enc = EncryptionScheme; |
| 75 | type Kh = HashScheme; |
| 76 | type Cs = ChecksumScheme; |
| 77 | |
| 78 | /// A generous scan deadline: a scan walks every index file, so a real store |
| 79 | /// needs far more than the shared user-request deadline. Ample for the test. |
| 80 | const SCAN_WAIT: Wait = Wait { |
| 81 | max_wait: Duration::from_secs(120), |
| 82 | check_interval: constant::CHECK_INTERVAL, |
| 83 | }; |
| 84 | |
| 85 | |
| 86 | #[test] |
| 87 | fn main() -> Outcome<()> { |
| 88 | log_set_level!("warn"); |
| 89 | let outcome = run_all(); |
| 90 | log_finish_wait!(); |
| 91 | outcome |
| 92 | } |
| 93 | |
| 94 | fn run_all() -> Outcome<()> { |
| 95 | res!(scenario_orphan_shrink_and_untouched()); |
| 96 | res!(scenario_unchunked_only()); |
| 97 | Ok(()) |
| 98 | } |
| 99 | |
| 100 | // A 32-byte at-rest key, as the gateway holds in a key file. Fixed here so the |
| 101 | // test is deterministic; NEVER hardcode a real key -- the tool reads it from a |
| 102 | // path at runtime. |
| 103 | fn test_key() -> [u8; 32] { [0x5au8; 32] } |
| 104 | |
| 105 | fn schemes_input() -> Outcome<RestSchemesInput<Enc, Kh, Kh, Cs>> { |
| 106 | let aes_gcm = res!(EncryptionScheme::new_aes_256_gcm_with_key(&test_key()[..])); |
| 107 | let crc32 = ChecksumScheme::new_crc32(); |
| 108 | Ok(RestSchemesInput::new( |
| 109 | Some(aes_gcm), |
| 110 | None::<HashScheme>, |
| 111 | None::<HashScheme>, |
| 112 | Some(crc32), |
| 113 | )) |
| 114 | } |
| 115 | |
| 116 | fn fresh_root(name: &str) -> Outcome<PathBuf> { |
| 117 | let _ = fs::remove_dir_all(name); |
| 118 | res!(fs::create_dir_all(name)); |
| 119 | Ok(res!(Path::new(name).canonicalize())) |
| 120 | } |
| 121 | |
| 122 | fn empty_root(name: &str) -> Outcome<PathBuf> { |
| 123 | let _ = fs::remove_dir_all(name); |
| 124 | res!(fs::create_dir_all(name)); |
| 125 | Ok(res!(Path::new(name).canonicalize())) |
| 126 | } |
| 127 | |
| 128 | /// Every regular file under `root`, keyed by path, mapped to its exact bytes. |
| 129 | /// The synthetic stores are small, so exact bytes are kept rather than a hash. |
| 130 | fn snapshot(root: &Path) -> Outcome<BTreeMap<PathBuf, Vec<u8>>> { |
| 131 | let mut out = BTreeMap::new(); |
| 132 | let mut stack = vec![root.to_path_buf()]; |
| 133 | while let Some(dir) = stack.pop() { |
| 134 | let rd = match fs::read_dir(&dir) { |
| 135 | Ok(rd) => rd, |
| 136 | Err(_) => continue, |
| 137 | }; |
| 138 | for entry in rd { |
| 139 | let entry = res!(entry); |
| 140 | let path = entry.path(); |
| 141 | let meta = res!(entry.metadata()); |
| 142 | if meta.is_dir() { |
| 143 | stack.push(path); |
| 144 | } else if meta.is_file() { |
| 145 | out.insert(path.clone(), res!(fs::read(&path))); |
| 146 | } |
| 147 | } |
| 148 | } |
| 149 | Ok(out) |
| 150 | } |
| 151 | |
| 152 | /// The base configuration for the synthetic stores: encrypted, several zones, |
| 153 | /// small files and a low chunk threshold so values chunk and files roll over, |
| 154 | /// no per-zone overrides so everything lives under the store root (as the |
| 155 | /// gateway store does). Garbage collection is driven per-store at start. |
| 156 | fn base_cfg() -> Outcome<OzoneConfig> { |
| 157 | let mut cfg = res!(setup::default_cfg()); |
| 158 | cfg.num_zones = 3; |
| 159 | cfg.num_cbots_per_zone = 1; |
| 160 | cfg.num_fbots_per_zone = 1; |
| 161 | cfg.num_wbots_per_zone = 1; |
| 162 | cfg.num_igbots_per_zone = 1; |
| 163 | cfg.data_file_max_bytes = 4_000; |
| 164 | cfg.rest_chunk_threshold = 500; |
| 165 | cfg.rest_chunk_bytes = 128; |
| 166 | cfg.cache_size_limit_bytes = 40_000_000; |
| 167 | cfg.zone_overrides = BTreeMap::new(); |
| 168 | Ok(cfg) |
| 169 | } |
| 170 | |
| 171 | |
| 172 | /// The value length written in overwrite `round`. Strictly increasing, so every |
| 173 | /// round has a distinct length -> a distinct chunk geometry -> distinct chunk |
| 174 | /// keys, whether keys are random-ticket (today) or deterministic (post-fix). All |
| 175 | /// lengths are well over the chunk threshold, so every round chunks. |
| 176 | fn round_size(round: usize) -> usize { 2_000 + round * 200 } |
| 177 | |
| 178 | fn scenario_orphan_shrink_and_untouched() -> Outcome<()> { |
| 179 | |
| 180 | let src_root = res!(fresh_root("./test_db_migrate_src")); |
| 181 | let tgt_root = res!(empty_root("./test_db_migrate_tgt")); |
| 182 | let user = setup::Uid::default(); |
| 183 | |
| 184 | // Distinct, self-checking payload for a key so a readback can be compared. |
| 185 | let payload = |tag: &str, n: usize| -> Dat { |
| 186 | let seed = fmt!("{}:", tag); |
| 187 | let mut s = String::with_capacity(n); |
| 188 | while s.len() < n { s.push_str(&seed); } |
| 189 | s.truncate(n); |
| 190 | Dat::Str(s) |
| 191 | }; |
| 192 | |
| 193 | // ── Phase 1: build the source, populate it, shut it down cleanly. ── |
| 194 | { |
| 195 | let db = res!(setup::start_db( |
| 196 | src_root.clone(), |
| 197 | Some(res!(base_cfg())), |
| 198 | res!(schemes_input()), |
| 199 | None, |
| 200 | false, // gc OFF: nothing is collected, so orphans accumulate as they would pre-fix |
| 201 | true, // wipe |
| 202 | )); |
| 203 | thread::sleep(Duration::from_secs(1)); |
| 204 | |
| 205 | // (a) Several keyspaces, several keys each -- small unchunked Complete keys. |
| 206 | for (prefix, keys) in &[ |
| 207 | ("acct", 3usize), |
| 208 | ("sess", 4), |
| 209 | ("lease", 2), |
| 210 | ("sub", 2), |
| 211 | ("a", 3), |
| 212 | ] { |
| 213 | for i in 0..*keys { |
| 214 | let k = dat!(fmt!("{}:{:04}", prefix, i)); |
| 215 | res!(store_wait(db.api(), k, payload(prefix, 40), user)); |
| 216 | } |
| 217 | } |
| 218 | |
| 219 | // (b) The sbody trap: one small "sync"-like record holding only hashes, |
| 220 | // and several large "sbody" Complete keys that NOTHING references. |
| 221 | // They are independent live Complete keys; a reachability copy from |
| 222 | // `sync:` would drop them. |
| 223 | let mut sync_hashes = String::new(); |
| 224 | for i in 0..5usize { |
| 225 | let body = payload("sbody", 300 + i * 10); // under the chunk threshold => unchunked Complete |
| 226 | let hk = fmt!("sbody:acctZ:{:08x}", i as u64); |
| 227 | sync_hashes.push_str(&hk); |
| 228 | sync_hashes.push(','); |
| 229 | res!(store_wait(db.api(), dat!(hk.clone()), body, user)); |
| 230 | } |
| 231 | // The sync record's value contains only the hashes; it does not carry the bodies. |
| 232 | res!(store_wait(db.api(), dat!("sync:acctZ"), Dat::Str(sync_hashes), user)); |
| 233 | |
| 234 | // (c) A chunked value overwritten many times -> orphaned chunk keys. |
| 235 | // Each overwrite CHANGES THE GEOMETRY (a different value length, so a |
| 236 | // different data_len and part count), which changes the chunk keys |
| 237 | // under any deterministic scheme as well as under today's random |
| 238 | // ticket. So the prior generation's chunk-data keys are genuinely |
| 239 | // orphaned either way, and this test cannot go vacuous once the |
| 240 | // deterministic-key leak fix lands on the same branch. A same-size |
| 241 | // overwrite would reuse the chunk keys under deterministic keys and |
| 242 | // mint no orphans at all. |
| 243 | let chunk_key = dat!("chunked:big"); |
| 244 | for round in 0..30usize { |
| 245 | let sz = round_size(round); |
| 246 | res!(store_wait(db.api(), chunk_key.clone(), payload(&fmt!("v{:02}", round), sz), user)); |
| 247 | } |
| 248 | |
| 249 | // (d) A key that is written then deleted -> a live tombstone that a scan |
| 250 | // emits but get() reads as absent. The migration must drop it. |
| 251 | res!(store_wait(db.api(), dat!("acct:deleteme"), payload("gone", 40), user)); |
| 252 | res!(delete_wait(db.api(), dat!("acct:deleteme"), user)); |
| 253 | |
| 254 | thread::sleep(Duration::from_secs(2)); |
| 255 | res!(db.shutdown()); |
| 256 | } |
| 257 | thread::sleep(Duration::from_secs(1)); |
| 258 | |
| 259 | // Snapshot the source while it is closed -- this is what the migration must |
| 260 | // not disturb in any data or index file. |
| 261 | let before = res!(snapshot(&src_root)); |
| 262 | let src_bytes_before: u64 = before.values().map(|v| v.len() as u64).sum(); |
| 263 | |
| 264 | // ── Phase 2: open source read-only (gc off), create target with the |
| 265 | // source's config, migrate, verify. ── |
| 266 | let report = { |
| 267 | let src_db = res!(setup::start_db( |
| 268 | src_root.clone(), |
| 269 | None, // load the source's own config.jdat |
| 270 | res!(schemes_input()), |
| 271 | None, |
| 272 | false, // gc OFF: a read-only, non-collecting open |
| 273 | false, // do NOT wipe |
| 274 | )); |
| 275 | thread::sleep(Duration::from_secs(1)); |
| 276 | |
| 277 | // Prove, structurally, that the source actually holds orphaned chunk |
| 278 | // data BEFORE the migration -- otherwise the orphan-drop path this tool |
| 279 | // exists for would not be exercised and the test would be vacuous. Only |
| 280 | // one key ("chunked:big") ever chunked, so its current bunch key names |
| 281 | // exactly `live_parts` live chunk keys; every other chunk-data record on |
| 282 | // disk (index >= 1) is an orphan of a superseded generation. |
| 283 | let live_parts = res!(live_bunch_num_parts(src_db.api(), &dat!("chunked:big"))); |
| 284 | let total_chunk_records = res!(count_chunk_data_records(src_db.api())); |
| 285 | let orphans = total_chunk_records.saturating_sub(live_parts); |
| 286 | warn!(sync_log::stream(), |
| 287 | "Orphan check: {} chunk-data records on disk, {} referenced by the live \ |
| 288 | value, so {} orphaned chunk-data keys with no live bunch key.", |
| 289 | total_chunk_records, live_parts, orphans); |
| 290 | if orphans == 0 { |
| 291 | return Err(err!( |
| 292 | "No orphaned chunk-data keys were minted (total chunk-data records {} == \ |
| 293 | live parts {}); the orphan-drop path is not being exercised, so the test \ |
| 294 | would be vacuous. The overwrite geometry must change per round.", |
| 295 | total_chunk_records, live_parts; |
| 296 | Test, Data)); |
| 297 | } |
| 298 | |
| 299 | // Carry the source's config to the fresh target, unchanged. |
| 300 | let src_cfg = src_db.cfg().clone(); |
| 301 | let tgt_db = res!(setup::start_db( |
| 302 | tgt_root.clone(), |
| 303 | Some(src_cfg), |
| 304 | res!(schemes_input()), |
| 305 | None, |
| 306 | false, // gc off during the bulk load |
| 307 | true, // fresh |
| 308 | )); |
| 309 | thread::sleep(Duration::from_secs(1)); |
| 310 | |
| 311 | let report = res!(migrate::migrate_live_set( |
| 312 | src_db.api(), |
| 313 | tgt_db.api(), |
| 314 | user, |
| 315 | None, // use each store's own (encrypted) default schemes |
| 316 | SCAN_WAIT, |
| 317 | )); |
| 318 | |
| 319 | // Print the report so the run carries the evidence. |
| 320 | warn!(sync_log::stream(), "{}", report.summary()); |
| 321 | |
| 322 | res!(tgt_db.shutdown()); |
| 323 | res!(src_db.shutdown()); |
| 324 | report |
| 325 | }; |
| 326 | thread::sleep(Duration::from_secs(1)); |
| 327 | |
| 328 | // ── Assertions on the report. ── |
| 329 | if !report.verified_ok { |
| 330 | return Err(err!("Migration did not verify."; Test, Data)); |
| 331 | } |
| 332 | // Present keys: acct(3)+sess(4)+lease(2)+sub(2)+a(3) = 14, + 5 sbody + 1 sync |
| 333 | // + 1 chunked = 21. The deleted key is a tombstone and must NOT be copied. |
| 334 | let expect_copied = 3 + 4 + 2 + 2 + 3 + 5 + 1 + 1; |
| 335 | if report.copied != expect_copied { |
| 336 | return Err(err!( |
| 337 | "Expected to copy {} present keys, copied {}.", expect_copied, report.copied; |
| 338 | Test, Data, Mismatch)); |
| 339 | } |
| 340 | if report.tombstones < 1 { |
| 341 | return Err(err!( |
| 342 | "Expected at least one tombstone (the deleted key) to be dropped, saw {}.", |
| 343 | report.tombstones; |
| 344 | Test, Data, Mismatch)); |
| 345 | } |
| 346 | // The sbody trap: every unreferenced sbody Complete key survived. |
| 347 | match report.per_prefix.get("sbody") { |
| 348 | Some(5) => (), |
| 349 | other => return Err(err!( |
| 350 | "Expected 5 unreferenced sbody Complete keys copied, per-prefix says {:?}. \ |
| 351 | Unreferenced Complete keys were dropped -- the data-loss case.", other; |
| 352 | Test, Data, Mismatch)), |
| 353 | } |
| 354 | // The deleted key must not appear under acct: acct should be exactly 3. |
| 355 | match report.per_prefix.get("acct") { |
| 356 | Some(3) => (), |
| 357 | other => return Err(err!( |
| 358 | "Expected acct prefix count 3 (the deleted key dropped), got {:?}.", other; |
| 359 | Test, Data, Mismatch)), |
| 360 | } |
| 361 | // The orphan-shrink property: the target is dramatically smaller than the |
| 362 | // source, because ~29 superseded generations of chunk data were dropped. |
| 363 | if !(report.target_bytes * 4 < report.source_bytes) { |
| 364 | return Err(err!( |
| 365 | "Expected the target to be far smaller than the source (orphans dropped): \ |
| 366 | source {} bytes, target {} bytes.", report.source_bytes, report.target_bytes; |
| 367 | Test, Data)); |
| 368 | } |
| 369 | |
| 370 | // ── The backup-safety property: no data or index file of the source changed. ── |
| 371 | let after = res!(snapshot(&src_root)); |
| 372 | let mut changed_data_index: Vec<String> = Vec::new(); |
| 373 | let mut other_changes: Vec<String> = Vec::new(); |
| 374 | for (path, bytes_before) in &before { |
| 375 | match after.get(path) { |
| 376 | Some(bytes_after) if bytes_after == bytes_before => (), |
| 377 | Some(_) => { |
| 378 | let name = path.to_string_lossy().to_string(); |
| 379 | if name.ends_with(".dat") || name.ends_with(".ind") { |
| 380 | changed_data_index.push(name); |
| 381 | } else { |
| 382 | other_changes.push(fmt!("{} (modified)", name)); |
| 383 | } |
| 384 | }, |
| 385 | None => other_changes.push(fmt!("{} (removed)", path.to_string_lossy())), |
| 386 | } |
| 387 | } |
| 388 | for path in after.keys() { |
| 389 | if !before.contains_key(path) { |
| 390 | let name = path.to_string_lossy().to_string(); |
| 391 | if name.ends_with(".dat") || name.ends_with(".ind") { |
| 392 | // A newly created (empty) live file is additive and does not |
| 393 | // touch existing data; note it but do not fail on it. |
| 394 | other_changes.push(fmt!("{} (new)", name)); |
| 395 | } else { |
| 396 | other_changes.push(fmt!("{} (new)", name)); |
| 397 | } |
| 398 | } |
| 399 | } |
| 400 | warn!(sync_log::stream(), |
| 401 | "Source after migration: {} data/index files changed in place; other changes: {:?}. \ |
| 402 | Source bytes before: {}.", |
| 403 | changed_data_index.len(), other_changes, src_bytes_before); |
| 404 | if !changed_data_index.is_empty() { |
| 405 | return Err(err!( |
| 406 | "A read-only migration modified {} of the source's data/index files in place: {:?}. \ |
| 407 | The rename-based runbook keeps the source as the only backup, so this would corrupt \ |
| 408 | the backup and is a blocker.", |
| 409 | changed_data_index.len(), changed_data_index; |
| 410 | Test, Data)); |
| 411 | } |
| 412 | |
| 413 | // ── Phase 3: reopen the target and confirm the live values read back. ── |
| 414 | { |
| 415 | let tgt_db = res!(setup::start_db( |
| 416 | tgt_root.clone(), |
| 417 | None, |
| 418 | res!(schemes_input()), |
| 419 | None, |
| 420 | false, |
| 421 | false, |
| 422 | )); |
| 423 | thread::sleep(Duration::from_secs(1)); |
| 424 | |
| 425 | // The chunked value reads back byte-identical to what was last written. |
| 426 | let expect = payload("v29", round_size(29)); |
| 427 | match res!(tgt_db.api().get_wait(&dat!("chunked:big"), None)) { |
| 428 | Some((v, _)) => { |
| 429 | if res!(v.as_bytes()) != res!(expect.as_bytes()) { |
| 430 | return Err(err!( |
| 431 | "Target chunked value does not match the last written value."; |
| 432 | Test, Data, Mismatch)); |
| 433 | } |
| 434 | }, |
| 435 | None => return Err(err!( |
| 436 | "Target is missing the chunked live value."; Test, Data, Missing)), |
| 437 | } |
| 438 | // Every sbody body survived and reads back. |
| 439 | for i in 0..5usize { |
| 440 | let expect = payload("sbody", 300 + i * 10); |
| 441 | let hk = fmt!("sbody:acctZ:{:08x}", i as u64); |
| 442 | match res!(tgt_db.api().get_wait(&dat!(hk.clone()), None)) { |
| 443 | Some((v, _)) => { |
| 444 | if res!(v.as_bytes()) != res!(expect.as_bytes()) { |
| 445 | return Err(err!( |
| 446 | "Target sbody value {} does not match.", hk; Test, Data, Mismatch)); |
| 447 | } |
| 448 | }, |
| 449 | None => return Err(err!( |
| 450 | "Target dropped an unreferenced sbody key: {}.", hk; Test, Data, Missing)), |
| 451 | } |
| 452 | } |
| 453 | // The deleted key reads as absent in the target, as in the source. |
| 454 | if res!(tgt_db.api().get_wait(&dat!("acct:deleteme"), None)).is_some() { |
| 455 | return Err(err!( |
| 456 | "Target holds the deleted key; a tombstone should have been dropped."; |
| 457 | Test, Data)); |
| 458 | } |
| 459 | res!(tgt_db.shutdown()); |
| 460 | } |
| 461 | thread::sleep(Duration::from_secs(1)); |
| 462 | |
| 463 | let _ = fs::remove_dir_all(&src_root); |
| 464 | let _ = fs::remove_dir_all(&tgt_root); |
| 465 | Ok(()) |
| 466 | } |
| 467 | |
| 468 | |
| 469 | fn scenario_unchunked_only() -> Outcome<()> { |
| 470 | |
| 471 | let src_root = res!(fresh_root("./test_db_migrate_u_src")); |
| 472 | let tgt_root = res!(empty_root("./test_db_migrate_u_tgt")); |
| 473 | let user = setup::Uid::default(); |
| 474 | |
| 475 | let payload = |tag: &str| -> Dat { Dat::Str(fmt!("{}-value-under-threshold", tag)) }; |
| 476 | |
| 477 | { |
| 478 | let db = res!(setup::start_db( |
| 479 | src_root.clone(), |
| 480 | Some(res!(base_cfg())), |
| 481 | res!(schemes_input()), |
| 482 | None, |
| 483 | false, |
| 484 | true, |
| 485 | )); |
| 486 | thread::sleep(Duration::from_secs(1)); |
| 487 | |
| 488 | // A handful of keys, each overwritten many times with gc off, so the |
| 489 | // source carries many superseded records the migration should drop. All |
| 490 | // values are well under the chunk threshold, so nothing chunks. |
| 491 | for i in 0..6usize { |
| 492 | let k = dat!(fmt!("rec:{:04}", i)); |
| 493 | for round in 0..20usize { |
| 494 | res!(store_wait(db.api(), k.clone(), payload(&fmt!("r{:02}", round)), user)); |
| 495 | } |
| 496 | } |
| 497 | thread::sleep(Duration::from_secs(2)); |
| 498 | res!(db.shutdown()); |
| 499 | } |
| 500 | thread::sleep(Duration::from_secs(1)); |
| 501 | |
| 502 | let report = { |
| 503 | let src_db = res!(setup::start_db( |
| 504 | src_root.clone(), None, res!(schemes_input()), None, false, false)); |
| 505 | thread::sleep(Duration::from_secs(1)); |
| 506 | let src_cfg = src_db.cfg().clone(); |
| 507 | let tgt_db = res!(setup::start_db( |
| 508 | tgt_root.clone(), Some(src_cfg), res!(schemes_input()), None, false, true)); |
| 509 | thread::sleep(Duration::from_secs(1)); |
| 510 | |
| 511 | let report = res!(migrate::migrate_live_set( |
| 512 | src_db.api(), tgt_db.api(), user, None, SCAN_WAIT)); |
| 513 | warn!(sync_log::stream(), "{}", report.summary()); |
| 514 | res!(tgt_db.shutdown()); |
| 515 | res!(src_db.shutdown()); |
| 516 | report |
| 517 | }; |
| 518 | thread::sleep(Duration::from_secs(1)); |
| 519 | |
| 520 | if !report.verified_ok { |
| 521 | return Err(err!("Unchunked migration did not verify."; Test, Data)); |
| 522 | } |
| 523 | if report.copied != 6 { |
| 524 | return Err(err!( |
| 525 | "Expected 6 live unchunked keys, copied {}.", report.copied; Test, Data, Mismatch)); |
| 526 | } |
| 527 | // No chunk records at all: the target is one current record per key, so it |
| 528 | // is smaller than a source carrying 20 generations of each. |
| 529 | if !(report.target_bytes < report.source_bytes) { |
| 530 | return Err(err!( |
| 531 | "Expected the unchunked target ({} bytes) smaller than the source ({} bytes).", |
| 532 | report.target_bytes, report.source_bytes; Test, Data)); |
| 533 | } |
| 534 | |
| 535 | let _ = fs::remove_dir_all(&src_root); |
| 536 | let _ = fs::remove_dir_all(&tgt_root); |
| 537 | Ok(()) |
| 538 | } |
| 539 | |
| 540 | |
| 541 | // The number of chunks the CURRENT value at `key` references, read from its |
| 542 | // bunch key. Errors if the key is not a chunked (bunch-key) value. |
| 543 | fn live_bunch_num_parts( |
| 544 | api: &oxedyne_fe2o3_o3db_sync::api::OzoneApi<{ setup::UID_LEN }, setup::Uid, Enc, Kh, Kh, Cs>, |
| 545 | key: &Dat, |
| 546 | ) |
| 547 | -> Outcome<usize> |
| 548 | { |
| 549 | let resp = res!(api.fetch_using_schemes(key, None)); |
| 550 | let enc = api.schemes().encrypter(); |
| 551 | match res!(resp.recv_daticle(enc, None)) { |
| 552 | (Some((Dat::Tup5u64(tup), _)), _) => Ok(PartKey(tup).num_parts() as usize), |
| 553 | (other, _) => Err(err!( |
| 554 | "Expected a chunked value (a bunch key Dat::Tup5u64) at {:?}, got {:?}.", |
| 555 | key, other; |
| 556 | Test, Data, Mismatch)), |
| 557 | } |
| 558 | } |
| 559 | |
| 560 | // Counts every chunk-data record (a stored key whose chunk index is >= 1) |
| 561 | // physically present across the store's index files. Mirrors the scan bot's |
| 562 | // index walk, but counts the chunk-data keys the scan elides. Used to prove |
| 563 | // orphaned chunks exist before a migration. |
| 564 | fn count_chunk_data_records( |
| 565 | api: &oxedyne_fe2o3_o3db_sync::api::OzoneApi<{ setup::UID_LEN }, setup::Uid, Enc, Kh, Kh, Cs>, |
| 566 | ) |
| 567 | -> Outcome<usize> |
| 568 | { |
| 569 | let csummer = api.schemes().checksummer().clone(); |
| 570 | let zdirs = res!(api.get_zone_dirs()); |
| 571 | let mut count = 0usize; |
| 572 | for (_zind, zdir) in &zdirs { |
| 573 | let mut fnums = Vec::new(); |
| 574 | for entry in res!(fs::read_dir(&zdir.dir)) { |
| 575 | let entry = res!(entry); |
| 576 | let path = entry.path(); |
| 577 | if !path.is_file() || ZoneDir::is_gc_temp_file(&path) { |
| 578 | continue; |
| 579 | } |
| 580 | if let Ok((fnum, FileType::Index)) = ZoneDir::ozone_file_number_and_type(&path) { |
| 581 | fnums.push(fnum); |
| 582 | } |
| 583 | } |
| 584 | for fnum in fnums { |
| 585 | let (_, file) = match zdir.open_ozone_file(fnum, &FileType::Index, &FileAccess::Reading) { |
| 586 | Ok(pair) => pair, |
| 587 | Err(_) => continue, // A file collected away between listing and opening. |
| 588 | }; |
| 589 | let mut reader = BufReader::new(file); |
| 590 | loop { |
| 591 | let key = match res!(StoredKey::<{ setup::UID_LEN }, setup::Uid>::load( |
| 592 | &mut reader, csummer.clone(), |
| 593 | )) { |
| 594 | None => break, |
| 595 | Some((skey, _, _)) => skey.into_key(), |
| 596 | }; |
| 597 | // Consume the index entry that follows the key, to reach the next. |
| 598 | match res!(StoredIndex::read(&mut reader, fnum, csummer.clone())) { |
| 599 | (Some(_), _) => (), |
| 600 | (None, _) => break, |
| 601 | } |
| 602 | if let Some(c) = key.index() { |
| 603 | if c >= 1 { |
| 604 | count += 1; |
| 605 | } |
| 606 | } |
| 607 | } |
| 608 | } |
| 609 | } |
| 610 | Ok(count) |
| 611 | } |
| 612 | |
| 613 | // Stores one pair and waits for every write part to be acknowledged. |
| 614 | fn store_wait( |
| 615 | api: &oxedyne_fe2o3_o3db_sync::api::OzoneApi<{ setup::UID_LEN }, setup::Uid, Enc, Kh, Kh, Cs>, |
| 616 | k: Dat, |
| 617 | v: Dat, |
| 618 | user: setup::Uid, |
| 619 | ) |
| 620 | -> Outcome<()> |
| 621 | { |
| 622 | let resp = res!(api.store(k, v, user)); |
| 623 | // The count, then each record written and durable: a writer's error is this test's failure. |
| 624 | res!(resp.recv_store_ack()); |
| 625 | Ok(()) |
| 626 | } |
| 627 | |
| 628 | // Deletes one key and waits for the tombstone write to be acknowledged. |
| 629 | fn delete_wait( |
| 630 | api: &oxedyne_fe2o3_o3db_sync::api::OzoneApi<{ setup::UID_LEN }, setup::Uid, Enc, Kh, Kh, Cs>, |
| 631 | k: Dat, |
| 632 | user: setup::Uid, |
| 633 | ) |
| 634 | -> Outcome<()> |
| 635 | { |
| 636 | let resp = api.responder(); |
| 637 | res!(api.delete_using_responder(&k, user, None, resp.clone())); |
| 638 | // The delete writes a single tombstone record; wait for its acknowledgement. |
| 639 | res!(resp.recv_delete_ack()); |
| 640 | Ok(()) |
| 641 | } |