oxedyne/fe2o3/fe2o3_o3db_sync/tests/two_handles.rs
6.2 KiB, 1 run
created by r1870400018:61251, 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 | //! Two handles open on one store at once, as Oregami's `voice` command runs beside its forge. |
| 2 | //! |
| 3 | //! Reproduces the multi-process fault found on 2026-09-23 (M3 in the ore lane's store timeouts |
| 4 | //! report): a zone survey hands its writer any incomplete data file, the live file of a handle |
| 5 | //! that is still running included, without claiming it, so both handles append to one file. Each |
| 6 | //! writer takes a record's offset from the file length it last saw, and the other handle's appends |
| 7 | //! make that stale. Ignored until live files are made exclusive across handles, which is a |
| 8 | //! decision still to be taken; run it with `--ignored` to see what is lost. |
| 9 | |
| 10 | use oxedyne_fe2o3_core::prelude::*; |
| 11 | use oxedyne_fe2o3_hash::{ |
| 12 | csum::ChecksumScheme, |
| 13 | hash::HashScheme, |
| 14 | }; |
| 15 | use oxedyne_fe2o3_iop_db::api::Database; |
| 16 | use oxedyne_fe2o3_jdat::prelude::*; |
| 17 | use oxedyne_fe2o3_o3db_sync::{ |
| 18 | O3db, |
| 19 | base::cfg::OzoneConfig, |
| 20 | data::core::RestSchemesInput, |
| 21 | test::setup::{ |
| 22 | self, |
| 23 | Uid, |
| 24 | UID_LEN, |
| 25 | }, |
| 26 | }; |
| 27 | |
| 28 | use std::{ |
| 29 | collections::BTreeMap, |
| 30 | path::{ |
| 31 | Path, |
| 32 | PathBuf, |
| 33 | }, |
| 34 | }; |
| 35 | |
| 36 | type TestDb = O3db< |
| 37 | { UID_LEN }, |
| 38 | Uid, |
| 39 | (), |
| 40 | HashScheme, |
| 41 | HashScheme, |
| 42 | ChecksumScheme, |
| 43 | >; |
| 44 | |
| 45 | #[test] |
| 46 | #[ignore] |
| 47 | fn main() -> Outcome<()> { |
| 48 | log_set_level!("error"); |
| 49 | let outcome = run(); |
| 50 | log_finish_wait!(); |
| 51 | outcome |
| 52 | } |
| 53 | |
| 54 | fn run() -> Outcome<()> { |
| 55 | let interleaved = res!(interleaved()); |
| 56 | let beside = res!(beside_a_running_store()); |
| 57 | msg!("Interleaved writes from two handles: {}", interleaved); |
| 58 | msg!("A second handle opened, written and closed beside a first: {}", beside); |
| 59 | if interleaved.lost() + beside.lost() > 0 { |
| 60 | return Err(err!( |
| 61 | "Two handles on one store lost writes. Interleaved: {}. Beside a running \ |
| 62 | store: {}.", interleaved, beside; |
| 63 | Test, Data, Missing)); |
| 64 | } |
| 65 | Ok(()) |
| 66 | } |
| 67 | |
| 68 | /// Both handles open, writing in turn. |
| 69 | fn interleaved() -> Outcome<Tally> { |
| 70 | let root = res!(fresh("./test_db_two_handles_interleaved")); |
| 71 | let a = res!(open(&root)); |
| 72 | let b = res!(open(&root)); |
| 73 | let mut written = Vec::new(); |
| 74 | for i in 0..40usize { |
| 75 | for (h, db) in [("a", &a), ("b", &b)] { |
| 76 | let (k, v) = kv(h, i); |
| 77 | res!(db.insert(k.clone(), v.clone(), Uid::default(), None)); |
| 78 | written.push((k, v)); |
| 79 | } |
| 80 | } |
| 81 | res!(a.close()); |
| 82 | res!(b.close()); |
| 83 | verify(&root, &written) |
| 84 | } |
| 85 | |
| 86 | /// The forge's pattern: a store held open, a command opening a second handle beside it to write |
| 87 | /// a few records and close, and the first writing on. |
| 88 | fn beside_a_running_store() -> Outcome<Tally> { |
| 89 | let root = res!(fresh("./test_db_two_handles_beside")); |
| 90 | let forge = res!(open(&root)); |
| 91 | let mut written = Vec::new(); |
| 92 | for i in 0..10usize { |
| 93 | let (k, v) = kv("forge", i); |
| 94 | res!(forge.insert(k.clone(), v.clone(), Uid::default(), None)); |
| 95 | written.push((k, v)); |
| 96 | } |
| 97 | { |
| 98 | let cli = res!(open(&root)); |
| 99 | for i in 0..5usize { |
| 100 | let (k, v) = kv("cli", i); |
| 101 | res!(cli.insert(k.clone(), v.clone(), Uid::default(), None)); |
| 102 | written.push((k, v)); |
| 103 | } |
| 104 | res!(cli.close()); |
| 105 | } |
| 106 | for i in 10..20usize { |
| 107 | let (k, v) = kv("forge", i); |
| 108 | res!(forge.insert(k.clone(), v.clone(), Uid::default(), None)); |
| 109 | written.push((k, v)); |
| 110 | } |
| 111 | res!(forge.close()); |
| 112 | verify(&root, &written) |
| 113 | } |
| 114 | |
| 115 | #[derive(Debug, Default)] |
| 116 | struct Tally { |
| 117 | right: usize, |
| 118 | wrong: usize, // read back as some other value |
| 119 | missing: usize, |
| 120 | failed: usize, // the read itself failed |
| 121 | first: Option<String>, // what the first bad read found |
| 122 | } |
| 123 | |
| 124 | impl Tally { |
| 125 | fn lost(&self) -> usize { self.wrong + self.missing + self.failed } |
| 126 | } |
| 127 | |
| 128 | impl std::fmt::Display for Tally { |
| 129 | fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { |
| 130 | ok!(write!(f, "{} of {} keys read back right, {} wrong, {} missing, {} failed", |
| 131 | self.right, self.right + self.lost(), self.wrong, self.missing, self.failed)); |
| 132 | match &self.first { |
| 133 | Some(first) => write!(f, "; first: {}", first), |
| 134 | None => Ok(()), |
| 135 | } |
| 136 | } |
| 137 | } |
| 138 | |
| 139 | /// Reopens the store alone and reads every key written to it. |
| 140 | fn verify(root: &Path, written: &[(Dat, Dat)]) -> Outcome<Tally> { |
| 141 | let db = res!(open(root)); |
| 142 | let mut tally = Tally::default(); |
| 143 | for (k, v) in written { |
| 144 | let bad = match db.get(k, None) { |
| 145 | Ok(Some((got, _))) if &got == v => { |
| 146 | tally.right += 1; |
| 147 | None |
| 148 | }, |
| 149 | Ok(Some((got, _))) => { |
| 150 | tally.wrong += 1; |
| 151 | Some(fmt!("{:?} read back as {:?}", k, got)) |
| 152 | }, |
| 153 | Ok(None) => { |
| 154 | tally.missing += 1; |
| 155 | Some(fmt!("{:?} is missing", k)) |
| 156 | }, |
| 157 | Err(e) => { |
| 158 | tally.failed += 1; |
| 159 | Some(fmt!("{:?} failed: {:?}", k, e)) |
| 160 | }, |
| 161 | }; |
| 162 | if tally.first.is_none() { |
| 163 | tally.first = bad; |
| 164 | } |
| 165 | } |
| 166 | res!(db.close()); |
| 167 | Ok(tally) |
| 168 | } |
| 169 | |
| 170 | /// Values of differing lengths, so that a record read from another's offset cannot pass for it. |
| 171 | fn kv(handle: &str, i: usize) -> (Dat, Dat) { |
| 172 | ( |
| 173 | dat!(fmt!("two handles {} {}", handle, i)), |
| 174 | dat!(fmt!("{}:{}:{}", handle, i, "x".repeat(i % 7 + 3))), |
| 175 | ) |
| 176 | } |
| 177 | |
| 178 | fn config() -> Outcome<OzoneConfig> { |
| 179 | let mut cfg = res!(setup::default_cfg()); |
| 180 | cfg.num_zones = 1; |
| 181 | cfg.num_wbots_per_zone = 1; |
| 182 | // Large enough that nothing rolls over: every record goes to the one live file. |
| 183 | cfg.data_file_max_bytes = 1_000_000; |
| 184 | cfg.zone_overrides = BTreeMap::new(); |
| 185 | Ok(cfg) |
| 186 | } |
| 187 | |
| 188 | fn schemes() -> RestSchemesInput<(), HashScheme, HashScheme, ChecksumScheme> { |
| 189 | RestSchemesInput::new( |
| 190 | None::<()>, |
| 191 | None::<HashScheme>, |
| 192 | None::<HashScheme>, |
| 193 | Some(ChecksumScheme::new_crc32()), |
| 194 | ) |
| 195 | } |
| 196 | |
| 197 | fn fresh(dir: &str) -> Outcome<PathBuf> { |
| 198 | let _ = std::fs::remove_dir_all(dir); |
| 199 | res!(std::fs::create_dir_all(dir)); |
| 200 | Ok(res!(Path::new(dir).canonicalize())) |
| 201 | } |
| 202 | |
| 203 | fn open(root: &Path) -> Outcome<TestDb> { |
| 204 | let mut db = res!(TestDb::new(root.to_path_buf(), Some(res!(config())), schemes(), Uid::default())); |
| 205 | res!(db.start("test")); |
| 206 | Ok(db) |
| 207 | } |