Oregami
Repositories/oxedyne/fe2o3

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

20.1 KiB, 95 runs

created by r1870400018:757, 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::constant,
4 bots::{
5 base::bot_deps::*,
6 worker::worker_deps::*,
7 },
8 data::{
9 cache::{
10 MetaLocation,
11 },
12 core::{
13 Key,
14 Value,
15 },
16 },
17 file::{
18 core::{
19 FileAccess,
20 FileType,
21 },
22 fcache::{
23 FileCache,
24 FileCacheEntry,
25 FileCacheIndex,
26 },
27 floc::{
28 FileLocation,
29 FileNum,
30 },
31 stored::StoredKey,
32 },
33};
34
35use oxedyne_fe2o3_iop_db::api::Meta;
36use oxedyne_fe2o3_jdat::{
37 prelude::*,
38 id::NumIdDat,
39};
40use oxedyne_fe2o3_iop_hash::csum::Checksummer;
41
42use std::{
43 fs::File,
44 io::{
45 BufReader,
46 Read,
47 Seek,
48 SeekFrom,
49 },
50 sync::{
51 Arc,
52 RwLock,
53 },
54};
55
56#[derive(Clone, Debug)]
57pub enum ReadResult<
58 const UIDL: usize,
59 UID: NumIdDat<UIDL>,
60> {
61 // Filebot issued
62 None,
63 Location(MetaLocation<UIDL, UID>, bool),
64 // Cachebot issued
65 Value(Vec<u8>, Meta<UIDL, UID>),
66 Deleted(Meta<UIDL, UID>),
67}
68
69pub struct ReaderBot<
70 const UIDL: usize,
71 UID: NumIdDat<UIDL>,
72 ENC: Encrypter,
73 KH: Hasher,
74 PR: Hasher,
75 CS: Checksummer,
76>{
77 // Identity
78 wind: WorkerInd,
79 wtyp: WorkerType,
80 // Bot
81 sem: Semaphore,
82 errc: Arc<Mutex<usize>>,
83 log_stream_id: String,
84 // Config
85 zdir: ZoneDir,
86 // Comms
87 chan_in: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>,
88 // API
89 api: OzoneApi<UIDL, UID, ENC, KH, PR, CS>,
90 // State
91 active: bool,
92 fcache: FileCache,
93 inited: bool,
94}
95
96impl<
97 const UIDL: usize,
98 UID: NumIdDat<UIDL> + 'static,
99 ENC: Encrypter + 'static,
100 KH: Hasher + 'static,
101 PR: Hasher,
102 CS: Checksummer,
103>
104 WorkerBot<UIDL, UID, ENC, KH, PR, CS> for ReaderBot<UIDL, UID, ENC, KH, PR, CS>
105{
106 workerbot_methods!();
107}
108
109impl<
110 const UIDL: usize,
111 UID: NumIdDat<UIDL> + 'static,
112 ENC: Encrypter + 'static,
113 KH: Hasher + 'static,
114 PR: Hasher,
115 CS: Checksummer,
116>
117 OzoneBot<UIDL, UID, ENC, KH, PR, CS> for ReaderBot<UIDL, UID, ENC, KH, PR, CS>
118{
119 ozonebot_methods!();
120}
121
122impl<
123 const UIDL: usize,
124 UID: NumIdDat<UIDL> + 'static,
125 ENC: Encrypter + 'static,
126 KH: Hasher + 'static,
127 PR: Hasher,
128 CS: Checksummer,
129>
130 Bot<{ BID_LEN }, Bid, OzoneMsg<UIDL, UID, ENC, KH>> for ReaderBot<UIDL, UID, ENC, KH, PR, CS>
131{
132 bot_methods!();
133
134 fn go(&mut self) {
135
136 sync_log::set_stream(self.log_stream_id());
137
138 if self.no_init() { return; }
139 self.now_listening();
140 loop {
141 if self.wind().b() < self.cfg().num_bots_per_zone((&self).wtyp()) {
142 if self.listen().must_end() { break; }
143 } else {
144 // This bot is to be terminated. Forward incoming messages to the remaining bots of
145 // this type.
146 }
147 }
148 }
149
150 fn listen(&mut self) -> LoopBreak {
151 match self.chan_in().recv() {
152 Err(e) => self.err_cannot_receive(err!(e,
153 "{}: Waiting for message.", self.ozid();
154 IO, Channel)),
155 Ok(msg) => {
156 if let Some(msg) = self.listen_worker(msg) {
157 match msg {
158 // COMMAND
159 OzoneMsg::FileReplaced(fnum, typ) => self.drop_cached_file(fnum, &typ),
160 // WORK
161 OzoneMsg::Read(key, cbpind, resp_r1) => {
162 let result = self.read(key, cbpind);
163 self.respond(result, &resp_r1);
164 },
165 _ => return self.listen_more(msg),
166 }
167 }
168 },
169 }
170 LoopBreak(false)
171 }
172}
173
174impl<
175 const UIDL: usize,
176 UID: NumIdDat<UIDL> + 'static,
177 ENC: Encrypter + 'static,
178 KH: Hasher + 'static,
179 PR: Hasher,
180 CS: Checksummer,
181>
182 ReaderBot<UIDL, UID, ENC, KH, PR, CS>
183{
184 pub fn new(
185 args: ZoneWorkerInitArgs<UIDL, UID, ENC, KH, PR, CS>,
186 )
187 -> Self
188 {
189 Self {
190 // Identity
191 wind: args.wind,
192 wtyp: args.wtyp,
193 // Bot
194 sem: args.sem,
195 errc: Arc::new(Mutex::new(0)),
196 log_stream_id: args.log_stream_id,
197 // Config
198 zdir: ZoneDir::default(),
199 // Comms
200 chan_in: args.chan_in,
201 // API
202 api: args.api,
203 // State
204 active: false,
205 fcache: FileCache::new(constant::FILE_CACHE_EXPIRY_SECS),
206 inited: false,
207 }
208 }
209
210 fn ref_file_cache(&self) -> &FileCache { &self.fcache }
211 fn mut_file_cache(&mut self) -> &mut FileCache { &mut self.fcache }
212
213 /// Forgets any open handle on the given file, so that the next read of it opens the path
214 /// afresh. A collection renames a new file over an old one, which leaves a cached handle on
215 /// the unlinked inode: correct bytes for the file that was, at offsets that belong to the file
216 /// that is. Dropping the entry is the whole of the repair, and it keeps the read path free of
217 /// any per-read check for the same thing.
218 fn drop_cached_file(
219 &mut self,
220 fnum: FileNum,
221 typ: &FileType,
222 ) {
223 let k = FileCacheIndex { fnum, typ: typ.clone() };
224 if self.mut_file_cache().mut_map().remove(&k).is_some() {
225 trace!(sync_log::stream(), "{}: Dropped the cached handle on {:?} file {}, which a \
226 collection has replaced.", self.ozid(), typ, fnum);
227 }
228 }
229
230 fn get_file(
231 &mut self,
232 fnum: FileNum,
233 typ: &FileType,
234 )
235 -> Outcome<Arc<RwLock<File>>>
236 {
237 let k = FileCacheIndex { fnum, typ: typ.clone() };
238 // If the cache has the file and its not stale, return it.
239 let mut delete = false;
240 if let Some(FileCacheEntry{ t, file }) = self.ref_file_cache().ref_map().get(&k) {
241 if t.elapsed() < *self.ref_file_cache().expiry() {
242 return Ok(Arc::clone(file));
243 } else {
244 delete = true;
245 }
246 }
247 if delete {
248 self.mut_file_cache().mut_map().remove(&k);
249 }
250 self.open_file(fnum, typ)
251 }
252
253 fn open_file(
254 &mut self,
255 fnum: FileNum,
256 typ: &FileType,
257 )
258 -> Outcome<Arc<RwLock<File>>>
259 {
260 let (_, file) = res!(self.zdir().open_ozone_file(
261 fnum,
262 typ,
263 &FileAccess::Reading,
264 ));
265 let file_locked = Arc::new(RwLock::new(file));
266 let len = self.ref_file_cache().len();
267 if len < constant::MAX_CACHED_FILES {
268 self.mut_file_cache().insert(fnum, typ, file_locked.clone());
269 }
270 Ok(file_locked)
271 }
272
273 /// Retrieves a value from the database.
274 /// 2. Asks the key-selected cbot for the value or file location, sending a new responder resp_r2.
275 /// 3. The cbot accesses its cache.
276 /// 4. In the case where only the file location is available, the cbot sends the read request (including resp_r2) to the file-selected fbot.
277 /// 5. The fbot either responds immediately giving the rbot permission to read the file because it is not being garbage collected, incrementing the file state reader count, or else adds the request to a buffer so that permission can be granted later when garbage collection is complete.
278 /// 6. The rbot waits to receive either the value (via the cbot) or the file location (via the fbot) through resp_r2. If garbage collection has just been performed, there is a chance that the value was updated during the process. A flag in the returned value message allows the caller to decide if they want to try the read again, or accept the possibility of an old value.
279 /// 7. If necessary the rbot reads the file location.
280 /// 8. Once reading is complete, a finish message is sent to the read channel of the file's fbot.
281 /// 9. The fbot decrements the reader count for the file state.
282 /// 10.The rbot returns the value to the caller using resp_r1.
283 ///
284 fn read(
285 &mut self,
286 key: Key,
287 cbpind: usize,
288 )
289 -> Outcome<OzoneMsg<UIDL, UID, ENC, KH>>
290 {
291 let cind = key.index();
292
293 // A location names a record by an offset into one generation of its file, and a collection
294 // makes a new generation. A read that waited behind a collection, or was handed an offset
295 // the collection has just moved, can reach an offset that now holds another record, and
296 // with records of one size that record passes its own checksums: the read returned another
297 // key's value, or an older version of its own, about once in three collections
298 // (2026-09-23). So a record is returned only if it is the one the cache named, its key
299 // and stamp matching, and anything else -- a failed checksum included -- is retried with a
300 // fresh location from the cbot, which a collection re-anchors before it renames the file.
301 // The retry is bounded, so that a file a supersession burst keeps collecting cannot spin
302 // a reader for ever; each attempt holds the reader-count pin (see the `ReadFileRequest`
303 // arm in bot_file.rs), so the attempts converge as soon as the file settles.
304 for attempt in 0..constant::MAX_READ_ATTEMPTS {
305
306 // <2> Send read request to cbot.
307 let resp_r2 = Responder::new(Some(self.ozid()));
308 let cbots = res!(self.cbots());
309 let bot = res!(cbots.get_bot(cbpind));
310 res!(bot.send(OzoneMsg::ReadCache(key.clone(), resp_r2.clone())));
311
312 // <6> We receive either the value or the file location from the cbot or fbot.
313 let (floc, meta, postgc) = match resp_r2.recv_timeout(constant::BOT_REQUEST_TIMEOUT) {
314 Err(e) => return Err(err!(e,
315 "While waiting on value or location from cbot or fbot.";
316 IO, Channel, Read)),
317 Ok(OzoneMsg::ReadResult(readres)) => {
318 match readres {
319 // <10> Return result to caller via resp_r1.
320 ReadResult::None |
321 ReadResult::Deleted(_) =>
322 return Ok(OzoneMsg::Value(Value::new(
323 None,
324 cind,
325 false,
326 ))),
327 // <10> Return result to caller via resp_r1.
328 ReadResult::Value(val, meta) => {
329 // All values are wrapped inside a Daticle
330 let (dat, _) = res!(Dat::from_bytes(&val));
331 return Ok(OzoneMsg::Value(Value::new(
332 Some((dat, meta)),
333 cind,
334 false,
335 )));
336 },
337 ReadResult::Location(mloc, postgc) => {
338 let floc = *mloc.file_location();
339 let meta = mloc.meta_move();
340 (floc, meta, postgc)
341 },
342 }
343 },
344 Ok(msg) => return Err(err!(
345 "Unrecognised response from cbot to read request: {:?}", msg;
346 Bug, Invalid, Input)),
347 };
348
349 // <7> Read the value from the file location. If the value was cached, it has been
350 // returned above already.
351 let vlen = floc.val().len as usize;
352 let fnum = floc.file_number();
353 // A `postgc == true` offset was remapped through the collector's move map, so it belongs
354 // to the inode the rename put behind the path, and any handle this rbot still holds is
355 // the pre-rename inode -- unlinked but open -- on which the new offset lands off a record
356 // boundary. The `FileReplaced` notice that would drop it is queued behind this read on a
357 // channel the rbot cannot drain while parked inside `read`, so it arrives too late;
358 // dropping here reopens the live file. A `postgc == false` offset is NOT dropped on
359 // before its first read: it may be an un-remapped offset that is correct in the handle
360 // still cached, and forcing it onto a reopened inode would read a different record that
361 // could pass its own checksum. The drop for a `false` read happens only after its
362 // checksum fails, and then only paired with a fresh re-fetch below -- never a reopen of
363 // the same offset.
364 if postgc {
365 self.drop_cached_file(fnum, &FileType::Data);
366 }
367 let result = self.read_checked(floc, &key, &meta);
368
369 // <8> Advise the fbot that reading has finished so it can decrement the reader count it
370 // took when it handed back the location. The count has to come down whether the read
371 // worked or not: the fbot will not collect a file whose reader count is above zero, so
372 // a read that returned early -- a checksum mismatch is the one that arrives in bursts
373 // -- left a count that never came down, and with it a file that could never be
374 // collected again for the life of the process. This is sent on every attempt,
375 // matching the increment the fbot takes on every location it hands back.
376 let bots = res!(self.fbots());
377 let (bot, _) = bots.choose_bot(&ChooseBot::ByFile(fnum));
378 res!(bot.send(OzoneMsg::ReadFinished(fnum)));
379
380 match result {
381 Ok(rec) => {
382 // The value, without its checksum, follows the key.
383 let klen = try_into!(usize, floc.klen);
384 let end = klen + vlen - res!(self.api().schms.checksummer().len());
385 // All values are wrapped inside a Dat::BU64.
386 let (dat, _) = res!(Dat::from_bytes(&rec[klen..end]));
387 return Ok(OzoneMsg::Value(Value::new(
388 Some((dat, meta)),
389 cind,
390 postgc,
391 )));
392 },
393 Err(e) => {
394 // The record read is not the one named: its offset and the handle it was read
395 // through disagree on which generation of the file they belong to, or the
396 // offset was taken before a collection moved the record. The final attempt
397 // never returns bytes it could not confirm, so it fails loud.
398 if attempt + 1 == constant::MAX_READ_ATTEMPTS {
399 let what = fmt!("{}: While reading {:?} file {} (postgc {}, attempt {} of \
400 {}).", self.ozid(), FileType::Data, fnum, postgc,
401 attempt + 1, constant::MAX_READ_ATTEMPTS);
402 // Tagged with its cause where that is a record other than the one named,
403 // since `Error::tags` reports only the outermost error's tags.
404 return Err(if e.tags().contains(&ErrTag::Mismatch) {
405 err!(e, "{}", what; Data, Mismatch)
406 } else {
407 err!(e, "{}", what; IO, File, Read)
408 });
409 }
410 trace!(sync_log::stream(), "{}: Retrying a read of {:?} file {} at {} \
411 (postgc {}, attempt {}): {}", self.ozid(), FileType::Data, fnum,
412 floc.start, postgc, attempt + 1, e);
413 // Drop the handle and loop. The next attempt fetches a fresh location from
414 // the cbot, which the collection re-anchored to the current generation before
415 // it renamed the file (`cache_data_file`'s cbot cache update completes inside
416 // the call at collect_file step 6, before the rename at step 7), then reads
417 // THAT offset against the reopened live inode.
418 self.drop_cached_file(fnum, &FileType::Data);
419 },
420 }
421 }
422
423 // The loop returns on the first clean read and on a spent retry budget, so this is only
424 // reached if the budget is zero, which the constant forbids.
425 Err(err!(
426 "{}: Read of a value exhausted its retry budget without a result.", self.ozid();
427 Bug, IO, Read))
428 }
429
430 /// Reads the record at the given location, key and value in one read, and returns it only if
431 /// it is the record the cache named: it carries the key asked for with the stamp the cache has
432 /// for it, and its value's checksum holds. A key alone would pass an older version of the
433 /// same key. Split out of `read` only so that the caller can report the read finished to the
434 /// fbot on the way out, whichever way this goes.
435 fn read_checked(
436 &mut self,
437 floc: FileLocation,
438 key: &Key,
439 meta: &Meta<UIDL, UID>,
440 )
441 -> Outcome<Vec<u8>>
442 {
443 let rec = res!(self.read_from_file(floc));
444 let klen = try_into!(usize, floc.klen);
445 let csummer = self.api().schemes().checksummer().clone();
446 let csum_len = res!(csummer.len());
447 if !res!(StoredKey::<UIDL, UID>::holds(&rec[..klen], key.as_bytes(), meta, csum_len)) {
448 return Err(err!(
449 "{}: The record at {} in data file {} is not the one its cache bot named, {:?} \
450 stamped {:?}.", self.ozid(), floc.start, floc.file_number(), key, meta.time;
451 Data, Mismatch));
452 }
453 res!(csummer.verify(&rec[klen..]));
454 Ok(rec)
455 }
456
457 fn read_from_file(
458 &mut self,
459 floc: FileLocation,
460 )
461 -> Outcome<Vec<u8>>
462 {
463 let locked_file = res!(self.get_file(floc.file_number(), &FileType::Data));
464 let mut file_write = lock_write!(locked_file, // seek requires mutability
465 "{}: While trying to read from the cached data file number {}.",
466 self.ozid(), floc.file_number(),
467 );
468
469 match file_write.seek(SeekFrom::Start(floc.keyval().start)) {
470 Err(e) => return Err(err!(e,
471 "{}: attempt to move to position {} in data file {}.",
472 self.ozid(), floc.keyval().start, floc.file_number();
473 IO, File, Seek)),
474 Ok(actual_pos) => {
475 if actual_pos != floc.keyval().start {
476 return Err(err!(
477 "{}: attempt to move to position {} in data file {} \
478 but only moved to {}.",
479 self.ozid(), floc.keyval().start, floc.file_number(), actual_pos;
480 IO, File, Seek));
481 }
482 let mut v = vec![0; floc.keyval().len as usize];
483 let file_clone = res!((*file_write).try_clone());
484 let mut reader = BufReader::new(file_clone);
485 match reader.read(&mut v) {
486 Err(e) => {
487 return Err(err!(e,
488 "{}: attempt to read {} bytes from position {} in data file {}.",
489 self.ozid(), floc.keyval().len, floc.keyval().start, floc.file_number();
490 IO, File, Read));
491 },
492 Ok(actually_read) => {
493 if actually_read != floc.keyval().len as usize {
494 return Err(err!(
495 "{:?}: attempt to read {} bytes from position {} \
496 in data file {}, but only read {} bytes.",
497 self.ozid(), floc.keyval().len, floc.keyval().start, floc.file_number(),
498 actually_read;
499 IO, File, Read));
500 }
501 return Ok(v);
502 },
503 }
504 }
505 }
506 }
507}