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 | |
| 27 | use oxedyne_fe2o3_core::prelude::*; |
| 28 | use oxedyne_fe2o3_hash::{ |
| 29 | csum::ChecksumScheme, |
| 30 | hash::HashScheme, |
| 31 | }; |
| 32 | use oxedyne_fe2o3_jdat::prelude::*; |
| 33 | use oxedyne_fe2o3_o3db_sync::{ |
| 34 | O3db, |
| 35 | data::core::RestSchemesInput, |
| 36 | test::setup::{ |
| 37 | self, |
| 38 | Uid, |
| 39 | UID_LEN, |
| 40 | }, |
| 41 | }; |
| 42 | |
| 43 | use 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. |
| 57 | const CLOSE_TIMEOUT: Duration = Duration::from_secs(30); |
| 58 | |
| 59 | fn the_key() -> Dat { |
| 60 | dat!("a value that has to survive a shared close") |
| 61 | } |
| 62 | |
| 63 | fn the_value() -> Dat { |
| 64 | dat!(42u8) |
| 65 | } |
| 66 | |
| 67 | #[test] |
| 68 | fn main() -> Outcome<()> { |
| 69 | log_set_level!("warn"); |
| 70 | let outcome = run(); |
| 71 | log_finish_wait!(); |
| 72 | outcome |
| 73 | } |
| 74 | |
| 75 | fn 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. |
| 164 | type TestDb = O3db< |
| 165 | { UID_LEN }, |
| 166 | Uid, |
| 167 | (), |
| 168 | HashScheme, |
| 169 | HashScheme, |
| 170 | ChecksumScheme, |
| 171 | >; |
| 172 | |
| 173 | /// Keeps whatever is already in the store. |
| 174 | fn 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 | } |