Oregami
Repositories/oxedyne/fe2o3

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

1use crate::{
2 prelude::*,
3 bots::base::bot_deps::*,
4 comm::channels::BotChannels,
5};
6
7use oxedyne_fe2o3_core::{
8 prelude::*,
9};
10use oxedyne_fe2o3_jdat::id::NumIdDat;
11
12use std::{
13 sync::Arc,
14};
15
16/// Listens internally, and possibly on the wire, for database commands.
17pub 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
37impl<
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
185impl<
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
198impl<
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}