oxedyne/fe2o3/fe2o3_o3db_sync/tests/chunked_value.rs
10.3 KiB, 10 runs
created by r1870400018:13997, 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 | //! A value too large for one record must come back as the value that was stored. |
| 2 | //! |
| 3 | //! A value whose encoding exceeds the chunk size is not written whole: it is split, each chunk is |
| 4 | //! stored under its own part key, and the key the caller used holds a `Tup5u64` part key pointing |
| 5 | //! at them. On the way back out, `fetch_chunks` collects the chunks, rejoins the bytes, decrypts |
| 6 | //! them and decodes them -- and hands back the caller's original `Dat`, fully formed. |
| 7 | //! |
| 8 | //! `get_wait` then treated that `Dat` as though it were still raw bytes and tried to decode it a |
| 9 | //! second time, so a chunked value only survived the round trip if it happened to be a byte string |
| 10 | //! -- and even then it came back wrong, its payload decoded as though it were an encoding. Every |
| 11 | //! other kind fell through to a catch-all and returned `Unexpected Dat ... returned`. |
| 12 | //! |
| 13 | //! Nothing in the suite stored a value large enough to be chunked, so the read path was never run. |
| 14 | //! The values that grow past the threshold in practice are the accumulating ones -- a ledger, an |
| 15 | //! append-only list, a document -- which is to say the ones a caller least expects to lose. This |
| 16 | //! test stores each of the kinds a caller actually stores, at a size that forces chunking, and |
| 17 | //! insists on getting back what it put in. |
| 18 | //! |
| 19 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 20 | //! Anthropic Claude |
| 21 | |
| 22 | use oxedyne_fe2o3_core::{ |
| 23 | prelude::*, |
| 24 | alt::Override, |
| 25 | rand::Rand, |
| 26 | }; |
| 27 | use oxedyne_fe2o3_crypto::enc::EncryptionScheme; |
| 28 | use oxedyne_fe2o3_hash::{ |
| 29 | csum::ChecksumScheme, |
| 30 | hash::HashScheme, |
| 31 | }; |
| 32 | use oxedyne_fe2o3_iop_db::api::{ |
| 33 | Database, |
| 34 | RestSchemesOverride, |
| 35 | }; |
| 36 | use oxedyne_fe2o3_jdat::{ |
| 37 | prelude::*, |
| 38 | file::JdatMapFile, |
| 39 | }; |
| 40 | use oxedyne_fe2o3_o3db_sync::{ |
| 41 | base::cfg::OzoneConfig, |
| 42 | data::core::RestSchemesInput, |
| 43 | test::setup, |
| 44 | }; |
| 45 | |
| 46 | use std::{ |
| 47 | path::Path, |
| 48 | thread, |
| 49 | time::Duration, |
| 50 | }; |
| 51 | |
| 52 | |
| 53 | // Enough records that the encoded list comfortably exceeds any chunk size the |
| 54 | // database is configured with, so the write path is certain to have split it. |
| 55 | const RECORDS: usize = 400; |
| 56 | |
| 57 | |
| 58 | /// A list of maps, which is the shape of every accumulating value a caller keeps under one key: |
| 59 | /// a ledger, an audit log, a list of bindings. |
| 60 | fn ledger() -> Dat { |
| 61 | let mut out = Vec::with_capacity(RECORDS); |
| 62 | for i in 0..RECORDS { |
| 63 | out.push(create_dat_ordmap(vec![ |
| 64 | (dat!("account_id"), dat!(fmt!("acct_{:08}", i))), |
| 65 | (dat!("kind"), dat!(if i == 0 { "topup" } else { "spend" })), |
| 66 | (dat!("delta_minor"), Dat::I64(if i == 0 { 5_000 } else { -1 })), |
| 67 | (dat!("ref"), dat!(fmt!("mail:someone{}@example.com", i))), |
| 68 | ])); |
| 69 | } |
| 70 | Dat::List(out) |
| 71 | } |
| 72 | |
| 73 | /// A long string, the other everyday large value: a document, a brief, a blob of text. |
| 74 | fn document() -> Dat { |
| 75 | let mut s = String::new(); |
| 76 | for i in 0..RECORDS { |
| 77 | s.push_str(&fmt!("line {} of a document long enough to be split across chunks.\n", i)); |
| 78 | } |
| 79 | dat!(s) |
| 80 | } |
| 81 | |
| 82 | /// A large byte string. This is the one kind the old code did not reject outright -- it decoded |
| 83 | /// the payload a second time, as though the bytes were themselves an encoding -- so it came back |
| 84 | /// as something else entirely, or failed. Silent corruption is worse than the error, and this |
| 85 | /// insists on the bytes. |
| 86 | fn blob() -> Dat { |
| 87 | let mut v = vec![0u8; 64 * 1024]; |
| 88 | Rand::fill_u8(&mut v); |
| 89 | Dat::BU32(v) |
| 90 | } |
| 91 | |
| 92 | pub fn test_chunked_value(_filter: &'static str) -> Outcome<()> { |
| 93 | |
| 94 | res!(std::fs::create_dir_all("./test_db_chunked_value")); |
| 95 | let db_root = res!(Path::new("./test_db_chunked_value").canonicalize()); |
| 96 | |
| 97 | let mut enckey = [0u8; 32]; |
| 98 | Rand::fill_u8(&mut enckey); |
| 99 | let aes_gcm = res!(EncryptionScheme::new_aes_256_gcm_with_key(&enckey[..])); |
| 100 | let crc32 = ChecksumScheme::new_crc32(); |
| 101 | let schms2: RestSchemesOverride<EncryptionScheme, HashScheme> = |
| 102 | RestSchemesOverride::default().set_encrypter(Override::Default(aes_gcm.clone())); |
| 103 | let schms2 = Some(&schms2); |
| 104 | let user = setup::Uid::default(); |
| 105 | let schms_input = RestSchemesInput::new( |
| 106 | Some(aes_gcm.clone()), |
| 107 | None::<HashScheme>, |
| 108 | None::<HashScheme>, |
| 109 | Some(crc32.clone()), |
| 110 | ); |
| 111 | |
| 112 | let mut cfg = res!(setup::default_cfg()); |
| 113 | cfg.sync_on_write = true; |
| 114 | cfg.zone_overrides = DaticleMap::new(); |
| 115 | |
| 116 | let cases: Vec<(&str, Dat)> = vec![ |
| 117 | ("chunked:ledger", ledger()), |
| 118 | ("chunked:document", document()), |
| 119 | ("chunked:blob", blob()), |
| 120 | ]; |
| 121 | |
| 122 | // ── Session 1: store values large enough to be chunked ─────── |
| 123 | { |
| 124 | test!(sync_log::stream(), "Session 1: storing values that exceed the chunk size."); |
| 125 | let mut db = res!(setup::start_db( |
| 126 | db_root.clone(), |
| 127 | Some(cfg.clone()), |
| 128 | schms_input.clone(), |
| 129 | None, |
| 130 | true, |
| 131 | true, // wipe: start from nothing |
| 132 | )); |
| 133 | |
| 134 | for (k, v) in &cases { |
| 135 | res!(db.insert(dat!(*k), v.clone(), user, schms2)); |
| 136 | } |
| 137 | |
| 138 | // Read them back in the same session, before any restart: the fault is in the read path, |
| 139 | // not in what reaches the disk, so it shows up immediately. |
| 140 | for (k, v) in &cases { |
| 141 | match res!(db.get(&dat!(*k), schms2)) { |
| 142 | Some((got, _)) => if got != *v { |
| 143 | return Err(err!( |
| 144 | "The value stored under {:?} did not survive the round trip. A value \ |
| 145 | larger than the chunk size is split on the way in and rejoined on the \ |
| 146 | way out, and what came back is not what went in.", k; |
| 147 | Test, Invalid, Data)); |
| 148 | }, |
| 149 | None => return Err(err!( |
| 150 | "The key {:?} was just written and cannot be read back.", k; |
| 151 | Test, Missing, Data)), |
| 152 | } |
| 153 | } |
| 154 | |
| 155 | res!(db.shutdown()); |
| 156 | } |
| 157 | |
| 158 | thread::sleep(Duration::from_secs(1)); |
| 159 | |
| 160 | // ── Session 2: and they must survive a restart ─────────────── |
| 161 | { |
| 162 | test!(sync_log::stream(), "Session 2: the chunked values must still read back."); |
| 163 | let db = res!(setup::start_db( |
| 164 | db_root.clone(), |
| 165 | Some(cfg.clone()), |
| 166 | schms_input.clone(), |
| 167 | None, |
| 168 | true, |
| 169 | false, // do not wipe: read what session 1 left |
| 170 | )); |
| 171 | |
| 172 | for (k, v) in &cases { |
| 173 | match res!(db.get(&dat!(*k), schms2)) { |
| 174 | Some((got, _)) => if got != *v { |
| 175 | return Err(err!( |
| 176 | "After a restart, the value under {:?} came back changed.", k; |
| 177 | Test, Invalid, Data)); |
| 178 | }, |
| 179 | None => return Err(err!( |
| 180 | "After a restart, the chunked value under {:?} is gone.", k; |
| 181 | Test, Missing, Data)), |
| 182 | } |
| 183 | } |
| 184 | |
| 185 | res!(db.shutdown()); |
| 186 | } |
| 187 | |
| 188 | // ── Session 3: the chunk configuration changes underneath them ── |
| 189 | // |
| 190 | // An operator raising a store's chunking threshold -- which is what a store carrying a |
| 191 | // too-small one needs -- must not thereby orphan everything already written under the old |
| 192 | // one. It does not, and this is the reason: a chunked value's geometry (how many chunks, how |
| 193 | // big) travels in the part key stored with the value, so the reader reconstructs it from what |
| 194 | // it finds, never from what the configuration currently says. The claim is worth a test |
| 195 | // rather than an assurance, because a wrong answer here is an unreadable production store. |
| 196 | { |
| 197 | test!(sync_log::stream(), "Session 3: the chunk configuration is raised on an \ |
| 198 | existing store; what was written under the old one must still read."); |
| 199 | |
| 200 | // Change the chunk geometry the store writes under. The values already on disk were split |
| 201 | // into 64-byte chunks; from now on the store splits into much larger ones, so anything |
| 202 | // read back that took its geometry from the configuration rather than from the value |
| 203 | // would be reassembled at the wrong stride and come back as rubbish. |
| 204 | // |
| 205 | // (The threshold is raised only as far as the store's own validation allows -- it must |
| 206 | // stay under 80% of the data file size, which these test files make small. That check is |
| 207 | // the same one that will refuse an ill-judged production setting.) |
| 208 | let cfg_path = OzoneConfig::config_path(&db_root); |
| 209 | let mut stored = res!(<OzoneConfig as JdatMapFile>::load(&cfg_path)); |
| 210 | stored.rest_chunk_bytes = 500; |
| 211 | res!(stored.save(&cfg_path, " ", true)); |
| 212 | |
| 213 | let db = res!(setup::start_db( |
| 214 | db_root.clone(), |
| 215 | Some(cfg.clone()), |
| 216 | schms_input.clone(), |
| 217 | None, |
| 218 | true, |
| 219 | false, // do not wipe: read what the earlier sessions left |
| 220 | )); |
| 221 | |
| 222 | for (k, v) in &cases { |
| 223 | match res!(db.get(&dat!(*k), schms2)) { |
| 224 | Some((got, _)) => if got != *v { |
| 225 | return Err(err!( |
| 226 | "The value under {:?} was chunked under one configuration and read back \ |
| 227 | under another, and it came back changed. A value's chunk geometry must \ |
| 228 | come from the part key stored with it, not from the current settings.", k; |
| 229 | Test, Invalid, Data)); |
| 230 | }, |
| 231 | None => return Err(err!( |
| 232 | "Changing the chunk geometry lost the value under {:?}, which was written \ |
| 233 | under the old one.", k; |
| 234 | Test, Missing, Data)), |
| 235 | } |
| 236 | } |
| 237 | |
| 238 | // And a value written under the NEW geometry reads back too, so the two coexist in one |
| 239 | // store rather than the store having to be all of one or all of the other. |
| 240 | let mut db = db; |
| 241 | let fresh = ledger(); |
| 242 | res!(db.insert(dat!("chunked:after"), fresh.clone(), user, schms2)); |
| 243 | match res!(db.get(&dat!("chunked:after"), schms2)) { |
| 244 | Some((got, _)) => if got != fresh { |
| 245 | return Err(err!( |
| 246 | "A value written under the new chunk geometry did not survive the round \ |
| 247 | trip."; Test, Invalid, Data)); |
| 248 | }, |
| 249 | None => return Err(err!( |
| 250 | "A value written under the new chunk geometry cannot be read back."; |
| 251 | Test, Missing, Data)), |
| 252 | } |
| 253 | |
| 254 | res!(db.shutdown()); |
| 255 | } |
| 256 | |
| 257 | Ok(()) |
| 258 | } |