Oregami
Repositories/oxedyne/fe2o3

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

1use 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
18use oxedyne_fe2o3_bot::{
19 bot::{
20 Bot,
21 LoopBreak,
22 },
23};
24use oxedyne_fe2o3_core::{
25 channels::Simplex,
26 thread::Semaphore,
27};
28use oxedyne_fe2o3_jdat::id::NumIdDat;
29
30use std::{
31 path::Path,
32};
33
34#[macro_export]
35macro_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]
53macro_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
64pub 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
164pub 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}