Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/src/bots/bot_super.rs

39.9 KiB, 143 runs

created by r1870400018:743, 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::OzoneBotType,
6 index::{
7 BotPoolInd,
8 ZoneInd,
9 },
10 },
11 bots::{
12 base::{
13 bot_deps::*,
14 handles::{
15 BotHandles,
16 Handle,
17 },
18 },
19 // Solo bots
20 bot_config::ConfigBot,
21 bot_server::ServerBot,
22 // Other bots
23 bot_zone::{
24 ZoneBot,
25 ZoneState,
26 },
27 worker::{
28 bot::ZoneWorkerInitArgs,
29 bot_cache::CacheBot,
30 bot_file::FileBot,
31 bot_initgc::InitGarbageBot,
32 bot_reader::ReaderBot,
33 bot_scan::ScanBot,
34 bot_writer::WriterBot,
35 worker_deps::*,
36 },
37 },
38 test::hooks,
39};
40
41use oxedyne_fe2o3_bot::Bot;
42use oxedyne_fe2o3_core::{
43 channels::{
44 simplex,
45 Recv,
46 },
47 thread::thread_channel,
48};
49use oxedyne_fe2o3_jdat::id::NumIdDat;
50
51use std::{
52 panic::{
53 self,
54 AssertUnwindSafe,
55 },
56 sync::Arc,
57 time::{
58 Duration,
59 Instant,
60 },
61 thread,
62};
63
64/// Manages the files and caches for an ozone database.
65#[derive(Debug)]
66pub struct Supervisor<
67 const UIDL: usize,
68 UID: NumIdDat<UIDL>,
69 ENC: Encrypter,
70 KH: Hasher,
71 PR: Hasher,
72 CS: Checksummer,
73>{
74 // Bot
75 sem: Semaphore,
76 errc: Arc<Mutex<usize>>,
77 log_stream_id: String,
78 // Comms
79 chan_in: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>,
80 chan_out: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>, // to the Master.
81 // API
82 api: OzoneApi<UIDL, UID, ENC, KH, PR, CS>,
83 // State
84 handles: BotHandles<UIDL, UID, ENC, KH>,
85 inited: bool,
86 trep: Instant,
87 zstats: Vec<ZoneState>,
88}
89
90impl<
91 const UIDL: usize,
92 UID: NumIdDat<UIDL> + 'static,
93 ENC: Encrypter + 'static,
94 KH: Hasher + 'static,
95 PR: Hasher + 'static,
96 CS: Checksummer + 'static,
97>
98 Bot<{ BID_LEN }, Bid, OzoneMsg<UIDL, UID, ENC, KH>> for
99 Supervisor<UIDL, UID, ENC, KH, PR, CS>
100{
101 bot_methods!();
102
103 fn go(&mut self) {
104
105 sync_log::set_stream(self.log_stream_id());
106
107 if self.no_init() { return; }
108 // The master is blocked in `O3db::start` until it hears one or the other of these. A
109 // panic while bringing the database up is a failed start like any other, so the bots
110 // already up are stopped: unwound out of this thread, it left them running with nothing
111 // to stop them, and `start` returned with them still over the directory (2026-09-23).
112 let brought_up = match panic::catch_unwind(AssertUnwindSafe(|| self.bring_up())) {
113 Ok(result) => result,
114 Err(_) => Err(err!(
115 "{}: The supervisor panicked while bringing the database up.", self.ozid();
116 Init, Thread, Panic)),
117 };
118 if let Err(e) = brought_up {
119 self.error(e.clone());
120 if let Err(e2) = self.chan_out.send(OzoneMsg::Error(e)) {
121 self.err_cannot_send(err!(e2,
122 "{}: Telling the master the database did not start.", self.ozid();
123 IO, Channel));
124 }
125 // A start that failed leaves no bots running over the directory.
126 let requester = fmt!("{} after a failed start", self.ozid());
127 let result = self.shutdown(requester, None);
128 self.result(&result);
129 return;
130 }
131 self.now_listening();
132 loop {
133 if self.listen().must_end() { break; }
134 }
135 }
136
137 fn listen(&mut self) -> LoopBreak {
138 let now = Instant::now();
139 if now.duration_since(self.trep) >= constant::HEALTH_CHECK_INTERVAL {
140 let result = self.report_health();
141 self.result(&result);
142 self.trep = now;
143 }
144 match self.chan_in().recv() {
145 Err(e) => self.err_cannot_receive(err!(e,
146 "{}: Waiting for message.", self.ozid();
147 IO, Channel)),
148 Ok(msg) => match msg {
149 OzoneMsg::ClearCache(_) |
150 OzoneMsg::DumpCacheRequest(_) |
151 OzoneMsg::DumpFiles(_) |
152 OzoneMsg::DumpFileStatesRequest(_) |
153 OzoneMsg::GcControl(_, _) |
154 OzoneMsg::GetZoneDir(_) |
155 //OzoneMsg::GetUsers(_) |
156 OzoneMsg::NewLiveFile(_, _)
157 => {
158 match self.chans().fwd_to_all_zones(msg) {
159 Err(e) => self.error(e),
160 Ok(_) => (),
161 }
162 },
163 OzoneMsg::OzoneStateRequest(resp) => {
164 self.respond(Ok(OzoneMsg::OzoneStateResponse(
165 self.ozone_state().clone())), &resp);
166 },
167 OzoneMsg::Shutdown(ozid, resp) => {
168 if let OzoneBotId::Master(_) = ozid {
169 let result = self.shutdown(fmt!("{}", ozid), Some(&resp));
170 self.result(&result);
171 return LoopBreak(true);
172 } else {
173 self.respond(Err(err!(
174 "{} attempted to shut down database, but only the Master \
175 can do this.", ozid;
176 Unauthorised)), &resp);
177 }
178 },
179 OzoneMsg::ZoneState(z, zstat) => {
180 //debug!(sync_log::stream(), "{}: zone {} state received: {:?}",self.ozid(),z,zstat);
181 if z+1 > self.ozone_state().len() {
182 self.error(err!(
183 "{}: The ZoneInd for a zone state update, {}, exceeds the \
184 number of zbot slots {}.",
185 self.ozid(), z, self.ozone_state().len();
186 Bug, Mismatch, Size));
187 } else {
188 self.ozone_state_mut()[z] = zstat;
189 }
190 },
191 _ => return self.listen_more(msg),
192 },
193 }
194 LoopBreak(false)
195 }
196}
197
198impl<
199 const UIDL: usize,
200 UID: NumIdDat<UIDL> + 'static,
201 ENC: Encrypter + 'static,
202 KH: Hasher + 'static,
203 PR: Hasher + 'static,
204 CS: Checksummer + 'static,
205>
206 OzoneBot<UIDL, UID, ENC, KH, PR, CS> for Supervisor<UIDL, UID, ENC, KH, PR, CS>
207{
208 ozonebot_methods!();
209}
210
211impl<
212 const UIDL: usize,
213 UID: NumIdDat<UIDL> + 'static,
214 ENC: Encrypter + 'static,
215 KH: Hasher + 'static,
216 PR: Hasher + 'static,
217 CS: Checksummer + 'static,
218>
219 Supervisor<UIDL, UID, ENC, KH, PR, CS>
220{
221 pub fn new(
222 args: BotInitArgs<UIDL, UID, ENC, KH, PR, CS>,
223 chan_out: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>,
224 )
225 -> Self
226 {
227 let handles = BotHandles::new(&args.api.cfg);
228 let zstats = vec![ZoneState::default(); args.api.cfg.num_zones()];
229 Self {
230 // Bot
231 sem: args.sem,
232 errc: Arc::new(Mutex::new(0)),
233 log_stream_id: args.log_stream_id,
234 // Comms
235 chan_in: args.chan_in,
236 chan_out,
237 // API
238 api: args.api,
239 // State
240 handles,
241 inited: false,
242 trep: Instant::now(),
243 zstats,
244 }
245 }
246
247 pub fn schemes(&self) -> &RestSchemes<ENC, KH, PR, CS> { &self.api.schms }
248 pub fn handles(&self) -> &BotHandles<UIDL, UID, ENC, KH> { &self.handles }
249 pub fn chans_mut(&mut self) -> &mut BotChannels<UIDL, UID, ENC, KH> { &mut self.api.chans }
250 pub fn handles_mut(&mut self) -> &mut BotHandles<UIDL, UID, ENC, KH> { &mut self.handles }
251
252 fn ozone_state(&self) -> &Vec<ZoneState> { &self.zstats }
253 fn ozone_state_mut(&mut self) -> &mut Vec<ZoneState> { &mut self.zstats }
254
255// pub fn start_db(&mut self) -> Outcome<()> {
256//
257// info!(sync_log::stream(), "{}: Starting database...", self.label());
258//
259// let nz = self.cfg().num_zones();
260// let ns = self.cfg().num_sbots();
261//
262// res!(thread::scope(|s| -> Outcome<()> {
263//
264// // Arise bots!
265// for xz in 0..nz {
266// let zind = ZoneInd::new(xz);
267// for wtyp in [
268// WorkerType::Cache,
269// WorkerType::File,
270// WorkerType::InitGarbage,
271// WorkerType::Reader,
272// WorkerType::Writer,
273// ] {
274// for b in 0..self.cfg().num_bots_per_zone(&wtyp) {
275// let bpind = BotPoolInd::new(b);
276// let wind = WorkerInd::new(zind, bpind);
277// let handle = res!(self.start_new_worker(s, &wtyp, &wind));
278// let chan = res!(handle.some_chan());
279// { res!(self.chans_mut().set_worker_bot(&wtyp, &wind, chan)); }
280// { res!(self.handles_mut().set_worker_bot(&wtyp, &wind, handle)); }
281// }
282// }
283// }
284//
285// // The zbots are initialised with cbot channels.
286// for xz in 0..nz {
287// let zind = ZoneInd::new(xz);
288// let handle = res!(self.start_new_zbot(s, &zind));
289// let chan = res!(handle.some_chan());
290// { res!(self.chans_mut().set_zbot(&zind, chan)); }
291// { res!(self.handles_mut().set_zbot(&zind, handle)); }
292// }
293//
294// // Start cfgbot.
295// let handle = res!(self.start_new_cfgbot(s));
296// let chan = res!(handle.some_chan());
297// { self.chans_mut().set_cfg(chan); }
298// { self.handles_mut().set_cfg(handle); }
299//
300// // Start srvbots.
301// for xs in 0..ns {
302// let bpind = BotPoolInd::new(xs);
303// let handle = res!(self.start_new_srvbot(s, bpind));
304// let chan = res!(handle.some_chan());
305// { res!(self.chans_mut().set_sbot(&bpind, chan)); }
306// { res!(self.handles_mut().set_sbot(&bpind, handle)); }
307// }
308//
309// // Broadcast updated BotChannels.
310// let resp = Responder::new(Some(self.ozid()));
311// let msg = OzoneMsg::Channels(self.chans().clone(), resp.clone());
312// res!(self.chans().cfg().send(msg.clone()));
313// res!(self.chans().fwd_to_all_zones(msg.clone()));
314// for _ in 0..nz+1 {
315// match res!(resp.recv_timeout(constant::BOT_REQUEST_TIMEOUT)) {
316// OzoneMsg::ChannelsReceived(ozid) => {
317// info!(sync_log::stream(), "{}: Channels received by {}", self.ozid(), ozid);
318// },
319// m => self.error(err!(
320// "Received {:?}, expecting ChannelsReceived confirmation.", m;
321// Channel, Read, Unexpected)),
322// }
323// }
324//
325// res!(self.chans().cfg().send(OzoneMsg::ZoneInitTrigger));
326//
327// Ok(())
328//
329// }));
330//
331// Ok(())
332// }
333
334// pub fn start_db(&mut self) -> Outcome<()> {
335// info!(sync_log::stream(), "{}: Starting database...", self.label());
336//
337// let nz = self.cfg().num_zones();
338// let ns = self.cfg().num_sbots();
339//
340// // Collect all worker handles and channels before mutating `self`.
341// let mut worker_handles = Vec::new();
342// let mut worker_channels = Vec::new();
343//
344// // Collect all zbot handles and channels before mutating `self`.
345// let mut zbot_handles = Vec::new();
346// let mut zbot_channels = Vec::new();
347//
348// // Collect all srvbot handles and channels before mutating `self`.
349// let mut srvbot_handles = Vec::new();
350// let mut srvbot_channels = Vec::new();
351//
352// let (cfg_handle, cfg_chan) = res!(thread::scope(|s| -> Outcome<(
353// Handle<'s, UIDL, UID, ENC, KH, ()>,
354// Simplex<OzoneMsg<UIDL, UID, ENC, KH>>,
355// )> {
356// // Arise bots!
357// for xz in 0..nz {
358// let zind = ZoneInd::new(xz);
359// for wtyp in [
360// WorkerType::Cache,
361// WorkerType::File,
362// WorkerType::InitGarbage,
363// WorkerType::Reader,
364// WorkerType::Writer,
365// ] {
366// for b in 0..self.cfg().num_bots_per_zone(&wtyp) {
367// let bpind = BotPoolInd::new(b);
368// let wind = WorkerInd::new(zind, bpind);
369// let handle = res!(self.start_new_worker(s, &wtyp, &wind));
370// let chan = res!(handle.some_chan());
371// worker_handles.push((wtyp.clone(), wind.clone(), handle));
372// worker_channels.push((wtyp.clone(), wind.clone(), chan));
373// }
374// }
375// }
376//
377// // The zbots are initialised with cbot channels.
378// for xz in 0..nz {
379// let zind = ZoneInd::new(xz);
380// let handle = res!(self.start_new_zbot(s, &zind));
381// let chan = res!(handle.some_chan());
382// zbot_handles.push((zind, handle));
383// zbot_channels.push((zind, chan));
384// }
385//
386// // Start cfgbot.
387// let handle = res!(self.start_new_cfgbot(s));
388// let chan = res!(handle.some_chan());
389// let cfg_handle = handle;
390// let cfg_chan = chan;
391//
392// // Start srvbots.
393// for xs in 0..ns {
394// let bpind = BotPoolInd::new(xs);
395// let handle = res!(self.start_new_srvbot(s, bpind));
396// let chan = res!(handle.some_chan());
397// srvbot_handles.push((bpind, handle));
398// srvbot_channels.push((bpind, chan));
399// }
400//
401// Ok((cfg_handle, cfg_chan))
402// }));
403//
404// // Now that the immutable borrows are no longer active, mutate `self`.
405// for (wtyp, wind, handle) in worker_handles {
406// res!(self.handles_mut().set_worker_bot(&wtyp, &wind, handle));
407// }
408// for (wtyp, wind, chan) in worker_channels {
409// res!(self.chans_mut().set_worker_bot(&wtyp, &wind, chan));
410// }
411//
412// for (zind, handle) in zbot_handles {
413// res!(self.handles_mut().set_zbot(&zind, handle));
414// }
415// for (zind, chan) in zbot_channels {
416// res!(self.chans_mut().set_zbot(&zind, chan));
417// }
418//
419// self.chans_mut().set_cfg(cfg_chan);
420// self.handles_mut().set_cfg(cfg_handle);
421//
422// for (bpind, handle) in srvbot_handles {
423// res!(self.handles_mut().set_sbot(&bpind, handle));
424// }
425// for (bpind, chan) in srvbot_channels {
426// res!(self.chans_mut().set_sbot(&bpind, chan));
427// }
428//
429// // Broadcast updated BotChannels.
430// let resp = Responder::new(Some(self.ozid()));
431// let msg = OzoneMsg::Channels(self.chans().clone(), resp.clone());
432// res!(self.chans().cfg().send(msg.clone()));
433// res!(self.chans().fwd_to_all_zones(msg.clone()));
434// for _ in 0..nz + 1 {
435// match res!(resp.recv_timeout(constant::BOT_REQUEST_TIMEOUT)) {
436// OzoneMsg::ChannelsReceived(ozid) => {
437// info!(sync_log::stream(), "{}: Channels received by {}", self.ozid(), ozid);
438// }
439// m => self.error(err!(
440// "Received {:?}, expecting ChannelsReceived confirmation.", m;
441// Channel, Read, Unexpected)),
442// }
443// }
444//
445// res!(self.chans().cfg().send(OzoneMsg::ZoneInitTrigger));
446//
447// Ok(())
448// }
449
450// pub fn start_db(&mut self) -> Outcome<()> {
451// info!(sync_log::stream(), "{}: Starting database...", self.label());
452//
453// let nz = self.cfg().num_zones();
454// let ns = self.cfg().num_sbots();
455//
456// res!(thread::scope(|s| -> Outcome<()> {
457// // Start worker bots first
458// for xz in 0..nz {
459// let zind = ZoneInd::new(xz);
460// for wtyp in [
461// WorkerType::Cache,
462// WorkerType::File,
463// WorkerType::InitGarbage,
464// WorkerType::Reader,
465// WorkerType::Writer,
466// ] {
467// for b in 0..self.cfg().num_bots_per_zone(&wtyp) {
468// let bpind = BotPoolInd::new(b);
469// let wind = WorkerInd::new(zind, bpind);
470// let handle = res!(self.start_new_worker(s, &wtyp, &wind));
471// let chan = res!(handle.some_chan());
472// res!(self.chans_mut().set_worker_bot(&wtyp, &wind, chan));
473// res!(self.handles_mut().set_worker_bot(&wtyp, &wind, handle));
474// }
475// }
476// }
477//
478// // Start zone bots
479// for xz in 0..nz {
480// let zind = ZoneInd::new(xz);
481// let handle = res!(self.start_new_zbot(s, &zind));
482// let chan = res!(handle.some_chan());
483// res!(self.chans_mut().set_zbot(&zind, chan));
484// res!(self.handles_mut().set_zbot(&zind, handle));
485// }
486//
487// // Start config bot
488// let cfg_handle = res!(self.start_new_cfgbot(s));
489// let cfg_chan = res!(cfg_handle.some_chan());
490// self.chans_mut().set_cfg(cfg_chan);
491// self.handles_mut().set_cfg(cfg_handle);
492//
493// // Start server bots
494// for xs in 0..ns {
495// let bpind = BotPoolInd::new(xs);
496// let handle = res!(self.start_new_srvbot(s, bpind));
497// let chan = res!(handle.some_chan());
498// res!(self.chans_mut().set_sbot(&bpind, chan));
499// res!(self.handles_mut().set_sbot(&bpind, handle));
500// }
501//
502// // Broadcast updated BotChannels
503// let resp = Responder::new(Some(self.ozid()));
504// let msg = OzoneMsg::Channels(self.chans().clone(), resp.clone());
505// res!(self.chans().cfg().send(msg.clone()));
506// res!(self.chans().fwd_to_all_zones(msg.clone()));
507//
508// for _ in 0..nz + 1 {
509// match res!(resp.recv_timeout(constant::BOT_REQUEST_TIMEOUT)) {
510// OzoneMsg::ChannelsReceived(ozid) => {
511// info!(sync_log::stream(), "{}: Channels received by {}", self.ozid(), ozid);
512// }
513// m => {
514// return Err(err!(
515// "Received {:?}, expecting ChannelsReceived confirmation.", m;
516// Channel, Read, Unexpected
517// ));
518// }
519// }
520// }
521//
522// res!(self.chans().cfg().send(OzoneMsg::ZoneInitTrigger));
523//
524// Ok(())
525// }));
526//
527// Ok(())
528// }
529
530 pub fn start_db(&mut self) -> Outcome<()> {
531
532 info!(sync_log::stream(), "{}: Starting database...", self.label());
533
534 let nz = self.cfg().num_zones();
535 let ns = self.cfg().num_sbots();
536
537 info!(sync_log::stream(), "{}: Starting worker bots...", self.label());
538
539 for xz in 0..nz {
540 let zind = ZoneInd::new(xz);
541 for wtyp in [
542 WorkerType::Cache,
543 WorkerType::File,
544 WorkerType::InitGarbage,
545 WorkerType::Reader,
546 WorkerType::Scan,
547 WorkerType::Writer,
548 ] {
549 let nb = self.cfg().num_bots_per_zone(&wtyp);
550 info!(sync_log::stream(), "{}: Starting {} {:?} bots...", self.label(), nb+1, wtyp);
551 for b in 0..nb {
552 let bpind = BotPoolInd::new(b);
553 let wind = WorkerInd::new(zind, bpind);
554
555 // Create worker bot inline instead of in a separate method
556 let chan = simplex();
557 let (semaphore, sentinel) = thread_channel();
558 let wg = self.handles().wait_end_ref().clone();
559 let api = OzoneApi::new(
560 OzoneBotId::new_worker(&wtyp, &wind),
561 self.db_root().to_path_buf(),
562 self.cfg().clone(),
563 self.chans().clone(),
564 self.schemes().clone(),
565 );
566
567 let args = ZoneWorkerInitArgs {
568 wind: wind.clone(),
569 wtyp: wtyp.clone(),
570 sem: semaphore,
571 log_stream_id: self.log_stream_id(),
572 chan_in: chan.clone(),
573 api,
574 };
575
576 let mut bot: Box<dyn WorkerBot<UIDL, UID, ENC, KH, PR, CS>> = match wtyp {
577 WorkerType::Cache => Box::new(CacheBot::new(args)),
578 WorkerType::File => Box::new(FileBot::new(args)),
579 WorkerType::InitGarbage => Box::new(InitGarbageBot::new(args)),
580 WorkerType::Reader => Box::new(ReaderBot::new(args)),
581 WorkerType::Scan => Box::new(ScanBot::new(args)),
582 WorkerType::Writer => Box::new(res!(WriterBot::new(args))),
583 };
584
585 res!(bot.init());
586 let builder = thread::Builder::new()
587 .name(bot.id().to_string())
588 .stack_size(constant::STACK_SIZE);
589
590 let ozid = bot.ozid().clone();
591 res!(builder.spawn(move || {
592 bot.go();
593 drop(wg);
594 }));
595 let handle = Handle::new(
596 Some(ozid),
597 sentinel,
598 Some(chan.clone()),
599 );
600
601 res!(self.chans_mut().set_worker_bot(&wtyp, &wind, chan));
602 res!(self.handles_mut().set_worker_bot(&wtyp, &wind, handle));
603 }
604 }
605 }
606
607 info!(sync_log::stream(), "{}: Starting {} zone bots...", self.label(), nz+1);
608 for xz in 0..nz {
609 let zind = ZoneInd::new(xz);
610 let chan = simplex();
611 let (semaphore, sentinel) = thread_channel();
612 let wg = self.handles().wait_end_ref().clone();
613 let api = OzoneApi::new(
614 OzoneBotId::ZoneBot(Bid::randef(), zind),
615 self.db_root().to_path_buf(),
616 self.cfg().clone(),
617 self.chans().clone(),
618 self.schemes().clone(),
619 );
620
621 let args = BotInitArgs {
622 sem: semaphore,
623 log_stream_id: self.log_stream_id(),
624 chan_in: chan.clone(),
625 api,
626 };
627
628 let mut bot = ZoneBot::new(args, zind);
629 res!(bot.init());
630 let builder = thread::Builder::new()
631 .name(bot.id().to_string())
632 .stack_size(constant::STACK_SIZE);
633
634 let ozid = bot.ozid().clone();
635 res!(builder.spawn(move || {
636 bot.go();
637 drop(wg);
638 }));
639 let handle = Handle::new(
640 Some(ozid),
641 sentinel,
642 Some(chan.clone()),
643 );
644
645 res!(self.chans_mut().set_zbot(&zind, chan));
646 res!(self.handles_mut().set_zbot(&zind, handle));
647 }
648
649 info!(sync_log::stream(), "{}: Starting config bot...", self.label());
650 let cfg_chan = simplex();
651 let (cfg_semaphore, cfg_sentinel) = thread_channel();
652 let cfg_wg = self.handles().wait_end_ref().clone();
653 let cfg_api = OzoneApi::new(
654 OzoneBotId::ConfigBot(Bid::randef()),
655 self.db_root().to_path_buf(),
656 self.cfg().clone(),
657 self.chans().clone(),
658 self.schemes().clone(),
659 );
660
661 let cfg_args = BotInitArgs {
662 sem: cfg_semaphore,
663 log_stream_id: self.log_stream_id(),
664 chan_in: cfg_chan.clone(),
665 api: cfg_api,
666 };
667
668 let mut cfg_bot = ConfigBot::new(cfg_args);
669 res!(cfg_bot.init());
670 let cfg_builder = thread::Builder::new()
671 .name(cfg_bot.id().to_string())
672 .stack_size(constant::STACK_SIZE);
673
674 let cfg_ozid = cfg_bot.ozid().clone();
675 res!(cfg_builder.spawn(move || {
676 cfg_bot.go();
677 drop(cfg_wg);
678 }));
679 let cfg_handle = Handle::new(
680 Some(cfg_ozid),
681 cfg_sentinel,
682 Some(cfg_chan.clone()),
683 );
684
685 self.chans_mut().set_cfg(cfg_chan);
686 self.handles_mut().set_cfg(cfg_handle);
687
688 info!(sync_log::stream(), "{}: Starting {} server bots...", self.label(), ns+1);
689 for xs in 0..ns {
690 let bpind = BotPoolInd::new(xs);
691 let chan = simplex();
692 let (semaphore, sentinel) = thread_channel();
693 let wg = self.handles().wait_end_ref().clone();
694 let api = OzoneApi::new(
695 OzoneBotId::ServerBot(Bid::randef(), bpind),
696 self.db_root().to_path_buf(),
697 self.cfg().clone(),
698 self.chans().clone(),
699 self.schemes().clone(),
700 );
701
702 let args = BotInitArgs {
703 sem: semaphore,
704 log_stream_id: self.log_stream_id(),
705 chan_in: chan.clone(),
706 api,
707 };
708
709 let mut bot = ServerBot::new(args);
710 res!(bot.init());
711 let builder = thread::Builder::new()
712 .name(bot.id().to_string())
713 .stack_size(constant::STACK_SIZE);
714
715 let ozid = bot.ozid().clone();
716 res!(builder.spawn(move || {
717 bot.go();
718 drop(wg);
719 }));
720 let handle = Handle::new(
721 Some(ozid),
722 sentinel,
723 Some(chan.clone()),
724 );
725
726 res!(self.chans_mut().set_sbot(&bpind, chan));
727 res!(self.handles_mut().set_sbot(&bpind, handle));
728 }
729
730 info!(sync_log::stream(), "{}: Broadcasting updated BotChannels...", self.label());
731 // Broadcast updated BotChannels
732 let resp = Responder::new(Some(self.ozid()));
733 let msg = OzoneMsg::Channels(self.chans().clone(), resp.clone());
734 res!(self.chans().cfg().send(msg.clone()));
735 res!(self.chans().fwd_to_all_zones(msg.clone()));
736
737 trace!(sync_log::stream(), "expecting {} channel receipt confirmations", nz+1);
738 // Start-up work, held to the control deadline: a zone bot confirms only once each of its
739 // workers has, and on a starved machine a thread just spawned may not run for a while.
740 for i in 0..nz + 1 {
741 match res!(resp.recv_timeout(constant::CONTROL_REQUEST_TIMEOUT)) {
742 OzoneMsg::ChannelsReceived(ozid) => {
743 info!(sync_log::stream(), "{}: Channels received by {}", self.ozid(), ozid);
744 }
745 m => {
746 return Err(err!(
747 "Received {:?}, expecting ChannelsReceived confirmation.", m;
748 Channel, Read, Unexpected
749 ));
750 }
751 }
752 trace!(sync_log::stream(), "i = {}", i+1);
753 }
754 info!(sync_log::stream(), "{}: All channel updates received by config bot and zone bots.", self.label());
755
756 Ok(())
757 }
758
759 /// Starts every bot, has every zone survey its files and load its caches, hands the master the
760 /// channels the bots listen on, and tells it the database is ready once every zone is.
761 ///
762 /// The master used to sleep for a second and read whatever had arrived by then, so a start that
763 /// took longer -- a starved machine spawning a few dozen threads -- left it holding channels no
764 /// bot reads, and its first request timed out on nothing. It now waits on `OzoneMsg::Ready`.
765 fn bring_up(&mut self) -> Outcome<()> {
766 res!(self.start_db());
767 hooks::supervisor_panic();
768
769 // 1. The zones start surveying at once, and answer on `zones` when they are done.
770 let zones = Responder::new(Some(self.ozid()));
771 res!(self.chans().cfg().send(OzoneMsg::ZoneInitTrigger(zones.clone())));
772
773 // 2. Meanwhile, the master is given the channels.
774 hooks::publish_delay();
775 let resp = Responder::new(Some(self.ozid()));
776 res!(self.chan_out.send(OzoneMsg::Channels(self.chans().clone(), resp.clone())));
777 match res!(resp.recv_timeout(constant::CONTROL_REQUEST_TIMEOUT)) {
778 OzoneMsg::ChannelsReceived(_) => (),
779 m => return Err(err!(
780 "{}: Received {:?}, expecting ChannelsReceived confirmation.", self.ozid(), m;
781 Channel, Unexpected)),
782 }
783
784 // 3. Every zone has to be ready before the database is. The first failure is reported
785 // as it arrives rather than after the rest, since the database cannot start either way.
786 let nz = self.cfg().num_zones();
787 let chan = match zones.channel() {
788 Some(chan) => chan,
789 None => return Err(err!(
790 "{}: The responder for zone readiness has no channel.", self.ozid();
791 Bug, Missing)),
792 };
793 let begun = Instant::now();
794 let mut ready = 0;
795 while ready < nz {
796 let left = constant::CONTROL_REQUEST_TIMEOUT.saturating_sub(begun.elapsed());
797 if left.is_zero() {
798 return Err(err!(
799 "{}: Only {} of {} zones finished surveying their files and loading their \
800 caches within {:?}.", self.ozid(), ready, nz, constant::CONTROL_REQUEST_TIMEOUT;
801 Init, Timeout));
802 }
803 match chan.recv_timeout(left) {
804 Recv::Empty => (), // Out of time, which the next pass reports.
805 Recv::Result(Err(e)) => return Err(err!(e,
806 "{}: While waiting for the zones to report themselves ready.", self.ozid();
807 Channel, Read)),
808 Recv::Result(Ok(OzoneMsg::Ok)) => ready += 1,
809 Recv::Result(Ok(OzoneMsg::Error(e))) => return Err(err!(e,
810 "{}: A zone could not be initialised.", self.ozid();
811 Init)),
812 Recv::Result(Ok(m)) => return Err(err!(
813 "{}: Received {:?}, expecting a zone to report itself ready.", self.ozid(), m;
814 Channel, Unexpected)),
815 }
816 }
817
818 res!(self.chan_out.send(OzoneMsg::Ready));
819 info!(sync_log::stream(), "{}: Ozone database start up complete.", self.label());
820 Ok(())
821 }
822
823// fn start_new_worker(
824// &'s self,
825// scope: &'s thread::Scope<'s, '_>,
826// wtyp: &WorkerType,
827// wind: &WorkerInd,
828// )
829// -> Outcome<Handle<'s, UIDL, UID, ENC, KH, ()>>
830// {
831// let chan = simplex();
832// let (semaphore, sentinel) = thread_channel();
833// let wg = self.handles().wait_end_ref().clone();
834// let api = OzoneApi::new(
835// OzoneBotId::new_worker(wtyp, wind),
836// self.db_root().to_path_buf(),
837// self.cfg().clone(),
838// self.chans().clone(),
839// self.schemes().clone(),
840// );
841// let args = ZoneWorkerInitArgs {
842// // Identity
843// wind: wind.clone(),
844// wtyp: wtyp.clone(),
845// // Bot
846// sem: semaphore,
847// // Comms
848// chan_in: chan.clone(),
849// // API
850// api,
851// };
852//
853// let mut bot: Box<dyn WorkerBot<UIDL, UID, ENC, KH, PR, CS>> = match wtyp {
854// WorkerType::Cache => Box::new(CacheBot::new(args)),
855// WorkerType::File => Box::new(FileBot::new(args)),
856// WorkerType::InitGarbage => Box::new(InitGarbageBot::new(args)),
857// WorkerType::Reader => Box::new(ReaderBot::new(args)),
858// WorkerType::Writer => Box::new(WriterBot::new(args)),
859// };
860// res!(bot.init());
861// let builder = thread::Builder::new()
862// .name(bot.id().to_string())
863// .stack_size(constant::STACK_SIZE);
864//
865// Ok(Handle::new(
866// Some(bot.ozid().clone()),
867// res!(builder.spawn_scoped(scope, move || {
868// bot.go();
869// drop(wg);
870// })),
871// sentinel,
872// Some(chan),
873// ))
874// }
875//
876// /// A zbot is neither a solo nor worker bot, since there is one per zone.
877// fn start_new_zbot(
878// &'s self,
879// scope: &'s thread::Scope<'s, '_>,
880// zind: &ZoneInd,
881// )
882// -> Outcome<Handle<'s, UIDL, UID, ENC, KH, ()>>
883// {
884// let chan = simplex();
885// let (semaphore, sentinel) = thread_channel();
886// let wg = self.handles().wait_end_ref().clone();
887// let api = OzoneApi::new(
888// OzoneBotId::ZoneBot(Bid::randef(), *zind),
889// self.db_root().to_path_buf(),
890// self.cfg().clone(),
891// self.chans().clone(),
892// self.schemes().clone(),
893// );
894// let args = BotInitArgs {
895// // Bot
896// sem: semaphore,
897// // Comms
898// chan_in: chan.clone(),
899// // API
900// api,
901// };
902// let mut bot = ZoneBot::new(args, *zind);
903// res!(bot.init());
904// let builder = thread::Builder::new()
905// .name(bot.id().to_string())
906// .stack_size(constant::STACK_SIZE);
907//
908// Ok(Handle::new(
909// Some(bot.ozid().clone()),
910// res!(builder.spawn_scoped(scope, move || {
911// bot.go();
912// drop(wg);
913// })),
914// sentinel,
915// Some(chan),
916// ))
917// }
918//
919// fn start_new_cfgbot(
920// &'s self,
921// scope: &'s thread::Scope<'s, '_>,
922// )
923// -> Outcome<Handle<'s, UIDL, UID, ENC, KH, ()>>
924// {
925// let chan = simplex();
926// let (semaphore, sentinel) = thread_channel();
927// let wg = self.handles().wait_end_ref().clone();
928// let api = OzoneApi::new(
929// OzoneBotId::ConfigBot(Bid::randef()),
930// self.db_root().to_path_buf(),
931// self.cfg().clone(),
932// self.chans().clone(),
933// self.schemes().clone(),
934// );
935// let args = BotInitArgs {
936// // Bot
937// sem: semaphore,
938// // Comms
939// chan_in: chan.clone(),
940// // API
941// api,
942// };
943// let mut bot = ConfigBot::new(args);
944// res!(bot.init());
945// let builder = thread::Builder::new()
946// .name(bot.ozid().to_string())
947// .stack_size(constant::STACK_SIZE);
948//
949// Ok(Handle::new(
950// Some(bot.ozid().clone()),
951// res!(builder.spawn_scoped(scope, move || {
952// bot.go();
953// drop(wg);
954// })),
955// sentinel,
956// Some(chan),
957// ))
958// }
959//
960// fn start_new_srvbot(
961// &'s self,
962// scope: &'s thread::Scope<'s, '_>,
963// bpind: BotPoolInd,
964// )
965// -> Outcome<Handle<'s, UIDL, UID, ENC, KH, ()>>
966// {
967// let chan = simplex();
968// let (semaphore, sentinel) = thread_channel();
969// let wg = self.handles().wait_end_ref().clone();
970// let api = OzoneApi::new(
971// OzoneBotId::ServerBot(Bid::randef(), bpind),
972// self.db_root().to_path_buf(),
973// self.cfg().clone(),
974// self.chans().clone(),
975// self.schemes().clone(),
976// );
977// let args = BotInitArgs {
978// // Bot
979// sem: semaphore,
980// // Comms
981// chan_in: chan.clone(),
982// // API
983// api,
984// };
985// let mut bot = ServerBot::new(args);
986// res!(bot.init());
987// let builder = thread::Builder::new()
988// .name(bot.ozid().to_string())
989// .stack_size(constant::STACK_SIZE);
990//
991// Ok(Handle::new(
992// Some(bot.ozid().clone()),
993// res!(builder.spawn_scoped(scope, move || {
994// bot.go();
995// drop(wg);
996// })),
997// sentinel,
998// Some(chan),
999// ))
1000// }
1001
1002 fn report_health(&self) -> Outcome<()> {
1003 let (expected, unresponsive) =
1004 res!(self.handles.get_unresponsive_bots(constant::PING_TIMEOUT));
1005 if unresponsive.len() > 0 {
1006 warn!(sync_log::stream(), "{} out of {} bots are unresponsive after {:?}:",
1007 unresponsive.len(), expected, constant::PING_TIMEOUT);
1008 for ozid in &unresponsive {
1009 warn!(sync_log::stream(), " {:?} is unresponsive", ozid);
1010 }
1011 let dead = res!(self.handles.get_dead_bots());
1012 if dead.len() > 0 {
1013 error!(sync_log::stream(), err!(
1014 "{} out of {} bots are dead:", dead.len(), expected;
1015 Thread, Missing));
1016 for ozid in &unresponsive {
1017 fault!(" {:?} is dead", ozid);
1018 }
1019 }
1020 } else {
1021 info!(sync_log::stream(), "{}: Ozone database health is good.", self.ozid());
1022 }
1023 Ok(())
1024 }
1025
1026 /// Gracefully shuts down the database. The bots are finished in order, each kind once those
1027 /// sending it work have ended (`BotChannels::finish_all`), and `resp` is answered when they
1028 /// have been, or when `constant::SHUTDOWN_MAX_WAIT` runs out, whichever comes first. What is
1029 /// left is then finished in the same order, however long the bots take to end, up to
1030 /// `constant::CONTROL_REQUEST_TIMEOUT` for each kind: finished early, the ones still waiting
1031 /// on them could only time out (2026-09-24). The master's wait for every thread to end
1032 /// covers this, since this thread is one of them.
1033 pub fn shutdown(
1034 &self,
1035 requester: String,
1036 resp: Option<&Responder<UIDL, UID, ENC, KH>>,
1037 )
1038 -> Outcome<()>
1039 {
1040 warn!(sync_log::stream(), "{}: Shutdown requested by {}, commencing...", self.label(), requester);
1041 let ended = |typs: &[WorkerType]| self.handles().ended(typs);
1042 let begun = Instant::now();
1043 let left = self.chans().finish_all(&ended, begun + constant::SHUTDOWN_MAX_WAIT);
1044 thread::sleep(Duration::from_secs(1));
1045 self.handles().report_status();
1046 if let Some(resp) = resp {
1047 self.respond(left.clone().map(|_| OzoneMsg::Ok), resp);
1048 }
1049 let mut left = res!(left);
1050 let unfinished = left.is_some();
1051 while let Some(stage) = left {
1052 left = res!(self.chans().finish_from(
1053 stage, &ended, Instant::now() + constant::CONTROL_REQUEST_TIMEOUT));
1054 if let Some(stage) = left {
1055 self.error(err!(
1056 "{}: Shutdown: Bots of stage {} had not ended after {:?}, so the rest are \
1057 finished without waiting for them.", self.ozid(), stage,
1058 constant::CONTROL_REQUEST_TIMEOUT;
1059 Thread, Timeout));
1060 left = res!(self.chans().finish_from(stage, |_: &[WorkerType]| true, Instant::now()));
1061 }
1062 }
1063 if unfinished {
1064 self.handles().report_status();
1065 }
1066 Ok(())
1067 }
1068}
1069