oxedyne/fe2o3/fe2o3_o3db_sync/src/comm/channels.rs
21.8 KiB, 78 runs
created by r1870400018:765, 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 | id::{ |
| 5 | OzoneBotId, |
| 6 | //OzoneBotType, |
| 7 | }, |
| 8 | index::{ |
| 9 | BotPoolInd, |
| 10 | WorkerInd, |
| 11 | ZoneInd, |
| 12 | }, |
| 13 | }, |
| 14 | bots::{ |
| 15 | worker::bot::WorkerType, |
| 16 | }, |
| 17 | comm::msg::OzoneMsg, |
| 18 | }; |
| 19 | |
| 20 | use oxedyne_fe2o3_core::{ |
| 21 | channels::{ |
| 22 | simplex, |
| 23 | Simplex, |
| 24 | }, |
| 25 | }; |
| 26 | use oxedyne_fe2o3_jdat::id::NumIdDat; |
| 27 | |
| 28 | use std::{ |
| 29 | ops::{ |
| 30 | Index, |
| 31 | IndexMut, |
| 32 | }, |
| 33 | time::{ |
| 34 | Duration, |
| 35 | Instant, |
| 36 | }, |
| 37 | }; |
| 38 | |
| 39 | use rand::Rng; |
| 40 | |
| 41 | const FINISH_CHECK_INTERVAL: Duration = Duration::from_millis(5); // finer than CHECK_INTERVAL: a close waits on it |
| 42 | |
| 43 | // The order in which a close finishes each zone's workers. A stage is sent its `Finish` only |
| 44 | // once every worker of the stage before it has ended, since those are the ones that send it work |
| 45 | // and wait on its answers. Finished together, a read or collection still queued behind the |
| 46 | // `Finish` of its reader or collector met a cache bot that had already ended, and waited out |
| 47 | // `BOT_REQUEST_TIMEOUT` for an answer nobody would send; and a record a syncer released late, a |
| 48 | // failed barrier's error among them, reached a cache bot that had ended, so its caller waited out |
| 49 | // the durability deadline (2026-09-24). |
| 50 | const FINISH_ORDER: [&[WorkerType]; 4] = [ |
| 51 | &[WorkerType::Reader, WorkerType::Scan, WorkerType::InitGarbage], // ask the others |
| 52 | &[WorkerType::Writer], // release to the caches |
| 53 | &[WorkerType::Cache], // tell the file bots |
| 54 | &[WorkerType::File], |
| 55 | ]; |
| 56 | |
| 57 | |
| 58 | #[derive(Clone, Debug, Eq, PartialEq)] |
| 59 | pub enum PoolType { |
| 60 | Cache, |
| 61 | File, |
| 62 | InitGarbage, |
| 63 | Reader, |
| 64 | Scan, |
| 65 | Writer, |
| 66 | Zone, |
| 67 | Server, |
| 68 | } |
| 69 | |
| 70 | impl From<&WorkerType> for PoolType { |
| 71 | fn from(wtyp: &WorkerType) -> Self { |
| 72 | match wtyp { |
| 73 | WorkerType::Cache => PoolType::Cache, |
| 74 | WorkerType::File => PoolType::File, |
| 75 | WorkerType::InitGarbage => PoolType::InitGarbage, |
| 76 | WorkerType::Reader => PoolType::Reader, |
| 77 | WorkerType::Scan => PoolType::Scan, |
| 78 | WorkerType::Writer => PoolType::Writer, |
| 79 | } |
| 80 | } |
| 81 | } |
| 82 | |
| 83 | #[derive(Clone, Debug)] |
| 84 | pub enum ChooseBot { |
| 85 | Randomly, |
| 86 | ByFile(u32), |
| 87 | } |
| 88 | |
| 89 | #[derive(Clone, Debug)] |
| 90 | pub struct ChannelPool< |
| 91 | const UIDL: usize, |
| 92 | UID: NumIdDat<UIDL>, |
| 93 | ENC: Encrypter, |
| 94 | KH: Hasher, |
| 95 | > { |
| 96 | typ: PoolType, |
| 97 | pool: Vec<Simplex<OzoneMsg<UIDL, UID, ENC, KH>>>, |
| 98 | } |
| 99 | |
| 100 | impl< |
| 101 | const UIDL: usize, |
| 102 | UID: NumIdDat<UIDL> + 'static, |
| 103 | ENC: Encrypter + 'static, |
| 104 | KH: Hasher + 'static, |
| 105 | > |
| 106 | ChannelPool<UIDL, UID, ENC, KH> |
| 107 | { |
| 108 | pub fn new(typ: &PoolType, n: usize) -> Self { |
| 109 | let mut pool = Vec::new(); |
| 110 | for _ in 0..n { |
| 111 | pool.push(simplex()); |
| 112 | } |
| 113 | Self { |
| 114 | typ: typ.clone(), |
| 115 | pool: pool, |
| 116 | } |
| 117 | } |
| 118 | |
| 119 | pub fn make(typ: &PoolType, pool: Vec<Simplex<OzoneMsg<UIDL, UID, ENC, KH>>>) -> Self { |
| 120 | Self { |
| 121 | typ: typ.clone(), |
| 122 | pool: pool, |
| 123 | } |
| 124 | } |
| 125 | |
| 126 | fn check_index(&self, ind: usize) -> Outcome<()> { |
| 127 | if ind > self.pool.len() { |
| 128 | return Err(err!("Index {} exceeds pool size {}.", ind, self.pool.len(); Index, TooBig)); |
| 129 | } |
| 130 | Ok(()) |
| 131 | } |
| 132 | |
| 133 | pub fn len(&self) -> usize { self.pool.len() } |
| 134 | |
| 135 | pub fn get_bot(&self, ind: usize) -> Outcome<&Simplex<OzoneMsg<UIDL, UID, ENC, KH>>> { |
| 136 | res!(self.check_index(ind)); |
| 137 | Ok(&self.pool[ind]) |
| 138 | } |
| 139 | |
| 140 | pub fn choose_bot( |
| 141 | &self, |
| 142 | how: &ChooseBot, |
| 143 | ) |
| 144 | -> (&Simplex<OzoneMsg<UIDL, UID, ENC, KH>>, BotPoolInd) |
| 145 | { |
| 146 | let n = self.pool.len(); |
| 147 | let i = match how { |
| 148 | ChooseBot::Randomly => { |
| 149 | let mut rng = rand::thread_rng(); |
| 150 | rng.gen_range(0..n) |
| 151 | }, |
| 152 | ChooseBot::ByFile(fnum) => (*fnum as usize) % n, |
| 153 | }; |
| 154 | (&self.pool[i], BotPoolInd::new(i)) |
| 155 | } |
| 156 | |
| 157 | pub fn set_bot(&mut self, ind: usize, chan: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>) -> Outcome<()> { |
| 158 | res!(self.check_index(ind)); |
| 159 | self.pool[ind] = chan; |
| 160 | Ok(()) |
| 161 | } |
| 162 | |
| 163 | pub fn finish_all(&self) -> Outcome<()> { |
| 164 | for chan in &self.pool { |
| 165 | res!(chan.send(OzoneMsg::Finish)); |
| 166 | } |
| 167 | Ok(()) |
| 168 | } |
| 169 | |
| 170 | pub fn msg_count(&self) -> Vec<usize> { |
| 171 | let mut queues = Vec::new(); |
| 172 | for chan in &self.pool { |
| 173 | queues.push(chan.len()); |
| 174 | } |
| 175 | queues |
| 176 | } |
| 177 | |
| 178 | pub fn msg_count_total(&self) -> usize { |
| 179 | let mut total: usize = 0; |
| 180 | for chan in &self.pool { |
| 181 | total += chan.len(); |
| 182 | } |
| 183 | total |
| 184 | } |
| 185 | |
| 186 | pub fn msg_count_non_zero(&self) -> bool { |
| 187 | let mut pending = false; |
| 188 | for chan in &self.pool { |
| 189 | pending = pending | (chan.len() > 0); |
| 190 | } |
| 191 | pending |
| 192 | } |
| 193 | |
| 194 | /// Returns the number of messages sent. |
| 195 | pub fn send_to_all(&self, msg: OzoneMsg<UIDL, UID, ENC, KH>) -> Outcome<usize> { |
| 196 | for chan in &self.pool { |
| 197 | res!(chan.send(msg.clone())); |
| 198 | } |
| 199 | Ok(self.len()) |
| 200 | } |
| 201 | } |
| 202 | |
| 203 | /// Channel message queue lengths for all worker bots in a zone. |
| 204 | #[derive(Clone, Debug)] |
| 205 | pub struct ZoneMsgCount { |
| 206 | cbots: Vec<usize>, |
| 207 | fbots: Vec<usize>, |
| 208 | igbots: Vec<usize>, |
| 209 | rbots: Vec<usize>, |
| 210 | scbots: Vec<usize>, |
| 211 | wbots: Vec<usize>, |
| 212 | } |
| 213 | |
| 214 | impl ZoneMsgCount { |
| 215 | pub fn total(&self) -> usize { |
| 216 | let mut total = 0; |
| 217 | total += self.cbots.iter().sum::<usize>(); |
| 218 | total += self.fbots.iter().sum::<usize>(); |
| 219 | total += self.igbots.iter().sum::<usize>(); |
| 220 | total += self.rbots.iter().sum::<usize>(); |
| 221 | total += self.scbots.iter().sum::<usize>(); |
| 222 | total += self.wbots.iter().sum::<usize>(); |
| 223 | total |
| 224 | } |
| 225 | } |
| 226 | |
| 227 | /// Channels for all worker bots in a zone. |
| 228 | #[derive(Clone, Debug)] |
| 229 | pub struct ZoneWorkerChannels< |
| 230 | const UIDL: usize, |
| 231 | UID: NumIdDat<UIDL>, |
| 232 | ENC: Encrypter, |
| 233 | KH: Hasher, |
| 234 | > { |
| 235 | cbots: ChannelPool<UIDL, UID, ENC, KH>, |
| 236 | fbots: ChannelPool<UIDL, UID, ENC, KH>, |
| 237 | igbots: ChannelPool<UIDL, UID, ENC, KH>, |
| 238 | rbots: ChannelPool<UIDL, UID, ENC, KH>, |
| 239 | scbots: ChannelPool<UIDL, UID, ENC, KH>, |
| 240 | wbots: ChannelPool<UIDL, UID, ENC, KH>, |
| 241 | } |
| 242 | |
| 243 | impl< |
| 244 | const UIDL: usize, |
| 245 | UID: NumIdDat<UIDL>, |
| 246 | ENC: Encrypter, |
| 247 | KH: Hasher, |
| 248 | > |
| 249 | Index<&WorkerType> for ZoneWorkerChannels<UIDL, UID, ENC, KH> { |
| 250 | type Output = ChannelPool<UIDL, UID, ENC, KH>; |
| 251 | |
| 252 | fn index(&self, typ: &WorkerType) -> &Self::Output { |
| 253 | match typ { |
| 254 | WorkerType::Cache => &self.cbots, |
| 255 | WorkerType::File => &self.fbots, |
| 256 | WorkerType::InitGarbage => &self.igbots, |
| 257 | WorkerType::Reader => &self.rbots, |
| 258 | WorkerType::Scan => &self.scbots, |
| 259 | WorkerType::Writer => &self.wbots, |
| 260 | } |
| 261 | } |
| 262 | } |
| 263 | |
| 264 | impl< |
| 265 | const UIDL: usize, |
| 266 | UID: NumIdDat<UIDL>, |
| 267 | ENC: Encrypter, |
| 268 | KH: Hasher, |
| 269 | > |
| 270 | IndexMut<&WorkerType> for ZoneWorkerChannels<UIDL, UID, ENC, KH> |
| 271 | { |
| 272 | fn index_mut(&mut self, typ: &WorkerType) -> &mut Self::Output { |
| 273 | match typ { |
| 274 | WorkerType::Cache => &mut self.cbots, |
| 275 | WorkerType::File => &mut self.fbots, |
| 276 | WorkerType::InitGarbage => &mut self.igbots, |
| 277 | WorkerType::Reader => &mut self.rbots, |
| 278 | WorkerType::Scan => &mut self.scbots, |
| 279 | WorkerType::Writer => &mut self.wbots, |
| 280 | } |
| 281 | } |
| 282 | } |
| 283 | |
| 284 | impl< |
| 285 | const UIDL: usize, |
| 286 | UID: NumIdDat<UIDL> + 'static, |
| 287 | ENC: Encrypter + 'static, |
| 288 | KH: Hasher + 'static, |
| 289 | > |
| 290 | ZoneWorkerChannels<UIDL, UID, ENC, KH> |
| 291 | { |
| 292 | /// Create a full set of worker channels for a zone, according to the given configuration. |
| 293 | pub fn new(cfg: &OzoneConfig) -> Self { |
| 294 | let nc = cfg.num_bots_per_zone(&WorkerType::Cache); |
| 295 | let nf = cfg.num_bots_per_zone(&WorkerType::File); |
| 296 | let nig = cfg.num_bots_per_zone(&WorkerType::InitGarbage); |
| 297 | let nr = cfg.num_bots_per_zone(&WorkerType::Reader); |
| 298 | let nsc = cfg.num_bots_per_zone(&WorkerType::Scan); |
| 299 | let nw = cfg.num_bots_per_zone(&WorkerType::Writer); |
| 300 | Self { |
| 301 | cbots: ChannelPool::new(&PoolType::Cache, nc), |
| 302 | fbots: ChannelPool::new(&PoolType::File, nf), |
| 303 | igbots: ChannelPool::new(&PoolType::InitGarbage,nig), |
| 304 | rbots: ChannelPool::new(&PoolType::Reader, nr), |
| 305 | scbots: ChannelPool::new(&PoolType::Scan, nsc), |
| 306 | wbots: ChannelPool::new(&PoolType::Writer, nw), |
| 307 | } |
| 308 | } |
| 309 | |
| 310 | fn msg_count_non_zero(&self) -> bool { |
| 311 | self.cbots.msg_count_non_zero() | |
| 312 | self.fbots.msg_count_non_zero() | |
| 313 | self.igbots.msg_count_non_zero() | |
| 314 | self.rbots.msg_count_non_zero() | |
| 315 | self.scbots.msg_count_non_zero() | |
| 316 | self.wbots.msg_count_non_zero() |
| 317 | } |
| 318 | |
| 319 | /// Finishes the zone's workers of the given types. |
| 320 | pub fn finish(&self, typs: &[WorkerType]) -> Outcome<()> { |
| 321 | for typ in typs { |
| 322 | res!(self[typ].finish_all()); |
| 323 | } |
| 324 | Ok(()) |
| 325 | } |
| 326 | |
| 327 | pub fn msg_count(&self) -> ZoneMsgCount { |
| 328 | ZoneMsgCount { |
| 329 | cbots: self.cbots.msg_count(), |
| 330 | fbots: self.fbots.msg_count(), |
| 331 | igbots: self.igbots.msg_count(), |
| 332 | rbots: self.rbots.msg_count(), |
| 333 | scbots: self.scbots.msg_count(), |
| 334 | wbots: self.wbots.msg_count(), |
| 335 | } |
| 336 | } |
| 337 | |
| 338 | pub fn total_bot_count(&self) -> usize { |
| 339 | let mut count = 0; |
| 340 | count += self.cbots.len(); |
| 341 | count += self.fbots.len(); |
| 342 | count += self.igbots.len(); |
| 343 | count += self.rbots.len(); |
| 344 | count += self.scbots.len(); |
| 345 | count += self.wbots.len(); |
| 346 | count |
| 347 | } |
| 348 | |
| 349 | pub fn send_to_all(&self, msg: OzoneMsg<UIDL, UID, ENC, KH>) -> Outcome<usize> { |
| 350 | let mut count = 0; |
| 351 | count += res!(self.cbots.send_to_all(msg.clone())); |
| 352 | count += res!(self.fbots.send_to_all(msg.clone())); |
| 353 | count += res!(self.igbots.send_to_all(msg.clone())); |
| 354 | count += res!(self.rbots.send_to_all(msg.clone())); |
| 355 | count += res!(self.scbots.send_to_all(msg.clone())); |
| 356 | count += res!(self.wbots.send_to_all(msg.clone())); |
| 357 | Ok(count) |
| 358 | } |
| 359 | } |
| 360 | |
| 361 | /// Message queue lengths for all channels. |
| 362 | #[derive(Clone, Debug)] |
| 363 | pub struct OzoneMsgCount { |
| 364 | nz: usize, |
| 365 | zwbots: Vec<ZoneMsgCount>, |
| 366 | zbots: Vec<usize>, |
| 367 | cfg: usize, |
| 368 | sbots: Vec<usize>, |
| 369 | sup: usize, |
| 370 | } |
| 371 | |
| 372 | impl OzoneMsgCount { |
| 373 | |
| 374 | pub fn total(&self) -> usize { |
| 375 | let mut total = 0; |
| 376 | self.zwbots.iter().for_each(|x| total += x.total()); |
| 377 | self.zbots.iter().for_each(|x| total += x); |
| 378 | total += self.cfg; |
| 379 | self.sbots.iter().for_each(|x| total += x); |
| 380 | total += self.sup; |
| 381 | total |
| 382 | } |
| 383 | |
| 384 | pub fn total_zone(&self) -> usize { |
| 385 | let mut total = 0; |
| 386 | self.zwbots.iter().for_each(|x| total += x.total()); |
| 387 | self.zbots.iter().for_each(|x| total += x); |
| 388 | total |
| 389 | } |
| 390 | } |
| 391 | |
| 392 | /// Channels for all bots in all zones. Rather than sharing references to these channels, clone them. Unlike `bots::base::handles::BotHandles`, this includes the `Supervisor`. |
| 393 | #[derive(Clone, Debug)] |
| 394 | pub struct BotChannels< |
| 395 | const UIDL: usize, |
| 396 | UID: NumIdDat<UIDL>, |
| 397 | ENC: Encrypter, |
| 398 | KH: Hasher, |
| 399 | > { |
| 400 | nz: usize, |
| 401 | zwbots: Vec<ZoneWorkerChannels<UIDL, UID, ENC, KH>>, |
| 402 | zbots: ChannelPool<UIDL, UID, ENC, KH>, |
| 403 | cfg: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>, |
| 404 | sbots: ChannelPool<UIDL, UID, ENC, KH>, |
| 405 | sup: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>, |
| 406 | } |
| 407 | |
| 408 | impl< |
| 409 | const UIDL: usize, |
| 410 | UID: NumIdDat<UIDL> + 'static, |
| 411 | ENC: Encrypter + 'static, |
| 412 | KH: Hasher + 'static, |
| 413 | > |
| 414 | BotChannels<UIDL, UID, ENC, KH> |
| 415 | { |
| 416 | /// Create a full set of functioning channels according to the given configuration. |
| 417 | pub fn new(cfg: &OzoneConfig) -> Self { |
| 418 | let nz = cfg.num_zones(); |
| 419 | let mut zwbots = Vec::new(); |
| 420 | for _ in 0..nz { |
| 421 | zwbots.push(ZoneWorkerChannels::new(cfg)); |
| 422 | } |
| 423 | Self { |
| 424 | nz, |
| 425 | zwbots, |
| 426 | zbots: ChannelPool::new(&PoolType::Zone, nz), |
| 427 | cfg: simplex(), |
| 428 | sbots: ChannelPool::new(&PoolType::Server, cfg.num_sbots()), |
| 429 | sup: simplex(), |
| 430 | } |
| 431 | } |
| 432 | |
| 433 | /// Returns channels for all worker pools, for the given zone. |
| 434 | pub fn get_all_workers_in_zone( |
| 435 | &self, |
| 436 | zind: &ZoneInd, |
| 437 | ) |
| 438 | -> Outcome<ZoneWorkerChannels<UIDL, UID, ENC, KH>> |
| 439 | { |
| 440 | res!(self.check_zone_index(**zind)); |
| 441 | Ok(self.zwbots[**zind].clone()) |
| 442 | } |
| 443 | |
| 444 | /// Returns channels for the worker pool of the given type, for the given zone. |
| 445 | pub fn get_workers_of_type_in_zone( |
| 446 | &self, |
| 447 | wtyp: &WorkerType, |
| 448 | zind: &ZoneInd, |
| 449 | ) |
| 450 | -> Outcome<ChannelPool<UIDL, UID, ENC, KH>> |
| 451 | { |
| 452 | res!(self.check_zone_index(**zind)); |
| 453 | Ok(self.zwbots[**zind][wtyp].clone()) |
| 454 | } |
| 455 | |
| 456 | /// Returns channels for all worker pools of the given type, across all zones. |
| 457 | pub fn get_all_workers_of_type( |
| 458 | &self, |
| 459 | wtyp: &WorkerType, |
| 460 | ) |
| 461 | -> Vec<ChannelPool<UIDL, UID, ENC, KH>> |
| 462 | { |
| 463 | let mut pools = Vec::new(); |
| 464 | for z in 0..self.nz { |
| 465 | pools.push(self.zwbots[z][wtyp].clone()); |
| 466 | } |
| 467 | pools |
| 468 | } |
| 469 | |
| 470 | pub fn all_zwbots(&self) -> &Vec<ZoneWorkerChannels<UIDL, UID, ENC, KH>> { &self.zwbots } |
| 471 | pub fn all_zbots(&self) -> &ChannelPool<UIDL, UID, ENC, KH> { &self.zbots } |
| 472 | pub fn cfg(&self) -> &Simplex<OzoneMsg<UIDL, UID, ENC, KH>> { &self.cfg } |
| 473 | pub fn all_sbots(&self) -> &ChannelPool<UIDL, UID, ENC, KH> { &self.sbots } |
| 474 | pub fn sup(&self) -> &Simplex<OzoneMsg<UIDL, UID, ENC, KH>> { &self.sup } |
| 475 | |
| 476 | pub fn get_sbot(&self, sind: &BotPoolInd) -> Outcome<&Simplex<OzoneMsg<UIDL, UID, ENC, KH>>> { |
| 477 | self.sbots.get_bot(**sind) |
| 478 | } |
| 479 | |
| 480 | pub fn get_zbot(&self, zind: &ZoneInd) -> Outcome<&Simplex<OzoneMsg<UIDL, UID, ENC, KH>>> { |
| 481 | self.zbots.get_bot(**zind) |
| 482 | } |
| 483 | |
| 484 | pub fn get_zwbots(&self, zind: &ZoneInd) -> Outcome<&ZoneWorkerChannels<UIDL, UID, ENC, KH>> { |
| 485 | if **zind > self.zwbots.len() { |
| 486 | return Err(err!( |
| 487 | "Index {} exceeds number of zones {}.", **zind, self.zwbots.len(); |
| 488 | Index, TooBig)); |
| 489 | } |
| 490 | Ok(&self.zwbots[**zind]) |
| 491 | } |
| 492 | |
| 493 | pub fn get_bot( |
| 494 | &self, |
| 495 | ozid: &OzoneBotId, |
| 496 | ) |
| 497 | -> Outcome<&Simplex<OzoneMsg<UIDL, UID, ENC, KH>>> |
| 498 | { |
| 499 | Ok(match ozid { |
| 500 | // Solo bots |
| 501 | OzoneBotId::ConfigBot(..) => self.cfg(), |
| 502 | OzoneBotId::Supervisor(..) => self.sup(), |
| 503 | OzoneBotId::ServerBot(_, bpind) => res!(self.get_sbot(bpind)), |
| 504 | OzoneBotId::ZoneBot(_, zind) => res!(self.get_zbot(zind)), |
| 505 | OzoneBotId::CacheBot(_, zind, bpind) => |
| 506 | res!(res!(self.get_zwbots(zind))[&WorkerType::Cache].get_bot(**bpind)), |
| 507 | OzoneBotId::FileBot(_, zind, bpind) => |
| 508 | res!(res!(self.get_zwbots(zind))[&WorkerType::File].get_bot(**bpind)), |
| 509 | OzoneBotId::InitGarbageBot(_, zind, bpind) => |
| 510 | res!(res!(self.get_zwbots(zind))[&WorkerType::InitGarbage].get_bot(**bpind)), |
| 511 | OzoneBotId::ReaderBot(_, zind, bpind) => |
| 512 | res!(res!(self.get_zwbots(zind))[&WorkerType::Reader].get_bot(**bpind)), |
| 513 | OzoneBotId::ScanBot(_, zind, bpind) => |
| 514 | res!(res!(self.get_zwbots(zind))[&WorkerType::Scan].get_bot(**bpind)), |
| 515 | OzoneBotId::WriterBot(_, zind, bpind) => |
| 516 | res!(res!(self.get_zwbots(zind))[&WorkerType::Writer].get_bot(**bpind)), |
| 517 | _ => return Err(err!( |
| 518 | "Cannot return channel for {:?}.", ozid; |
| 519 | Bug, Invalid)), |
| 520 | }) |
| 521 | } |
| 522 | |
| 523 | // Mutate |
| 524 | pub fn zbots_mut(&mut self) -> &mut ChannelPool<UIDL, UID, ENC, KH> { &mut self.zbots } |
| 525 | pub fn cfg_mut(&mut self) -> &mut Simplex<OzoneMsg<UIDL, UID, ENC, KH>> { &mut self.cfg } |
| 526 | pub fn sbots_mut(&mut self) -> &mut ChannelPool<UIDL, UID, ENC, KH> { &mut self.sbots } |
| 527 | pub fn sup_mut(&mut self) -> &mut Simplex<OzoneMsg<UIDL, UID, ENC, KH>> { &mut self.sup } |
| 528 | |
| 529 | pub fn set_sbot( |
| 530 | &mut self, |
| 531 | bpind: &BotPoolInd, |
| 532 | chan: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>, |
| 533 | ) |
| 534 | -> Outcome<()> |
| 535 | { |
| 536 | self.sbots.set_bot(**bpind, chan) |
| 537 | } |
| 538 | |
| 539 | pub fn set_zbot( |
| 540 | &mut self, |
| 541 | zind: &ZoneInd, |
| 542 | chan: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>, |
| 543 | ) |
| 544 | -> Outcome<()> |
| 545 | { |
| 546 | self.zbots.set_bot(**zind, chan) |
| 547 | } |
| 548 | |
| 549 | pub fn set_cfg(&mut self, chan: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>) { self.cfg = chan; } |
| 550 | pub fn set_sup(&mut self, chan: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>) { self.sup = chan; } |
| 551 | |
| 552 | pub fn set_worker_bot( |
| 553 | &mut self, |
| 554 | wtyp: &WorkerType, |
| 555 | wind: &WorkerInd, |
| 556 | chan: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>, |
| 557 | ) |
| 558 | -> Outcome<()> |
| 559 | { |
| 560 | res!(self.check_zone_index(wind.z())); |
| 561 | self.zwbots[wind.z()][wtyp].set_bot(wind.b(), chan) |
| 562 | } |
| 563 | |
| 564 | fn check_zone_index(&self, ind: usize) -> Outcome<()> { |
| 565 | if ind > self.nz { |
| 566 | return Err(err!( |
| 567 | "Zone index {} into BotChannels exceeds number of zones {}.", |
| 568 | ind, self.nz; |
| 569 | Index, TooBig)); |
| 570 | } |
| 571 | Ok(()) |
| 572 | } |
| 573 | |
| 574 | /// Sends a `Finish` to every bot but the supervisor: the servers first, then each zone's |
| 575 | /// workers stage by stage in `FINISH_ORDER`, then the zone bots and the config bot. Before |
| 576 | /// each stage it waits until `ended` says every worker of the stage before it has ended. It |
| 577 | /// stops waiting at `until` and returns the stage it was waiting on, already finished, for |
| 578 | /// `finish_from` to carry on from; `None` once every bot has been sent its `Finish`. |
| 579 | pub fn finish_all<F: Fn(&[WorkerType]) -> bool>( |
| 580 | &self, |
| 581 | ended: F, |
| 582 | until: Instant, |
| 583 | ) |
| 584 | -> Outcome<Option<usize>> |
| 585 | { |
| 586 | // Starve servers. |
| 587 | res!(self.sbots.send_to_all(OzoneMsg::Finish)); |
| 588 | warn!(sync_log::stream(), "Shutdown: Completion request sent to server, finishing the \ |
| 589 | other bots in order, waiting up to {:?} for them.", |
| 590 | until.saturating_duration_since(Instant::now())); |
| 591 | res!(self.finish_stage(0)); |
| 592 | self.finish_from(0, ended, until) |
| 593 | } |
| 594 | |
| 595 | /// Carries on the `finish_all` that stopped waiting on `stage`. An `ended` that is always |
| 596 | /// true finishes everything left at once. |
| 597 | pub fn finish_from<F: Fn(&[WorkerType]) -> bool>( |
| 598 | &self, |
| 599 | stage: usize, |
| 600 | ended: F, |
| 601 | until: Instant, |
| 602 | ) |
| 603 | -> Outcome<Option<usize>> |
| 604 | { |
| 605 | let mut stage = stage; |
| 606 | while stage + 1 < FINISH_ORDER.len() { |
| 607 | let left = until.saturating_duration_since(Instant::now()); |
| 608 | let (start, timed_out) = res!(oxedyne_fe2o3_core::time::wait_for_true( |
| 609 | FINISH_CHECK_INTERVAL.min(left), |
| 610 | left, |
| 611 | || ended(FINISH_ORDER[stage]), |
| 612 | )); |
| 613 | if timed_out { |
| 614 | // Counted, not listed: a channel can only be read by taking its messages, and |
| 615 | // taken to be listed here, the records the writers had just released were |
| 616 | // destroyed, and their callers waited out the durability deadline for writes that |
| 617 | // had landed (2026-09-23). |
| 618 | warn!(sync_log::stream(), "Shutdown: The {:?} bots had not all ended after {:?}, \ |
| 619 | so those after them wait; pending: {:?}", |
| 620 | FINISH_ORDER[stage], start.elapsed(), self.msg_count()); |
| 621 | return Ok(Some(stage)); |
| 622 | } |
| 623 | stage += 1; |
| 624 | res!(self.finish_stage(stage)); |
| 625 | } |
| 626 | res!(self.zbots.finish_all()); |
| 627 | res!(self.cfg().send(OzoneMsg::Finish)); |
| 628 | warn!(sync_log::stream(), "Shutdown: Every bot has been sent its completion request."); |
| 629 | Ok(None) |
| 630 | } |
| 631 | |
| 632 | fn finish_stage(&self, stage: usize) -> Outcome<()> { |
| 633 | for zone in &self.zwbots { |
| 634 | res!(zone.finish(FINISH_ORDER[stage])); |
| 635 | } |
| 636 | Ok(()) |
| 637 | } |
| 638 | |
| 639 | pub fn fwd_to_all_zones(&self, msg: OzoneMsg<UIDL, UID, ENC, KH>) -> Outcome<usize> { |
| 640 | self.zbots.send_to_all(msg) |
| 641 | } |
| 642 | |
| 643 | pub fn send_to_all(&self, msg: OzoneMsg<UIDL, UID, ENC, KH>) -> Outcome<usize> { |
| 644 | let mut count = 0; |
| 645 | for zwbot in &self.zwbots { |
| 646 | count += res!(zwbot.send_to_all(msg.clone())); |
| 647 | } |
| 648 | count += res!(self.zbots.send_to_all(msg.clone())); |
| 649 | { |
| 650 | res!(self.cfg.send(msg.clone())); |
| 651 | count += 1; |
| 652 | } |
| 653 | count += res!(self.sbots.send_to_all(msg.clone())); |
| 654 | { |
| 655 | res!(self.sup.send(msg.clone())); |
| 656 | count += 1; |
| 657 | } |
| 658 | Ok(count) |
| 659 | } |
| 660 | |
| 661 | pub fn msg_count(&self) -> OzoneMsgCount { |
| 662 | let mut zone_counts = Vec::new(); |
| 663 | for zone in &self.zwbots { |
| 664 | zone_counts.push(zone.msg_count()); |
| 665 | } |
| 666 | OzoneMsgCount { |
| 667 | nz: self.nz, |
| 668 | zwbots: zone_counts, |
| 669 | zbots: self.zbots.msg_count(), |
| 670 | cfg: self.cfg.len(), |
| 671 | sbots: self.sbots.msg_count(), |
| 672 | sup: self.sup.len(), |
| 673 | } |
| 674 | } |
| 675 | |
| 676 | pub fn dump_pending_messages( |
| 677 | lines: Vec<String>, // obtain using drain_messages |
| 678 | label: &str, |
| 679 | z: Option<usize>, |
| 680 | b: Option<usize>, |
| 681 | ) { |
| 682 | match (z, b) { |
| 683 | (Some(z), Some(b)) => debug!(sync_log::stream(), " Z{} B{} {} messages ({}):", z, b, label, lines.len()), |
| 684 | (Some(z), None) => debug!(sync_log::stream(), " Z{} {} messages ({}):", z, label, lines.len()), |
| 685 | (None, None) => debug!(sync_log::stream(), " {} messages ({}):", label, lines.len()), |
| 686 | _ => (), |
| 687 | } |
| 688 | if lines.len() > 0 { |
| 689 | for line in lines { |
| 690 | debug!(sync_log::stream(), " {}", line); |
| 691 | } |
| 692 | } |
| 693 | } |
| 694 | } |