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
| 1 | use 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 | |
| 35 | use oxedyne_fe2o3_iop_db::api::Meta; |
| 36 | use oxedyne_fe2o3_jdat::{ |
| 37 | prelude::*, |
| 38 | id::NumIdDat, |
| 39 | }; |
| 40 | use oxedyne_fe2o3_iop_hash::csum::Checksummer; |
| 41 | |
| 42 | use 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)] |
| 57 | pub 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 | |
| 69 | pub 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 | |
| 96 | impl< |
| 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 | |
| 109 | impl< |
| 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 | |
| 122 | impl< |
| 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 | |
| 174 | impl< |
| 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 | } |