oxedyne/fe2o3/fe2o3_o3db_sync/src/bots/base/bot.rs
6.2 KiB, 38 runs
created by r1870400018:731, 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 | cfg::OzoneConfig, |
| 5 | id::{ |
| 6 | Bid, |
| 7 | BID_LEN, |
| 8 | OzoneBotId, |
| 9 | }, |
| 10 | }, |
| 11 | comm::{ |
| 12 | channels::BotChannels, |
| 13 | msg::OzoneMsg, |
| 14 | response::Responder, |
| 15 | }, |
| 16 | }; |
| 17 | |
| 18 | use oxedyne_fe2o3_bot::{ |
| 19 | bot::{ |
| 20 | Bot, |
| 21 | LoopBreak, |
| 22 | }, |
| 23 | }; |
| 24 | use oxedyne_fe2o3_core::{ |
| 25 | channels::Simplex, |
| 26 | thread::Semaphore, |
| 27 | }; |
| 28 | use oxedyne_fe2o3_jdat::id::NumIdDat; |
| 29 | |
| 30 | use std::{ |
| 31 | path::Path, |
| 32 | }; |
| 33 | |
| 34 | #[macro_export] |
| 35 | macro_rules! bot_methods { () => { |
| 36 | fn id(&self) -> Bid { self.ozid().bid() } |
| 37 | fn errc(&self) -> &Arc<Mutex<usize>> { &self.errc } |
| 38 | fn chan_in(&self) -> &Simplex<OzoneMsg<UIDL, UID, ENC, KH>> { &self.chan_in } |
| 39 | fn label(&self) -> String { fmt!("{}", self.ozid()) } |
| 40 | fn err_count_warning(&self) -> usize { constant::BOT_ERR_COUNT_WARNING } |
| 41 | fn log_stream_id(&self) -> String { self.log_stream_id.clone() } |
| 42 | fn set_chan_in(&mut self, chan_in: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>) { |
| 43 | self.chan_in = chan_in; |
| 44 | } |
| 45 | fn init(&mut self) -> Outcome<()> { |
| 46 | info!(sync_log::stream(), "{:?}: Initialising.", self.ozid()); |
| 47 | self.inited = true; |
| 48 | Ok(()) |
| 49 | } |
| 50 | } } |
| 51 | |
| 52 | #[macro_export] |
| 53 | macro_rules! ozonebot_methods { () => { |
| 54 | fn api(&self) -> &OzoneApi<UIDL, UID, ENC, KH, PR, CS> { &self.api } |
| 55 | fn api_mut(&mut self) -> &mut OzoneApi<UIDL, UID, ENC, KH, PR, CS> { &mut self.api } |
| 56 | fn ozid(&self) -> &OzoneBotId { &self.api.ozid } |
| 57 | fn db_root(&self) -> &Path { &self.api.db_root } |
| 58 | fn cfg(&self) -> &OzoneConfig { &self.api.cfg } |
| 59 | fn chans(&self) -> &BotChannels<UIDL, UID, ENC, KH> { &self.api.chans } |
| 60 | fn inited(&self) -> bool { self.inited } |
| 61 | fn set_chans(&mut self, chans: BotChannels<UIDL, UID, ENC, KH>) { self.api.chans = chans } |
| 62 | } } |
| 63 | |
| 64 | pub trait OzoneBot< |
| 65 | const UIDL: usize, |
| 66 | UID: NumIdDat<UIDL> + 'static, |
| 67 | ENC: Encrypter + 'static, |
| 68 | KH: Hasher + 'static, |
| 69 | PR: Hasher, |
| 70 | CS: Checksummer, |
| 71 | >: |
| 72 | Bot<{ BID_LEN }, Bid, OzoneMsg<UIDL, UID, ENC, KH>> |
| 73 | { |
| 74 | // Required. |
| 75 | fn api(&self) -> &OzoneApi<UIDL, UID, ENC, KH, PR, CS>; |
| 76 | fn api_mut(&mut self) -> &mut OzoneApi<UIDL, UID, ENC, KH, PR, CS>; |
| 77 | fn ozid(&self) -> &OzoneBotId; |
| 78 | fn db_root(&self) -> &Path; |
| 79 | fn cfg(&self) -> &OzoneConfig; |
| 80 | fn chans(&self) -> &BotChannels<UIDL, UID, ENC, KH>; |
| 81 | fn inited(&self) -> bool; |
| 82 | fn set_chans(&mut self, chans: BotChannels<UIDL, UID, ENC, KH>); |
| 83 | |
| 84 | // Provided. |
| 85 | fn no_init(&self) -> bool { |
| 86 | if !self.inited() { |
| 87 | error!(sync_log::stream(), err!( |
| 88 | "Attempt to start {} before running init()", self.label(); |
| 89 | Init, Missing)); |
| 90 | return true; |
| 91 | } |
| 92 | false |
| 93 | } |
| 94 | fn respond( |
| 95 | &self, |
| 96 | result: Outcome<OzoneMsg<UIDL, UID, ENC, KH>>, |
| 97 | resp: &Responder<UIDL, UID, ENC, KH>, |
| 98 | ) { |
| 99 | match resp.channel() { |
| 100 | None => return, |
| 101 | Some(simplex) => { |
| 102 | let msg = match result { |
| 103 | Err(e) => OzoneMsg::Error(e), |
| 104 | Ok(m) => m, |
| 105 | }; |
| 106 | let err_msg = format!( |
| 107 | "While trying to return a msg {:?} via a responder for ticket {}", |
| 108 | msg, resp.ticket(), |
| 109 | ); |
| 110 | if let Err(e) = simplex.send(msg) { |
| 111 | self.err_cannot_send(err!(e, "{}", err_msg; Channel, Write)); |
| 112 | } |
| 113 | }, |
| 114 | } |
| 115 | } |
| 116 | |
| 117 | /// Message handling common to all Ozone bots. |
| 118 | fn listen_more(&mut self, msg: OzoneMsg<UIDL, UID, ENC, KH>) -> LoopBreak { |
| 119 | match msg { |
| 120 | OzoneMsg::Finish => { |
| 121 | trace!(sync_log::stream(), "{}: Finish message received, finishing now.", self.ozid()); |
| 122 | return LoopBreak(true); |
| 123 | }, |
| 124 | OzoneMsg::Ready => info!(sync_log::stream(), "{} ready to receive messages now.", self.ozid()), |
| 125 | OzoneMsg::Channels(chans, resp) => { |
| 126 | self.set_chans(chans); |
| 127 | match self.chans().get_bot(self.ozid()) { |
| 128 | Err(e) => self.error(e), |
| 129 | Ok(chan) => { |
| 130 | if self.chan_in().len() > 0 { |
| 131 | BotChannels::<UIDL, UID, ENC, KH>::dump_pending_messages( |
| 132 | self.chan_in().drain_messages(), |
| 133 | &fmt!("Updating {} channel, clearing out existing", self.ozid()), |
| 134 | None, |
| 135 | None, |
| 136 | ); |
| 137 | } |
| 138 | let chan_clone = chan.clone(); |
| 139 | self.set_chan_in(chan_clone); |
| 140 | }, |
| 141 | } |
| 142 | self.respond(Ok(OzoneMsg::ChannelsReceived(self.ozid().clone())), &resp); |
| 143 | trace!(sync_log::stream(), "{}: Channel update received.", self.ozid()); |
| 144 | }, |
| 145 | OzoneMsg::Ping(id, resp) => { |
| 146 | // A ping reports the bot's error tally as well as its liveness, so a |
| 147 | // caller can tell a healthy bot from one that is running but failing. |
| 148 | let errs = match self.error_count() { |
| 149 | Ok(n) => n, |
| 150 | Err(_) => usize::MAX, |
| 151 | }; |
| 152 | if let Err(e) = resp.send(OzoneMsg::Pong(self.ozid().clone(), errs)) { |
| 153 | self.err_cannot_send(err!(e, |
| 154 | "Attempt to return a ping from {:?} failed", id; |
| 155 | IO, Channel)); |
| 156 | } |
| 157 | }, |
| 158 | _ => error!(sync_log::stream(), err!("{}: Message {:?} not recognised.", self.ozid(), msg; Invalid, Input)), |
| 159 | } |
| 160 | LoopBreak(false) |
| 161 | } |
| 162 | } |
| 163 | |
| 164 | pub struct BotInitArgs< |
| 165 | const UIDL: usize, |
| 166 | UID: NumIdDat<UIDL>, |
| 167 | ENC: Encrypter, |
| 168 | KH: Hasher, |
| 169 | PR: Hasher, |
| 170 | CS: Checksummer, |
| 171 | >{ |
| 172 | // Bot |
| 173 | pub sem: Semaphore, |
| 174 | pub log_stream_id: String, |
| 175 | // Comms |
| 176 | pub chan_in: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>, |
| 177 | // API |
| 178 | pub api: OzoneApi<UIDL, UID, ENC, KH, PR, CS>, |
| 179 | } |