Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/src/bots/base/handles.rs

19.1 KiB, 129 runs

created by r1870400018:735, 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 bots::worker::bot::WorkerType,
4 base::{
5 id::OzoneBotId,
6 index::{
7 BotPoolInd,
8 WorkerInd,
9 ZoneInd,
10 },
11 },
12 comm::{
13 msg::OzoneMsg,
14 response::{
15 Responder,
16 Wait,
17 },
18 },
19};
20
21use oxedyne_fe2o3_core::{
22 channels::Simplex,
23 thread::Sentinel,
24};
25use oxedyne_fe2o3_jdat::id::NumIdDat;
26
27use std::{
28 time::Duration,
29};
30
31use crossbeam_utils::sync::WaitGroup;
32
33
34#[derive(Debug)]
35pub struct Handle<
36 const UIDL: usize,
37 UID: NumIdDat<UIDL>,
38 ENC: Encrypter,
39 KH: Hasher,
40> {
41 ozid: Option<OzoneBotId>,
42 sentinel: Sentinel,
43 chan: Option<Simplex<OzoneMsg<UIDL, UID, ENC, KH>>>,
44}
45
46impl<
47 const UIDL: usize,
48 UID: NumIdDat<UIDL> + 'static,
49 ENC: Encrypter + 'static,
50 KH: Hasher + 'static,
51>
52 Default for Handle<UIDL, UID, ENC, KH>
53{
54 fn default() -> Self {
55 Self {
56 ozid: None,
57 sentinel: Sentinel::default(),
58 chan: None,
59 }
60 }
61}
62
63impl<
64 const UIDL: usize,
65 UID: NumIdDat<UIDL> + 'static,
66 ENC: Encrypter + 'static,
67 KH: Hasher + 'static,
68>
69 Handle<UIDL, UID, ENC, KH>
70{
71 pub fn new(
72 ozid: Option<OzoneBotId>,
73 sentinel: Sentinel,
74 chan: Option<Simplex<OzoneMsg<UIDL, UID, ENC, KH>>>,
75 )
76 -> Self
77 {
78 Self {
79 ozid,
80 sentinel,
81 chan,
82 }
83 }
84
85 pub fn ozid(&self) -> &Option<OzoneBotId> { &self.ozid }
86 pub fn sentinel(&self) -> &Sentinel { &self.sentinel }
87 pub fn chan(&self) -> &Option<Simplex<OzoneMsg<UIDL, UID, ENC, KH>>> {
88 &self.chan
89 }
90 pub fn some_ozid(&self) -> Outcome<OzoneBotId> {
91 match &self.ozid {
92 Some(ozid) => Ok(ozid.clone()),
93 None => Err(err!(
94 "Handle contains no id as expected.";
95 Identifier, Missing)),
96 }
97 }
98 pub fn some_chan(&self) -> Outcome<Simplex<OzoneMsg<UIDL, UID, ENC, KH>>> {
99 match &self.chan {
100 Some(chan) => Ok(chan.clone()),
101 None => Err(err!(
102 "Handle contains no channel as expected.";
103 Channel, Missing)),
104 }
105 }
106}
107
108/// Contains handles for all the bots (except the `Supervisor`), for use by the `Supervisor`.
109#[derive(Debug)]
110pub struct BotHandles<
111 const UIDL: usize,
112 UID: NumIdDat<UIDL>,
113 ENC: Encrypter,
114 KH: Hasher,
115> {
116 nz: usize,
117 cbots: Vec<Vec<Handle<UIDL, UID, ENC, KH>>>,
118 fbots: Vec<Vec<Handle<UIDL, UID, ENC, KH>>>,
119 igbots: Vec<Vec<Handle<UIDL, UID, ENC, KH>>>,
120 rbots: Vec<Vec<Handle<UIDL, UID, ENC, KH>>>,
121 scbots: Vec<Vec<Handle<UIDL, UID, ENC, KH>>>,
122 wbots: Vec<Vec<Handle<UIDL, UID, ENC, KH>>>,
123 zbots: Vec<Handle<UIDL, UID, ENC, KH>>,
124 cfg: Handle<UIDL, UID, ENC, KH>,
125 sbots: Vec<Handle<UIDL, UID, ENC, KH>>,
126 wait_init: WaitGroup,
127 wait_end: WaitGroup,
128}
129
130impl<
131 const UIDL: usize,
132 UID: NumIdDat<UIDL> + 'static,
133 ENC: Encrypter + 'static,
134 KH: Hasher + 'static,
135>
136 Default for BotHandles<UIDL, UID, ENC, KH>
137{
138 fn default() -> Self {
139 Self {
140 nz: 0,
141 cbots: Vec::new(),
142 fbots: Vec::new(),
143 igbots: Vec::new(),
144 rbots: Vec::new(),
145 scbots: Vec::new(),
146 wbots: Vec::new(),
147 zbots: Vec::new(),
148 cfg: Handle::default(),
149 sbots: Vec::new(),
150 wait_init: WaitGroup::new(),
151 wait_end: WaitGroup::new(),
152 }
153 }
154}
155
156impl<
157 const UIDL: usize,
158 UID: NumIdDat<UIDL> + 'static,
159 ENC: Encrypter + 'static,
160 KH: Hasher + 'static,
161>
162 BotHandles<UIDL, UID, ENC, KH>
163{
164 /// Create a full set of empty handles according to the given configuration.
165 pub fn new(cfg: &OzoneConfig) -> Self {
166 let nz = cfg.num_zones();
167 let ns = cfg.num_sbots();
168 let nc = cfg.num_bots_per_zone(&WorkerType::Cache);
169 let nf = cfg.num_bots_per_zone(&WorkerType::File);
170 let nig = cfg.num_bots_per_zone(&WorkerType::InitGarbage);
171 let nr = cfg.num_bots_per_zone(&WorkerType::Reader);
172 let nsc = cfg.num_bots_per_zone(&WorkerType::Scan);
173 let nw = cfg.num_bots_per_zone(&WorkerType::Writer);
174 let mut cbots = Vec::new();
175 for _ in 0..nz {
176 let mut bots = Vec::new();
177 for _ in 0..nc {
178 bots.push(Handle::<UIDL, UID, ENC, KH>::default());
179 }
180 cbots.push(bots);
181 }
182 let mut fbots = Vec::new();
183 for _ in 0..nz {
184 let mut bots = Vec::new();
185 for _ in 0..nf {
186 bots.push(Handle::<UIDL, UID, ENC, KH>::default());
187 }
188 fbots.push(bots);
189 }
190 let mut igbots = Vec::new();
191 for _ in 0..nz {
192 let mut bots = Vec::new();
193 for _ in 0..nig {
194 bots.push(Handle::<UIDL, UID, ENC, KH>::default());
195 }
196 igbots.push(bots);
197 }
198 let mut rbots = Vec::new();
199 for _ in 0..nz {
200 let mut bots = Vec::new();
201 for _ in 0..nr {
202 bots.push(Handle::<UIDL, UID, ENC, KH>::default());
203 }
204 rbots.push(bots);
205 }
206 let mut scbots = Vec::new();
207 for _ in 0..nz {
208 let mut bots = Vec::new();
209 for _ in 0..nsc {
210 bots.push(Handle::<UIDL, UID, ENC, KH>::default());
211 }
212 scbots.push(bots);
213 }
214 let mut wbots = Vec::new();
215 for _ in 0..nz {
216 let mut bots = Vec::new();
217 for _ in 0..nw {
218 bots.push(Handle::<UIDL, UID, ENC, KH>::default());
219 }
220 wbots.push(bots);
221 }
222 let mut zbots = Vec::new();
223 for _ in 0..nz {
224 zbots.push(Handle::<UIDL, UID, ENC, KH>::default());
225 }
226 let mut sbots = Vec::new();
227 for _ in 0..ns {
228 sbots.push(Handle::<UIDL, UID, ENC, KH>::default());
229 }
230 Self {
231 nz,
232 cbots,
233 fbots,
234 igbots,
235 rbots,
236 scbots,
237 wbots,
238 zbots,
239 sbots,
240 ..Default::default()
241 }
242 }
243
244 // Use
245 pub fn all_cbots(&self) -> &Vec<Vec<Handle<UIDL, UID, ENC, KH>>> { &self.cbots }
246 pub fn all_fbots(&self) -> &Vec<Vec<Handle<UIDL, UID, ENC, KH>>> { &self.fbots }
247 pub fn all_igbots(&self) -> &Vec<Vec<Handle<UIDL, UID, ENC, KH>>> { &self.igbots }
248 pub fn all_rbots(&self) -> &Vec<Vec<Handle<UIDL, UID, ENC, KH>>> { &self.rbots }
249 /// Returns the per-zone scan-bot handles.
250 pub fn all_scbots(&self) -> &Vec<Vec<Handle<UIDL, UID, ENC, KH>>> { &self.scbots }
251 pub fn all_wbots(&self) -> &Vec<Vec<Handle<UIDL, UID, ENC, KH>>> { &self.wbots }
252 pub fn all_zbots(&self) -> &Vec<Handle<UIDL, UID, ENC, KH>> { &self.zbots }
253 pub fn cfg(&self) -> &Handle<UIDL, UID, ENC, KH> { &self.cfg }
254 pub fn all_sbots(&self) -> &Vec<Handle<UIDL, UID, ENC, KH>> { &self.sbots }
255 pub fn wait_init_ref(&self) -> &WaitGroup { &self.wait_init }
256 pub fn wait_end_ref(&self) -> &WaitGroup { &self.wait_end }
257
258 pub fn get_zbot(&self, zind: &ZoneInd) -> Outcome<&Handle<UIDL, UID, ENC, KH>> {
259 res!(self.check_zone_index(**zind));
260 Ok(&self.zbots[**zind])
261 }
262
263 // Mutate
264 pub fn cbots_mut(&mut self) -> &mut Vec<Vec<Handle<UIDL, UID, ENC, KH>>> { &mut self.cbots }
265 pub fn fbots_mut(&mut self) -> &mut Vec<Vec<Handle<UIDL, UID, ENC, KH>>> { &mut self.fbots }
266 pub fn igbots_mut(&mut self) -> &mut Vec<Vec<Handle<UIDL, UID, ENC, KH>>> { &mut self.igbots }
267 pub fn rbots_mut(&mut self) -> &mut Vec<Vec<Handle<UIDL, UID, ENC, KH>>> { &mut self.rbots }
268 /// Returns a mutable reference to the per-zone scan-bot handles.
269 pub fn scbots_mut(&mut self) -> &mut Vec<Vec<Handle<UIDL, UID, ENC, KH>>> { &mut self.scbots }
270 pub fn wbots_mut(&mut self) -> &mut Vec<Vec<Handle<UIDL, UID, ENC, KH>>> { &mut self.wbots }
271 pub fn zbots_mut(&mut self) -> &mut Vec<Handle<UIDL, UID, ENC, KH>> { &mut self.zbots }
272 pub fn cfg_mut(&mut self) -> &mut Handle<UIDL, UID, ENC, KH> { &mut self.cfg }
273 pub fn sbots_mut(&mut self) -> &mut Vec<Handle<UIDL, UID, ENC, KH>> { &mut self.sbots }
274
275 pub fn set_sbot(&mut self, bpind: &BotPoolInd, hand: Handle<UIDL, UID, ENC, KH>) -> Outcome<()> {
276 self.sbots[**bpind] = hand;
277 Ok(())
278 }
279 pub fn set_zbot(&mut self, zind: &ZoneInd, hand: Handle<UIDL, UID, ENC, KH>) -> Outcome<()> {
280 res!(self.check_zone_index(**zind));
281 self.zbots[**zind] = hand;
282 Ok(())
283 }
284 pub fn set_cfg(&mut self, hand: Handle<UIDL, UID, ENC, KH>) { self.cfg = hand; }
285
286 /// Have all the workers of the given types, in every zone, ended?
287 pub fn ended(&self, typs: &[WorkerType]) -> bool {
288 typs.iter().all(|typ| {
289 let pools = match typ {
290 WorkerType::Cache => &self.cbots,
291 WorkerType::File => &self.fbots,
292 WorkerType::InitGarbage => &self.igbots,
293 WorkerType::Reader => &self.rbots,
294 WorkerType::Scan => &self.scbots,
295 WorkerType::Writer => &self.wbots,
296 };
297 pools.iter().all(|zone| zone.iter().all(|h| h.sentinel().is_finished()))
298 })
299 }
300
301 pub fn wait_init(self) {
302 self.wait_init.wait();
303 }
304 pub fn wait_end(self) {
305 self.wait_end.wait();
306 }
307
308 pub fn set_worker_bot(
309 &mut self,
310 wtyp: &WorkerType,
311 wind: &WorkerInd,
312 hand: Handle<UIDL, UID, ENC, KH>,
313 )
314 -> Outcome<()>
315 {
316 res!(self.check_zone_index(wind.z()));
317 match wtyp {
318 WorkerType::Cache => self.cbots[wind.z()][wind.b()] = hand,
319 WorkerType::File => self.fbots[wind.z()][wind.b()] = hand,
320 WorkerType::InitGarbage => self.igbots[wind.z()][wind.b()] = hand,
321 WorkerType::Reader => self.rbots[wind.z()][wind.b()] = hand,
322 WorkerType::Scan => self.scbots[wind.z()][wind.b()] = hand,
323 WorkerType::Writer => self.wbots[wind.z()][wind.b()] = hand,
324 }
325 Ok(())
326 }
327
328 fn check_zone_index(&self, ind: usize) -> Outcome<()> {
329 if ind > self.nz {
330 return Err(err!(
331 "Zone index {} into BotHandles exceeds number of zones {}.",
332 ind, self.nz;
333 Bug, Excessive));
334 }
335 Ok(())
336 }
337
338 pub fn report_status(&self) {
339 for z in 0..self.nz {
340 for pool in [
341 &self.cbots[z],
342 &self.fbots[z],
343 &self.igbots[z],
344 &self.rbots[z],
345 &self.scbots[z],
346 &self.wbots[z],
347 ] {
348 for h in pool {
349 if !h.sentinel().is_finished() {
350 msg!("{:?} bot is not finished", h.ozid());
351 }
352 }
353 }
354 }
355 for h in &self.zbots {
356 if !h.sentinel().is_finished() {
357 msg!("{:?} bot is not finished", h.ozid());
358 }
359 }
360 for h in &self.sbots {
361 if !h.sentinel().is_finished() {
362 msg!("{:?} bot is not finished", h.ozid());
363 }
364 }
365 for h in [&self.cfg] {
366 if !h.sentinel().is_finished() {
367 msg!("{:?} bot is not finished", h.ozid());
368 }
369 }
370 }
371
372 /// Returns the ids of bots threads that have finished.
373 pub fn get_dead_bots(&self) -> Outcome<Vec<OzoneBotId>> {
374
375 let mut finished = Vec::new();
376
377 for z in 0..self.nz {
378 for pool in [
379 &self.cbots[z],
380 &self.fbots[z],
381 &self.igbots[z],
382 &self.rbots[z],
383 &self.scbots[z],
384 &self.wbots[z],
385 ] {
386 for h in pool {
387 if h.sentinel().is_finished() {
388 finished.push(res!(h.some_ozid().clone()));
389 }
390 }
391 }
392 }
393 for h in &self.zbots {
394 if h.sentinel().is_finished() {
395 finished.push(res!(h.some_ozid().clone()));
396 }
397 }
398 for h in &self.sbots {
399 if h.sentinel().is_finished() {
400 finished.push(res!(h.some_ozid().clone()));
401 }
402 }
403 for h in [&self.cfg] {
404 if h.sentinel().is_finished() {
405 finished.push(res!(h.some_ozid().clone()));
406 }
407 }
408
409 Ok(finished)
410 }
411
412 /// Returns the ids of bots that fail to respond to a ping within the specified timeout.
413 ///
414 /// # Arguments
415 /// * `timeout` - The duration to wait for a response from each bot.
416 ///
417 /// Returns a list of `OzoneBotIds` for bots that did not respond in time.
418 pub fn get_unresponsive_bots(
419 &self,
420 timeout: Duration,
421 )
422 -> Outcome<(usize, Vec<OzoneBotId>)>
423 {
424 let wait = Wait {
425 max_wait: timeout.clone(),
426 check_interval: constant::CHECK_INTERVAL,
427 };
428 let resp = Responder::new(None);
429 let mut all_bot_ids = Vec::new();
430
431 // Send pings to all bots.
432 for handle in self.iter() {
433 if let Some(chan) = handle.chan() {
434 let ozid = res!(handle.some_ozid());
435 match chan.send(OzoneMsg::Ping(ozid.clone(), resp.clone())) {
436 Ok(_) => {
437 all_bot_ids.push(ozid);
438 }
439 Err(e) => error!(sync_log::stream(), err!(e,
440 "While sending ping to bot {:?}", ozid;
441 Channel, Write))
442 }
443 }
444 }
445 let expected = all_bot_ids.len();
446
447 // Track responses and build set of responsive bots.
448 let (_, responsive) = res!(resp.recv_pongs(wait));
449 if responsive.len() > expected {
450 error!(sync_log::stream(), err!(
451 "Expecting {} messages via responder, received {} after \
452 {:?}.", responsive.len(), expected, timeout;
453 Input, Mismatch, Size));
454 }
455
456 // Collect unresponsive bot IDs by comparing against those that responded.
457 let mut unresponsive = Vec::new();
458 for ozid in &all_bot_ids {
459 if !responsive.contains(ozid) {
460 unresponsive.push(ozid.clone());
461 }
462 }
463
464 Ok((expected, unresponsive))
465 }
466}
467
468/// Immutable iterator over all bot handles in a BotHandles collection.
469#[derive(Debug)]
470pub struct BotHandlesIter<
471 'iter,
472 const UIDL: usize,
473 UID: NumIdDat<UIDL>,
474 ENC: Encrypter,
475 KH: Hasher,
476> {
477 handles: &'iter BotHandles<UIDL, UID, ENC, KH>,
478 zone_index: usize,
479 pool_type: usize, // Index into the pool types (cbots, fbots, etc.).
480 bot_index: usize, // Index within the current pool.
481 stage: IterStage, // Tracks which group of bots we're iterating over.
482}
483
484/// Tracks the current stage of iteration through different bot types.
485#[derive(Debug, PartialEq)]
486enum IterStage {
487 Workers, // Iterating through worker bot pools (cbots, fbots, etc.).
488 ZoneBots, // Iterating through zone bots.
489 StoreBots, // Iterating through store bots.
490 Config, // Iterating through the config bot.
491 Done, // Iteration complete.
492}
493
494impl<
495 'iter,
496 const UIDL: usize,
497 UID: NumIdDat<UIDL> + 'static,
498 ENC: Encrypter + 'static,
499 KH: Hasher + 'static,
500>
501 BotHandlesIter<'iter, UIDL, UID, ENC, KH>
502{
503 fn new(handles: &'iter BotHandles<UIDL, UID, ENC, KH>) -> Self {
504 Self {
505 handles,
506 zone_index: 0,
507 pool_type: 0,
508 bot_index: 0,
509 stage: IterStage::Workers,
510 }
511 }
512
513 /// Returns the next worker bot handle, if any remain in the current zone.
514 fn next_worker(&mut self) -> Option<&'iter Handle<UIDL, UID, ENC, KH>> {
515 let pools = [
516 self.handles.all_cbots(),
517 self.handles.all_fbots(),
518 self.handles.all_igbots(),
519 self.handles.all_rbots(),
520 self.handles.all_scbots(),
521 self.handles.all_wbots(),
522 ];
523
524 // Ensure we haven't exceeded available pools.
525 if self.pool_type >= pools.len() {
526 return None;
527 }
528
529 let current_pool = &pools[self.pool_type][self.zone_index];
530
531 // If we've exhausted the current pool.
532 if self.bot_index >= current_pool.len() {
533 self.bot_index = 0;
534 self.pool_type += 1;
535 return self.next_worker();
536 }
537
538 // If we've exhausted the current zone.
539 if self.zone_index >= self.handles.nz {
540 self.zone_index = 0;
541 self.pool_type += 1;
542 return self.next_worker();
543 }
544
545 let handle = &current_pool[self.bot_index];
546 self.bot_index += 1;
547 Some(handle)
548 }
549}
550
551impl<
552 'iter,
553 const UIDL: usize,
554 UID: NumIdDat<UIDL> + 'static,
555 ENC: Encrypter + 'static,
556 KH: Hasher + 'static,
557>
558 Iterator for BotHandlesIter<'iter, UIDL, UID, ENC, KH>
559{
560 type Item = &'iter Handle<UIDL, UID, ENC, KH>;
561
562 fn next(&mut self) -> Option<Self::Item> {
563 match self.stage {
564 IterStage::Workers => {
565 if let Some(handle) = self.next_worker() {
566 return Some(handle);
567 }
568 self.stage = IterStage::ZoneBots;
569 self.zone_index = 0;
570 self.next()
571 }
572 IterStage::ZoneBots => {
573 if self.zone_index < self.handles.nz {
574 let handle = &self.handles.all_zbots()[self.zone_index];
575 self.zone_index += 1;
576 Some(handle)
577 } else {
578 self.stage = IterStage::StoreBots;
579 self.bot_index = 0;
580 self.next()
581 }
582 }
583 IterStage::StoreBots => {
584 if self.bot_index < self.handles.all_sbots().len() {
585 let handle = &self.handles.all_sbots()[self.bot_index];
586 self.bot_index += 1;
587 Some(handle)
588 } else {
589 self.stage = IterStage::Config;
590 self.next()
591 }
592 }
593 IterStage::Config => {
594 self.stage = IterStage::Done;
595 Some(self.handles.cfg())
596 }
597 IterStage::Done => None,
598 }
599 }
600}
601
602impl<
603 const UIDL: usize,
604 UID: NumIdDat<UIDL> + 'static,
605 ENC: Encrypter + 'static,
606 KH: Hasher + 'static,
607>
608 BotHandles<UIDL, UID, ENC, KH>
609{
610 /// Returns an iterator over references to all bot handles.
611 pub fn iter<'iter>(&'iter self) -> BotHandlesIter<'iter, UIDL, UID, ENC, KH> {
612 BotHandlesIter::new(self)
613 }
614}