oxedyne/fe2o3/fe2o3_o3db_sync/tests/rollover.rs
19.3 KiB, 80 runs
created by r1870400018:18290, 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 | //! Ozone data file rollover and supersession accounting test. |
| 2 | //! |
| 3 | //! A zone seals its live data file when the next record would push it past |
| 4 | //! `data_file_max_bytes`, and the writer bot asks its zone bot for the next |
| 5 | //! file number. Every record written into a file is recorded in that file's |
| 6 | //! `FileState.dmap`, and when the record is superseded by a later write of the |
| 7 | //! same key the entry is flipped from `Cur` to `Old`, which is what tells the |
| 8 | //! garbage collector how much of the file is reclaimable. |
| 9 | //! |
| 10 | //! This test drives many rollovers with a tiny file size limit and repeatedly |
| 11 | //! overwrites a small set of keys, so that supersessions land in files that |
| 12 | //! have already been sealed. It then asserts the accounting invariant: |
| 13 | //! |
| 14 | //! * one `dmap` entry per record written, |
| 15 | //! * exactly one `Cur` entry per distinct key, |
| 16 | //! * every other entry `Old`, |
| 17 | //! * an old-record counter that agrees with the map, |
| 18 | //! * not one error logged by any bot. |
| 19 | //! |
| 20 | //! When the zone hands a writer a file number that is already in use, the |
| 21 | //! receiving file bot replaces that file's `FileState` wholesale and the |
| 22 | //! entries recorded so far are lost. Every later supersession of one of those |
| 23 | //! records then fails in `FileState::register_old` with "a data entry starting |
| 24 | //! at position N in the FileState was not found", garbage is never registered, |
| 25 | //! and the store grows without bound. The invariant above catches that. |
| 26 | //! |
| 27 | //! The remaining phases cover a restart, which takes over the incomplete live |
| 28 | //! file left behind and rebuilds the caches from the index files, and garbage |
| 29 | //! collection, measured against the same churn run with collection disabled. |
| 30 | //! |
| 31 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 32 | //! Anthropic Claude |
| 33 | |
| 34 | use oxedyne_fe2o3_core::{ |
| 35 | prelude::*, |
| 36 | alt::Override, |
| 37 | }; |
| 38 | use oxedyne_fe2o3_crypto::enc::EncryptionScheme; |
| 39 | use oxedyne_fe2o3_hash::{ |
| 40 | csum::ChecksumScheme, |
| 41 | hash::HashScheme, |
| 42 | }; |
| 43 | use oxedyne_fe2o3_iop_db::api::{ |
| 44 | Database, |
| 45 | RestSchemesOverride, |
| 46 | }; |
| 47 | use oxedyne_fe2o3_jdat::prelude::*; |
| 48 | use oxedyne_fe2o3_o3db_sync::{ |
| 49 | base::{ |
| 50 | cfg::OzoneConfig, |
| 51 | constant, |
| 52 | index::WorkerInd, |
| 53 | }, |
| 54 | data::core::RestSchemesInput, |
| 55 | file::state::{ |
| 56 | DataState, |
| 57 | FileStateMap, |
| 58 | }, |
| 59 | test::setup, |
| 60 | O3db, |
| 61 | }; |
| 62 | |
| 63 | use std::{ |
| 64 | collections::BTreeMap, |
| 65 | fs, |
| 66 | path::{ |
| 67 | Path, |
| 68 | PathBuf, |
| 69 | }, |
| 70 | thread, |
| 71 | time::Duration, |
| 72 | }; |
| 73 | |
| 74 | const NKEYS: usize = 40; // distinct keys churned through the rollover |
| 75 | const NOVER: usize = 12; // overwrites of each key after its first write |
| 76 | const GC_ROUNDS: usize = 40; // overwrite rounds in the collection phase |
| 77 | |
| 78 | pub fn test_rollover(_filter: &'static str) -> Outcome<()> { |
| 79 | |
| 80 | let db_root = res!(canonical_dir("./test_db_rollover")); |
| 81 | |
| 82 | let enckey = [0x7bu8; 32]; |
| 83 | let aes_gcm = res!(EncryptionScheme::new_aes_256_gcm_with_key(&enckey[..])); |
| 84 | let crc32 = ChecksumScheme::new_crc32(); |
| 85 | let schms2: RestSchemesOverride<EncryptionScheme, HashScheme> = |
| 86 | RestSchemesOverride::default() |
| 87 | .set_encrypter(Override::Default(aes_gcm.clone())); |
| 88 | let schms2 = Some(&schms2); |
| 89 | let user = setup::Uid::default(); |
| 90 | |
| 91 | let schms_input = RestSchemesInput::new( |
| 92 | Some(aes_gcm.clone()), |
| 93 | None::<HashScheme>, |
| 94 | None::<HashScheme>, |
| 95 | Some(crc32.clone()), |
| 96 | ); |
| 97 | |
| 98 | // One writer bot per zone is the simple case: a rollover can only ever |
| 99 | // collide with the writer's own live file. |
| 100 | res!(supersession_accounting( |
| 101 | "one writer per zone", |
| 102 | &db_root, |
| 103 | &schms_input, |
| 104 | schms2, |
| 105 | user, |
| 106 | 1, |
| 107 | )); |
| 108 | // Two writer bots per zone is the shipped default, and is the case where a |
| 109 | // reused file number can be handed to a writer other than the one that |
| 110 | // already holds it. |
| 111 | res!(supersession_accounting( |
| 112 | "two writers per zone", |
| 113 | &db_root, |
| 114 | &schms_input, |
| 115 | schms2, |
| 116 | user, |
| 117 | 2, |
| 118 | )); |
| 119 | |
| 120 | // A restart takes over the incomplete live file left behind, which is the |
| 121 | // other way the zone counter could be left at a number already in use. |
| 122 | res!(restart_keeps_file_numbers_unique( |
| 123 | &db_root, |
| 124 | &schms_input, |
| 125 | schms2, |
| 126 | user, |
| 127 | )); |
| 128 | |
| 129 | res!(gc_reclaims_sealed_files( |
| 130 | &db_root, |
| 131 | &schms_input, |
| 132 | schms2, |
| 133 | user, |
| 134 | )); |
| 135 | |
| 136 | test!(sync_log::stream(), "Rollover test passed."); |
| 137 | Ok(()) |
| 138 | } |
| 139 | |
| 140 | /// Tiny data files, so rollovers happen every few records, and garbage |
| 141 | /// collection under the caller's control. |
| 142 | fn rollover_cfg(nwbots: u16) -> Outcome<OzoneConfig> { |
| 143 | let mut cfg = res!(setup::default_cfg()); |
| 144 | cfg.num_zones = 1; |
| 145 | cfg.num_cbots_per_zone = 2; |
| 146 | cfg.num_fbots_per_zone = 2; |
| 147 | cfg.num_igbots_per_zone = 2; |
| 148 | cfg.num_wbots_per_zone = nwbots; |
| 149 | cfg.data_file_max_bytes = 4_000; |
| 150 | cfg.rest_chunk_threshold = 3_000; |
| 151 | cfg.rest_chunk_bytes = 1_000; |
| 152 | cfg.zone_overrides = mapdat!{ |
| 153 | 1u16 => mapdat!{ "dir" => "", "max_size" => 100_000_000u64 }, |
| 154 | }.get_map().unwrap(); |
| 155 | Ok(cfg) |
| 156 | } |
| 157 | |
| 158 | /// Writes `NKEYS` keys `NOVER + 1` times each through many rollovers, then |
| 159 | /// checks that every record written is still accounted for in some file state, |
| 160 | /// and that exactly one record per key is current. |
| 161 | fn supersession_accounting( |
| 162 | label: &str, |
| 163 | db_root: &PathBuf, |
| 164 | schms_input: &RestSchemesInput< |
| 165 | EncryptionScheme, |
| 166 | HashScheme, |
| 167 | HashScheme, |
| 168 | ChecksumScheme, |
| 169 | >, |
| 170 | schms2: Option<&RestSchemesOverride<EncryptionScheme, HashScheme>>, |
| 171 | user: setup::Uid, |
| 172 | nwbots: u16, |
| 173 | ) |
| 174 | -> Outcome<()> |
| 175 | { |
| 176 | test!(sync_log::stream(), "+--- rollover: {} ---", label); |
| 177 | |
| 178 | let cfg = res!(rollover_cfg(nwbots)); |
| 179 | let db = res!(setup::start_db( |
| 180 | db_root.clone(), |
| 181 | Some(cfg.clone()), |
| 182 | schms_input.clone(), |
| 183 | None, |
| 184 | false, // gc off: nothing is retired, so the accounting is exact. |
| 185 | true, // wipe: start from an empty zone directory. |
| 186 | )); |
| 187 | |
| 188 | thread::sleep(Duration::from_millis(200)); |
| 189 | |
| 190 | let mut nwrites = 0usize; |
| 191 | for round in 0..(NOVER + 1) { |
| 192 | for k in 0..NKEYS { |
| 193 | res!(db.insert( |
| 194 | dat!(fmt!("rk{:03}", k)), |
| 195 | dat!(fmt!("rv{:03}r{:02}", k, round)), |
| 196 | user, |
| 197 | schms2, |
| 198 | )); |
| 199 | nwrites += 1; |
| 200 | } |
| 201 | } |
| 202 | |
| 203 | // Let the trailing UpdateData and ScheduleOld messages drain. |
| 204 | thread::sleep(Duration::from_millis(500)); |
| 205 | |
| 206 | // Read every key back, to confirm the rollover left the data readable. |
| 207 | for k in 0..NKEYS { |
| 208 | let key = dat!(fmt!("rk{:03}", k)); |
| 209 | match res!(db.get(&key, schms2)) { |
| 210 | Some((val, _meta)) => { |
| 211 | let expected = dat!(fmt!("rv{:03}r{:02}", k, NOVER)); |
| 212 | if val != expected { |
| 213 | return Err(err!( |
| 214 | "{}: key {} round-tripped as {:?}, expected {:?}.", |
| 215 | label, k, val, expected; |
| 216 | Test, Mismatch)); |
| 217 | } |
| 218 | }, |
| 219 | None => return Err(err!( |
| 220 | "{}: key {} not found after {} writes.", label, k, nwrites; |
| 221 | Test, Missing)), |
| 222 | } |
| 223 | } |
| 224 | |
| 225 | let states = res!(db.api().collect_file_states(constant::USER_REQUEST_WAIT)); |
| 226 | let (nfiles, ncur, nold, noldcnt) = count_data_states(&states); |
| 227 | |
| 228 | test!(sync_log::stream(), |
| 229 | "{}: {} writes over {} tracked files: {} current, {} old, {} tracked \ |
| 230 | total, {} counted old.", |
| 231 | label, nwrites, nfiles, ncur, nold, ncur + nold, noldcnt); |
| 232 | |
| 233 | if ncur + nold != nwrites { |
| 234 | return Err(err!( |
| 235 | "{}: {} records were written but the file states track {} \ |
| 236 | ({} current + {} old). Entries lost from a FileState can never be \ |
| 237 | flagged old, so their bytes can never be reclaimed.", |
| 238 | label, nwrites, ncur + nold, ncur, nold; |
| 239 | Test, Mismatch, Data)); |
| 240 | } |
| 241 | if ncur != NKEYS { |
| 242 | return Err(err!( |
| 243 | "{}: {} distinct keys were written but {} records are flagged \ |
| 244 | current.", label, NKEYS, ncur; |
| 245 | Test, Mismatch, Data)); |
| 246 | } |
| 247 | if noldcnt != nold { |
| 248 | return Err(err!( |
| 249 | "{}: {} records are flagged old in the record maps but the \ |
| 250 | old-record counters total {}. The two disagree, so the \ |
| 251 | reclaimable byte total the collector trusts is wrong.", |
| 252 | label, nold, noldcnt; |
| 253 | Test, Mismatch, Data)); |
| 254 | } |
| 255 | res!(assert_no_bot_errors(&db, label)); |
| 256 | |
| 257 | res!(db.shutdown()); |
| 258 | thread::sleep(Duration::from_millis(200)); |
| 259 | |
| 260 | test!(sync_log::stream(), "+--- rollover: {} : passed ---", label); |
| 261 | Ok(()) |
| 262 | } |
| 263 | |
| 264 | /// Churns, shuts the database down, restarts it over the files left behind and |
| 265 | /// churns again. On restart a writer bot takes over the incomplete live file, |
| 266 | /// so the zone live file counter must be left above it: if the counter is |
| 267 | /// wound back below a file already in use, the first rollover after the restart |
| 268 | /// hands that number out again and the file's record entries are lost. |
| 269 | fn restart_keeps_file_numbers_unique( |
| 270 | db_root: &PathBuf, |
| 271 | schms_input: &RestSchemesInput< |
| 272 | EncryptionScheme, |
| 273 | HashScheme, |
| 274 | HashScheme, |
| 275 | ChecksumScheme, |
| 276 | >, |
| 277 | schms2: Option<&RestSchemesOverride<EncryptionScheme, HashScheme>>, |
| 278 | user: setup::Uid, |
| 279 | ) |
| 280 | -> Outcome<()> |
| 281 | { |
| 282 | test!(sync_log::stream(), "+--- rollover: restart ---"); |
| 283 | |
| 284 | let cfg = res!(rollover_cfg(1)); |
| 285 | |
| 286 | // First run, from empty. |
| 287 | let db = res!(setup::start_db( |
| 288 | db_root.clone(), |
| 289 | Some(cfg.clone()), |
| 290 | schms_input.clone(), |
| 291 | None, |
| 292 | false, // gc off. |
| 293 | true, // wipe. |
| 294 | )); |
| 295 | thread::sleep(Duration::from_millis(200)); |
| 296 | let n1 = res!(churn(&db, "sk", user, schms2, NOVER + 1)); |
| 297 | thread::sleep(Duration::from_millis(500)); |
| 298 | res!(assert_no_bot_errors(&db, "restart: first run")); |
| 299 | res!(db.shutdown()); |
| 300 | thread::sleep(Duration::from_millis(500)); |
| 301 | |
| 302 | // Second run, over the files the first left behind. |
| 303 | let db = res!(setup::start_db( |
| 304 | db_root.clone(), |
| 305 | Some(cfg.clone()), |
| 306 | schms_input.clone(), |
| 307 | None, |
| 308 | false, // gc off. |
| 309 | false, // no wipe: this is the point of the phase. |
| 310 | )); |
| 311 | thread::sleep(Duration::from_millis(500)); |
| 312 | res!(assert_no_bot_errors(&db, "restart: after reload")); |
| 313 | let n2 = res!(churn(&db, "sk", user, schms2, NOVER + 1)); |
| 314 | thread::sleep(Duration::from_millis(500)); |
| 315 | |
| 316 | let states = res!(db.api().collect_file_states(constant::USER_REQUEST_WAIT)); |
| 317 | let (nfiles, ncur, nold, noldcnt) = count_data_states(&states); |
| 318 | |
| 319 | test!(sync_log::stream(), |
| 320 | "restart: {} writes before and {} after, {} tracked files: {} current, \ |
| 321 | {} old, {} tracked total, {} counted old.", |
| 322 | n1, n2, nfiles, ncur, nold, ncur + nold, noldcnt); |
| 323 | |
| 324 | if ncur + nold != n1 + n2 { |
| 325 | return Err(err!( |
| 326 | "restart: {} records were written across the two runs but the file \ |
| 327 | states track {} ({} current + {} old).", |
| 328 | n1 + n2, ncur + nold, ncur, nold; |
| 329 | Test, Mismatch, Data)); |
| 330 | } |
| 331 | if ncur != NKEYS { |
| 332 | return Err(err!( |
| 333 | "restart: {} distinct keys were written but {} records are flagged \ |
| 334 | current.", NKEYS, ncur; |
| 335 | Test, Mismatch, Data)); |
| 336 | } |
| 337 | if noldcnt != nold { |
| 338 | return Err(err!( |
| 339 | "restart: {} records are flagged old in the record maps but the \ |
| 340 | old-record counters total {}.", nold, noldcnt; |
| 341 | Test, Mismatch, Data)); |
| 342 | } |
| 343 | res!(assert_no_bot_errors(&db, "restart: second run")); |
| 344 | |
| 345 | res!(db.shutdown()); |
| 346 | thread::sleep(Duration::from_millis(200)); |
| 347 | |
| 348 | test!(sync_log::stream(), "+--- rollover: restart : passed ---"); |
| 349 | Ok(()) |
| 350 | } |
| 351 | |
| 352 | fn churn< |
| 353 | ENC: oxedyne_fe2o3_iop_crypto::enc::Encrypter + 'static, |
| 354 | KH: oxedyne_fe2o3_iop_hash::api::Hasher + 'static, |
| 355 | PR: oxedyne_fe2o3_iop_hash::api::Hasher + 'static, |
| 356 | CS: oxedyne_fe2o3_iop_hash::csum::Checksummer + 'static, |
| 357 | >( |
| 358 | db: &O3db<{ setup::UID_LEN }, setup::Uid, ENC, KH, PR, CS>, |
| 359 | prefix: &str, |
| 360 | user: setup::Uid, |
| 361 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 362 | rounds: usize, |
| 363 | ) |
| 364 | -> Outcome<usize> |
| 365 | { |
| 366 | let mut nwrites = 0; |
| 367 | for round in 0..rounds { |
| 368 | for k in 0..NKEYS { |
| 369 | res!(db.insert( |
| 370 | dat!(fmt!("{}{:03}", prefix, k)), |
| 371 | dat!(fmt!("{}v{:03}r{:02}", prefix, k, round)), |
| 372 | user, |
| 373 | schms2, |
| 374 | )); |
| 375 | nwrites += 1; |
| 376 | } |
| 377 | } |
| 378 | Ok(nwrites) |
| 379 | } |
| 380 | |
| 381 | /// Heavy overwrite churn over many sealed data files, run twice: once with |
| 382 | /// garbage collection off and once with it on. The store with collection on |
| 383 | /// must be a fraction of the size of the store without it. The run with |
| 384 | /// collection off is the oracle: it is the same workload measured with the only |
| 385 | /// mechanism that can reclaim bytes disabled, so the comparison cannot be |
| 386 | /// satisfied by anything other than reclamation actually happening. |
| 387 | fn gc_reclaims_sealed_files( |
| 388 | db_root: &PathBuf, |
| 389 | schms_input: &RestSchemesInput< |
| 390 | EncryptionScheme, |
| 391 | HashScheme, |
| 392 | HashScheme, |
| 393 | ChecksumScheme, |
| 394 | >, |
| 395 | schms2: Option<&RestSchemesOverride<EncryptionScheme, HashScheme>>, |
| 396 | user: setup::Uid, |
| 397 | ) |
| 398 | -> Outcome<()> |
| 399 | { |
| 400 | test!(sync_log::stream(), "+--- rollover: garbage collection ---"); |
| 401 | |
| 402 | let uncollected = res!(churn_and_measure( |
| 403 | db_root, schms_input, schms2, user, false)); |
| 404 | let collected = res!(churn_and_measure( |
| 405 | db_root, schms_input, schms2, user, true)); |
| 406 | |
| 407 | test!(sync_log::stream(), |
| 408 | "gc: {} bytes of data files with collection off, {} with it on.", |
| 409 | uncollected, collected); |
| 410 | |
| 411 | // The live set is a fortieth of everything written, so a working collector |
| 412 | // should leave well under a quarter of the uncollected store. Anything at |
| 413 | // or above it means sealed files are not being reclaimed. |
| 414 | if collected * 4 >= uncollected { |
| 415 | return Err(err!( |
| 416 | "Garbage collection reclaimed little or nothing: the same churn \ |
| 417 | left {} bytes of data files with collection on against {} with it \ |
| 418 | off.", collected, uncollected; |
| 419 | Test, Mismatch, Data)); |
| 420 | } |
| 421 | |
| 422 | test!(sync_log::stream(), "+--- rollover: garbage collection : passed ---"); |
| 423 | Ok(()) |
| 424 | } |
| 425 | |
| 426 | /// The database is freshly wiped first, so the returned byte total covers only |
| 427 | /// this run's churn. |
| 428 | fn churn_and_measure( |
| 429 | db_root: &PathBuf, |
| 430 | schms_input: &RestSchemesInput< |
| 431 | EncryptionScheme, |
| 432 | HashScheme, |
| 433 | HashScheme, |
| 434 | ChecksumScheme, |
| 435 | >, |
| 436 | schms2: Option<&RestSchemesOverride<EncryptionScheme, HashScheme>>, |
| 437 | user: setup::Uid, |
| 438 | gc_on: bool, |
| 439 | ) |
| 440 | -> Outcome<u64> |
| 441 | { |
| 442 | let label = if gc_on { "collection on" } else { "collection off" }; |
| 443 | let cfg = res!(rollover_cfg(1)); |
| 444 | let db = res!(setup::start_db( |
| 445 | db_root.clone(), |
| 446 | Some(cfg.clone()), |
| 447 | schms_input.clone(), |
| 448 | None, |
| 449 | gc_on, |
| 450 | true, // wipe. |
| 451 | )); |
| 452 | |
| 453 | thread::sleep(Duration::from_millis(200)); |
| 454 | |
| 455 | // Enough churn that the superseded copies vastly outweigh the live set. |
| 456 | let mut nwrites = 0usize; |
| 457 | for round in 0..GC_ROUNDS { |
| 458 | for k in 0..NKEYS { |
| 459 | res!(db.insert( |
| 460 | dat!(fmt!("gk{:03}", k)), |
| 461 | dat!(fmt!("gv{:03}r{:02}", k, round)), |
| 462 | user, |
| 463 | schms2, |
| 464 | )); |
| 465 | nwrites += 1; |
| 466 | } |
| 467 | } |
| 468 | |
| 469 | // Give the gbots time to finish transcribing. |
| 470 | thread::sleep(Duration::from_secs(2)); |
| 471 | |
| 472 | let bytes = res!(zone_data_bytes(&db_root)); |
| 473 | let states = res!(db.api().collect_file_states(constant::USER_REQUEST_WAIT)); |
| 474 | let (nfiles, ncur, nold, _) = count_data_states(&states); |
| 475 | |
| 476 | test!(sync_log::stream(), |
| 477 | "gc {}: {} writes, {} bytes of data files remain over {} tracked \ |
| 478 | files, {} current and {} old records tracked.", |
| 479 | label, nwrites, bytes, nfiles, ncur, nold); |
| 480 | |
| 481 | if nold > nwrites - NKEYS { |
| 482 | return Err(err!( |
| 483 | "gc {}: more records are flagged old ({}) than were superseded ({}).", |
| 484 | label, nold, nwrites - NKEYS; |
| 485 | Test, Mismatch, Data)); |
| 486 | } |
| 487 | res!(assert_no_bot_errors(&db, label)); |
| 488 | |
| 489 | res!(db.shutdown()); |
| 490 | thread::sleep(Duration::from_millis(200)); |
| 491 | |
| 492 | Ok(bytes) |
| 493 | } |
| 494 | |
| 495 | /// Fails if any bot in the database has logged an error. The flag-as-old |
| 496 | /// failure this test exists for is logged by a file bot rather than returned |
| 497 | /// to the caller, so this is how the log storm itself is asserted away. |
| 498 | fn assert_no_bot_errors< |
| 499 | ENC: oxedyne_fe2o3_iop_crypto::enc::Encrypter + 'static, |
| 500 | KH: oxedyne_fe2o3_iop_hash::api::Hasher + 'static, |
| 501 | PR: oxedyne_fe2o3_iop_hash::api::Hasher + 'static, |
| 502 | CS: oxedyne_fe2o3_iop_hash::csum::Checksummer + 'static, |
| 503 | >( |
| 504 | db: &O3db<{ setup::UID_LEN }, setup::Uid, ENC, KH, PR, CS>, |
| 505 | label: &str, |
| 506 | ) |
| 507 | -> Outcome<()> |
| 508 | { |
| 509 | let (errs, nbots) = res!(db.api().bot_error_count(constant::USER_REQUEST_WAIT)); |
| 510 | test!(sync_log::stream(), "{}: {} bots reported {} errors.", label, nbots, errs); |
| 511 | if errs > 0 { |
| 512 | return Err(err!( |
| 513 | "{}: the {} bots of the database logged {} errors during the run. \ |
| 514 | A supersession that cannot be registered is logged, not returned, \ |
| 515 | so any error here means the rollover bookkeeping is still wrong.", |
| 516 | label, nbots, errs; |
| 517 | Test, Mismatch, Data)); |
| 518 | } |
| 519 | Ok(()) |
| 520 | } |
| 521 | |
| 522 | /// Across every file bot in every zone: the number of tracked files, the |
| 523 | /// current and old record counts, and the old count the file states track. |
| 524 | fn count_data_states( |
| 525 | states: &BTreeMap<WorkerInd, FileStateMap>, |
| 526 | ) |
| 527 | -> (usize, usize, usize, usize) |
| 528 | { |
| 529 | let mut nfiles = 0; |
| 530 | let mut ncur = 0; |
| 531 | let mut nold = 0; |
| 532 | let mut noldcnt = 0; |
| 533 | for (_wind, fstates) in states { |
| 534 | for (_fnum, fstat) in fstates.map() { |
| 535 | nfiles += 1; |
| 536 | noldcnt += fstat.get_old_count(); |
| 537 | for (_start, dstat) in fstat.data_map() { |
| 538 | match dstat { |
| 539 | DataState::Cur => ncur += 1, |
| 540 | DataState::Old => nold += 1, |
| 541 | } |
| 542 | } |
| 543 | } |
| 544 | } |
| 545 | (nfiles, ncur, nold, noldcnt) |
| 546 | } |
| 547 | |
| 548 | fn zone_data_bytes(db_root: &Path) -> Outcome<u64> { |
| 549 | let mut total = 0u64; |
| 550 | res!(walk_data_files(db_root, &mut total)); |
| 551 | Ok(total) |
| 552 | } |
| 553 | |
| 554 | fn walk_data_files(dir: &Path, total: &mut u64) -> Outcome<()> { |
| 555 | for entry in res!(fs::read_dir(dir)) { |
| 556 | let entry = res!(entry); |
| 557 | let path = entry.path(); |
| 558 | if path.is_dir() { |
| 559 | res!(walk_data_files(&path, total)); |
| 560 | } else if path.extension().map(|e| e == constant::DATA_FILE_EXT).unwrap_or(false) { |
| 561 | let meta = res!(entry.metadata()); |
| 562 | *total += meta.len(); |
| 563 | } |
| 564 | } |
| 565 | Ok(()) |
| 566 | } |
| 567 | |
| 568 | /// Creates the directory if it does not exist. |
| 569 | fn canonical_dir(p: &str) -> Outcome<PathBuf> { |
| 570 | match Path::new(p).canonicalize() { |
| 571 | Ok(path) => Ok(path), |
| 572 | Err(_) => { |
| 573 | res!(fs::create_dir_all(p)); |
| 574 | match Path::new(p).canonicalize() { |
| 575 | Ok(path) => Ok(path), |
| 576 | Err(e) => Err(err!(e, "Cannot canonicalise {:?}.", p; IO, Path)), |
| 577 | } |
| 578 | }, |
| 579 | } |
| 580 | } |