oxedyne/fe2o3/fe2o3_o3db_sync/src/db.rs
22.3 KiB, 137 runs
created by r1870400018:787, 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 | use crate::{ |
| 2 | prelude::*, |
| 3 | base::{ |
| 4 | constant, |
| 5 | id::{ |
| 6 | Bid, |
| 7 | OzoneBotId, |
| 8 | }, |
| 9 | }, |
| 10 | bots::{ |
| 11 | base::{ |
| 12 | bot::{ |
| 13 | BotInitArgs, |
| 14 | OzoneBot, |
| 15 | }, |
| 16 | handles::Handle, |
| 17 | }, |
| 18 | bot_super::Supervisor, |
| 19 | }, |
| 20 | comm::{ |
| 21 | channels::BotChannels, |
| 22 | msg::OzoneMsg, |
| 23 | response::Responder, |
| 24 | }, |
| 25 | data::{ |
| 26 | core::{ |
| 27 | RestSchemes, |
| 28 | RestSchemesInput, |
| 29 | }, |
| 30 | }, |
| 31 | file::core::find_files, |
| 32 | }; |
| 33 | |
| 34 | use oxedyne_fe2o3_bot::Bot; |
| 35 | use oxedyne_fe2o3_core::{ |
| 36 | channels::{ |
| 37 | simplex, |
| 38 | Simplex, |
| 39 | Recv, |
| 40 | }, |
| 41 | path::NormalPath, |
| 42 | rand::RanDef, |
| 43 | thread::{ |
| 44 | Sentinel, |
| 45 | thread_channel, |
| 46 | }, |
| 47 | }; |
| 48 | use oxedyne_fe2o3_jdat::{ |
| 49 | prelude::*, |
| 50 | cfg::Config, |
| 51 | file::JdatMapFile, |
| 52 | id::NumIdDat, |
| 53 | }; |
| 54 | use oxedyne_fe2o3_namex::id::{ |
| 55 | InNamex, |
| 56 | NamexId, |
| 57 | }; |
| 58 | use oxedyne_fe2o3_text::string::Stringer; |
| 59 | |
| 60 | use std::{ |
| 61 | fs, |
| 62 | path::{ |
| 63 | Path, |
| 64 | PathBuf, |
| 65 | }, |
| 66 | sync::{ |
| 67 | Arc, |
| 68 | Mutex, |
| 69 | RwLock, |
| 70 | mpsc, |
| 71 | }, |
| 72 | thread, |
| 73 | time::{ |
| 74 | Duration, |
| 75 | Instant, |
| 76 | }, |
| 77 | }; |
| 78 | |
| 79 | use crossbeam_utils::sync::WaitGroup; |
| 80 | |
| 81 | |
| 82 | /// The main Ozone database struct. |
| 83 | /// |
| 84 | /// # Storage specifications |
| 85 | /// The user can change the various data transformation schemes in three ways. |
| 86 | /// ## Invocation |
| 87 | /// Upon invocation, rest schemes conforming to the traits `oxedyne_fe2o3_iop_hash::csum::Checksummer`, |
| 88 | /// `oxedyne_fe2o3_iop_hash::api::Hasher`, `oxedyne_fe2o3_iop_crypto::enc::Encrypter` and |
| 89 | /// `oxedyne_fe2o3_iop_crypto::sign::Signer` can be given to the `O3db` instance. When a scheme is not |
| 90 | /// provided, a hardwired default is set. |
| 91 | /// ## Configuration |
| 92 | /// Default schemes can be overridden at invocation or upon any subsequent configuration file |
| 93 | /// changes. Schemes are limited to those provided in `oxedyne_fe2o3_hash` and `oxedyne_fe2o3_crypto`. |
| 94 | /// ## Per value basis |
| 95 | /// Finally, schemes for storage of data at rest can be set explicitely for any key-value pair, |
| 96 | /// overriding invocation or default schemes. |
| 97 | /// |
| 98 | /// # Directory layout |
| 99 | /// The example below shows a database that has been used with 3 and 5 zones. The database is |
| 100 | /// invoked with the absolute path to the db_root, and an optional `OzoneConfig` which is used if |
| 101 | /// a configuration file is not found. |
| 102 | /// |
| 103 | /// ```ignore |
| 104 | /// |
| 105 | /// /../my_o3db <- db_root with db_name, aka db_container |
| 106 | /// ├── config.jdat |
| 107 | /// ├── 003_zone <- zone_root |
| 108 | /// │ ├── zone_001 <- zone_dir |
| 109 | /// │ │ ├── 000_000_001.dat |
| 110 | /// │ │ ├── 000_000_001.ind |
| 111 | /// │ │ ├── 000_000_002.dat |
| 112 | /// │ │ └── 000_000_002.ind |
| 113 | /// │ ├── zone_002 <- zone_dir |
| 114 | /// │ │ ├── 000_000_001.dat |
| 115 | /// │ │ ├── 000_000_001.ind |
| 116 | /// │ │ ├── 000_000_002.dat |
| 117 | /// │ │ └── 000_000_002.ind |
| 118 | /// │ └── zone_003 <- zone_dir |
| 119 | /// │ ├── 000_000_001.dat |
| 120 | /// │ ├── 000_000_001.ind |
| 121 | /// │ ├── 000_000_002.dat |
| 122 | /// │ └── 000_000_002.ind |
| 123 | /// └── 005_zone <- zone_root |
| 124 | /// ├── zone_002 <- zone_dir |
| 125 | /// │ ├── 000_000_001.dat |
| 126 | /// │ └── 000_000_001.ind |
| 127 | /// ├── zone_003 <- zone_dir |
| 128 | /// │ ├── 000_000_001.dat |
| 129 | /// │ └── 000_000_001.ind |
| 130 | /// ├── zone_004 <- zone_dir |
| 131 | /// │ ├── 000_000_001.dat |
| 132 | /// │ └── 000_000_001.ind |
| 133 | /// └── zone_005 <- zone_dir |
| 134 | /// ├── 000_000_001.dat |
| 135 | /// └── 000_000_001.ind |
| 136 | /// |
| 137 | /// a_zone_container_dir <- zone_container |
| 138 | /// └── 005_zone <- zone_root |
| 139 | /// └─── zone_001 <- zone_dir |
| 140 | /// ├── 000_000_001.dat |
| 141 | /// └── 000_000_001.ind |
| 142 | /// ``` |
| 143 | #[derive(Clone, Debug)] |
| 144 | pub struct O3db< |
| 145 | const UIDL: usize, // User identifier byte length. |
| 146 | UID: NumIdDat<UIDL>, // User identifier. |
| 147 | ENC: Encrypter, // Symmetric encryption of data at rest. |
| 148 | KH: Hasher, // Hashes database keys. |
| 149 | PR: Hasher, // Pseudo-randomiser hash to distribute cache data. |
| 150 | CS: Checksummer, // Checks integrity of data at rest. |
| 151 | >{ |
| 152 | db_root: PathBuf, |
| 153 | chan_inbox: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>, |
| 154 | api: OzoneApi<UIDL, UID, ENC, KH, PR, CS>, |
| 155 | closing: Arc<Mutex<Closing>>, |
| 156 | } |
| 157 | |
| 158 | #[derive(Debug)] |
| 159 | struct Closing { |
| 160 | wg: Option<WaitGroup>, |
| 161 | done: bool, |
| 162 | } |
| 163 | |
| 164 | impl< |
| 165 | const UIDL: usize, |
| 166 | UID: NumIdDat<UIDL> + 'static, |
| 167 | ENC: Encrypter + 'static, |
| 168 | KH: Hasher + 'static, |
| 169 | PR: Hasher + 'static, |
| 170 | CS: Checksummer + 'static, |
| 171 | > |
| 172 | O3db<UIDL, UID, ENC, KH, PR, CS> |
| 173 | { |
| 174 | /// Create a new Ozone database instance. Some validation is performed, but the database is |
| 175 | /// not properly activated until `O3db::start` is called. |
| 176 | pub fn new<P: Into<PathBuf>>( |
| 177 | db_root: P, |
| 178 | cfg_opt: Option<OzoneConfig>, |
| 179 | schms_input: RestSchemesInput<ENC, KH, PR, CS>, |
| 180 | _uid_template: UID, |
| 181 | ) |
| 182 | -> Outcome<Self> |
| 183 | { |
| 184 | // Check constants. |
| 185 | res!(OzoneConfig::check_constants()); |
| 186 | |
| 187 | let db_root = db_root.into(); |
| 188 | if !db_root.exists() { |
| 189 | warn!(sync_log::stream(), "Ozone database root directory {:?} does not exist, attempting to create...", |
| 190 | db_root); |
| 191 | res!(fs::create_dir_all(&db_root)); |
| 192 | info!(sync_log::stream(), "{:?} created.", db_root); |
| 193 | } |
| 194 | |
| 195 | let cfg_path = OzoneConfig::config_path(&db_root); |
| 196 | let mut cfg = if cfg_path.is_file() { |
| 197 | res!(<OzoneConfig as JdatMapFile>::load(&cfg_path)) |
| 198 | } else { |
| 199 | match cfg_opt { |
| 200 | Some(cfg) => { |
| 201 | res!(cfg.save(&cfg_path, " ", true)); |
| 202 | warn!(sync_log::stream(), |
| 203 | "Configuration file {:?} saved, using default configuration provided.", |
| 204 | cfg_path, |
| 205 | ); |
| 206 | cfg |
| 207 | } |
| 208 | None => return Err(err!( |
| 209 | "You must supply an OzoneConfig."; |
| 210 | Input, Missing)), |
| 211 | } |
| 212 | }; |
| 213 | |
| 214 | // Check configuration. |
| 215 | res!(cfg.check_and_fix()); |
| 216 | |
| 217 | // File system. |
| 218 | let zone_root = cfg.zone_root(&db_root); |
| 219 | res!(fs::create_dir_all(&zone_root)); |
| 220 | |
| 221 | let schms = RestSchemes::from(schms_input); |
| 222 | |
| 223 | let chans = BotChannels::new(&cfg); // This is the original from which all derive. |
| 224 | let ozid = OzoneBotId::Master(Bid::randef()); |
| 225 | |
| 226 | let api = OzoneApi::new( |
| 227 | ozid, |
| 228 | db_root.clone(), |
| 229 | cfg, |
| 230 | chans, |
| 231 | schms, |
| 232 | ); |
| 233 | |
| 234 | Ok(Self { |
| 235 | db_root, |
| 236 | chan_inbox: simplex(), |
| 237 | api, |
| 238 | closing: Arc::new(Mutex::new(Closing { |
| 239 | wg: None, |
| 240 | done: false, |
| 241 | })), |
| 242 | }) |
| 243 | } |
| 244 | |
| 245 | pub fn db_root(&self) -> &Path { &self.db_root } |
| 246 | pub fn api(&self) -> &OzoneApi<UIDL, UID, ENC, KH, PR, CS> { &self.api } |
| 247 | pub fn api_mut(&mut self) -> &mut OzoneApi<UIDL, UID, ENC, KH, PR, CS> { &mut self.api } |
| 248 | |
| 249 | /// Thread-safe mutable sharing of the API. |
| 250 | pub fn share_api(self) -> Arc<RwLock<OzoneApi<UIDL, UID, ENC, KH, PR, CS>>> { |
| 251 | Arc::new(RwLock::new(self.api)) |
| 252 | } |
| 253 | |
| 254 | pub fn updated_api(&mut self) -> Outcome<&mut OzoneApi<UIDL, UID, ENC, KH, PR, CS>> { |
| 255 | res!(self.update()); |
| 256 | Ok(&mut self.api) |
| 257 | } |
| 258 | |
| 259 | // Convenience. |
| 260 | pub fn ozid(&self) -> &OzoneBotId { &self.api.ozid } |
| 261 | pub fn cfg(&self) -> &OzoneConfig { &self.api.cfg } |
| 262 | pub fn chans(&self) -> &BotChannels<UIDL, UID, ENC, KH> { &self.api.chans } |
| 263 | pub fn schemes(&self) -> &RestSchemes<ENC, KH, PR, CS> { &self.api.schms } |
| 264 | pub fn responder(&self) -> Responder<UIDL, UID, ENC, KH> { Responder::new(Some(&self.ozid())) } |
| 265 | pub fn no_responder() -> Responder<UIDL, UID, ENC, KH> { Responder::none(None) } |
| 266 | |
| 267 | pub fn update(&mut self) -> Outcome<()> { |
| 268 | let ozid = self.api.ozid.clone(); |
| 269 | Self::drain(&self.chan_inbox, &ozid, &mut self.api.chans, &mut self.api.cfg) |
| 270 | } |
| 271 | |
| 272 | fn drain( |
| 273 | inbox: &Simplex<OzoneMsg<UIDL, UID, ENC, KH>>, |
| 274 | ozid: &OzoneBotId, |
| 275 | chans: &mut BotChannels<UIDL, UID, ENC, KH>, |
| 276 | cfg: &mut OzoneConfig, |
| 277 | ) |
| 278 | -> Outcome<()> |
| 279 | { |
| 280 | loop { // loop to ensure we get the latest BotChannels |
| 281 | match inbox.try_recv() { |
| 282 | Recv::Empty => break, |
| 283 | Recv::Result(Err(e)) => { |
| 284 | return Err(e); |
| 285 | }, |
| 286 | Recv::Result(Ok(msg)) => match msg { |
| 287 | OzoneMsg::Channels(new_chans, resp) => { |
| 288 | *chans = new_chans; |
| 289 | res!(resp.send( |
| 290 | OzoneMsg::ChannelsReceived(ozid.clone())) |
| 291 | ); |
| 292 | }, |
| 293 | OzoneMsg::Config(new_cfg) => { |
| 294 | *cfg = new_cfg; |
| 295 | }, |
| 296 | _ => { |
| 297 | return Err(err!( |
| 298 | "{}: Unrecognised channel update message: {:?}.", |
| 299 | ozid, msg; |
| 300 | Invalid, Input, Channel)); |
| 301 | }, |
| 302 | } |
| 303 | } |
| 304 | } |
| 305 | Ok(()) |
| 306 | } |
| 307 | |
| 308 | /// Start the Ozone database. Returns once every zone has surveyed its files and loaded its |
| 309 | /// caches, bounded by `constant::CONTROL_REQUEST_TIMEOUT`. A start that fails returns once the |
| 310 | /// bots it brought up have stopped, so the directory can be opened again at once. |
| 311 | pub fn start< |
| 312 | S: Into<String>, |
| 313 | >( |
| 314 | &mut self, |
| 315 | log_stream_id: S, |
| 316 | ) |
| 317 | -> Outcome<Handle< UIDL, UID, ENC, KH>> |
| 318 | { |
| 319 | let log_stream_id = log_stream_id.into(); |
| 320 | sync_log::set_stream(log_stream_id.clone()); |
| 321 | |
| 322 | for line in constant::SPLASH.split("\n") { |
| 323 | info!(sync_log::stream(), "{}", line); |
| 324 | } |
| 325 | for line in Stringer::new(fmt!("{:?}", self.schemes())).to_lines(" ") { |
| 326 | info!(sync_log::stream(), "{}", line); |
| 327 | } |
| 328 | // Write config to a file now that we have a directory structure. |
| 329 | res!(self.cfg().write_config_file(self.db_root())); |
| 330 | |
| 331 | // Create and start the supervisor. |
| 332 | let (semaphore, sentinel) = thread_channel(); |
| 333 | let api = OzoneApi::new( |
| 334 | OzoneBotId::Supervisor(Bid::randef()), |
| 335 | self.db_root.clone(), |
| 336 | self.cfg().clone(), |
| 337 | self.chans().clone(), |
| 338 | self.schemes().clone(), |
| 339 | ); |
| 340 | let args = BotInitArgs { |
| 341 | // Bot |
| 342 | sem: semaphore, |
| 343 | log_stream_id, |
| 344 | // Comms |
| 345 | chan_in: self.chans().sup().clone(), |
| 346 | // API |
| 347 | api, |
| 348 | }; |
| 349 | let mut sup = Supervisor::new( |
| 350 | args, |
| 351 | self.chan_inbox.clone(), |
| 352 | ); |
| 353 | res!(sup.init()); // Starts all the other bots. |
| 354 | // One clone for the database side, however many handles share it, and |
| 355 | // one for the supervisor thread to drop when it has ended. |
| 356 | let wg_end = sup.handles().wait_end_ref().clone(); |
| 357 | { |
| 358 | let mut closing = lock_mutex!(self.closing, |
| 359 | "Taking the shutdown record while starting the database."); |
| 360 | closing.wg = Some(sup.handles().wait_end_ref().clone()); |
| 361 | closing.done = false; |
| 362 | } |
| 363 | |
| 364 | let sup_ozid = sup.ozid().clone(); |
| 365 | |
| 366 | let builder = thread::Builder::new() |
| 367 | .name(sup_ozid.to_string()) |
| 368 | .stack_size(constant::STACK_SIZE); |
| 369 | res!(builder.spawn(move || { |
| 370 | sup.go(); |
| 371 | drop(wg_end); |
| 372 | })); |
| 373 | |
| 374 | let handle = Handle::new( |
| 375 | Some(sup_ozid), |
| 376 | sentinel.clone(), |
| 377 | Some(self.chans().sup().clone()), |
| 378 | ); |
| 379 | |
| 380 | if let Err(e) = self.await_ready(&sentinel) { |
| 381 | // A supervisor that failed to start the database stops its bots itself. One still |
| 382 | // starting is asked to once it is done, since nobody will use what it starts. |
| 383 | let stop = OzoneMsg::Shutdown(self.ozid().clone(), Responder::none(Some(self.ozid()))); |
| 384 | if let Err(e2) = self.chans().sup().send(stop) { |
| 385 | error!(sync_log::stream(), err!(e2, |
| 386 | "{}: Asking the supervisor to stop after a failed start.", self.ozid(); |
| 387 | Channel, Write)); |
| 388 | } |
| 389 | // Nothing is left for a close to do. |
| 390 | let wg = { |
| 391 | let mut closing = lock_mutex!(self.closing, |
| 392 | "Taking the shutdown record after a failed start."); |
| 393 | closing.done = true; |
| 394 | closing.wg.take() |
| 395 | }; |
| 396 | // The caller may open the directory again as soon as this returns, as Oregami's |
| 397 | // forge does on its next request, so it returns once the bots are gone. It returned |
| 398 | // while they were still stopping, and the next open could bring up a second set over |
| 399 | // the same files (2026-09-23). Their stopping is start-up work, held to the control |
| 400 | // deadline like the rest of it. |
| 401 | if sentinel.was_interrupted() { |
| 402 | // A panic while bringing the database up is caught and stops the bots like any |
| 403 | // failed start, so this is a panic in stopping them, and nothing is left that |
| 404 | // would stop the rest. |
| 405 | return Err(err!(e, |
| 406 | "{}: The database did not start, and its supervisor panicked while stopping \ |
| 407 | the bots it had brought up, some of which may still be running.", self.ozid(); |
| 408 | Init, Thread, Panic)); |
| 409 | } |
| 410 | if let Some(wg) = wg { |
| 411 | if !res!(Self::await_stopped(wg, constant::CONTROL_REQUEST_TIMEOUT)) { |
| 412 | return Err(err!(e, |
| 413 | "{}: The database did not start, and the bots it brought up had not all \ |
| 414 | stopped {:?} later.", self.ozid(), constant::CONTROL_REQUEST_TIMEOUT; |
| 415 | Init, Timeout)); |
| 416 | } |
| 417 | } |
| 418 | return Err(e); |
| 419 | } |
| 420 | |
| 421 | info!(sync_log::stream(), "Database initialisation and activation complete."); |
| 422 | |
| 423 | Ok(handle) |
| 424 | } |
| 425 | |
| 426 | /// Takes the channels the supervisor hands over and waits for it to say every zone is ready. |
| 427 | /// |
| 428 | /// This was a one-second sleep, after which the caller read whatever had arrived. A slower |
| 429 | /// start -- a starved machine spawning a few dozen threads -- left this handle with the |
| 430 | /// channels it was built with, which no bot reads, and its first request timed out on nothing |
| 431 | /// (2026-09-23). A get issued before a zone had loaded its cache could also find nothing. |
| 432 | fn await_ready(&mut self, sentinel: &Sentinel) -> Outcome<()> { |
| 433 | let begun = Instant::now(); |
| 434 | loop { |
| 435 | match self.chan_inbox.recv_timeout(constant::CHECK_INTERVAL) { |
| 436 | Recv::Empty => { |
| 437 | if sentinel.is_finished() { |
| 438 | return Err(err!( |
| 439 | "{}: The supervisor stopped before the database was ready.", |
| 440 | self.ozid(); |
| 441 | Init, Thread)); |
| 442 | } |
| 443 | if begun.elapsed() > constant::CONTROL_REQUEST_TIMEOUT { |
| 444 | return Err(err!( |
| 445 | "{}: The database was not ready within {:?} of starting. A zone \ |
| 446 | surveying a very large store can take long; this deadline is \ |
| 447 | constant::CONTROL_REQUEST_TIMEOUT.", self.ozid(), |
| 448 | constant::CONTROL_REQUEST_TIMEOUT; |
| 449 | Init, Timeout)); |
| 450 | } |
| 451 | }, |
| 452 | Recv::Result(Err(e)) => return Err(err!(e, |
| 453 | "{}: While waiting for the supervisor to start the database.", self.ozid(); |
| 454 | Init, Channel, Read)), |
| 455 | Recv::Result(Ok(msg)) => match msg { |
| 456 | OzoneMsg::Channels(chans, resp) => { |
| 457 | self.api.chans = chans; |
| 458 | res!(resp.send(OzoneMsg::ChannelsReceived(self.api.ozid.clone()))); |
| 459 | }, |
| 460 | OzoneMsg::Config(cfg) => self.api.cfg = cfg, |
| 461 | OzoneMsg::Ready => return Ok(()), |
| 462 | OzoneMsg::Error(e) => return Err(err!(e, |
| 463 | "{}: The database did not start.", self.ozid(); |
| 464 | Init)), |
| 465 | msg => return Err(err!( |
| 466 | "{}: Unexpected message while the database was starting: {:?}.", |
| 467 | self.ozid(), msg; |
| 468 | Init, Channel, Unexpected)), |
| 469 | }, |
| 470 | } |
| 471 | } |
| 472 | } |
| 473 | |
| 474 | /// Waits for every thread holding a clone of the wait group to let it go, for at most |
| 475 | /// `within`. `WaitGroup::wait` has no deadline, so the waiting is done by a thread of its own. |
| 476 | fn await_stopped(wg: WaitGroup, within: Duration) -> Outcome<bool> { |
| 477 | let (tx, rx) = mpsc::channel(); |
| 478 | let waiter = res!(thread::Builder::new() |
| 479 | .name(fmt!("o3db-stopping")) |
| 480 | .spawn(move || { |
| 481 | wg.wait(); |
| 482 | // A caller whose wait ran out has gone, and there is no one else to tell. |
| 483 | let _ = tx.send(()); |
| 484 | })); |
| 485 | match rx.recv_timeout(within) { |
| 486 | Ok(()) => match waiter.join() { |
| 487 | // Joined, so that it is not itself still running when the caller looks. |
| 488 | Ok(()) => Ok(true), |
| 489 | Err(_) => Err(err!( |
| 490 | "The thread waiting for the database's bots to stop panicked."; |
| 491 | Thread, Panic)), |
| 492 | }, |
| 493 | Err(_) => Ok(false), |
| 494 | } |
| 495 | } |
| 496 | |
| 497 | /// Find all data and index files of the existing database. |
| 498 | pub fn find_all_data_files(&self) -> Outcome<Vec<PathBuf>> { |
| 499 | |
| 500 | let mut found_files = Vec::new(); |
| 501 | |
| 502 | let cur_dir = res!(std::env::current_dir()); |
| 503 | info!(sync_log::stream(), "The current directory is {}", cur_dir.display()); |
| 504 | |
| 505 | let db_root = &self.db_root; |
| 506 | |
| 507 | info!(sync_log::stream(), "Searching for all data and index files in {:?}", db_root); |
| 508 | |
| 509 | if db_root.exists() && db_root.is_dir() { |
| 510 | let files = res!(find_files(&db_root)); |
| 511 | for file in files { |
| 512 | found_files.push(file); |
| 513 | } |
| 514 | } |
| 515 | |
| 516 | for (zind_dat, zone_dat) in self.cfg().zone_overrides() { |
| 517 | if let Ok(Some(Dat::Str(dir))) = zone_dat.map_get(&dat!("dir")) { |
| 518 | let dir = db_root.join(dir).normalise(); |
| 519 | info!(sync_log::stream(), "Searching for all data and index files in zone {:?} override {:?}", |
| 520 | zind_dat, dir); |
| 521 | let files = res!(find_files(&dir)); |
| 522 | for file in files { |
| 523 | found_files.push(file); |
| 524 | } |
| 525 | } |
| 526 | } |
| 527 | |
| 528 | Ok(found_files) |
| 529 | } |
| 530 | |
| 531 | /// Gracefully shut down the database, including the supervisor. |
| 532 | /// |
| 533 | /// Consumes the handle, which is the right shape when there is only one. |
| 534 | /// Where the database is shared -- an `Arc`, or a clone held by a worker |
| 535 | /// thread -- use [`Self::close`], which asks the same of the supervisor |
| 536 | /// through a borrow. |
| 537 | pub fn shutdown(self) -> Outcome<()> { |
| 538 | self.close() |
| 539 | } |
| 540 | |
| 541 | pub fn close(&self) -> Outcome<()> { |
| 542 | // Held for the whole of the shutdown, so that a second caller waits |
| 543 | // here and finds the work already done rather than doing it again. |
| 544 | let mut closing = lock_mutex!(self.closing, |
| 545 | "Taking the shutdown record while closing the database."); |
| 546 | if closing.done { |
| 547 | return Ok(()); |
| 548 | } |
| 549 | |
| 550 | // The latest channel set the supervisor has broadcast, since the |
| 551 | // shutdown request goes down whichever channel is current. Into local |
| 552 | // copies: `&self` cannot write them back into the handle, and after a |
| 553 | // shutdown there is nothing left for them to be useful to. |
| 554 | let mut chans = self.api.chans.clone(); |
| 555 | let mut cfg = self.api.cfg.clone(); |
| 556 | res!(Self::drain(&self.chan_inbox, &self.api.ozid, &mut chans, &mut cfg)); |
| 557 | |
| 558 | let self_id = self.ozid(); |
| 559 | let resp = self.responder(); |
| 560 | if let Err(e) = chans.sup().send( |
| 561 | OzoneMsg::Shutdown(self_id.clone(), resp.clone()) |
| 562 | ) { |
| 563 | return Err(err!(e, |
| 564 | "{}: Cannot send shutdown request to supervisor.", self_id; |
| 565 | Channel, Write)); |
| 566 | } |
| 567 | warn!(sync_log::stream(), "Shutdown: Waiting for response from supervisor..."); |
| 568 | match res!(resp.recv_timeout(constant::USER_REQUEST_TIMEOUT)) { |
| 569 | OzoneMsg::Error(e) => return Err(err!(e, |
| 570 | "{}: The supervisor had a problem during shutdown.", self_id; |
| 571 | Thread)), |
| 572 | OzoneMsg::Ok => (), |
| 573 | msg => return Err(err!( |
| 574 | "{}: Unexpected response from supervisor during shutdown: {:?}", self_id, msg; |
| 575 | Channel, Unexpected)), |
| 576 | } |
| 577 | warn!(sync_log::stream(), "Shutdown: Succesfully completed by supervisor, waiting for final \ |
| 578 | verification of termination of all threads..."); |
| 579 | // Taken rather than cloned: `WaitGroup::wait` counts the clones that are |
| 580 | // left, so waiting on a copy while the original lives never returns. |
| 581 | if let Some(wg) = closing.wg.take() { |
| 582 | wg.wait(); |
| 583 | } |
| 584 | closing.done = true; |
| 585 | warn!(sync_log::stream(), "Shutdown: Verified."); |
| 586 | Ok(()) |
| 587 | } |
| 588 | |
| 589 | } |
| 590 | |
| 591 | impl< |
| 592 | const UIDL: usize, |
| 593 | UID: NumIdDat<UIDL> + 'static, |
| 594 | ENC: Encrypter + 'static, |
| 595 | KH: Hasher + 'static, |
| 596 | PR: Hasher + 'static, |
| 597 | CS: Checksummer + 'static, |
| 598 | > |
| 599 | InNamex for O3db<UIDL, UID, ENC, KH, PR, CS> |
| 600 | { |
| 601 | fn name_id(&self) -> Outcome<NamexId> { |
| 602 | NamexId::try_from(constant::NAMEX_ID) |
| 603 | } |
| 604 | } |