Oregami
Repositories/oxedyne/fe2o3

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

1use 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
34use oxedyne_fe2o3_bot::Bot;
35use 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};
48use oxedyne_fe2o3_jdat::{
49 prelude::*,
50 cfg::Config,
51 file::JdatMapFile,
52 id::NumIdDat,
53};
54use oxedyne_fe2o3_namex::id::{
55 InNamex,
56 NamexId,
57};
58use oxedyne_fe2o3_text::string::Stringer;
59
60use 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
79use 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)]
144pub 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)]
159struct Closing {
160 wg: Option<WaitGroup>,
161 done: bool,
162}
163
164impl<
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
591impl<
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}