oxedyne/fe2o3/fe2o3_o3db_sync/src/bots/bot_server.rs
8.1 KiB, 45 runs
created by r1870400018:741, 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 | bots::base::bot_deps::*, |
| 4 | comm::channels::BotChannels, |
| 5 | }; |
| 6 | |
| 7 | use oxedyne_fe2o3_core::{ |
| 8 | prelude::*, |
| 9 | }; |
| 10 | use oxedyne_fe2o3_jdat::id::NumIdDat; |
| 11 | |
| 12 | use std::{ |
| 13 | sync::Arc, |
| 14 | }; |
| 15 | |
| 16 | /// Listens internally, and possibly on the wire, for database commands. |
| 17 | pub struct ServerBot< |
| 18 | const UIDL: usize, |
| 19 | UID: NumIdDat<UIDL>, |
| 20 | ENC: Encrypter, |
| 21 | KH: Hasher, |
| 22 | PR: Hasher, |
| 23 | CS: Checksummer, |
| 24 | >{ |
| 25 | // Bot |
| 26 | sem: Semaphore, |
| 27 | errc: Arc<Mutex<usize>>, |
| 28 | log_stream_id: String, |
| 29 | // Comms |
| 30 | chan_in: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>, |
| 31 | // API |
| 32 | api: OzoneApi<UIDL, UID, ENC, KH, PR, CS>, |
| 33 | // State |
| 34 | inited: bool, |
| 35 | } |
| 36 | |
| 37 | impl< |
| 38 | const UIDL: usize, |
| 39 | UID: NumIdDat<UIDL> + 'static, |
| 40 | ENC: Encrypter + 'static, |
| 41 | KH: Hasher + 'static, |
| 42 | PR: Hasher, |
| 43 | CS: Checksummer, |
| 44 | > |
| 45 | Bot<{ BID_LEN }, Bid, OzoneMsg<UIDL, UID, ENC, KH>> for ServerBot<UIDL, UID, ENC, KH, PR, CS> |
| 46 | { |
| 47 | bot_methods!(); |
| 48 | |
| 49 | fn go(&mut self) { |
| 50 | |
| 51 | sync_log::set_stream(self.log_stream_id()); |
| 52 | |
| 53 | if self.no_init() { return; } |
| 54 | info!(sync_log::stream(), "{}: Listening for database requests.", self.ozid()); |
| 55 | self.now_listening(); |
| 56 | loop { |
| 57 | if self.listen().must_end() { break; } |
| 58 | } |
| 59 | } |
| 60 | |
| 61 | fn listen(&mut self) -> LoopBreak { |
| 62 | // INTERNAL |
| 63 | // Block until a message arrives. The server bot has no periodic maintenance of its |
| 64 | // own, so it must sleep rather than poll the channel, otherwise an idle database burns |
| 65 | // a CPU core. A shutdown is delivered as an `OzoneMsg::Finish` on this same channel, |
| 66 | // which wakes the blocking receive immediately. |
| 67 | match self.chan_in().recv() { |
| 68 | Err(e) => self.err_cannot_receive(err!(e, |
| 69 | "{}: Waiting for message on internal channel.", self.ozid(); |
| 70 | IO, Channel)), |
| 71 | Ok(msg) => match msg { |
| 72 | OzoneMsg::Get { key, schms2, resp } => { |
| 73 | match self.api().get_wait(&key, schms2.as_ref()) { |
| 74 | Err(e) => { |
| 75 | // The caller is waiting on this answer, and a failure only logged |
| 76 | // reached it as a timeout that named no cause. |
| 77 | let e = err!(e, |
| 78 | "{}: While trying to get value for key {:?}", self.ozid(), key; |
| 79 | Data, Read); |
| 80 | self.error(e.clone()); |
| 81 | self.respond(Err(e), &resp); |
| 82 | }, |
| 83 | Ok(result) => match resp.send(OzoneMsg::GetResult(result)) { |
| 84 | Err(e) => self.err_cannot_send(err!(e, |
| 85 | "{}: While sending an OzoneMsg::GetResult back via a responder.", |
| 86 | self.ozid(); |
| 87 | Data, Channel)), |
| 88 | Ok(()) => (), |
| 89 | }, |
| 90 | } |
| 91 | }, |
| 92 | OzoneMsg::Put { key, val, user, schms2, resp } => { |
| 93 | debug!(sync_log::stream(), "Store key: {:?}",key); |
| 94 | let caller = resp.clone(); |
| 95 | match self.api().store_dat_using_responder( |
| 96 | key, |
| 97 | val, |
| 98 | user, |
| 99 | schms2.as_ref(), |
| 100 | resp, |
| 101 | ) { |
| 102 | Err(e) => { |
| 103 | // As for a get: the caller hears the cause, not a timeout. |
| 104 | let e = err!(e, |
| 105 | "{}: While trying to put value.", self.ozid(); |
| 106 | Data, Write); |
| 107 | self.error(e.clone()); |
| 108 | self.respond(Err(e), &caller); |
| 109 | }, |
| 110 | Ok(_nchunks) => (), |
| 111 | } |
| 112 | }, |
| 113 | _ => if self.listen_more(msg).must_end() { |
| 114 | return LoopBreak(true); |
| 115 | }, |
| 116 | // TODO one for OzoneMsg::Delete? |
| 117 | }, |
| 118 | } |
| 119 | |
| 120 | //// EXTERNAL |
| 121 | //match self.sock.recv_from(&mut self.buf) { // Receive udp packet, non-blocking. |
| 122 | // Err(e) => { |
| 123 | // match self.timer.write() { |
| 124 | // Err(e) => self.error(err!( |
| 125 | // "While locking timer for writing: {}.", e), ErrTag::Poisoned)), |
| 126 | // Ok(mut unlocked_timer) => { unlocked_timer.update(); }, |
| 127 | // } |
| 128 | // //self.timer.update(); |
| 129 | // match e.kind() { |
| 130 | // io::ErrorKind::WouldBlock | io::ErrorKind::InvalidInput => (), |
| 131 | // _ => self.err_cannot_receive(Error::from(e), |
| 132 | // errmsg!("Waiting for message on external UDP socket")), |
| 133 | // } |
| 134 | // }, |
| 135 | // Ok((n, src_addr)) => { |
| 136 | // let mut buf_clone = [0u8; constant::UDP_BUFFER_SIZE]; |
| 137 | // for i in 0..n { |
| 138 | // buf_clone[i] = self.buf[i]; |
| 139 | // } |
| 140 | // let state = wire::ServerProcessorEnv::<POWH, SGN> { |
| 141 | // buf: buf_clone, |
| 142 | // n, |
| 143 | // src_addr, |
| 144 | // // Comms |
| 145 | // //wschms: WireSchemes<WENC, WCS, POWH, SGN, HS>, |
| 146 | // //buf: [u8; constant::UDP_BUFFER_SIZE], |
| 147 | // //chan: Simplex<Msg>, |
| 148 | // //chans: BotChannels, |
| 149 | // protoref: self.protoref.clone(), // Arc |
| 150 | // timer: self.timer.clone(), // Arc |
| 151 | // // Schemes. |
| 152 | // schmdb: self.schmdb.clone(), // Arc |
| 153 | // // Keys. |
| 154 | // pack_sigkeys: self.pack_sigkeys.clone(), // Arc |
| 155 | // // Declared source address protection. |
| 156 | // agrd: self.agrd.clone(), // Arc |
| 157 | // // User protection. |
| 158 | // ugrd: self.ugrd.clone(), // Arc |
| 159 | // // Packet validation. |
| 160 | // packval: self.packval.clone(), |
| 161 | // gpzparams: self.gpzparams.clone(), |
| 162 | // // Message assembly. |
| 163 | // massembler: self.massembler.clone(), // Arc |
| 164 | // ma_params: self.ma_params.clone(), |
| 165 | // // Database configuration values. |
| 166 | // time_horiz: self.cfg().server_pow_time_horiz_secs, |
| 167 | // accept_unknown: self.cfg().server_accept_unknown_users, |
| 168 | // }; |
| 169 | // task::spawn(state.process( |
| 170 | // //&mut self RingTimer<{ constant::REQ_TIMER_LEN }>, |
| 171 | // )); |
| 172 | // }, |
| 173 | //} // Receive udp packet. |
| 174 | |
| 175 | //// Message assembly garbage collection. |
| 176 | //if self.ma_gc_last.elapsed() > self.ma_gc_int { |
| 177 | // self.massembler.message_assembly_garbage_collection(&self.ma_params); |
| 178 | // self.ma_gc_last = Instant::now(); |
| 179 | //} |
| 180 | |
| 181 | LoopBreak(false) |
| 182 | } |
| 183 | } |
| 184 | |
| 185 | impl< |
| 186 | const UIDL: usize, |
| 187 | UID: NumIdDat<UIDL> + 'static, |
| 188 | ENC: Encrypter + 'static, |
| 189 | KH: Hasher + 'static, |
| 190 | PR: Hasher, |
| 191 | CS: Checksummer, |
| 192 | > |
| 193 | OzoneBot<UIDL, UID, ENC, KH, PR, CS> for ServerBot<UIDL, UID, ENC, KH, PR, CS> |
| 194 | { |
| 195 | ozonebot_methods!(); |
| 196 | } |
| 197 | |
| 198 | impl< |
| 199 | const UIDL: usize, |
| 200 | UID: NumIdDat<UIDL> + 'static, |
| 201 | ENC: Encrypter + 'static, |
| 202 | KH: Hasher + 'static, |
| 203 | PR: Hasher, |
| 204 | CS: Checksummer, |
| 205 | > |
| 206 | ServerBot<UIDL, UID, ENC, KH, PR, CS> |
| 207 | { |
| 208 | pub fn new( |
| 209 | args: BotInitArgs<UIDL, UID, ENC, KH, PR, CS>, |
| 210 | ) |
| 211 | -> Self |
| 212 | { |
| 213 | Self { |
| 214 | // Bot |
| 215 | sem: args.sem, |
| 216 | errc: Arc::new(Mutex::new(0)), |
| 217 | log_stream_id: args.log_stream_id, |
| 218 | // Comms |
| 219 | chan_in: args.chan_in, |
| 220 | // API |
| 221 | api: args.api, |
| 222 | // State |
| 223 | inited: false, |
| 224 | } |
| 225 | } |
| 226 | |
| 227 | |
| 228 | } |