oxedyne/fe2o3/fe2o3_o3db_sync/tests/crash_restart_durability.rs
12.8 KiB, 1 run
created by r1870400018:38052, 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 | //! Crash-restart durability of the data-file rebuild (BUG B). |
| 2 | //! |
| 3 | //! # Why this exists |
| 4 | //! |
| 5 | //! An o3db zone keeps an append-only data file and an index file beside it. On |
| 6 | //! restart the index is read to populate the caches; when the index is missing |
| 7 | //! or disagrees with the data file, the zone instead rebuilds the index by |
| 8 | //! walking the data file (`InitGarbageBot::init_cache_data_file`). Under the |
| 9 | //! pre-durability default (never fsync) a crash routinely leaves the data and |
| 10 | //! index flushed to different points, so this rebuild path is the one a real |
| 11 | //! restart falls into. |
| 12 | //! |
| 13 | //! The property a log-structured store owes its caller is that an interrupted |
| 14 | //! final append costs the one record it interrupted and nothing else. The |
| 15 | //! index-driven path keeps that property (see `scan_torn_tail`), but the data |
| 16 | //! rebuild used to `res!`-abort the whole per-file walk on the first torn |
| 17 | //! record. That turned one torn record into a rebuild that returns an error, |
| 18 | //! which fails zone initialisation, which fails `Store::open`: a single torn |
| 19 | //! tail took the WHOLE store offline, and on the full (~1 MiB) live file of the |
| 20 | //! Ochre repro it was ~3k records' worth of data that could not be read back. |
| 21 | //! |
| 22 | //! This test reproduces that shape with many small unchunked records across |
| 23 | //! several rollovers, then simulates a crash by damaging the tail of one sealed |
| 24 | //! data file and removing its index (forcing the data rebuild). Two crash |
| 25 | //! shapes are covered, because an interrupted append leaves either: |
| 26 | //! |
| 27 | //! - a torn key: a prefix of the next record reached disk (here a lone cache |
| 28 | //! hash), so `StoredKey::load` reads the hash and then meets the end of the |
| 29 | //! file loading the key; |
| 30 | //! - a torn value: the key and the value's length header are whole but the |
| 31 | //! value body was cut off. `StoredValue::count` walks a value by seeking |
| 32 | //! over its declared length rather than reading it, so this is invisible at |
| 33 | //! the key/value decode and shows up only as a record that claims to end |
| 34 | //! past the file. |
| 35 | //! |
| 36 | //! Before the fix either shape aborts the rebuild and the store will not reopen |
| 37 | //! at all; after it, the rebuild truncates the torn tail and recovers every |
| 38 | //! other record in that file and every record in every other file. |
| 39 | //! |
| 40 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 41 | //! Anthropic Claude |
| 42 | |
| 43 | use oxedyne_fe2o3_core::{ |
| 44 | prelude::*, |
| 45 | alt::Override, |
| 46 | }; |
| 47 | use oxedyne_fe2o3_crypto::enc::EncryptionScheme; |
| 48 | use oxedyne_fe2o3_hash::{ |
| 49 | csum::ChecksumScheme, |
| 50 | hash::HashScheme, |
| 51 | }; |
| 52 | use oxedyne_fe2o3_iop_db::api::{ |
| 53 | Database, |
| 54 | RestSchemesOverride, |
| 55 | ScanOpts, |
| 56 | }; |
| 57 | use oxedyne_fe2o3_jdat::prelude::*; |
| 58 | use oxedyne_fe2o3_o3db_sync::{ |
| 59 | base::cfg::OzoneConfig, |
| 60 | data::core::RestSchemesInput, |
| 61 | test::setup, |
| 62 | }; |
| 63 | |
| 64 | use std::{ |
| 65 | fs, |
| 66 | io::Write, |
| 67 | path::{ |
| 68 | Path, |
| 69 | PathBuf, |
| 70 | }, |
| 71 | thread, |
| 72 | time::Duration, |
| 73 | }; |
| 74 | |
| 75 | /// How the last append to a data file was interrupted. |
| 76 | #[derive(Clone, Copy, Debug)] |
| 77 | enum Damage { |
| 78 | /// A prefix of the next record (a lone cache hash) reached disk: the key |
| 79 | /// decode runs off the end of the file. |
| 80 | TornKey, |
| 81 | /// The value's length header is whole but its body was cut short: the |
| 82 | /// record claims to end past the file. |
| 83 | TornValue, |
| 84 | } |
| 85 | |
| 86 | /// Every file under `root` with the given extension, with its byte length. |
| 87 | fn files_with_ext(root: &Path, ext: &str) -> Outcome<Vec<(PathBuf, u64)>> { |
| 88 | let mut found = Vec::new(); |
| 89 | let mut stack = vec![root.to_path_buf()]; |
| 90 | while let Some(dir) = stack.pop() { |
| 91 | for entry in res!(fs::read_dir(&dir)) { |
| 92 | let path = res!(entry).path(); |
| 93 | if path.is_dir() { |
| 94 | stack.push(path); |
| 95 | } else if path.extension().map(|e| e == ext).unwrap_or(false) { |
| 96 | let len = res!(fs::metadata(&path)).len(); |
| 97 | found.push((path, len)); |
| 98 | } |
| 99 | } |
| 100 | } |
| 101 | found.sort(); |
| 102 | Ok(found) |
| 103 | } |
| 104 | |
| 105 | fn run_case(damage: Damage) -> Outcome<()> { |
| 106 | |
| 107 | let written: usize = 300; |
| 108 | |
| 109 | let dirname = fmt!("./test_db_crash_restart_{:?}", damage).to_lowercase(); |
| 110 | let db_root = res!(Path::new(&dirname).canonicalize().or_else(|_| { |
| 111 | ok!(fs::create_dir_all(&dirname)); |
| 112 | Path::new(&dirname).canonicalize() |
| 113 | })); |
| 114 | |
| 115 | let enckey = [0x37u8; 32]; |
| 116 | let aes_gcm = res!(EncryptionScheme::new_aes_256_gcm_with_key(&enckey[..])); |
| 117 | let crc32 = ChecksumScheme::new_crc32(); |
| 118 | let schms2: RestSchemesOverride<EncryptionScheme, HashScheme> = |
| 119 | RestSchemesOverride::default() |
| 120 | .set_encrypter(Override::Default(aes_gcm.clone())); |
| 121 | let schms2 = Some(&schms2); |
| 122 | let user = setup::Uid::default(); |
| 123 | |
| 124 | let schms_input = RestSchemesInput::new( |
| 125 | Some(aes_gcm.clone()), |
| 126 | None::<HashScheme>, |
| 127 | None::<HashScheme>, |
| 128 | Some(crc32.clone()), |
| 129 | ); |
| 130 | |
| 131 | // Start from the test config, but exercise the real default durability |
| 132 | // policy (a bounded group commit). The file size is shrunk so a few |
| 133 | // hundred ~300 B records force many rollovers and leave several full sealed |
| 134 | // files; the records stay well under `rest_chunk_threshold`, so every write |
| 135 | // takes the single unchunked path, exactly as the Ochre repro did. |
| 136 | let mut cfg = res!(setup::default_cfg()); |
| 137 | cfg.num_zones = 2; |
| 138 | cfg.num_cbots_per_zone = 2; |
| 139 | cfg.num_wbots_per_zone = 1; |
| 140 | cfg.num_igbots_per_zone = 2; |
| 141 | cfg.data_file_max_bytes = 8_000; |
| 142 | cfg.rest_chunk_threshold = 1_500; // 300 B values stay unchunked. |
| 143 | cfg.rest_chunk_bytes = 64; |
| 144 | cfg.sync_interval_ms = OzoneConfig::default().sync_interval_ms; // the real default floor. |
| 145 | cfg.zone_overrides = mapdat!{ |
| 146 | 1u16 => mapdat!{ "dir" => "", "max_size" => 10_000_000u64 }, |
| 147 | 2u16 => mapdat!{ "dir" => "", "max_size" => 10_000_000u64 }, |
| 148 | }.get_map().unwrap_or_default(); |
| 149 | |
| 150 | test!(sync_log::stream(), "+--- crash-restart durability: {:?} ---", damage); |
| 151 | |
| 152 | // ── 1. Write many small records over several rollovers, then stop cleanly ── |
| 153 | let db = res!(setup::start_db( |
| 154 | db_root.clone(), |
| 155 | Some(cfg.clone()), |
| 156 | schms_input.clone(), |
| 157 | None, |
| 158 | false, // gc off: the rebuild path is what is under test. |
| 159 | true, // wipe |
| 160 | )); |
| 161 | thread::sleep(Duration::from_millis(500)); |
| 162 | let val_body = "x".repeat(300); |
| 163 | for i in 0..written { |
| 164 | res!(db.insert( |
| 165 | dat!(fmt!("crash:{:04}", i)), |
| 166 | dat!(fmt!("{}:{:04}", val_body, i)), |
| 167 | user, |
| 168 | schms2, |
| 169 | )); |
| 170 | } |
| 171 | thread::sleep(Duration::from_millis(500)); |
| 172 | res!(db.shutdown()); |
| 173 | thread::sleep(Duration::from_millis(500)); |
| 174 | |
| 175 | // ── 2. Simulate the crash on one sealed data file and remove its index so |
| 176 | // the restart must rebuild it from the data and meet the torn tail. A |
| 177 | // sealed (not the highest-numbered, live) file is chosen so the tear |
| 178 | // is unambiguously a crash artefact rather than a writer appending. ── |
| 179 | let dats = res!(files_with_ext(&db_root, "dat")); |
| 180 | if dats.len() < 3 { |
| 181 | return Err(err!( |
| 182 | "Only {} data files were written; the config should have forced \ |
| 183 | several rollovers, so this test would prove nothing.", dats.len(); |
| 184 | Test, Size)); |
| 185 | } |
| 186 | let live = dats[dats.len() - 1].0.clone(); |
| 187 | let mut victim: Option<(PathBuf, u64)> = None; |
| 188 | for (p, len) in &dats { |
| 189 | if *p == live { |
| 190 | continue; |
| 191 | } |
| 192 | if victim.as_ref().map(|(_, n)| *len > *n).unwrap_or(true) { |
| 193 | victim = Some((p.clone(), *len)); |
| 194 | } |
| 195 | } |
| 196 | let (victim_path, victim_len) = res!(victim.ok_or_else(|| err!( |
| 197 | "No sealed data file was found to damage."; |
| 198 | Test, Missing))); |
| 199 | test!(sync_log::stream(), |
| 200 | "Damaging sealed data file {:?} ({} bytes) by {:?} and removing its index.", |
| 201 | victim_path, victim_len, damage); |
| 202 | |
| 203 | match damage { |
| 204 | Damage::TornKey => { |
| 205 | // Append a lone cache hash (CACHE_HASH_BYTES = 4): the key decode |
| 206 | // reads it, then meets EOF loading the key itself. |
| 207 | let mut file = res!(fs::OpenOptions::new().append(true).open(&victim_path)); |
| 208 | res!(file.write_all(&[0u8; 4])); |
| 209 | }, |
| 210 | Damage::TornValue => { |
| 211 | // Cut the last few bytes so the final value's body is short while |
| 212 | // its length header survives. |
| 213 | let file = res!(fs::OpenOptions::new().write(true).open(&victim_path)); |
| 214 | res!(file.set_len(victim_len.saturating_sub(9))); |
| 215 | }, |
| 216 | } |
| 217 | // Remove the matching index file to force the data-rebuild path. |
| 218 | let mut ind_path = victim_path.clone(); |
| 219 | ind_path.set_extension("ind"); |
| 220 | if ind_path.is_file() { |
| 221 | res!(fs::remove_file(&ind_path)); |
| 222 | } |
| 223 | |
| 224 | // ── 3. Reopen. Before the fix the rebuild aborts and this very open fails; |
| 225 | // after the fix it truncates the torn tail and opens. ── |
| 226 | let db = match setup::start_db( |
| 227 | db_root.clone(), |
| 228 | Some(cfg.clone()), |
| 229 | schms_input.clone(), |
| 230 | None, |
| 231 | false, // gc off |
| 232 | false, // keep what is there |
| 233 | ) { |
| 234 | Ok(db) => db, |
| 235 | Err(e) => return Err(err!(e, |
| 236 | "The store would not reopen after a single {:?} in one data file. A \ |
| 237 | torn tail on the rebuild path must cost that one record, not the whole \ |
| 238 | store: this is the pre-fix abort-the-whole-file behaviour.", damage; |
| 239 | Test, Data)), |
| 240 | }; |
| 241 | thread::sleep(Duration::from_millis(500)); |
| 242 | |
| 243 | // ── 4. Scan the store. A scan answers from the index files on disk, so it |
| 244 | // is the true measure of the rebuild: before the fix the aborted |
| 245 | // rebuild never writes the torn file's index, so every record in that |
| 246 | // file is invisible to a scan for the life of the store (and the |
| 247 | // rebuild re-aborts on every later restart); after the fix the index |
| 248 | // is rebuilt from the records the crash did not take, so they are |
| 249 | // scannable again. A `get` by key would hide this difference, because |
| 250 | // the rebuild streams each record into the caches as it reads it, so |
| 251 | // the records ahead of the tear are cached either way -- it is the |
| 252 | // index, and therefore the scan, that the whole-file abort destroys. ── |
| 253 | let scanned = match db.scan(&ScanOpts::with_str_prefix("crash:"), schms2) { |
| 254 | Ok(rows) => rows, |
| 255 | Err(e) => { |
| 256 | let _ = db.shutdown(); |
| 257 | return Err(err!(e, |
| 258 | "[{:?}] The prefix scan failed after the crash-restart.", damage; |
| 259 | Test, Data)); |
| 260 | }, |
| 261 | }; |
| 262 | let scanned_n = scanned.len(); |
| 263 | let scan_lost = written.saturating_sub(scanned_n); |
| 264 | test!(sync_log::stream(), |
| 265 | "[{:?}] scan returned {} of {} records; {} not scannable.", |
| 266 | damage, scanned_n, written, scan_lost); |
| 267 | |
| 268 | // The property: an interrupted final append costs that one record, not a |
| 269 | // whole file's worth. One tear was injected, so at most a small constant |
| 270 | // number of records may drop out of the scan; a file's worth missing means |
| 271 | // the whole-file abort left the torn file with no rebuilt index. |
| 272 | if scan_lost > 3 { |
| 273 | let _ = db.shutdown(); |
| 274 | return Err(err!( |
| 275 | "[{:?}] {} of {} records are missing from the scan after a single torn \ |
| 276 | tail. The aborted rebuild left the torn data file with no index, so a \ |
| 277 | whole file's records became invisible to a scan -- the pre-fix \ |
| 278 | whole-file blast radius.", |
| 279 | damage, scan_lost, written; |
| 280 | Test, Data)); |
| 281 | } |
| 282 | |
| 283 | // And every scanned key must read back cleanly: a torn record left behind as |
| 284 | // a poison pill would error here (the 2026-08-24 "Mismatch detected" scan |
| 285 | // failure), whereas the fix removes it. |
| 286 | for (k, _, _) in &scanned { |
| 287 | match db.get(k, schms2) { |
| 288 | Ok(Some(_)) => (), |
| 289 | Ok(None) => { |
| 290 | let _ = db.shutdown(); |
| 291 | return Err(err!( |
| 292 | "[{:?}] Scan returned key {:?} but a get of it found nothing.", damage, k; |
| 293 | Test, Data)); |
| 294 | }, |
| 295 | Err(e) => { |
| 296 | let _ = db.shutdown(); |
| 297 | return Err(err!(e, |
| 298 | "[{:?}] A scanned key {:?} would not read back -- a torn record \ |
| 299 | left as a poison pill.", damage, k; |
| 300 | Test, Data)); |
| 301 | }, |
| 302 | } |
| 303 | } |
| 304 | |
| 305 | res!(db.shutdown()); |
| 306 | thread::sleep(Duration::from_millis(500)); |
| 307 | Ok(()) |
| 308 | } |
| 309 | |
| 310 | fn run() -> Outcome<()> { |
| 311 | test!(sync_log::stream(), "+---------------------------------------------+"); |
| 312 | test!(sync_log::stream(), "| CRASH-RESTART DURABILITY (data rebuild) |"); |
| 313 | test!(sync_log::stream(), "+---------------------------------------------+"); |
| 314 | res!(run_case(Damage::TornKey)); |
| 315 | res!(run_case(Damage::TornValue)); |
| 316 | test!(sync_log::stream(), "Crash-restart durability test passed."); |
| 317 | Ok(()) |
| 318 | } |
| 319 | |
| 320 | #[test] |
| 321 | fn main() -> Outcome<()> { |
| 322 | log_set_level!("trace"); |
| 323 | let outcome = run(); |
| 324 | log_finish_wait!(); |
| 325 | outcome |
| 326 | } |