oxedyne/fe2o3/fe2o3_o3db_sync/tests/chunk_leak.rs
25.1 KiB, 3 runs
created by r1870400018:38082, 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 | //! Chunked-value overwrite and delete must not leak chunk records. |
| 2 | //! |
| 3 | //! A value whose encoding exceeds the chunk threshold is split: the bytes live in chunk records |
| 4 | //! and the user's key holds a `Tup5u64` bunch key naming them. Garbage collection is |
| 5 | //! supersession-based -- a record becomes reclaimable only when the same key is written again -- |
| 6 | //! so when the chunk records were keyed by a random per-operation ticket, overwriting or deleting |
| 7 | //! such a value superseded the bunch key but left every chunk record under an unreachable key. |
| 8 | //! Nothing rewrote those keys, so nothing flagged them old, so the collector never saw them: every |
| 9 | //! overwrite leaked a whole value's worth of chunk bytes, permanently. On the gateway this reached |
| 10 | //! ~19 GB. |
| 11 | //! |
| 12 | //! The fix derives the chunk set identifier from the key rather than from a random ticket, so an |
| 13 | //! overwrite writes the new chunks under the same keys as the old and the ordinary supersession |
| 14 | //! path reclaims them; and the delete path reads the current bunch key and tombstones each chunk. |
| 15 | //! Reads are unaffected: a value's chunk keys are reconstructed from the set_id stored in its bunch |
| 16 | //! key, so values written under the old random scheme still read after the switch. |
| 17 | //! |
| 18 | //! These tests force chunking with an explicit, prod-like threshold (not the drifting dev |
| 19 | //! default). `overwrite_reclaims_chunks` hammers one key -- which is also the concurrent |
| 20 | //! same-key supersession that exposed the accounting race the fix closes -- and insists the |
| 21 | //! footprint stays near the one-value baseline. `delete_reclaims_chunks` insists a deleted |
| 22 | //! chunked value's chunks are reclaimed, not only its bunch key tombstoned. `shrink_is_bounded` |
| 23 | //! insists repeated overwrites at a smaller geometry reach a bounded steady state. |
| 24 | //! `old_random_key_value_still_reads` writes a value under a simulated pre-upgrade (random) |
| 25 | //! set_id and insists it still reads, before and after a normal overwrite. `churn_survives_restart` |
| 26 | //! insists the on-disk accounting a heavy same-key churn leaves behind is self-consistent across a |
| 27 | //! restart. |
| 28 | //! |
| 29 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 30 | //! Anthropic Claude |
| 31 | |
| 32 | use oxedyne_fe2o3_core::{ |
| 33 | prelude::*, |
| 34 | alt::Override, |
| 35 | rand::Rand, |
| 36 | }; |
| 37 | use oxedyne_fe2o3_crypto::enc::EncryptionScheme; |
| 38 | use oxedyne_fe2o3_hash::{ |
| 39 | csum::ChecksumScheme, |
| 40 | hash::HashScheme, |
| 41 | }; |
| 42 | use oxedyne_fe2o3_iop_db::api::{ |
| 43 | Database, |
| 44 | RestSchemesOverride, |
| 45 | }; |
| 46 | use oxedyne_fe2o3_jdat::prelude::*; |
| 47 | use oxedyne_fe2o3_o3db_sync::{ |
| 48 | base::constant, |
| 49 | data::core::RestSchemesInput, |
| 50 | test::setup, |
| 51 | }; |
| 52 | |
| 53 | use std::{ |
| 54 | fs, |
| 55 | path::{ |
| 56 | Path, |
| 57 | PathBuf, |
| 58 | }, |
| 59 | thread, |
| 60 | time::Duration, |
| 61 | }; |
| 62 | |
| 63 | const NOVERWRITES: usize = 20; // overwrites of the one chunked key |
| 64 | const VALUE_BYTES: usize = 6_000; // comfortably over the threshold, several chunks |
| 65 | const SMALL_BYTES: usize = 3_200; // still over the threshold, but fewer chunks |
| 66 | |
| 67 | /// A byte string of the given size, filled deterministically so each version is distinct on disk. |
| 68 | fn value_of(seed: u8, len: usize) -> Dat { |
| 69 | let mut v = vec![0u8; len]; |
| 70 | for (i, b) in v.iter_mut().enumerate() { |
| 71 | *b = seed.wrapping_add((i % 251) as u8); |
| 72 | } |
| 73 | Dat::BU32(v) |
| 74 | } |
| 75 | |
| 76 | fn chunky_value(seed: u8) -> Dat { value_of(seed, VALUE_BYTES) } |
| 77 | |
| 78 | pub fn test_chunk_leak(_filter: &'static str) -> Outcome<()> { |
| 79 | |
| 80 | let db_root = res!(canonical_dir("./test_db_chunk_leak")); |
| 81 | |
| 82 | let mut enckey = [0u8; 32]; |
| 83 | Rand::fill_u8(&mut enckey); |
| 84 | let aes_gcm = res!(EncryptionScheme::new_aes_256_gcm_with_key(&enckey[..])); |
| 85 | let crc32 = ChecksumScheme::new_crc32(); |
| 86 | let schms2: RestSchemesOverride<EncryptionScheme, HashScheme> = |
| 87 | RestSchemesOverride::default().set_encrypter(Override::Default(aes_gcm.clone())); |
| 88 | let schms2 = Some(&schms2); |
| 89 | let user = setup::Uid::default(); |
| 90 | let schms_input = RestSchemesInput::new( |
| 91 | Some(aes_gcm.clone()), |
| 92 | None::<HashScheme>, |
| 93 | None::<HashScheme>, |
| 94 | Some(crc32.clone()), |
| 95 | ); |
| 96 | |
| 97 | // Force chunking deterministically: an explicit threshold the value exceeds, a small chunk |
| 98 | // size so a value is several chunks, and a data file small enough that files seal and garbage |
| 99 | // collection runs -- but large enough to satisfy the threshold-to-file-size ratio check. |
| 100 | let mut cfg = res!(setup::default_cfg()); |
| 101 | cfg.data_file_max_bytes = 16_000; |
| 102 | cfg.rest_chunk_threshold = 3_000; |
| 103 | cfg.rest_chunk_bytes = 1_000; |
| 104 | cfg.sync_on_write = true; |
| 105 | cfg.zone_overrides = DaticleMap::new(); |
| 106 | |
| 107 | res!(overwrite_reclaims_chunks(&db_root, &cfg, &schms_input, schms2, user)); |
| 108 | res!(delete_reclaims_chunks(&db_root, &cfg, &schms_input, schms2, user)); |
| 109 | res!(shrink_is_bounded(&db_root, &cfg, &schms_input, schms2, user)); |
| 110 | res!(old_random_key_value_still_reads(&db_root, &cfg, &schms_input, schms2, user)); |
| 111 | res!(churn_survives_restart(&db_root, &cfg, &schms_input, schms2, user)); |
| 112 | |
| 113 | Ok(()) |
| 114 | } |
| 115 | |
| 116 | /// Overwriting a chunked value many times must leave the data footprint near the one-value |
| 117 | /// baseline, not growing by a value's worth of chunks each time. Hammering one key is also the |
| 118 | /// concurrent same-key supersession that exposed the accounting race the fix closes: a bounded |
| 119 | /// footprint here means every superseded chunk was flagged old and collected without a |
| 120 | /// mis-accounting abort. A genuine accounting inconsistency would abort collection in a bot |
| 121 | /// thread and the footprint would grow, tripping the bound below. |
| 122 | fn overwrite_reclaims_chunks( |
| 123 | db_root: &PathBuf, |
| 124 | cfg: &oxedyne_fe2o3_o3db_sync::base::cfg::OzoneConfig, |
| 125 | schms_input: &RestSchemesInput<EncryptionScheme, HashScheme, HashScheme, ChecksumScheme>, |
| 126 | schms2: Option<&RestSchemesOverride<EncryptionScheme, HashScheme>>, |
| 127 | user: setup::Uid, |
| 128 | ) |
| 129 | -> Outcome<()> |
| 130 | { |
| 131 | test!(sync_log::stream(), "+--- chunk leak: overwrite reclaims chunks ---"); |
| 132 | |
| 133 | let db = res!(setup::start_db( |
| 134 | db_root.clone(), |
| 135 | Some(cfg.clone()), |
| 136 | schms_input.clone(), |
| 137 | None, |
| 138 | true, // gc on |
| 139 | true, // wipe |
| 140 | )); |
| 141 | |
| 142 | let key = dat!("chunked:churned"); |
| 143 | |
| 144 | // First write, then settle, and measure the single-value baseline. |
| 145 | let (exists, nchunks) = res!(db.insert(key.clone(), chunky_value(0), user, schms2)); |
| 146 | if nchunks < 2 { |
| 147 | return Err(err!( |
| 148 | "The test value was not chunked ({} chunk(s)); the threshold or value size is wrong, \ |
| 149 | so the test would not exercise the leak.", nchunks; |
| 150 | Test, Invalid, Configuration)); |
| 151 | } |
| 152 | if exists { |
| 153 | return Err(err!("The churned key unexpectedly already existed on a fresh store."; |
| 154 | Test, Invalid, Data)); |
| 155 | } |
| 156 | thread::sleep(Duration::from_secs(2)); |
| 157 | let baseline = res!(zone_data_bytes(db_root)); |
| 158 | test!(sync_log::stream(), |
| 159 | "chunk leak: one {}-byte value is {} chunks, {} bytes of data files at baseline.", |
| 160 | VALUE_BYTES, nchunks, baseline); |
| 161 | |
| 162 | // Overwrite the same key many times. Each overwrite writes new chunk bytes under the same |
| 163 | // key-derived chunk keys; the prior value's chunks must be reclaimed rather than left behind. |
| 164 | for i in 1..=NOVERWRITES { |
| 165 | res!(db.insert(key.clone(), chunky_value(i as u8), user, schms2)); |
| 166 | } |
| 167 | thread::sleep(Duration::from_secs(3)); |
| 168 | let after = res!(zone_data_bytes(db_root)); |
| 169 | |
| 170 | // The value must still read back as the last one written. |
| 171 | match res!(db.get(&key, schms2)) { |
| 172 | Some((got, _)) => if got != chunky_value(NOVERWRITES as u8) { |
| 173 | return Err(err!( |
| 174 | "After {} overwrites the chunked value came back changed.", NOVERWRITES; |
| 175 | Test, Invalid, Data)); |
| 176 | }, |
| 177 | None => return Err(err!( |
| 178 | "After {} overwrites the chunked value cannot be read back.", NOVERWRITES; |
| 179 | Test, Missing, Data)), |
| 180 | } |
| 181 | |
| 182 | test!(sync_log::stream(), |
| 183 | "chunk leak: after {} overwrites, {} bytes of data files (baseline {}, ratio {:.1}x).", |
| 184 | NOVERWRITES, after, baseline, after as f64 / baseline as f64); |
| 185 | |
| 186 | // Without reclamation the footprint grows by a value's worth of chunks per overwrite -- about |
| 187 | // NOVERWRITES times the baseline. With it, the collector reclaims the superseded chunks and |
| 188 | // only the current value survives, a small multiple of the baseline. The bound sits well |
| 189 | // below the leak and above the reclaimed steady state. |
| 190 | let bound = baseline.saturating_mul(6); |
| 191 | if after >= bound { |
| 192 | return Err(err!( |
| 193 | "Chunk leak on overwrite: {} bytes of data files after {} overwrites of a {}-chunk \ |
| 194 | value, against a one-value baseline of {} bytes. The superseded chunk records are not \ |
| 195 | being reclaimed -- the footprint is growing with the overwrite count instead of \ |
| 196 | staying bounded.", after, NOVERWRITES, nchunks, baseline; |
| 197 | Test, Mismatch, Data)); |
| 198 | } |
| 199 | |
| 200 | res!(db.shutdown()); |
| 201 | thread::sleep(Duration::from_millis(200)); |
| 202 | test!(sync_log::stream(), "+--- chunk leak: overwrite reclaims chunks : passed ---"); |
| 203 | Ok(()) |
| 204 | } |
| 205 | |
| 206 | /// Deleting a chunked value must reclaim its chunk records, not only tombstone the bunch key. |
| 207 | fn delete_reclaims_chunks( |
| 208 | db_root: &PathBuf, |
| 209 | cfg: &oxedyne_fe2o3_o3db_sync::base::cfg::OzoneConfig, |
| 210 | schms_input: &RestSchemesInput<EncryptionScheme, HashScheme, HashScheme, ChecksumScheme>, |
| 211 | schms2: Option<&RestSchemesOverride<EncryptionScheme, HashScheme>>, |
| 212 | user: setup::Uid, |
| 213 | ) |
| 214 | -> Outcome<()> |
| 215 | { |
| 216 | test!(sync_log::stream(), "+--- chunk leak: delete reclaims chunks ---"); |
| 217 | |
| 218 | let db = res!(setup::start_db( |
| 219 | db_root.clone(), |
| 220 | Some(cfg.clone()), |
| 221 | schms_input.clone(), |
| 222 | None, |
| 223 | true, // gc on |
| 224 | true, // wipe |
| 225 | )); |
| 226 | |
| 227 | // Enough chunked values under distinct keys that the reclaimed chunk bytes dominate the fixed |
| 228 | // per-key tombstone residual, so the drop on delete is unambiguous. |
| 229 | let nkeys = 16usize; |
| 230 | for k in 0..nkeys { |
| 231 | let (_, nchunks) = res!(db.insert(dat!(fmt!("del:{:03}", k)), chunky_value(k as u8), user, schms2)); |
| 232 | if nchunks < 2 { |
| 233 | return Err(err!( |
| 234 | "The delete-case value was not chunked ({} chunk(s)).", nchunks; |
| 235 | Test, Invalid, Configuration)); |
| 236 | } |
| 237 | } |
| 238 | thread::sleep(Duration::from_secs(3)); |
| 239 | let with_values = res!(zone_data_bytes(db_root)); |
| 240 | |
| 241 | // Delete them all. |
| 242 | for k in 0..nkeys { |
| 243 | let existed = res!(db.delete(&dat!(fmt!("del:{:03}", k)), user, schms2)); |
| 244 | if !existed { |
| 245 | return Err(err!("Delete reported the chunked key {} absent.", k; Test, Missing, Data)); |
| 246 | } |
| 247 | } |
| 248 | thread::sleep(Duration::from_secs(4)); |
| 249 | let after_delete = res!(zone_data_bytes(db_root)); |
| 250 | |
| 251 | // Every deleted key reads back as absent. |
| 252 | for k in 0..nkeys { |
| 253 | if res!(db.get(&dat!(fmt!("del:{:03}", k)), schms2)).is_some() { |
| 254 | return Err(err!("Deleted chunked key {} still reads back a value.", k; |
| 255 | Test, Invalid, Data)); |
| 256 | } |
| 257 | } |
| 258 | |
| 259 | test!(sync_log::stream(), |
| 260 | "chunk leak: {} bytes of data files with {} chunked values present, {} after deleting all.", |
| 261 | with_values, nkeys, after_delete); |
| 262 | |
| 263 | // Deletion that reclaims chunks leaves only small tombstones plus whatever sits in still-live |
| 264 | // files, so the footprint after deleting every value must fall well below the footprint with |
| 265 | // them present. Without reclamation the chunk bytes remain and the footprint barely moves. |
| 266 | if after_delete * 2 >= with_values { |
| 267 | return Err(err!( |
| 268 | "Chunk leak on delete: deleting {} chunked values took the data footprint from {} \ |
| 269 | bytes only to {} bytes. The chunk records of a deleted value are not being reclaimed.", |
| 270 | nkeys, with_values, after_delete; |
| 271 | Test, Mismatch, Data)); |
| 272 | } |
| 273 | |
| 274 | res!(db.shutdown()); |
| 275 | thread::sleep(Duration::from_millis(200)); |
| 276 | test!(sync_log::stream(), "+--- chunk leak: delete reclaims chunks : passed ---"); |
| 277 | Ok(()) |
| 278 | } |
| 279 | |
| 280 | /// Repeatedly overwriting a chunked value with a smaller (still chunked) one must reach a bounded |
| 281 | /// steady state, not grow with the overwrite count. |
| 282 | /// |
| 283 | /// A chunk key encodes the value's geometry (its length and part count) as well as the key-derived |
| 284 | /// set_id, so an overwrite that *changes* the geometry -- a grow or a shrink across a chunk |
| 285 | /// boundary -- writes its chunks under different keys and does not supersede the prior geometry's |
| 286 | /// chunks: that one prior value leaks and is left for the offline orphan sweep. But overwrites at |
| 287 | /// the *same* geometry share chunk keys and reclaim normally, so a run of same-size overwrites is |
| 288 | /// bounded regardless of length. This test shrinks once and then holds the smaller size, and |
| 289 | /// insists the steady state after many small overwrites is no worse than after a few. |
| 290 | fn shrink_is_bounded( |
| 291 | db_root: &PathBuf, |
| 292 | cfg: &oxedyne_fe2o3_o3db_sync::base::cfg::OzoneConfig, |
| 293 | schms_input: &RestSchemesInput<EncryptionScheme, HashScheme, HashScheme, ChecksumScheme>, |
| 294 | schms2: Option<&RestSchemesOverride<EncryptionScheme, HashScheme>>, |
| 295 | user: setup::Uid, |
| 296 | ) |
| 297 | -> Outcome<()> |
| 298 | { |
| 299 | test!(sync_log::stream(), "+--- chunk leak: shrink is bounded ---"); |
| 300 | |
| 301 | let db = res!(setup::start_db( |
| 302 | db_root.clone(), |
| 303 | Some(cfg.clone()), |
| 304 | schms_input.clone(), |
| 305 | None, |
| 306 | true, |
| 307 | true, |
| 308 | )); |
| 309 | |
| 310 | let key = dat!("chunked:shrinking"); |
| 311 | |
| 312 | // A large chunked value, then shrink to a smaller (still chunked) one. |
| 313 | let (_, big_chunks) = res!(db.insert(key.clone(), value_of(0, VALUE_BYTES), user, schms2)); |
| 314 | let (_, small_chunks) = res!(db.insert(key.clone(), value_of(1, SMALL_BYTES), user, schms2)); |
| 315 | if big_chunks < 2 || small_chunks < 2 || small_chunks >= big_chunks { |
| 316 | return Err(err!( |
| 317 | "Shrink case geometry is wrong: big={} chunks, small={} chunks; the small value must \ |
| 318 | still be chunked and have fewer chunks than the big one.", big_chunks, small_chunks; |
| 319 | Test, Invalid, Configuration)); |
| 320 | } |
| 321 | |
| 322 | // Warm up to a steady state at the smaller geometry, measure, then run three times as many |
| 323 | // more overwrites and measure again. Same-size overwrites share chunk keys, so the collector |
| 324 | // reclaims each prior small value: past the warm-up the footprint must not grow with the |
| 325 | // count. Comparing two post-warm-up points (rather than an early point against a late one) |
| 326 | // isolates a genuine linear leak from the one-off ramp to steady state and the single big |
| 327 | // value leaked at the shrink transition. |
| 328 | for i in 2..=11 { |
| 329 | res!(db.insert(key.clone(), value_of(i as u8, SMALL_BYTES), user, schms2)); |
| 330 | } |
| 331 | thread::sleep(Duration::from_secs(3)); |
| 332 | let after_warm = res!(zone_data_bytes(db_root)); |
| 333 | |
| 334 | for i in 12..=41 { |
| 335 | res!(db.insert(key.clone(), value_of(i as u8, SMALL_BYTES), user, schms2)); |
| 336 | } |
| 337 | thread::sleep(Duration::from_secs(3)); |
| 338 | let after_more = res!(zone_data_bytes(db_root)); |
| 339 | |
| 340 | match res!(db.get(&key, schms2)) { |
| 341 | Some((got, _)) => if got != value_of(41, SMALL_BYTES) { |
| 342 | return Err(err!("After shrinking overwrites the value came back changed."; Test, Invalid, Data)); |
| 343 | }, |
| 344 | None => return Err(err!("After shrinking overwrites the value cannot be read back."; Test, Missing, Data)), |
| 345 | } |
| 346 | |
| 347 | test!(sync_log::stream(), |
| 348 | "chunk leak: shrink steady state {} bytes after 10 small overwrites, {} bytes after 40.", |
| 349 | after_warm, after_more); |
| 350 | |
| 351 | // Thirty more overwrites must not have grown the footprint materially: reclaim at the smaller |
| 352 | // geometry is keeping pace. A generous factor absorbs sealing and collection timing jitter |
| 353 | // while still catching a genuine linear leak (thirty un-reclaimed small values would be many |
| 354 | // times larger). |
| 355 | if after_more > after_warm * 3 / 2 { |
| 356 | return Err(err!( |
| 357 | "Chunk leak on shrink: the footprint grew from {} bytes after warm-up to {} bytes \ |
| 358 | after thirty more overwrites, so same-size overwrites are not reclaiming their \ |
| 359 | predecessors.", after_warm, after_more; |
| 360 | Test, Mismatch, Data)); |
| 361 | } |
| 362 | |
| 363 | res!(db.shutdown()); |
| 364 | thread::sleep(Duration::from_millis(200)); |
| 365 | test!(sync_log::stream(), "+--- chunk leak: shrink is bounded : passed ---"); |
| 366 | Ok(()) |
| 367 | } |
| 368 | |
| 369 | /// A chunked value written under the old scheme -- a random per-operation set_id -- must still read |
| 370 | /// after the switch to key-derived identifiers, because the reader reconstructs the chunk keys from |
| 371 | /// the set_id stored in the bunch key, never from a recomputed one. This forces a random set_id to |
| 372 | /// stand in for a pre-upgrade value, reads it, then overwrites it the ordinary way and reads again: |
| 373 | /// the value is correct throughout. (The first ordinary overwrite does not reclaim the old |
| 374 | /// random-keyed chunks -- a different key scheme -- which is the documented migration cost left for |
| 375 | /// the offline sweep; correctness of reads, not reclamation, is what this asserts.) |
| 376 | fn old_random_key_value_still_reads( |
| 377 | db_root: &PathBuf, |
| 378 | cfg: &oxedyne_fe2o3_o3db_sync::base::cfg::OzoneConfig, |
| 379 | schms_input: &RestSchemesInput<EncryptionScheme, HashScheme, HashScheme, ChecksumScheme>, |
| 380 | schms2: Option<&RestSchemesOverride<EncryptionScheme, HashScheme>>, |
| 381 | user: setup::Uid, |
| 382 | ) |
| 383 | -> Outcome<()> |
| 384 | { |
| 385 | test!(sync_log::stream(), "+--- chunk leak: old random-key value still reads ---"); |
| 386 | |
| 387 | let mut db = res!(setup::start_db( |
| 388 | db_root.clone(), |
| 389 | Some(cfg.clone()), |
| 390 | schms_input.clone(), |
| 391 | None, |
| 392 | true, |
| 393 | true, |
| 394 | )); |
| 395 | |
| 396 | let key = dat!("chunked:legacy"); |
| 397 | let legacy = value_of(7, VALUE_BYTES); |
| 398 | |
| 399 | // Write it as an earlier build did: chunks addressed by a random set_id, unrelated to the key. |
| 400 | // A key-derived write would use `chunk_set_id(key)`, so a random one guarantees the stored |
| 401 | // set_id differs from what the current scheme would compute -- exactly the pre-upgrade case. |
| 402 | let mut set_id_bytes = [0u8; 8]; |
| 403 | Rand::fill_u8(&mut set_id_bytes); |
| 404 | let random_set_id = u64::from_be_bytes(set_id_bytes); |
| 405 | |
| 406 | let resp = db.api().responder(); |
| 407 | let nchunks = res!(db.api().store_dat_using_responder_forcing_set_id( |
| 408 | key.clone(), |
| 409 | legacy.clone(), |
| 410 | user, |
| 411 | schms2, |
| 412 | resp.clone(), |
| 413 | random_set_id, |
| 414 | )); |
| 415 | if nchunks < 2 { |
| 416 | return Err(err!( |
| 417 | "The legacy value was not chunked ({} chunk(s)).", nchunks; Test, Invalid, Configuration)); |
| 418 | } |
| 419 | // Drain the write acknowledgements so the value is settled before it is read: the count, |
| 420 | // then every record written and durable. |
| 421 | res!(resp.recv_store_ack()); |
| 422 | thread::sleep(Duration::from_secs(1)); |
| 423 | |
| 424 | // It must read back correctly even though its set_id is not the key-derived one. |
| 425 | match res!(db.get(&key, schms2)) { |
| 426 | Some((got, _)) => if got != legacy { |
| 427 | return Err(err!( |
| 428 | "A chunked value written under a random (pre-upgrade) set_id did not read back \ |
| 429 | unchanged, so the reader is not reconstructing chunk keys from the stored set_id."; |
| 430 | Test, Invalid, Data)); |
| 431 | }, |
| 432 | None => return Err(err!( |
| 433 | "A chunked value written under a random (pre-upgrade) set_id cannot be read at all."; |
| 434 | Test, Missing, Data)), |
| 435 | } |
| 436 | |
| 437 | // Now overwrite it the ordinary (key-derived) way and read again: the new value is correct, and |
| 438 | // the old value remained readable right up to the overwrite. |
| 439 | let fresh = value_of(8, VALUE_BYTES); |
| 440 | res!(db.insert(key.clone(), fresh.clone(), user, schms2)); |
| 441 | thread::sleep(Duration::from_secs(1)); |
| 442 | match res!(db.get(&key, schms2)) { |
| 443 | Some((got, _)) => if got != fresh { |
| 444 | return Err(err!( |
| 445 | "After overwriting a legacy random-keyed value the new value did not read back."; |
| 446 | Test, Invalid, Data)); |
| 447 | }, |
| 448 | None => return Err(err!( |
| 449 | "After overwriting a legacy random-keyed value the key cannot be read."; |
| 450 | Test, Missing, Data)), |
| 451 | } |
| 452 | |
| 453 | res!(db.shutdown()); |
| 454 | thread::sleep(Duration::from_millis(200)); |
| 455 | test!(sync_log::stream(), "+--- chunk leak: old random-key value still reads : passed ---"); |
| 456 | Ok(()) |
| 457 | } |
| 458 | |
| 459 | /// A heavy same-key churn -- the concurrent supersession that exposed the accounting race -- must |
| 460 | /// leave the store self-consistent across a restart. A restart rebuilds every file state from what |
| 461 | /// is physically on disk, so if the churn had left the accounting inconsistent (records flagged old |
| 462 | /// twice, records superseded but never inserted, a data file out of step with its state) the |
| 463 | /// rebuilt store would drop or corrupt the value or fail to start. Reading the value back |
| 464 | /// unchanged, with a bounded footprint, after a restart is the black-box proof that the |
| 465 | /// reconciliation kept the on-disk accounting exact. |
| 466 | fn churn_survives_restart( |
| 467 | db_root: &PathBuf, |
| 468 | cfg: &oxedyne_fe2o3_o3db_sync::base::cfg::OzoneConfig, |
| 469 | schms_input: &RestSchemesInput<EncryptionScheme, HashScheme, HashScheme, ChecksumScheme>, |
| 470 | schms2: Option<&RestSchemesOverride<EncryptionScheme, HashScheme>>, |
| 471 | user: setup::Uid, |
| 472 | ) |
| 473 | -> Outcome<()> |
| 474 | { |
| 475 | test!(sync_log::stream(), "+--- chunk leak: churn survives restart ---"); |
| 476 | |
| 477 | let key = dat!("chunked:restart"); |
| 478 | let last = 40u8; |
| 479 | |
| 480 | { |
| 481 | let db = res!(setup::start_db( |
| 482 | db_root.clone(), |
| 483 | Some(cfg.clone()), |
| 484 | schms_input.clone(), |
| 485 | None, |
| 486 | true, |
| 487 | true, |
| 488 | )); |
| 489 | // Rapid, un-spaced overwrites of one key: the burst that races record insertion against |
| 490 | // sibling-chunk supersession and drives the reconciliation and the drain-gated collector. |
| 491 | for i in 0..=last { |
| 492 | res!(db.insert(key.clone(), chunky_value(i), user, schms2)); |
| 493 | } |
| 494 | thread::sleep(Duration::from_secs(3)); |
| 495 | res!(db.shutdown()); |
| 496 | } |
| 497 | thread::sleep(Duration::from_secs(1)); |
| 498 | |
| 499 | // Restart over what the churn left on disk. |
| 500 | let db = res!(setup::start_db( |
| 501 | db_root.clone(), |
| 502 | Some(cfg.clone()), |
| 503 | schms_input.clone(), |
| 504 | None, |
| 505 | true, |
| 506 | false, // do not wipe: read what the churn left |
| 507 | )); |
| 508 | |
| 509 | match res!(db.get(&key, schms2)) { |
| 510 | Some((got, _)) => if got != chunky_value(last) { |
| 511 | return Err(err!( |
| 512 | "After a heavy same-key churn and a restart, the value came back changed: the \ |
| 513 | on-disk accounting the churn left was not self-consistent."; |
| 514 | Test, Invalid, Data)); |
| 515 | }, |
| 516 | None => return Err(err!( |
| 517 | "After a heavy same-key churn and a restart, the value is gone."; |
| 518 | Test, Missing, Data)), |
| 519 | } |
| 520 | |
| 521 | let after = res!(zone_data_bytes(db_root)); |
| 522 | let one_value = res!(single_value_baseline(cfg, schms_input, schms2, user)); |
| 523 | test!(sync_log::stream(), |
| 524 | "chunk leak: after churn + restart, {} bytes of data files (one value ~{} bytes).", |
| 525 | after, one_value); |
| 526 | if after >= one_value.saturating_mul(8) { |
| 527 | return Err(err!( |
| 528 | "After a heavy same-key churn and a restart, the footprint is {} bytes against a \ |
| 529 | one-value baseline of ~{} bytes: the churn leaked rather than reclaiming.", |
| 530 | after, one_value; |
| 531 | Test, Mismatch, Data)); |
| 532 | } |
| 533 | |
| 534 | res!(db.shutdown()); |
| 535 | thread::sleep(Duration::from_millis(200)); |
| 536 | test!(sync_log::stream(), "+--- chunk leak: churn survives restart : passed ---"); |
| 537 | Ok(()) |
| 538 | } |
| 539 | |
| 540 | /// The settled on-disk size of a single chunked value, measured in a throwaway store, for use as a |
| 541 | /// baseline the churn footprint is judged against. |
| 542 | fn single_value_baseline( |
| 543 | cfg: &oxedyne_fe2o3_o3db_sync::base::cfg::OzoneConfig, |
| 544 | schms_input: &RestSchemesInput<EncryptionScheme, HashScheme, HashScheme, ChecksumScheme>, |
| 545 | schms2: Option<&RestSchemesOverride<EncryptionScheme, HashScheme>>, |
| 546 | user: setup::Uid, |
| 547 | ) |
| 548 | -> Outcome<u64> |
| 549 | { |
| 550 | let base_root = res!(canonical_dir("./test_db_chunk_leak_baseline")); |
| 551 | let db = res!(setup::start_db( |
| 552 | base_root.clone(), |
| 553 | Some(cfg.clone()), |
| 554 | schms_input.clone(), |
| 555 | None, |
| 556 | true, |
| 557 | true, |
| 558 | )); |
| 559 | res!(db.insert(dat!("chunked:one"), chunky_value(0), user, schms2)); |
| 560 | thread::sleep(Duration::from_secs(2)); |
| 561 | let bytes = res!(zone_data_bytes(&base_root)); |
| 562 | res!(db.shutdown()); |
| 563 | thread::sleep(Duration::from_millis(200)); |
| 564 | Ok(bytes) |
| 565 | } |
| 566 | |
| 567 | fn zone_data_bytes(db_root: &Path) -> Outcome<u64> { |
| 568 | let mut total = 0u64; |
| 569 | res!(walk_data_files(db_root, &mut total)); |
| 570 | Ok(total) |
| 571 | } |
| 572 | |
| 573 | fn walk_data_files(dir: &Path, total: &mut u64) -> Outcome<()> { |
| 574 | for entry in res!(fs::read_dir(dir)) { |
| 575 | let entry = res!(entry); |
| 576 | let path = entry.path(); |
| 577 | if path.is_dir() { |
| 578 | res!(walk_data_files(&path, total)); |
| 579 | } else if path.extension().map(|e| e == constant::DATA_FILE_EXT).unwrap_or(false) { |
| 580 | let meta = res!(entry.metadata()); |
| 581 | *total += meta.len(); |
| 582 | } |
| 583 | } |
| 584 | Ok(()) |
| 585 | } |
| 586 | |
| 587 | /// Creates the directory if it does not exist. |
| 588 | fn canonical_dir(p: &str) -> Outcome<PathBuf> { |
| 589 | match Path::new(p).canonicalize() { |
| 590 | Ok(path) => Ok(path), |
| 591 | Err(_) => { |
| 592 | res!(fs::create_dir_all(p)); |
| 593 | match Path::new(p).canonicalize() { |
| 594 | Ok(path) => Ok(path), |
| 595 | Err(e) => Err(err!(e, "Cannot canonicalise {:?}.", p; IO, Path)), |
| 596 | } |
| 597 | }, |
| 598 | } |
| 599 | } |