Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/src/bots/worker/bot_cache.rs

14.3 KiB, 69 runs

created by r1870400018:751, 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::{
4 base::bot_deps::*,
5 worker::{
6 bot_reader::ReadResult,
7 worker_deps::*,
8 },
9 },
10 data::{
11 cache::{
12 Cache,
13 ValueOrLocation,
14 },
15 core::Key,
16 },
17 file::{
18 floc::FileLocation,
19 stored::RecordDigest,
20 },
21 test::hooks,
22};
23
24use oxedyne_fe2o3_core::channels::Recv;
25use oxedyne_fe2o3_iop_db::api::Meta;
26use oxedyne_fe2o3_jdat::id::NumIdDat;
27
28use std::{
29 sync::Arc,
30 time::Instant,
31};
32
33#[derive(Debug)]
34pub struct CacheBot<
35 const UIDL: usize,
36 UID: NumIdDat<UIDL>,
37 ENC: Encrypter,
38 KH: Hasher,
39 PR: Hasher,
40 CS: Checksummer,
41>{
42 // Identity
43 wind: WorkerInd,
44 wtyp: WorkerType,
45 // Bot
46 sem: Semaphore,
47 errc: Arc<Mutex<usize>>,
48 log_stream_id: String,
49 // Config
50 zdir: ZoneDir,
51 // Comms
52 chan_in: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>,
53 // API
54 api: OzoneApi<UIDL, UID, ENC, KH, PR, CS>,
55 // State
56 active: bool,
57 cache: Cache<UIDL, UID>,
58 inited: bool,
59 trep: Instant,
60}
61
62impl<
63 const UIDL: usize,
64 UID: NumIdDat<UIDL> + 'static,
65 ENC: Encrypter + 'static,
66 KH: Hasher + 'static,
67 PR: Hasher,
68 CS: Checksummer,
69>
70 WorkerBot<UIDL, UID, ENC, KH, PR, CS> for CacheBot<UIDL, UID, ENC, KH, PR, CS>
71{
72 workerbot_methods!();
73}
74
75impl<
76 const UIDL: usize,
77 UID: NumIdDat<UIDL> + 'static,
78 ENC: Encrypter + 'static,
79 KH: Hasher + 'static,
80 PR: Hasher,
81 CS: Checksummer,
82>
83 OzoneBot<UIDL, UID, ENC, KH, PR, CS> for CacheBot<UIDL, UID, ENC, KH, PR, CS>
84{
85 ozonebot_methods!();
86}
87
88impl<
89 const UIDL: usize,
90 UID: NumIdDat<UIDL> + 'static,
91 ENC: Encrypter + 'static,
92 KH: Hasher + 'static,
93 PR: Hasher,
94 CS: Checksummer,
95>
96 Bot<{ BID_LEN }, Bid, OzoneMsg<UIDL, UID, ENC, KH>> for CacheBot<UIDL, UID, ENC, KH, PR, CS>
97{
98 bot_methods!();
99
100 fn go(&mut self) {
101
102 sync_log::set_stream(self.log_stream_id());
103
104 if self.no_init() { return; }
105 self.now_listening();
106 loop {
107 if self.wind().b() < self.cfg().num_bots_per_zone(self.wtyp()) {
108
109 if self.trep.elapsed() > self.cfg().zone_state_update_interval() {
110 // Automated state reporting.
111 self.trep = Instant::now();
112 if let Some(zbot) = self.zbot() {
113 if let Err(e) = zbot.send(
114 OzoneMsg::CacheSize(
115 self.wind().b(),
116 self.cache().get_size(),
117 self.cache().get_ancillary_size(),
118 )
119 ) {
120 self.result(&Err(err!(e,
121 "{}: Cannot send cache size update to zbot.", self.ozid();
122 Channel, Write)));
123 }
124 }
125 }
126
127 if self.listen().must_end() { break; }
128
129 } else {
130 // This bot is to be terminated. Forward incoming messages to the remaining bots of
131 // this type.
132 }
133 }
134 }
135
136 fn listen(&mut self) -> LoopBreak {
137 match self.chan_in().recv_timeout(self.cfg().zone_state_update_interval()) {
138 Recv::Result(Err(e)) => self.err_cannot_receive(err!(e,
139 "{}: Waiting for message.", self.ozid();
140 IO, Channel)),
141 Recv::Result(Ok(msg)) => {
142 if let Some(msg) = self.listen_worker(msg) {
143 match msg {
144 // COMMAND
145 OzoneMsg::ClearCache(resp) => {
146 self.cache_mut().clear_all_values();
147 self.respond(Ok(OzoneMsg::Ok), &resp);
148 },
149 OzoneMsg::SetCacheSizeLimit(size_lim) => {
150 self.cache_mut().set_lim(size_lim);
151 },
152 // WRITE
153 OzoneMsg::GcCacheUpdateRequest(buf, resp_g1) => {
154 let mut old_flocs = Vec::new();
155 for (key, floc, meta) in buf {
156 if let Some(old_floc) = self.cache_mut().reanchor(&key, &floc, &meta) {
157 match RecordDigest::new(&key, &meta) {
158 Ok(rid) => old_flocs.push((old_floc, rid)),
159 // Its move entry stays, and keeps the file from being
160 // collected again; the location itself is right.
161 Err(e) => self.error(err!(e,
162 "{}: Naming a record a collection re-anchored.", self.ozid();
163 Data, Encode)),
164 }
165 }
166 }
167 if let Err(e) = resp_g1.send(
168 OzoneMsg::GcCacheUpdateResponse(old_flocs)
169 ) {
170 self.err_cannot_send(err!(e,
171 "{}: Sending cache update response back to gbot.", self.ozid();
172 IO, Channel));
173 }
174 },
175 OzoneMsg::Insert(key, val, cind, floc, ilen, meta, resp_w1, unconfirmed) => {
176 let result = self.insert(key, val, cind, floc, ilen, meta, resp_w1, unconfirmed);
177 self.result(&result);
178 },
179 // READ
180 OzoneMsg::DumpCacheRequest(resp) => {
181 if let Err(e) = resp.send(OzoneMsg::DumpCacheResponse(
182 self.wind().clone(),
183 self.cache().clone(),
184 )) {
185 self.err_cannot_send(err!(e,
186 "{}: Responding to {:?} with cache dump.", self.ozid(), resp.ozid();
187 Data, IO, Channel));
188 }
189 },
190 //OzoneMsg::GetUsers(resp) => {
191 // let mut kuserdat = Vec::new();
192 // let prefix = constant::USER_DAT.byte_prefix();
193 // for (k, _) in self.cache().map() {
194 // if k.len() > UsrKindId::CODE_BYTE_LEN {
195 // if k[0] == Dat::USR_CODE {
196 // if UsrKindId::prefix_matches(prefix, &k[1..]) {
197 // let kdat = match Dat::from_bytes(&k) {
198 // Ok((dat, _)) => dat,
199 // _ => continue,
200 // };
201 // if let Dat::Usr(_, optboxdat) = &kdat {
202 // match optboxdat {
203 // Some(boxdat) => match **boxdat {
204 // Dat::U128(id) => {
205 // kuserdat.push((id, kdat.clone()));
206 // },
207 // dat => self.error(err!(
208 // "Custom Usr daticle should contain a \
209 // Dat::U128 but found {:?}.", dat,
210 // ), Bug, Invalid, Input)),
211 // },
212 // None => self.error(err!(
213 // "Custom Usr daticle should contain \
214 // something but None was found.",
215 // ), Bug, Invalid, Input)),
216 // }
217 // }
218 // }
219 // }
220 // }
221 // }
222 // self.respond(Ok(OzoneMsg::UserKeys(kuserdat)), &resp);
223 //},
224 OzoneMsg::ReadCache(key, resp_r2) => {
225 let result = self.read(key, resp_r2);
226 self.result(&result);
227 },
228 _ => return self.listen_more(msg),
229 }
230 }
231 },
232 Recv::Empty => (),
233 }
234 LoopBreak(false)
235 }
236
237}
238
239impl<
240 const UIDL: usize,
241 UID: NumIdDat<UIDL> + 'static,
242 ENC: Encrypter + 'static,
243 KH: Hasher + 'static,
244 PR: Hasher,
245 CS: Checksummer,
246>
247 CacheBot<UIDL, UID, ENC, KH, PR, CS>
248{
249 pub fn new(
250 args: ZoneWorkerInitArgs<UIDL, UID, ENC, KH, PR, CS>,
251 )
252 -> Self
253 {
254 let cache = Cache::new(Some(&args.api.ozid));
255 Self {
256 // Identity
257 wind: args.wind,
258 wtyp: args.wtyp,
259 // Bot
260 sem: args.sem,
261 errc: Arc::new(Mutex::new(0)),
262 log_stream_id: args.log_stream_id,
263 // Config
264 zdir: ZoneDir::default(),
265 // Comms
266 chan_in: args.chan_in,
267 // API
268 api: args.api,
269 // State
270 active: false,
271 cache,
272 inited: false,
273 trep: Instant::now(),
274 }
275 }
276
277 fn cache(&self) -> &Cache<UIDL, UID> { &self.cache }
278 fn cache_mut(&mut self) -> &mut Cache<UIDL, UID> { &mut self.cache }
279
280 pub fn max_file_len(&self) -> usize {
281 self.cfg().data_file_max_bytes as usize
282 }
283
284 pub fn activate(mut self) -> Self {
285 self.active = true;
286 self
287 }
288
289 pub fn insert(
290 &mut self,
291 key: Vec<u8>,
292 val: Option<Vec<u8>>,
293 cind: Option<usize>,
294 floc: FileLocation,
295 ilen: usize,
296 meta: Meta<UIDL, UID>,
297 resp_w1: Responder<UIDL, UID, ENC, KH>,
298 unconfirmed: Option<Error<ErrTag>>,
299 )
300 -> Outcome<()>
301 {
302 hooks::insert_delay();
303 // [12] Insert the data into the key-chosen zone cache.
304 let floc_new = floc.clone();
305 let floc_old_opt = match self.cache.insert(
306 key,
307 val,
308 floc,
309 meta,
310 ) {
311 Ok(floc_old_opt) => floc_old_opt,
312 Err(e) => {
313 // The caller is waiting on this answer. Only logged, a failure here reached it
314 // as an expired durability deadline, which says the write is on its way.
315 let e = err!(e,
316 "{}: A written record could not be entered in the cache.", self.ozid();
317 Data, Write, Unconfirmed);
318 self.respond(Err(e.clone()), &resp_w1);
319 return Err(e);
320 },
321 };
322
323 let key_present = floc_old_opt.is_some();
324
325 // [13] Inform the caller of successful file write and cache insertion, or of the barrier
326 // that failed after the write: told here, the caller can read its write back once it
327 // hears, as it can a confirmed one.
328 match (unconfirmed, cind) {
329 (Some(e), _) => self.respond(Err(e), &resp_w1),
330 (None, Some(cind)) => self.respond(Ok(OzoneMsg::KeyChunkExists(key_present, cind)), &resp_w1),
331 (None, None) => self.respond(Ok(OzoneMsg::KeyExists(key_present)), &resp_w1),
332 }
333 self.respond(Ok(OzoneMsg::Finish), &resp_w1);
334
335 // [14] Insert the new data into the file state data map, via a file-selected cbot.
336 let bots = res!(self.fbots());
337 let (bot, _) = bots.choose_bot(
338 &ChooseBot::ByFile(floc_new.file_number())
339 );
340 res!(bot.send(OzoneMsg::UpdateData {
341 floc_new,
342 ilen,
343 floc_old_opt,
344 from_id: self.ozid().clone(),
345 }));
346
347 Ok(())
348 }
349
350 pub fn read(
351 &mut self,
352 key: Key,
353 resp_r2: Responder<UIDL, UID, ENC, KH>,
354 )
355 -> Outcome<()>
356 {
357 // <3> The cbot accesses its cache.
358 match res!(self.cache.get(key.as_bytes())) {
359 Some(vloc) => {
360 match vloc {
361 ValueOrLocation::Location(mloc) => {
362 // <4> Only the file location is available, so send a request to the
363 // appropriate fbot, forwarding the responder.
364 let fnum = mloc.file_number();
365 let bots = res!(self.fbots());
366 let (bot, _) = bots.choose_bot(&ChooseBot::ByFile(fnum));
367 // The key goes with the location, so that a move entry at the same
368 // offset is taken only if it is this record's.
369 res!(bot.send(
370 OzoneMsg::ReadFileRequest(
371 fnum,
372 key.into_bytes(),
373 mloc,
374 resp_r2,
375 )));
376 },
377 // <6> Send result directly back to rbot.
378 ValueOrLocation::Value(val, meta) =>
379 res!(resp_r2.send(OzoneMsg::ReadResult(ReadResult::Value(val, meta)))),
380 ValueOrLocation::Deleted(meta) =>
381 res!(resp_r2.send(OzoneMsg::ReadResult(ReadResult::Deleted(meta)))),
382 }
383
384 },
385 // <6> Send result directly back to rbot.
386 None => res!(resp_r2.send(OzoneMsg::ReadResult(ReadResult::None))),
387 }
388 Ok(())
389 }
390}