Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/tests/store_timeouts.rs

17.4 KiB, 11 runs

created by r1870400018:61247, 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//! On a busy machine Oregami's store reported writes as failed that went on to land, and came up
2//! answering nothing (2026-09-23). Each fault is reproduced here without the machine having to be
3//! busy: `test::hooks` holds every durability barrier, or the supervisor, for longer than the
4//! deadline that used to expire. The hooks are process-wide, which is why this is a test binary
5//! of its own with a single test.
6
7use oxedyne_fe2o3_core::prelude::*;
8use oxedyne_fe2o3_hash::{
9 csum::ChecksumScheme,
10 hash::HashScheme,
11};
12use oxedyne_fe2o3_iop_db::api::Database;
13use oxedyne_fe2o3_jdat::prelude::*;
14use oxedyne_fe2o3_o3db_sync::{
15 O3db,
16 base::{
17 cfg::OzoneConfig,
18 constant,
19 },
20 comm::msg::OzoneMsg,
21 data::core::RestSchemesInput,
22 test::{
23 hooks,
24 setup::{
25 self,
26 Uid,
27 UID_LEN,
28 },
29 },
30};
31
32use std::{
33 collections::BTreeMap,
34 os::unix::fs::PermissionsExt,
35 path::{
36 Path,
37 PathBuf,
38 },
39 thread,
40 time::{
41 Duration,
42 Instant,
43 },
44};
45
46type TestDb = O3db<
47 { UID_LEN },
48 Uid,
49 (),
50 HashScheme,
51 HashScheme,
52 ChecksumScheme,
53>;
54
55#[test]
56fn main() -> Outcome<()> {
57 log_set_level!("warn");
58 let outcome = run();
59 // Whatever failed, the next binary must not inherit a slow store.
60 hooks::set_barrier_delay(Duration::ZERO);
61 hooks::set_publish_delay(Duration::ZERO);
62 hooks::set_supervisor_panic(false);
63 log_finish_wait!();
64 outcome
65}
66
67/// Every check runs, whichever fail, and all that fail are reported.
68fn run() -> Outcome<()> {
69 let slow_start = slow_start_hands_over_live_channels();
70 hooks::set_publish_delay(Duration::ZERO);
71 let slow_barrier = slow_barrier_still_acknowledges();
72 hooks::set_barrier_delay(Duration::ZERO);
73 let deadline = durability_deadline_reports_written();
74 hooks::set_barrier_delay(Duration::ZERO);
75 let writer = writer_failure_reaches_the_caller();
76 let zone = failed_zone_fails_the_start();
77 let panicked = panicked_supervisor_stops_its_bots();
78 hooks::set_supervisor_panic(false);
79 let failed: Vec<Error<ErrTag>> = [slow_start, slow_barrier, deadline, writer, zone, panicked]
80 .into_iter()
81 .filter_map(|r| r.err())
82 .collect();
83 match failed.len() {
84 0 => Ok(()),
85 1 => match failed.into_iter().next() {
86 Some(e) => Err(e),
87 None => Ok(()), // unreachable
88 },
89 n => Err(err!(
90 "{} checks failed: {:?}", n, failed;
91 Test)),
92 }
93}
94
95/// The supervisor takes longer to hand over its channels than the one second `start` used to
96/// sleep. The handle must hold the live channels when `start` returns, with no `updated_api`
97/// and no sleep: before, every bot but the supervisor was unreachable, and a ping heard one of
98/// them.
99fn slow_start_hands_over_live_channels() -> Outcome<()> {
100 let delay = Duration::from_millis(2_500);
101 let root = res!(fresh("./test_db_store_timeouts_start"));
102 hooks::set_publish_delay(delay);
103 let begun = Instant::now();
104 let opened = open(&root, res!(config()));
105 hooks::set_publish_delay(Duration::ZERO);
106 let db = res!(opened);
107 let took = begun.elapsed();
108 if took < delay {
109 return Err(err!(
110 "start returned after {:?}, before the supervisor had handed over its channels \
111 {:?} in.", took, delay;
112 Test, Timeout));
113 }
114
115 // Straight through the handle, as Oregami's wake does. A ping goes to every bot the handle
116 // holds a channel for and fails unless each of them answers.
117 let (_, pongs) = res!(db.api().ping_bots(constant::USER_REQUEST_WAIT));
118 if pongs.len() < 2 {
119 return Err(err!(
120 "Only {} bot answered a ping as soon as start returned.", pongs.len();
121 Test, Missing));
122 }
123
124 let key = dat!("reachable after a slow start");
125 res!(db.insert(key.clone(), dat!(7u8), Uid::default(), None));
126 match res!(db.get(&key, None)) {
127 Some((v, _)) => req!(v, dat!(7u8), "The value written after a slow start."),
128 None => return Err(err!(
129 "A value written after a slow start could not be read back."; Test, Missing)),
130 }
131 res!(db.close());
132 Ok(())
133}
134
135/// Every durability barrier takes longer than the user request deadline. A write used to fail
136/// at that deadline and land a moment later; now it is answered written at once and durable
137/// when the barrier completes. Writes queued behind the barrier share it, and none of them is
138/// held by it before being written.
139fn slow_barrier_still_acknowledges() -> Outcome<()> {
140 let delay = constant::USER_REQUEST_TIMEOUT + Duration::from_secs(1);
141 let root = res!(fresh("./test_db_store_timeouts_barrier"));
142 let mut cfg = res!(config());
143 cfg.sync_on_write = true;
144 let db = res!(open(&root, cfg));
145
146 hooks::set_barrier_delay(delay);
147 let begun = Instant::now();
148 let mut writers = Vec::new();
149 for i in 0..3u8 {
150 let db = db.clone();
151 writers.push(res!(thread::Builder::new()
152 .name(fmt!("writer {}", i))
153 .spawn(move || -> Outcome<Duration> {
154 let begun = Instant::now();
155 res!(db.insert(key(i), dat!(i), Uid::default(), None));
156 Ok(begun.elapsed())
157 })));
158 }
159 for (i, writer) in writers.into_iter().enumerate() {
160 match writer.join() {
161 Ok(Ok(took)) => if took < delay {
162 return Err(err!(
163 "Write {} was acknowledged in {:?}, before the {:?} barrier it waits on \
164 could have completed.", i, took, delay;
165 Test, Unexpected));
166 },
167 Ok(Err(e)) => return Err(err!(e,
168 "Write {} failed behind a {:?} barrier, which the user request deadline of \
169 {:?} does not measure.", i, delay, constant::USER_REQUEST_TIMEOUT;
170 Test, Write)),
171 Err(_) => return Err(err!("Writer thread {} panicked.", i; Test, Thread)),
172 }
173 }
174 hooks::set_barrier_delay(Duration::ZERO);
175 test!(sync_log::stream(), "Three writes behind {:?} barriers took {:?}.", delay, begun.elapsed());
176 for i in 0..3u8 {
177 match res!(db.get(&key(i), None)) {
178 Some((v, _)) => req!(v, dat!(i), "The value of a slowly synced write."),
179 None => return Err(err!("Write {} is not readable.", i; Test, Missing)),
180 }
181 }
182 res!(db.close());
183
184 // Reopened, to show the writes are in the files and not just the cache.
185 let db = res!(open(&root, res!(config())));
186 for i in 0..3u8 {
187 match res!(db.get(&key(i), None)) {
188 Some((v, _)) => req!(v, dat!(i), "The value of a slowly synced write, reopened."),
189 None => return Err(err!("Write {} did not survive a reopen.", i; Test, Missing)),
190 }
191 }
192 res!(db.close());
193 Ok(())
194}
195
196/// When the durability deadline itself expires, the caller is told the write is in the files
197/// but not yet confirmed durable, and that is true: the value is readable once the barrier
198/// completes. The deadline is passed in short here, since the real one is two minutes.
199fn durability_deadline_reports_written() -> Outcome<()> {
200 let delay = Duration::from_secs(4);
201 let durability = Duration::from_secs(2);
202 let root = res!(fresh("./test_db_store_timeouts_deadline"));
203 let mut cfg = res!(config());
204 cfg.sync_on_write = true;
205 let db = res!(open(&root, cfg));
206
207 hooks::set_barrier_delay(delay);
208 let begun = Instant::now();
209 let resp = res!(db.api().store(key(9), dat!(9u8), Uid::default()));
210 let n = match res!(resp.recv_timeout(constant::USER_REQUEST_TIMEOUT)) {
211 OzoneMsg::Chunks(n) => n,
212 msg => return Err(err!("Expected the record count, received {:?}.", msg;
213 Test, Unexpected)),
214 };
215 let outcome = resp.recv_write_acks(n, constant::USER_REQUEST_TIMEOUT, durability);
216 let took = begun.elapsed();
217 hooks::set_barrier_delay(Duration::ZERO);
218 match outcome {
219 Ok(acks) => return Err(err!(
220 "A write behind a {:?} barrier was confirmed durable within {:?}: {:?}.",
221 delay, durability, acks;
222 Test, Unexpected)),
223 Err(e) => {
224 let text = fmt!("{:?}", e);
225 if !text.contains("not confirmed durable") || !text.contains("has not failed") {
226 return Err(err!(
227 "An expired durability deadline must say the write is in the files but not \
228 yet durable, and it said: {}", text;
229 Test, Mismatch));
230 }
231 if took >= delay {
232 return Err(err!(
233 "The durability deadline of {:?} took {:?} to expire.", durability, took;
234 Test, Timeout));
235 }
236 if !e.tags().contains(&ErrTag::Unconfirmed) {
237 return Err(err!(
238 "An expired durability deadline is tagged {:?}, which does not say \
239 Unconfirmed: a caller cannot tell it from a write that never landed without \
240 reading the words.", e.tags();
241 Test, Mismatch));
242 }
243 },
244 }
245 // What the error said: the write lands when the disk completes it.
246 thread::sleep(delay.saturating_sub(begun.elapsed()) + Duration::from_millis(500));
247 match res!(db.get(&key(9), None)) {
248 Some((v, _)) => req!(v, dat!(9u8), "A write reported as not yet durable."),
249 None => return Err(err!(
250 "A write reported as written but not yet durable never became readable.";
251 Test, Missing)),
252 }
253 res!(db.close());
254 Ok(())
255}
256
257/// A writer that fails says why, where it used to leave the caller to time out on nothing. The
258/// zone directory is made read-only, so the write that needs a new live file cannot have one.
259fn writer_failure_reaches_the_caller() -> Outcome<()> {
260 let root = res!(fresh("./test_db_store_timeouts_writer"));
261 let mut cfg = res!(config());
262 cfg.num_zones = 1;
263 let db = res!(open(&root, cfg.clone()));
264 let big = || dat!(vec![0x5au8; 1_000]); // two of these overflow a 2,000 byte live file
265
266 res!(db.insert(key(1), big(), Uid::default(), None));
267 let zdir = res!(zone_dir(&root, &cfg));
268 res!(std::fs::set_permissions(&zdir, std::fs::Permissions::from_mode(0o555)));
269 let begun = Instant::now();
270 let outcome = db.insert(key(2), big(), Uid::default(), None);
271 let took = begun.elapsed();
272 res!(std::fs::set_permissions(&zdir, std::fs::Permissions::from_mode(0o755)));
273 match outcome {
274 Ok(_) => return Err(err!(
275 "A write needing a live file in a read-only zone directory succeeded."; Test, Unexpected)),
276 Err(e) => {
277 let text = fmt!("{:?}", e);
278 if text.contains("Failed to receive a message via responder") {
279 return Err(err!(
280 "A writer's failure reached the caller as a timeout: {}", text; Test, Mismatch));
281 }
282 if !text.contains("new live file") || !text.contains("ermission denied") {
283 return Err(err!(
284 "A writer's failure must name its cause, and it said: {}", text; Test, Mismatch));
285 }
286 if took >= constant::USER_REQUEST_TIMEOUT {
287 return Err(err!(
288 "The writer's failure took {:?} to arrive.", took; Test, Timeout));
289 }
290 },
291 }
292 // The directory is writable again, and so is the store.
293 res!(db.insert(key(3), big(), Uid::default(), None));
294 match res!(db.get(&key(3), None)) {
295 Some((v, _)) => req!(v, big(), "A write after the zone recovered."),
296 None => return Err(err!("A write after the zone recovered is missing."; Test, Missing)),
297 }
298 res!(db.close());
299 Ok(())
300}
301
302/// A zone that cannot be initialised fails the start. It used to be logged, and the store came
303/// up with that zone's writer holding no directory, so its writes landed in the working directory
304/// of whatever process had opened it. The failed start returns once the bots it brought up have
305/// stopped: it returned with 32 threads running where there had been 4, and a caller that opens
306/// the directory again at once, as Oregami's forge does on its next request, could have two sets
307/// of bots over one directory.
308fn failed_zone_fails_the_start() -> Outcome<()> {
309 let root = res!(fresh("./test_db_store_timeouts_zone"));
310 let mut cfg = res!(config());
311 cfg.zone_overrides = res!(mapdat!{
312 1u16 => mapdat!{
313 "dir" => "../test_db_store_timeouts_no_such_container",
314 "max_size" => 1_000_000u64,
315 },
316 }.get_map().ok_or_else(|| err!("The zone override is not a map."; Test, Bug)));
317 let before = res!(threads());
318 let begun = Instant::now();
319 let mut db = res!(TestDb::new(root.clone(), Some(cfg), schemes(), Uid::default()));
320 match db.start("test") {
321 Ok(_) => return Err(err!(
322 "A store whose zone directory does not exist started."; Test, Unexpected)),
323 Err(e) => {
324 let text = fmt!("{:?}", e);
325 if !text.contains("does not exist and must be created") {
326 return Err(err!(
327 "The failed start must name the missing directory, and it said: {}", text;
328 Test, Mismatch));
329 }
330 },
331 }
332 let took = begun.elapsed();
333 if took >= constant::USER_REQUEST_TIMEOUT {
334 return Err(err!("The failed start took {:?} to report.", took; Test, Timeout));
335 }
336 // A thread that has let go of its bot ends a moment later, so the count gets that moment.
337 let at_return = res!(threads());
338 let settled = Instant::now() + Duration::from_millis(100);
339 let mut after = at_return;
340 while after > before && Instant::now() < settled {
341 thread::sleep(Duration::from_millis(5));
342 after = res!(threads());
343 }
344 if after > before {
345 return Err(err!(
346 "A failed start returned with {} threads running, {} a moment later, where there were \
347 {} before it: the bots it brought up were still running.", at_return, after, before;
348 Test, Unexpected));
349 }
350 // Nothing is left running for a close to wait on.
351 let begun = Instant::now();
352 res!(db.close());
353 if begun.elapsed() >= Duration::from_secs(1) {
354 return Err(err!(
355 "Closing a store whose start failed took {:?}.", begun.elapsed(); Test, Timeout));
356 }
357 Ok(())
358}
359
360/// The supervisor panics once it has brought the bots up, before the database is ready. The
361/// start fails and returns once those bots have stopped. It returned at once, with every bot
362/// still running and nothing left that would stop them, so the directory could be opened again
363/// beside them.
364fn panicked_supervisor_stops_its_bots() -> Outcome<()> {
365 let root = res!(fresh("./test_db_store_timeouts_panic"));
366 let cfg = res!(config());
367 let before = res!(threads());
368 hooks::set_supervisor_panic(true);
369 let mut db = res!(TestDb::new(root.clone(), Some(cfg), schemes(), Uid::default()));
370 let started = db.start("test");
371 hooks::set_supervisor_panic(false);
372 if started.is_ok() {
373 return Err(err!(
374 "A store whose supervisor panicked while starting it started."; Test, Unexpected));
375 }
376 // A thread that has let go of its bot ends a moment later, so the count gets that moment.
377 let at_return = res!(threads());
378 let settled = Instant::now() + Duration::from_millis(100);
379 let mut after = at_return;
380 while after > before && Instant::now() < settled {
381 thread::sleep(Duration::from_millis(5));
382 after = res!(threads());
383 }
384 if after > before {
385 return Err(err!(
386 "A start whose supervisor panicked returned with {} threads running, {} a moment \
387 later, where there were {} before it: the bots it brought up were still running.",
388 at_return, after, before;
389 Test, Unexpected));
390 }
391 let begun = Instant::now();
392 res!(db.close());
393 if begun.elapsed() >= Duration::from_secs(1) {
394 return Err(err!(
395 "Closing a store whose start failed took {:?}.", begun.elapsed(); Test, Timeout));
396 }
397 Ok(())
398}
399
400fn key(i: u8) -> Dat {
401 dat!(fmt!("store timeouts key {}", i))
402}
403
404/// The threads this process is running.
405fn threads() -> Outcome<usize> {
406 Ok(res!(std::fs::read_dir("/proc/self/task")).count())
407}
408
409/// The crate's test configuration, with every zone inside the test's own directory.
410fn config() -> Outcome<OzoneConfig> {
411 let mut cfg = res!(setup::default_cfg());
412 cfg.zone_overrides = BTreeMap::new();
413 Ok(cfg)
414}
415
416fn schemes() -> RestSchemesInput<(), HashScheme, HashScheme, ChecksumScheme> {
417 RestSchemesInput::new(
418 None::<()>,
419 None::<HashScheme>,
420 None::<HashScheme>,
421 Some(ChecksumScheme::new_crc32()),
422 )
423}
424
425/// An empty directory of this test's own.
426fn fresh(dir: &str) -> Outcome<PathBuf> {
427 let _ = std::fs::remove_dir_all(dir);
428 res!(std::fs::create_dir_all(dir));
429 Ok(res!(Path::new(dir).canonicalize()))
430}
431
432fn open(root: &Path, cfg: OzoneConfig) -> Outcome<TestDb> {
433 let mut db = res!(TestDb::new(root.to_path_buf(), Some(cfg), schemes(), Uid::default()));
434 res!(db.start("test"));
435 Ok(db)
436}
437
438/// The directory of the first zone, which is the only one when there is one.
439fn zone_dir(root: &Path, cfg: &OzoneConfig) -> Outcome<PathBuf> {
440 let zdir = cfg.zone_root(root).join("zone_001");
441 if !zdir.is_dir() {
442 return Err(err!("Expected a zone directory at {:?}.", zdir; Test, Missing));
443 }
444 Ok(zdir)
445}