Oregami
Repositories/oxedyne/fe2o3

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
25use oxedyne_fe2o3_core::{
26 prelude::*,
27 alt::Override,
28};
29use oxedyne_fe2o3_crypto::enc::EncryptionScheme;
30use oxedyne_fe2o3_hash::{
31 csum::ChecksumScheme,
32 hash::HashScheme,
33};
34use oxedyne_fe2o3_iop_db::api::{
35 Database,
36 RestSchemesOverride,
37};
38use oxedyne_fe2o3_jdat::prelude::*;
39use oxedyne_fe2o3_o3db_sync::{
40 data::core::RestSchemesInput,
41 test::setup,
42};
43
44use std::{
45 path::Path,
46 thread,
47 time::Duration,
48};
49
50pub 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
114fn 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<()>
129where
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}