oxedyne/fe2o3/fe2o3_o3db_sync/tests/scan_torn_tail.rs
11.0 KiB, 1 run
created by r1870400018:35555, 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 | //! What a scan does when the last record in an index file is damaged. |
| 2 | //! |
| 3 | //! # Why this exists |
| 4 | //! |
| 5 | //! An o3db zone's index is an append-only log, and a process that dies without |
| 6 | //! a clean shutdown leaves whatever the last append had managed to put on disk. |
| 7 | //! That is not an exotic condition: `dev/verify_operators.mjs` restarts the |
| 8 | //! Daimond gateway with `SIGKILL` twice in one run, against the same store, and |
| 9 | //! every deployment ends the same way sooner or later. |
| 10 | //! |
| 11 | //! The property a log-structured store owes its caller is that damage at the |
| 12 | //! TAIL costs the records in that damage and nothing else, and this check was |
| 13 | //! written to find out whether o3db keeps it. **It does**, which is the finding |
| 14 | //! and the reason the check is worth keeping: |
| 15 | //! |
| 16 | //! - Damage to an INDEX file costs nothing at all, either shape. The index is |
| 17 | //! rebuilt from the data file beside it when the zone starts, so the damage is |
| 18 | //! gone before a reader sees it: 12 of 12 written records read back. |
| 19 | //! - Damage to a DATA file costs exactly the record it damaged. The scan still |
| 20 | //! returns all 12 keys, because a scan answers from the index; 11 of the 12 |
| 21 | //! values then read back and the twelfth does not. |
| 22 | //! |
| 23 | //! So a store that was stopped mid-write loses what was in flight and stays |
| 24 | //! readable otherwise. Where `store.scan_prefix("lic:") failed after 10.092 s: |
| 25 | //! Mismatch detected` came from on 2026-08-24 is therefore NOT o3db widening |
| 26 | //! one bad record into a dead prefix -- it is the caller doing that, and the |
| 27 | //! caller is `Store::scan_prefix` in the Daimond gateway, which `res!`s the |
| 28 | //! `get` of every key the scan returned. |
| 29 | //! |
| 30 | //! Two shapes are checked because an interrupted append leaves either: a record |
| 31 | //! cut short, and a record whose bytes are all there and wrong. |
| 32 | //! |
| 33 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 34 | //! Anthropic Claude |
| 35 | |
| 36 | use oxedyne_fe2o3_core::{ |
| 37 | prelude::*, |
| 38 | alt::Override, |
| 39 | }; |
| 40 | use oxedyne_fe2o3_crypto::enc::EncryptionScheme; |
| 41 | use oxedyne_fe2o3_hash::{ |
| 42 | csum::ChecksumScheme, |
| 43 | hash::HashScheme, |
| 44 | }; |
| 45 | use oxedyne_fe2o3_iop_db::api::{ |
| 46 | Database, |
| 47 | RestSchemesOverride, |
| 48 | ScanOpts, |
| 49 | }; |
| 50 | use oxedyne_fe2o3_jdat::prelude::*; |
| 51 | use oxedyne_fe2o3_o3db_sync::{ |
| 52 | data::core::RestSchemesInput, |
| 53 | test::setup, |
| 54 | }; |
| 55 | |
| 56 | use std::{ |
| 57 | fs, |
| 58 | path::{ |
| 59 | Path, |
| 60 | PathBuf, |
| 61 | }, |
| 62 | thread, |
| 63 | time::Duration, |
| 64 | }; |
| 65 | |
| 66 | /// How the last append was interrupted. |
| 67 | #[derive(Clone, Copy, Debug)] |
| 68 | enum Damage { |
| 69 | /// The write did not finish: the file is short by a few bytes. |
| 70 | CutShort, |
| 71 | /// The bytes arrived and are wrong, which is what a checksum is for. |
| 72 | Corrupted, |
| 73 | } |
| 74 | |
| 75 | /// Every file under `root` with the given extension. |
| 76 | fn files_with_ext(root: &Path, ext: &str) -> Outcome<Vec<PathBuf>> { |
| 77 | let mut found = Vec::new(); |
| 78 | let mut stack = vec![root.to_path_buf()]; |
| 79 | while let Some(dir) = stack.pop() { |
| 80 | for entry in res!(fs::read_dir(&dir)) { |
| 81 | let path = res!(entry).path(); |
| 82 | if path.is_dir() { |
| 83 | stack.push(path); |
| 84 | } else if path.extension().map(|e| e == ext).unwrap_or(false) { |
| 85 | found.push(path); |
| 86 | } |
| 87 | } |
| 88 | } |
| 89 | found.sort(); |
| 90 | Ok(found) |
| 91 | } |
| 92 | |
| 93 | fn run_case(target: &'static str, damage: Damage, written: usize) -> Outcome<()> { |
| 94 | |
| 95 | let dirname = fmt!("./test_db_scan_torn_{}_{:?}", target, damage).to_lowercase(); |
| 96 | let db_root = res!(Path::new(&dirname).canonicalize().or_else(|_| { |
| 97 | ok!(fs::create_dir_all(&dirname)); |
| 98 | Path::new(&dirname).canonicalize() |
| 99 | })); |
| 100 | |
| 101 | let enckey = [0x42u8; 32]; |
| 102 | let aes_gcm = res!(EncryptionScheme::new_aes_256_gcm_with_key(&enckey[..])); |
| 103 | let crc32 = ChecksumScheme::new_crc32(); |
| 104 | let schms2: RestSchemesOverride<EncryptionScheme, HashScheme> = |
| 105 | RestSchemesOverride::default() |
| 106 | .set_encrypter(Override::Default(aes_gcm.clone())); |
| 107 | let schms2 = Some(&schms2); |
| 108 | let user = setup::Uid::default(); |
| 109 | |
| 110 | let schms_input = RestSchemesInput::new( |
| 111 | Some(aes_gcm.clone()), |
| 112 | None::<HashScheme>, |
| 113 | None::<HashScheme>, |
| 114 | Some(crc32.clone()), |
| 115 | ); |
| 116 | |
| 117 | let mut cfg = res!(setup::default_cfg()); |
| 118 | cfg.num_zones = 1; |
| 119 | cfg.num_cbots_per_zone = 2; |
| 120 | cfg.num_igbots_per_zone = 2; |
| 121 | cfg.data_file_max_bytes = 200_000; |
| 122 | cfg.zone_overrides = mapdat!{ |
| 123 | 1u16 => mapdat!{ "dir" => "", "max_size" => 10_000_000u64 }, |
| 124 | }.get_map().unwrap_or_default(); |
| 125 | |
| 126 | test!(sync_log::stream(), "+--- torn tail: .{} {:?} ---", target, damage); |
| 127 | |
| 128 | // ── 1. Write a prefix's worth of records and stop cleanly ── |
| 129 | let mut db = res!(setup::start_db( |
| 130 | db_root.clone(), |
| 131 | Some(cfg.clone()), |
| 132 | schms_input.clone(), |
| 133 | None, |
| 134 | false, // gc off: an index rebuilt from the data file would repair the |
| 135 | // damage before the scan saw it, and what is under test here is |
| 136 | // the scan. |
| 137 | true, // wipe |
| 138 | )); |
| 139 | thread::sleep(Duration::from_secs(1)); |
| 140 | for i in 0..written { |
| 141 | res!(db.insert( |
| 142 | dat!(fmt!("lic:{:03}", i)), |
| 143 | dat!(fmt!("licence_{}", i)), |
| 144 | user, |
| 145 | schms2, |
| 146 | )); |
| 147 | } |
| 148 | thread::sleep(Duration::from_millis(500)); |
| 149 | res!(db.shutdown()); |
| 150 | thread::sleep(Duration::from_millis(500)); |
| 151 | |
| 152 | // ── 2. Damage the last record of the largest such file ── |
| 153 | let inds = res!(files_with_ext(&db_root, target)); |
| 154 | let mut biggest: Option<(PathBuf, u64)> = None; |
| 155 | for p in inds { |
| 156 | let len = res!(fs::metadata(&p)).len(); |
| 157 | if biggest.as_ref().map(|(_, n)| len > *n).unwrap_or(true) { |
| 158 | biggest = Some((p, len)); |
| 159 | } |
| 160 | } |
| 161 | let (path, len) = res!(biggest.ok_or_else(|| err!( |
| 162 | "No .{} file was written under {:?}, so there is nothing to damage and \ |
| 163 | this check would prove nothing.", target, db_root; |
| 164 | Test, Missing))); |
| 165 | test!(sync_log::stream(), "Damaging {:?} ({} bytes) by {:?}.", path, len, damage); |
| 166 | match damage { |
| 167 | Damage::CutShort => { |
| 168 | let file = res!(fs::OpenOptions::new().write(true).open(&path)); |
| 169 | res!(file.set_len(len.saturating_sub(9))); |
| 170 | }, |
| 171 | Damage::Corrupted => { |
| 172 | let mut bytes = res!(fs::read(&path)); |
| 173 | let n = bytes.len(); |
| 174 | if n < 4 { |
| 175 | return Err(err!("File {:?} is too small to damage.", path; |
| 176 | Test, Size)); |
| 177 | } |
| 178 | // The last four bytes of a record are its checksum, so flipping a |
| 179 | // bit just before them damages content the checksum covers. |
| 180 | bytes[n - 5] ^= 0xFF; |
| 181 | res!(fs::write(&path, &bytes)); |
| 182 | }, |
| 183 | } |
| 184 | |
| 185 | // ── 3. Reopen and scan ── |
| 186 | let db = res!(setup::start_db( |
| 187 | db_root.clone(), |
| 188 | Some(cfg.clone()), |
| 189 | schms_input.clone(), |
| 190 | None, |
| 191 | false, // gc off |
| 192 | false, // keep what is there |
| 193 | )); |
| 194 | thread::sleep(Duration::from_secs(1)); |
| 195 | |
| 196 | let outcome = db.scan(&ScanOpts::with_str_prefix("lic:"), schms2); |
| 197 | let rows = match outcome { |
| 198 | Ok(rows) => rows, |
| 199 | Err(e) => { |
| 200 | let _ = db.shutdown(); |
| 201 | return Err(err!(e, |
| 202 | "A {:?} tail in one .{} file made the WHOLE prefix scan fail. \ |
| 203 | Damage at the tail of an append-only log must cost the records \ |
| 204 | in that damage and no others: this is a store that cannot be \ |
| 205 | read at all after an unclean stop.", damage, target; |
| 206 | Test, Data)); |
| 207 | }, |
| 208 | }; |
| 209 | test!(sync_log::stream(), "Scan returned {} of {} after .{} {:?}.", |
| 210 | rows.len(), written, target, damage); |
| 211 | |
| 212 | // Every row that did come back must be intact and ours. A scan that |
| 213 | // survives damage by handing back rubbish would be worse than one that |
| 214 | // fails. |
| 215 | for (k, _, _) in &rows { |
| 216 | match k { |
| 217 | Dat::Str(s) => if !s.starts_with("lic:") { |
| 218 | let _ = db.shutdown(); |
| 219 | return Err(err!("Scan returned {:?}, which is not a 'lic:' key.", k; |
| 220 | Test, Mismatch)); |
| 221 | }, |
| 222 | other => { |
| 223 | let _ = db.shutdown(); |
| 224 | return Err(err!("Expected a Dat::Str key, got {:?}.", other; |
| 225 | Test, Mismatch)); |
| 226 | }, |
| 227 | } |
| 228 | } |
| 229 | if rows.len() == 0 { |
| 230 | let _ = db.shutdown(); |
| 231 | return Err(err!( |
| 232 | "The scan came back empty after {:?}. One damaged record cost every \ |
| 233 | one of the {} written, which is the same outage as an error and is \ |
| 234 | quieter about it.", damage, written; |
| 235 | Test, Data)); |
| 236 | } |
| 237 | |
| 238 | // And now read each value, because that is what the caller does. o3db's |
| 239 | // scan answers from the INDEX and hands every value back as `Dat::Empty`, |
| 240 | // so a scan alone never opens a data file and never checksums one -- |
| 241 | // `Store::scan_prefix` in the gateway scans for keys and then `get`s each, |
| 242 | // and it is the `get` that reads the record whose checksum failed on |
| 243 | // 2026-08-24. A check that stopped at the scan would be measuring the half |
| 244 | // of the path that cannot fail. |
| 245 | let mut read_ok = 0usize; |
| 246 | let mut read_bad = Vec::new(); |
| 247 | for (k, _, _) in &rows { |
| 248 | match db.get(k, schms2) { |
| 249 | Ok(Some(_)) => read_ok += 1, |
| 250 | Ok(None) => read_bad.push(fmt!("{:?}: gone", k)), |
| 251 | Err(e) => read_bad.push(fmt!("{:?}: {}", k, e.plain())), |
| 252 | } |
| 253 | } |
| 254 | test!(sync_log::stream(), "Read back {} of {} scanned keys after .{} {:?}; {} would not read.", |
| 255 | read_ok, rows.len(), target, damage, read_bad.len()); |
| 256 | |
| 257 | // The property: damage at the tail costs the records in that damage and no |
| 258 | // others. One damaged record may cost itself; it may not cost the store. |
| 259 | if read_ok + 1 < written { |
| 260 | let _ = db.shutdown(); |
| 261 | return Err(err!( |
| 262 | "After .{} {:?}, only {} of {} records could be read back. Damage at \ |
| 263 | the tail of an append-only log must cost the records in that damage \ |
| 264 | and no others. What would not read: {:?}", |
| 265 | target, damage, read_ok, written, read_bad; |
| 266 | Test, Data)); |
| 267 | } |
| 268 | |
| 269 | res!(db.shutdown()); |
| 270 | thread::sleep(Duration::from_millis(500)); |
| 271 | Ok(()) |
| 272 | } |
| 273 | |
| 274 | pub fn test_scan_torn_tail(_filter: &'static str) -> Outcome<()> { |
| 275 | test!(sync_log::stream(), "+---------------------------------------------+"); |
| 276 | test!(sync_log::stream(), "| SCAN OVER A TORN TAIL |"); |
| 277 | test!(sync_log::stream(), "+---------------------------------------------+"); |
| 278 | // The index first, which o3db rebuilds from the data file at start-up. Both |
| 279 | // shapes cost nothing, measured: this pair is here so that the repair stays |
| 280 | // repaired, and so the data pair below cannot be read as a general claim |
| 281 | // about damage. |
| 282 | res!(run_case("ind", Damage::CutShort, 12)); |
| 283 | res!(run_case("ind", Damage::Corrupted, 12)); |
| 284 | // And the data file, which is what the index is rebuilt FROM, so there is |
| 285 | // nothing behind it to repair it with. |
| 286 | res!(run_case("dat", Damage::CutShort, 12)); |
| 287 | res!(run_case("dat", Damage::Corrupted, 12)); |
| 288 | test!(sync_log::stream(), "Torn tail test passed."); |
| 289 | Ok(()) |
| 290 | } |