oxedyne/fe2o3/fe2o3_o3db_sync/tests/durability.rs
5.9 KiB, 12 runs
created by r1870400018:10994, 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 | //! Ozone durability barrier integration test. |
| 2 | //! |
| 3 | //! Exercises the new `sync_on_write` / `sync_every_n_writes` / |
| 4 | //! `sync_interval_ms` config knobs by: |
| 5 | //! |
| 6 | //! 1. Building a database with `sync_on_write = true`, writing a |
| 7 | //! handful of keys, reading them back, and asserting that every |
| 8 | //! write completed through the `sync_data` code path without |
| 9 | //! error. This is the strongest-guarantee policy and is the one |
| 10 | //! most operators will reach for. |
| 11 | //! 2. Repeating the exercise with `sync_every_n_writes = 3` and |
| 12 | //! with `sync_interval_ms = 50`, so all three policies of a |
| 13 | //! writer's syncer (`bots::worker::syncer`) run under test. |
| 14 | //! |
| 15 | //! The test does not attempt to measure actual disk sync (that |
| 16 | //! requires strace or a kernel tracepoint and is tied to the |
| 17 | //! filesystem). Instead it asserts the end-to-end correctness |
| 18 | //! property an operator cares about: that acknowledged writes are |
| 19 | //! readable through a fresh live file pair after each policy's sync |
| 20 | //! cadence has fired at least once. |
| 21 | //! |
| 22 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 23 | //! Anthropic Claude |
| 24 | |
| 25 | use oxedyne_fe2o3_core::{ |
| 26 | prelude::*, |
| 27 | alt::Override, |
| 28 | }; |
| 29 | use oxedyne_fe2o3_crypto::enc::EncryptionScheme; |
| 30 | use oxedyne_fe2o3_hash::{ |
| 31 | csum::ChecksumScheme, |
| 32 | hash::HashScheme, |
| 33 | }; |
| 34 | use oxedyne_fe2o3_iop_db::api::{ |
| 35 | Database, |
| 36 | RestSchemesOverride, |
| 37 | }; |
| 38 | use oxedyne_fe2o3_jdat::prelude::*; |
| 39 | use oxedyne_fe2o3_o3db_sync::{ |
| 40 | data::core::RestSchemesInput, |
| 41 | test::setup, |
| 42 | }; |
| 43 | |
| 44 | use std::{ |
| 45 | path::Path, |
| 46 | thread, |
| 47 | time::Duration, |
| 48 | }; |
| 49 | |
| 50 | pub fn test_durability(_filter: &'static str) -> Outcome<()> { |
| 51 | |
| 52 | let db_root = res!(Path::new("./test_db_durability").canonicalize().or_else(|_| { |
| 53 | ok!(std::fs::create_dir_all("./test_db_durability")); |
| 54 | Path::new("./test_db_durability").canonicalize() |
| 55 | })); |
| 56 | |
| 57 | let enckey = [0x5au8; 32]; |
| 58 | let aes_gcm = res!(EncryptionScheme::new_aes_256_gcm_with_key(&enckey[..])); |
| 59 | let crc32 = ChecksumScheme::new_crc32(); |
| 60 | let schms2: RestSchemesOverride<EncryptionScheme, HashScheme> = |
| 61 | RestSchemesOverride::default() |
| 62 | .set_encrypter(Override::Default(aes_gcm.clone())); |
| 63 | let schms2 = Some(&schms2); |
| 64 | let user = setup::Uid::default(); |
| 65 | |
| 66 | let schms_input = RestSchemesInput::new( |
| 67 | Some(aes_gcm.clone()), |
| 68 | None::<HashScheme>, |
| 69 | None::<HashScheme>, |
| 70 | Some(crc32.clone()), |
| 71 | ); |
| 72 | |
| 73 | // Three independent sub-tests: strongest-guarantee, group-commit |
| 74 | // by count, and group-commit by time window. |
| 75 | res!(run_with_policy( |
| 76 | "sync_on_write", |
| 77 | &db_root, |
| 78 | &schms_input, |
| 79 | schms2, |
| 80 | user, |
| 81 | |cfg| { |
| 82 | cfg.sync_on_write = true; |
| 83 | }, |
| 84 | 12, |
| 85 | )); |
| 86 | res!(run_with_policy( |
| 87 | "sync_every_n_writes=3", |
| 88 | &db_root, |
| 89 | &schms_input, |
| 90 | schms2, |
| 91 | user, |
| 92 | |cfg| { |
| 93 | cfg.sync_every_n_writes = 3; |
| 94 | }, |
| 95 | 9, |
| 96 | )); |
| 97 | res!(run_with_policy( |
| 98 | "sync_interval_ms=50", |
| 99 | &db_root, |
| 100 | &schms_input, |
| 101 | schms2, |
| 102 | user, |
| 103 | |cfg| { |
| 104 | cfg.sync_interval_ms = 50; |
| 105 | }, |
| 106 | 6, |
| 107 | )); |
| 108 | |
| 109 | test!(sync_log::stream(), |
| 110 | "Durability barrier test passed under every sync policy branch."); |
| 111 | Ok(()) |
| 112 | } |
| 113 | |
| 114 | fn run_with_policy<F>( |
| 115 | label: &str, |
| 116 | db_root: &std::path::PathBuf, |
| 117 | schms_input: &RestSchemesInput< |
| 118 | EncryptionScheme, |
| 119 | HashScheme, |
| 120 | HashScheme, |
| 121 | ChecksumScheme, |
| 122 | >, |
| 123 | schms2: Option<&RestSchemesOverride<EncryptionScheme, HashScheme>>, |
| 124 | user: setup::Uid, |
| 125 | apply: F, |
| 126 | count: u32, |
| 127 | ) |
| 128 | -> Outcome<()> |
| 129 | where |
| 130 | F: FnOnce(&mut oxedyne_fe2o3_o3db_sync::base::cfg::OzoneConfig), |
| 131 | { |
| 132 | test!(sync_log::stream(), "+--- durability: {} ---", label); |
| 133 | |
| 134 | let mut cfg = res!(setup::default_cfg()); |
| 135 | cfg.num_zones = 2; |
| 136 | cfg.num_cbots_per_zone = 1; |
| 137 | cfg.num_wbots_per_zone = 1; |
| 138 | cfg.num_igbots_per_zone = 1; |
| 139 | cfg.data_file_max_bytes = 200_000; |
| 140 | cfg.zone_overrides = mapdat!{ |
| 141 | 1u16 => mapdat!{ "dir" => "", "max_size" => 10_000_000u64 }, |
| 142 | 2u16 => mapdat!{ "dir" => "", "max_size" => 10_000_000u64 }, |
| 143 | }.get_map().unwrap(); |
| 144 | apply(&mut cfg); |
| 145 | |
| 146 | let mut db = res!(setup::start_db( |
| 147 | db_root.clone(), |
| 148 | Some(cfg.clone()), |
| 149 | schms_input.clone(), |
| 150 | None, |
| 151 | false, // gc off: keep the write path simple. |
| 152 | true, // wipe: every sub-test starts clean. |
| 153 | )); |
| 154 | |
| 155 | thread::sleep(Duration::from_millis(200)); |
| 156 | |
| 157 | for i in 0..count { |
| 158 | res!(db.insert( |
| 159 | dat!(fmt!("dkey:{:03}", i)), |
| 160 | dat!(fmt!("dval_{}", i)), |
| 161 | user, |
| 162 | schms2, |
| 163 | )); |
| 164 | } |
| 165 | |
| 166 | // Time-window sub-test needs to wait past the interval so the |
| 167 | // last few writes actually fire a sync before we stop. |
| 168 | thread::sleep(Duration::from_millis(120)); |
| 169 | |
| 170 | // Read every key back, to assert the write side acknowledged |
| 171 | // and the read side can see the value. If any sync step |
| 172 | // errored out, the write would have propagated that upstream |
| 173 | // and this loop would fail here. |
| 174 | for i in 0..count { |
| 175 | let key = dat!(fmt!("dkey:{:03}", i)); |
| 176 | let got = res!(db.get(&key, schms2)); |
| 177 | match got { |
| 178 | Some((val, _meta)) => { |
| 179 | let expected = dat!(fmt!("dval_{}", i)); |
| 180 | if val != expected { |
| 181 | return Err(err!( |
| 182 | "{}: round-tripped value mismatch: key={}, \ |
| 183 | got={:?}, expected={:?}", |
| 184 | label, i, val, expected; |
| 185 | Test, Mismatch)); |
| 186 | } |
| 187 | }, |
| 188 | None => return Err(err!( |
| 189 | "{}: key {} not found after durable write.", |
| 190 | label, i; |
| 191 | Test, Missing)), |
| 192 | } |
| 193 | } |
| 194 | |
| 195 | test!(sync_log::stream(), |
| 196 | "+--- durability: {} : {} keys round-tripped ---", |
| 197 | label, count); |
| 198 | Ok(()) |
| 199 | } |