Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/tests/shared_shutdown.rs

6.8 KiB, 36 runs

created by r1870400018:21570, 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//! Closing a database that more than one handle is holding.
2//!
3//! `O3db` is `Clone`, and until now the only way to close one was
4//! `shutdown(self)`, which consumes. That shape has two failures, and both of
5//! them are silent hangs rather than errors:
6//!
7//! 1. A handle inside an `Arc` -- which is how a server shares a store with the
8//! threads doing its slow work -- cannot be consumed at all, because
9//! `Arc::try_unwrap` never succeeds while another thread is holding one.
10//! 2. Even when a handle *can* be consumed, the shutdown ended on
11//! `WaitGroup::wait`, and a wait group counts its clones. Every clone of the
12//! database carried a clone of the group, so the shutdown waited for handles
13//! that were still alive and never returned.
14//!
15//! `O3db::close(&self)` is the answer to both. What is checked here is that it
16//! returns while a second handle is alive, that calling it twice is safe, and --
17//! the part that matters to whatever is storing data -- that a database closed
18//! this way opens again with its contents intact.
19//!
20//! Every wait in this file is bounded and reported as a failure. A test for a
21//! deadlock that deadlocks says nothing except that somebody has to press
22//! Ctrl-C.
23//!
24//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
25//! Anthropic Claude
26
27use oxedyne_fe2o3_core::prelude::*;
28use oxedyne_fe2o3_hash::{
29 csum::ChecksumScheme,
30 hash::HashScheme,
31};
32use oxedyne_fe2o3_jdat::prelude::*;
33use oxedyne_fe2o3_o3db_sync::{
34 O3db,
35 data::core::RestSchemesInput,
36 test::setup::{
37 self,
38 Uid,
39 UID_LEN,
40 },
41};
42
43use std::{
44 path::PathBuf,
45 sync::{
46 mpsc,
47 Arc,
48 },
49 thread,
50 time::Duration,
51};
52
53// How long a close may take before it is called a deadlock. A shutdown of an
54// idle database is a message and a thread join, and takes well under a second.
55// Thirty seconds is not a measurement, it is the difference between a failing
56// test and a hung one.
57const CLOSE_TIMEOUT: Duration = Duration::from_secs(30);
58
59fn the_key() -> Dat {
60 dat!("a value that has to survive a shared close")
61}
62
63fn the_value() -> Dat {
64 dat!(42u8)
65}
66
67#[test]
68fn main() -> Outcome<()> {
69 log_set_level!("warn");
70 let outcome = run();
71 log_finish_wait!();
72 outcome
73}
74
75fn run() -> Outcome<()> {
76
77 // A directory of this test's own. A fixture shared with another test is
78 // wiped by whichever of them gets there second.
79 let db_dir = PathBuf::from("./test_db_shared_shutdown");
80 let _ = std::fs::remove_dir_all(&db_dir);
81 res!(std::fs::create_dir_all(&db_dir));
82 let db_root = res!(db_dir.canonicalize());
83
84 // ------------------------------------------------ write, then close shared
85 {
86 let db = Arc::new(res!(open(&db_root)));
87 thread::sleep(Duration::from_secs(1));
88
89 let user = Default::default();
90 let resp = res!(db.api().store(the_key(), the_value(), user));
91 if let Err(e) = resp.recv_store_ack() {
92 return Err(err!(e,
93 "The database refused the write this test is built on."; Test, IO));
94 }
95 match res!(db.api().get_wait(&the_key(), None)) {
96 Some((v, _meta)) => req!(v, the_value(),
97 "The value did not come back before the database was closed."),
98 None => return Err(err!(
99 "The value was not there before the database was closed, so \
100 nothing this test goes on to check would mean anything."; Test, Missing)),
101 }
102
103 // A second handle, held for the whole of the close. This is the shape
104 // that used to hang: a clone alive means a clone of the wait group
105 // alive, and the wait group's wait counts clones.
106 let held = db.clone();
107
108 // The close itself, on a thread, so that a hang is a failure rather
109 // than a test that never ends.
110 let closing = db.clone();
111 let (tx, rx) = mpsc::channel();
112 res!(thread::Builder::new()
113 .name("closing".to_string())
114 .spawn(move || {
115 let _ = tx.send(closing.close());
116 }));
117 match rx.recv_timeout(CLOSE_TIMEOUT) {
118 Ok(result) => res!(result),
119 Err(_) => return Err(err!(
120 "A close through a shared handle did not return within {:?}, \
121 with a second handle alive. That is the wait group counting its \
122 own clones again.", CLOSE_TIMEOUT; Test, Timeout)),
123 }
124
125 // And a second close, which two threads both noticing a stop will make.
126 // It must be safe and it must not wait for anything.
127 let (tx, rx) = mpsc::channel();
128 res!(thread::Builder::new()
129 .name("closing-again".to_string())
130 .spawn(move || {
131 let _ = tx.send(held.close());
132 }));
133 match rx.recv_timeout(CLOSE_TIMEOUT) {
134 Ok(result) => res!(result),
135 Err(_) => return Err(err!(
136 "Closing an already closed database did not return within {:?}.",
137 CLOSE_TIMEOUT; Test, Timeout)),
138 }
139 }
140
141 // ------------------------------------------------ and it opens again
142 // The claim worth making: a store closed through a shared handle is a store
143 // that can be opened, not one that has to be repaired first.
144 {
145 let db = res!(open(&db_root));
146 thread::sleep(Duration::from_secs(1));
147 match res!(db.api().get_wait(&the_key(), None)) {
148 Some((v, _meta)) => req!(v, the_value(),
149 "A database closed through a shared handle reopened with the \
150 wrong value in it."),
151 None => return Err(err!(
152 "A database closed through a shared handle reopened without the \
153 value it was holding."; Test, Missing, Data)),
154 }
155 res!(db.shutdown());
156 }
157
158 let _ = std::fs::remove_dir_all(&db_dir);
159 Ok(())
160}
161
162/// No encryption, so that the second opening reads what the first one wrote
163/// without a key having to be carried between them.
164type TestDb = O3db<
165 { UID_LEN },
166 Uid,
167 (),
168 HashScheme,
169 HashScheme,
170 ChecksumScheme,
171>;
172
173/// Keeps whatever is already in the store.
174fn open(db_root: &PathBuf) -> Outcome<TestDb> {
175 let schms_input = RestSchemesInput::new(
176 None::<()>,
177 None::<HashScheme>,
178 None::<HashScheme>,
179 Some(ChecksumScheme::new_crc32()),
180 );
181 let mut cfg = res!(setup::default_cfg());
182 // Every zone in this test's own directory, not the shared container the default names.
183 cfg.zone_overrides = Default::default();
184 setup::start_db(
185 db_root.clone(),
186 Some(cfg),
187 schms_input,
188 None,
189 false,
190 // The config file is left alone: this test reopens the same database
191 // and a wipe would make it a different one.
192 false,
193 )
194}