oxedyne/fe2o3/fe2o3_o3db_sync/src/bots/worker/bot_file.rs
31.9 KiB, 156 runs
created by r1870400018:753, 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::{ |
| 4 | base::bot_deps::*, |
| 5 | worker::{ |
| 6 | bot_reader::ReadResult, |
| 7 | worker_deps::*, |
| 8 | }, |
| 9 | }, |
| 10 | file::{ |
| 11 | floc::{ |
| 12 | FileLocation, |
| 13 | FileNum, |
| 14 | }, |
| 15 | state::FileStateMap, |
| 16 | stored::RecordDigest, |
| 17 | }, |
| 18 | test::hooks, |
| 19 | }; |
| 20 | |
| 21 | use oxedyne_fe2o3_core::channels::Recv; |
| 22 | use oxedyne_fe2o3_jdat::id::NumIdDat; |
| 23 | |
| 24 | use std::{ |
| 25 | collections::BTreeMap, |
| 26 | fs::self, |
| 27 | sync::Arc, |
| 28 | time::Instant, |
| 29 | }; |
| 30 | |
| 31 | #[derive(Clone, Debug)] |
| 32 | pub enum GcControl { |
| 33 | On(bool), // switch gc on or off |
| 34 | Auto(bool), // set auto gc |
| 35 | Manual(FileNum), // file number |
| 36 | } |
| 37 | |
| 38 | #[derive(Debug)] |
| 39 | pub struct FileBot< |
| 40 | const UIDL: usize, |
| 41 | UID: NumIdDat<UIDL>, |
| 42 | ENC: Encrypter, |
| 43 | KH: Hasher, |
| 44 | PR: Hasher, |
| 45 | CS: Checksummer, |
| 46 | >{ |
| 47 | // Identity |
| 48 | wind: WorkerInd, |
| 49 | wtyp: WorkerType, |
| 50 | // Bot |
| 51 | sem: Semaphore, |
| 52 | errc: Arc<Mutex<usize>>, |
| 53 | log_stream_id: String, |
| 54 | // Config |
| 55 | zdir: ZoneDir, |
| 56 | // Comms |
| 57 | chan_in: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>, |
| 58 | // API |
| 59 | api: OzoneApi<UIDL, UID, ENC, KH, PR, CS>, |
| 60 | // State |
| 61 | active: bool, |
| 62 | auto_gc: bool, |
| 63 | gcbuf: BTreeMap<FileNum, Vec<OzoneMsg<UIDL, UID, ENC, KH>>>, |
| 64 | gc_on: bool, |
| 65 | inited: bool, |
| 66 | states: FileStateMap, |
| 67 | trep: Instant, |
| 68 | } |
| 69 | |
| 70 | impl< |
| 71 | const UIDL: usize, |
| 72 | UID: NumIdDat<UIDL> + 'static, |
| 73 | ENC: Encrypter + 'static, |
| 74 | KH: Hasher + 'static, |
| 75 | PR: Hasher, |
| 76 | CS: Checksummer, |
| 77 | > |
| 78 | WorkerBot<UIDL, UID, ENC, KH, PR, CS> for FileBot<UIDL, UID, ENC, KH, PR, CS> |
| 79 | { |
| 80 | workerbot_methods!(); |
| 81 | } |
| 82 | |
| 83 | impl< |
| 84 | const UIDL: usize, |
| 85 | UID: NumIdDat<UIDL> + 'static, |
| 86 | ENC: Encrypter + 'static, |
| 87 | KH: Hasher + 'static, |
| 88 | PR: Hasher, |
| 89 | CS: Checksummer, |
| 90 | > |
| 91 | OzoneBot<UIDL, UID, ENC, KH, PR, CS> for FileBot<UIDL, UID, ENC, KH, PR, CS> |
| 92 | { |
| 93 | ozonebot_methods!(); |
| 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 | Bot<{ BID_LEN }, Bid, OzoneMsg<UIDL, UID, ENC, KH>> for FileBot<UIDL, UID, ENC, KH, PR, CS> |
| 105 | { |
| 106 | bot_methods!(); |
| 107 | |
| 108 | fn go(&mut self) { |
| 109 | |
| 110 | sync_log::set_stream(self.log_stream_id()); |
| 111 | |
| 112 | if self.no_init() { return; } |
| 113 | self.now_listening(); |
| 114 | loop { |
| 115 | if self.wind().b() < self.cfg().num_bots_per_zone((&self).wtyp()) { |
| 116 | |
| 117 | if self.trep.elapsed() > self.cfg().zone_state_update_interval() { |
| 118 | // Automated state reporting. |
| 119 | self.trep = Instant::now(); |
| 120 | if let Some(zbot) = self.zbot() { |
| 121 | if let Err(e) = zbot.send( |
| 122 | OzoneMsg::ShardFileSize(self.wind().b(), self.states().get_size()) |
| 123 | ) { |
| 124 | self.result(&Err(err!(e, |
| 125 | "{}: Cannot send cache size update to zbot.", self.ozid(); |
| 126 | Channel, Write))); |
| 127 | } |
| 128 | } |
| 129 | } |
| 130 | |
| 131 | if self.listen().must_end() { break; } |
| 132 | |
| 133 | } else { |
| 134 | // This bot is to be terminated. Forward incoming messages to the remaining bots of |
| 135 | // this type. |
| 136 | } |
| 137 | } |
| 138 | } |
| 139 | |
| 140 | fn listen(&mut self) -> LoopBreak { |
| 141 | match self.chan_in().recv_timeout(self.cfg().zone_state_update_interval()) { |
| 142 | Recv::Result(Err(e)) => self.err_cannot_receive(err!(e, |
| 143 | "{}: Waiting for message.", self.ozid(); |
| 144 | IO, Channel)), |
| 145 | Recv::Result(Ok(msg)) => { |
| 146 | if let Some(msg) = self.listen_worker(msg) { |
| 147 | if self.listen_work(&msg) { |
| 148 | return self.listen_cmd(msg); |
| 149 | } |
| 150 | } |
| 151 | }, |
| 152 | Recv::Empty => (), |
| 153 | } |
| 154 | LoopBreak(false) |
| 155 | } |
| 156 | |
| 157 | } |
| 158 | |
| 159 | impl< |
| 160 | const UIDL: usize, |
| 161 | UID: NumIdDat<UIDL> + 'static, |
| 162 | ENC: Encrypter + 'static, |
| 163 | KH: Hasher + 'static, |
| 164 | PR: Hasher, |
| 165 | CS: Checksummer, |
| 166 | > |
| 167 | FileBot<UIDL, UID, ENC, KH, PR, CS> |
| 168 | { |
| 169 | pub fn new( |
| 170 | args: ZoneWorkerInitArgs<UIDL, UID, ENC, KH, PR, CS>, |
| 171 | ) |
| 172 | -> Self |
| 173 | { |
| 174 | Self { |
| 175 | // Identity |
| 176 | wind: args.wind, |
| 177 | wtyp: args.wtyp, |
| 178 | // Bot |
| 179 | sem: args.sem, |
| 180 | errc: Arc::new(Mutex::new(0)), |
| 181 | log_stream_id: args.log_stream_id, |
| 182 | // Config |
| 183 | zdir: ZoneDir::default(), |
| 184 | // Comms |
| 185 | chan_in: args.chan_in, |
| 186 | // API |
| 187 | api: args.api, |
| 188 | // State |
| 189 | active: false, |
| 190 | auto_gc: true, |
| 191 | gcbuf: BTreeMap::new(), |
| 192 | gc_on: false, |
| 193 | inited: false, |
| 194 | states: FileStateMap::default(), |
| 195 | trep: Instant::now(), |
| 196 | } |
| 197 | } |
| 198 | |
| 199 | fn states(&self) -> &FileStateMap { &self.states } |
| 200 | fn states_mut(&mut self) -> &mut FileStateMap { &mut self.states } |
| 201 | fn gc_buffer(&self) -> &BTreeMap<FileNum, Vec<OzoneMsg<UIDL, UID, ENC, KH>>> { &self.gcbuf } |
| 202 | fn gc_buffer_mut(&mut self) -> &mut BTreeMap<FileNum, Vec<OzoneMsg<UIDL, UID, ENC, KH>>> { &mut self.gcbuf } |
| 203 | |
| 204 | fn gc_auto_active(&self) -> bool { self.auto_gc } |
| 205 | |
| 206 | pub fn activate(mut self) -> Self { |
| 207 | self.active = true; |
| 208 | self |
| 209 | } |
| 210 | |
| 211 | pub fn listen_work( |
| 212 | &mut self, |
| 213 | msg: &OzoneMsg<UIDL, UID, ENC, KH>, |
| 214 | ) |
| 215 | -> bool |
| 216 | { |
| 217 | match msg { |
| 218 | OzoneMsg::ProcessGcBuffer(msgbox) => self.process_work(&*msgbox, true), |
| 219 | _ => self.process_work(msg, false), |
| 220 | } |
| 221 | |
| 222 | } |
| 223 | |
| 224 | /// Returns a flag indicating whether to keep listening. |
| 225 | pub fn process_work( |
| 226 | &mut self, |
| 227 | msg: &OzoneMsg<UIDL, UID, ENC, KH>, |
| 228 | processing_buffer: bool, |
| 229 | ) |
| 230 | -> bool |
| 231 | { |
| 232 | match msg { |
| 233 | // WRITE |
| 234 | OzoneMsg::ScheduleOld(floc, rid, from_id) => { |
| 235 | // [17] Schedule the old file location for deletion. |
| 236 | if !self.gc_active( |
| 237 | floc.file_number(), |
| 238 | &msg, |
| 239 | processing_buffer, |
| 240 | ) { |
| 241 | let result = self.schedule_deletion(floc, rid, from_id); |
| 242 | self.result(&result); |
| 243 | } |
| 244 | } |
| 245 | OzoneMsg::UpdateData { floc_new, ilen, floc_old_opt, from_id } => { |
| 246 | // [15] Add new data to the given live file state. |
| 247 | let result = self.update_data(floc_new, *ilen, floc_old_opt.as_ref(), from_id); |
| 248 | self.result(&result); |
| 249 | } |
| 250 | OzoneMsg::GcCompleted(fnum, new_fstat, size_dec) => { |
| 251 | match self.states_mut().get_state_mut(*fnum) { |
| 252 | Ok(fstat) => { |
| 253 | *fstat = new_fstat.clone(); |
| 254 | fstat.set_gc(false); |
| 255 | // Process buffer. |
| 256 | self.gc_active( |
| 257 | *fnum, |
| 258 | &OzoneMsg::None, |
| 259 | false, |
| 260 | ); |
| 261 | } |
| 262 | Err(e) => { |
| 263 | self.error(err!(e, |
| 264 | "{}: Cannot update file state for file {} after garbage collection \ |
| 265 | because it cannot be found in the file state map.", |
| 266 | self.ozid(), fnum; |
| 267 | Bug, Missing, Data)); |
| 268 | return false; |
| 269 | }, |
| 270 | }; |
| 271 | let result = self.states_mut().dec_size(*size_dec); |
| 272 | self.result(&result); |
| 273 | } |
| 274 | // READ |
| 275 | OzoneMsg::DumpFileStatesRequest(resp) => { |
| 276 | // TODO send back files currently being gc'd. |
| 277 | if let Err(e) = resp.send(OzoneMsg::DumpFileStatesResponse( |
| 278 | self.wind().clone(), |
| 279 | self.states().clone(), |
| 280 | )) { |
| 281 | self.err_cannot_send(err!(e, |
| 282 | "{}: Responding to {:?} with file states dump.", self.ozid(), resp.ozid(); |
| 283 | Data, IO, Channel)); |
| 284 | } |
| 285 | } |
| 286 | OzoneMsg::ReadFileRequest(fnum, kbyts, mloc, resp_r2) => { |
| 287 | if !self.gc_active( |
| 288 | *fnum, |
| 289 | &msg, |
| 290 | processing_buffer, |
| 291 | ) { |
| 292 | // <5> Increment the reader count and send the read result, an implicit |
| 293 | // permission to read, to the rbot. |
| 294 | let (result, msg2) = match self.states_mut().get_state_mut(*fnum) { |
| 295 | Ok(fstat) => { |
| 296 | let mut mloc2 = mloc.clone(); |
| 297 | let mut postgc = false; |
| 298 | // Only a file a collection has just moved records in has move entries, |
| 299 | // so the record is named only then. A read looks up its move and |
| 300 | // leaves the entry: it is kept for the record's supersession, which |
| 301 | // finding it gone flagged whatever now sits at the old offset as old. |
| 302 | if !fstat.no_pending_moves() { |
| 303 | match RecordDigest::new(kbyts, mloc.meta()) { |
| 304 | Ok(rid) => { |
| 305 | let dloc = mloc2.file_location().keyval(); |
| 306 | if let Some(new_start) = fstat.moved_to(&dloc, &rid) { |
| 307 | mloc2.new_start_position(new_start); |
| 308 | postgc = true; |
| 309 | } |
| 310 | }, |
| 311 | // The reader confirms the record it reads, so an unmapped |
| 312 | // location costs it a retry, never a wrong answer. |
| 313 | Err(e) => error!(sync_log::stream(), err!(e, |
| 314 | "Naming the record a read of file {} asks for.", fnum; |
| 315 | Data, Encode)), |
| 316 | } |
| 317 | } |
| 318 | // Increment the reader count whether or not the location was remapped. |
| 319 | // The count is the pin that keeps a file from being collected while a |
| 320 | // read of it is in flight (`schedule_deletion` will not start a |
| 321 | // collection unless `no_readers`), and a `postgc` read needs that pin |
| 322 | // as much as any other: the burst that carried this record can trip the |
| 323 | // trigger again at once, and a second collection renaming the file |
| 324 | // between here and the rbot's read would leave the just-handed-back |
| 325 | // offset pointing into a superseded inode -- a checksum failure the |
| 326 | // rbot's handle drop cannot repair, because the offset itself is stale. |
| 327 | // The rbot sends `ReadFinished` on every path, so this decrements |
| 328 | // cleanly; leaving it out here was also what drove the reader count |
| 329 | // below zero when a `postgc` read reported a finish it never counted. |
| 330 | ( |
| 331 | fstat.inc_readers(), |
| 332 | OzoneMsg::ReadResult(ReadResult::Location(mloc2, postgc)), |
| 333 | ) |
| 334 | }, |
| 335 | Err(e) => ( |
| 336 | Ok(()), |
| 337 | OzoneMsg::Error(err!(e, |
| 338 | "Read file request for file {}.", fnum; |
| 339 | Bug, Missing, Data)), |
| 340 | ), |
| 341 | }; |
| 342 | self.result(&result); |
| 343 | // <6> Send file location back to rbot, representing permission to perform a read. |
| 344 | self.respond(Ok(msg2), resp_r2); |
| 345 | } |
| 346 | } |
| 347 | OzoneMsg::ReadFinished(fnum) => { |
| 348 | if !self.gc_active( |
| 349 | *fnum, |
| 350 | &msg, |
| 351 | processing_buffer, |
| 352 | ) { |
| 353 | // <9> Decrement the reader count now that a read has completed. The last read |
| 354 | // finishing does not look at the file for collection (see `maybe_collect`). |
| 355 | let result = match self.states_mut().get_state_mut(*fnum) { |
| 356 | Ok(fstat) => { |
| 357 | let result = fstat.dec_readers(); |
| 358 | result |
| 359 | }, |
| 360 | Err(_) => { |
| 361 | warn!(sync_log::stream(), |
| 362 | "A read completion for file {} has been received, but the file state \ |
| 363 | no longer exists, ignoring.", fnum, |
| 364 | ); |
| 365 | Ok(()) |
| 366 | }, |
| 367 | }; |
| 368 | self.result(&result); |
| 369 | } |
| 370 | } |
| 371 | _ => return true, |
| 372 | } |
| 373 | false |
| 374 | } |
| 375 | |
| 376 | pub fn listen_cmd(&mut self, msg: OzoneMsg<UIDL, UID, ENC, KH>) -> LoopBreak { |
| 377 | match msg { |
| 378 | // COMMAND |
| 379 | OzoneMsg::NewFileStates(shard) => { |
| 380 | self.states = shard; |
| 381 | } |
| 382 | //Ok(Config(cfg, tik)) => { |
| 383 | // self.cfg = cfg; |
| 384 | // let msg = OzoneMsg::ConfigConfirm(self.ozid().clone(), tik); |
| 385 | // if let Err(e) = self.sup().send(msg.clone()) { |
| 386 | // self.err_cannot_send(fmt!("{:?}", msg)); |
| 387 | // } |
| 388 | //}, |
| 389 | OzoneMsg::CloseOldLiveFileState { |
| 390 | fnum_old, |
| 391 | fnum_new, |
| 392 | new_dat_size, |
| 393 | new_ind_size, |
| 394 | resp, |
| 395 | } => { |
| 396 | // The writer rolling over waits on this, and would otherwise wait out its deadline |
| 397 | // on a failure only logged here. |
| 398 | let caller = resp.clone(); |
| 399 | if let Err(e) = self.close_old_live_file_state( |
| 400 | fnum_old, |
| 401 | fnum_new, |
| 402 | new_dat_size, |
| 403 | new_ind_size, |
| 404 | resp, |
| 405 | ) { |
| 406 | self.error(e.clone()); |
| 407 | self.respond(Err(e), &caller); |
| 408 | } |
| 409 | }, |
| 410 | OzoneMsg::OpenNewLiveFileState { |
| 411 | fnum_new, |
| 412 | new_dat_size, |
| 413 | new_ind_size, |
| 414 | resp, |
| 415 | } => { |
| 416 | let result = self.open_new_live_file_state( |
| 417 | fnum_new, |
| 418 | new_dat_size, |
| 419 | new_ind_size, |
| 420 | ); |
| 421 | self.respond(result, &resp); |
| 422 | }, |
| 423 | OzoneMsg::GcControl(gc_ctrl, _) => { |
| 424 | match gc_ctrl { |
| 425 | GcControl::On(state) => self.gc_on = state, |
| 426 | GcControl::Auto(state) => self.auto_gc = state, |
| 427 | GcControl::Manual(_) => warn!(sync_log::stream(), "Manual gc not yet implemented."), |
| 428 | } |
| 429 | } |
| 430 | _ => return self.listen_more(msg), |
| 431 | } |
| 432 | LoopBreak(false) |
| 433 | } |
| 434 | |
| 435 | /// Capture read and write messages related to a file in a buffer while its garbage is being |
| 436 | /// collected. There are two indications that a file is in the process of garbage collection; |
| 437 | /// a flag `gc_active` in its file state, and a entry in the `gc_buffer` map keyed to the file |
| 438 | /// number. undergoing garbage collection. If so, the incoming message is appended to the |
| 439 | /// entry. Otherwise |
| 440 | fn gc_active( |
| 441 | &mut self, |
| 442 | fnum: FileNum, |
| 443 | msg: &OzoneMsg<UIDL, UID, ENC, KH>, |
| 444 | processing: bool, |
| 445 | ) |
| 446 | -> bool |
| 447 | { |
| 448 | if processing { return false; } |
| 449 | |
| 450 | // Borrow check work around: read gc flag first. |
| 451 | let flag = match self.states().get_state(fnum) { |
| 452 | Ok(fstat) => Some(fstat.gc_active()), |
| 453 | Err(_) => None, |
| 454 | }; |
| 455 | let mut process_buffer = false; |
| 456 | let gc_active = match self.gc_buffer_mut().get_mut(&fnum) { |
| 457 | Some(gcbuf) => { |
| 458 | match flag { |
| 459 | Some(gc_active) => { |
| 460 | if gc_active { |
| 461 | // Buffer the incoming message because the file garbage is being |
| 462 | // collected. |
| 463 | gcbuf.push(msg.clone()); |
| 464 | } else { |
| 465 | // Process the buffer messages because garbage collection has finished. |
| 466 | // This needs to be pushed outside this match scope because work_msg |
| 467 | // requires another mutable borrow of self. |
| 468 | process_buffer = true; |
| 469 | } |
| 470 | }, |
| 471 | None => self.error(err!( |
| 472 | "{}: A gc buffer exists for file {} but no file state exists. \ |
| 473 | Could not buffer the received {:?}.", self.ozid(), fnum, msg; |
| 474 | Bug, Missing, Data)), |
| 475 | } |
| 476 | true |
| 477 | }, |
| 478 | None => false, // No need to do anything. |
| 479 | }; |
| 480 | if process_buffer { |
| 481 | // Take the buffer out before replaying it. Replaying a message can start a |
| 482 | // fresh collection of the same file, which creates a new, empty buffer; removing |
| 483 | // the entry after the replay would take that new buffer with it, and every |
| 484 | // later message for the file would then be applied to a data file in the middle |
| 485 | // of being transcribed, with a second collector free to start on it as well. |
| 486 | if let Some(gcbuf) = self.gc_buffer_mut().remove(&fnum) { |
| 487 | for msg in gcbuf { |
| 488 | // Once a fresh collection has started, the rest of the replay belongs to |
| 489 | // its buffer rather than to the file being transcribed. |
| 490 | match self.gc_buffer_mut().get_mut(&fnum) { |
| 491 | Some(newbuf) => newbuf.push(msg), |
| 492 | None => { self.listen_work(&OzoneMsg::ProcessGcBuffer(Box::new(msg))); }, |
| 493 | } |
| 494 | } |
| 495 | } |
| 496 | } |
| 497 | gc_active |
| 498 | } |
| 499 | |
| 500 | fn schedule_deletion( |
| 501 | &mut self, |
| 502 | floc: &FileLocation, |
| 503 | rid: &RecordDigest, |
| 504 | from: &OzoneBotId, |
| 505 | ) |
| 506 | -> Outcome<()> |
| 507 | { |
| 508 | let self_id = self.ozid().clone(); |
| 509 | |
| 510 | // [17] Schedule the old file location for deletion. |
| 511 | // [17.1] Update the current data in the file state data map to old. |
| 512 | match self.states_mut().get_state_mut(floc.file_number()) { |
| 513 | Err(e) => return Err(err!(e, |
| 514 | "{:?}: Request from {:?} to delete {:?}.", self_id, from, floc; |
| 515 | Bug, Missing, Data)), |
| 516 | Ok(fstat) => { |
| 517 | // Perform mapping to new start position, resulting from scheduling messages which |
| 518 | // have backed up during previous garbage collection. Only this record's move is |
| 519 | // taken: another's at the same offset is waiting for its own supersession. |
| 520 | let mut floc2 = floc.clone(); |
| 521 | if let Some(new_start) = fstat.map_and_remove(&floc2.keyval(), rid) { |
| 522 | floc2.start = new_start; |
| 523 | } |
| 524 | |
| 525 | // Register data as old. |
| 526 | if let Err(e) = fstat.register_old(&floc2.keyval()) { |
| 527 | return Err(err!(e, "{:?}: file {}.", self_id, floc2.file_number(); Data)); |
| 528 | } |
| 529 | }, |
| 530 | } |
| 531 | |
| 532 | // [17.2] Check whether garbage collection should be triggered for the file. |
| 533 | self.maybe_collect(floc.file_number()) |
| 534 | } |
| 535 | |
| 536 | /// Starts a collection of the file if it is eligible now. A file deferred on any input to |
| 537 | /// eligibility is looked at again only when something calls this: a supersession, a seal, a |
| 538 | /// record landing in a sealed file. Until 2026-09-23 only a supersession did, and a file that |
| 539 | /// crossed the trigger while it was live, or before its last writes had drained, kept its |
| 540 | /// garbage until a later supersession happened to land in it: for a file whose remaining |
| 541 | /// records are never superseded, for ever. |
| 542 | /// |
| 543 | /// Two changes that can make a file eligible deliberately do not call this. The last read |
| 544 | /// finishing would start a collection in the middle of a burst of reads of the file, a chunked |
| 545 | /// value's, and a read queued behind the collection is replayed at its old offset with |
| 546 | /// nothing to remap it, since the collection re-anchored its cache entry and dropped the move. |
| 547 | /// With records of one size that offset holds another valid record, which the reader returned |
| 548 | /// until it confirmed key and stamp (2026-09-23): starting collections there failed the first |
| 549 | /// read of a same-key churn's value after a restart in 5 to 8 runs of 12. The reader now |
| 550 | /// retries such a read, but the trigger stays withdrawn until that is measured. Switching |
| 551 | /// collection on would hand every file a start-up load had found garbage in to the collectors |
| 552 | /// at once, and a read of a file waiting its turn waits with it. |
| 553 | fn maybe_collect(&mut self, fnum: FileNum) -> Outcome<()> { |
| 554 | let self_id = self.ozid().clone(); |
| 555 | if self.gc_on { |
| 556 | let mut gc_activated = false; |
| 557 | match self.states().get_state(fnum) { |
| 558 | Ok(fstat) => { |
| 559 | let oldvals = fstat.get_old_sum() as f64; |
| 560 | let datfilemax = self.cfg().data_file_max_bytes as f64; |
| 561 | let trigger = constant::OLD_DATA_PERCENT_GC_TRIGGER; |
| 562 | let oldfrac = 100.0 * (oldvals / datfilemax); |
| 563 | let eligible = |
| 564 | ((oldfrac > trigger) || fstat.is_all_data_old()) && |
| 565 | fstat.no_pending_moves() && |
| 566 | !fstat.is_live() && |
| 567 | !fstat.gc_active() && // Never set a second collector on the same file. |
| 568 | fstat.no_readers() && |
| 569 | self.gc_auto_active(); |
| 570 | // A sealed file may still have writes draining. A record's bytes reach the |
| 571 | // data file in `WriterBot::write` before its accounting does: the accounting |
| 572 | // travels writer -> cbot -> fbot as an `UpdateData`, while the seal travels |
| 573 | // writer -> fbot directly and can overtake it. So at the moment a sibling |
| 574 | // record's supersession trips this trigger, `dat_size` and `dmap` can still |
| 575 | // lag the physical file by the records whose `UpdateData` is in flight. |
| 576 | // Collecting then is doubly wrong: the snapshot's accounting disagrees with |
| 577 | // the file it transcribes (the `old_sum != old_size - new_size` abort), and a |
| 578 | // later in-flight `UpdateData` would insert a now-stale position into the |
| 579 | // rewritten file. Defer until the on-disk size equals the accounted size -- |
| 580 | // the point at which every write has drained. This is rare on the unchunked |
| 581 | // path (old bytes accrue a small record at a time, long after the file has |
| 582 | // sealed and drained) but routine for a chunked value, whose single overwrite |
| 583 | // both rolls a file mid-burst and supersedes a whole value's worth of records |
| 584 | // at once. The record whose landing completes the drain brings the file |
| 585 | // back here (`update_data`), so deferral only delays. |
| 586 | let drained = if eligible { |
| 587 | let mut dat_path = self.zdir().dir.clone(); |
| 588 | dat_path.push(ZoneDir::relative_file_path(&FileType::Data, fnum)); |
| 589 | match std::fs::metadata(&dat_path) { |
| 590 | Ok(m) => m.len() == fstat.get_data_file_size() as u64, |
| 591 | // Cannot confirm the file has drained, so do not collect it yet. |
| 592 | Err(_) => false, |
| 593 | } |
| 594 | } else { |
| 595 | false |
| 596 | }; |
| 597 | // Case (c) from `register_old`: once the file has fully drained, every parked |
| 598 | // supersession must have found its record. Any left over refers to a record |
| 599 | // that is not on disk -- a genuine accounting fault, not the write-path race -- |
| 600 | // so fail loudly rather than collect a file whose accounting is inconsistent. |
| 601 | if eligible && drained && !fstat.pending_old_empty() { |
| 602 | return Err(err!( |
| 603 | "{:?}: File {} has drained (on-disk size equals accounted size) yet \ |
| 604 | {} superseded record(s) were never inserted: {:?}. This is an \ |
| 605 | accounting inconsistency, not the write-path race.", |
| 606 | self_id, fnum, fstat.pending_old().len(), fstat.pending_old(); |
| 607 | Bug, Missing, Data)); |
| 608 | } |
| 609 | if eligible && drained { |
| 610 | // [18.1] Select a gbot to collect the garbage. |
| 611 | debug!(sync_log::stream(), "{}: Automated garbage collection for file {}", self_id, fnum); |
| 612 | let bots = res!(self.igbots()); |
| 613 | let (bot, _) = bots.choose_bot(&ChooseBot::Randomly); |
| 614 | if fstat.is_all_old() { |
| 615 | // [18.2] Just delete the data file and its index file if it has no current data. |
| 616 | for ftyp in [FileType::Data, FileType::Index] { |
| 617 | let mut path = self.zdir().dir.clone(); |
| 618 | path.push(ZoneDir::relative_file_path(&ftyp, fnum)); |
| 619 | if path.is_file() { |
| 620 | res!(fs::remove_file(path)); |
| 621 | } |
| 622 | } |
| 623 | debug!(sync_log::stream(), |
| 624 | "{}: All the data in file {} is old, the file has therefore been deleted.", |
| 625 | self_id, fnum, |
| 626 | ); |
| 627 | } else { |
| 628 | res!(bot.send(OzoneMsg::CollectGarbage { |
| 629 | fnum, |
| 630 | fstat: fstat.clone(), |
| 631 | fbot_index: self.wind().b(), |
| 632 | })); |
| 633 | // [18.3] Create a gc buffer entry. |
| 634 | self.gc_buffer_mut().insert(fnum, Vec::new()); |
| 635 | //fstat.set_gc(true); // [#] Moved out of scope due to borrow checker |
| 636 | gc_activated = true; |
| 637 | } |
| 638 | } |
| 639 | }, |
| 640 | Err(e) => return Err(err!(e, |
| 641 | "{:?}: Evaluating file {} for collection, which has no state.", self_id, fnum; |
| 642 | Bug, Missing, Data)), |
| 643 | } |
| 644 | // [#] Moved out to here due to borrow checker. |
| 645 | if gc_activated { |
| 646 | match self.states_mut().get_state_mut(fnum) { |
| 647 | Ok(fstat) => fstat.set_gc(true), |
| 648 | _ => (), // unreachable |
| 649 | } |
| 650 | } |
| 651 | } |
| 652 | Ok(()) |
| 653 | } |
| 654 | |
| 655 | fn update_data( |
| 656 | &mut self, |
| 657 | floc_new: &FileLocation, |
| 658 | ilen: usize, |
| 659 | floc_old_opt: Option<&(FileLocation, RecordDigest)>, |
| 660 | from: &OzoneBotId, |
| 661 | ) |
| 662 | -> Outcome<()> |
| 663 | { |
| 664 | // [15] Add new data to the given live file state. |
| 665 | match self.states_mut().insert_new(floc_new, ilen) { |
| 666 | Err(e) => return Err(err!(e, |
| 667 | "{:?}: Request from {:?} to insert {:?}.", self.ozid(), from, floc_new; |
| 668 | Data)), |
| 669 | Ok(()) => (), |
| 670 | }; |
| 671 | |
| 672 | // [16] Advise the appropriate fbot to schedule the old data for deletion in its file state |
| 673 | // data map. |
| 674 | if let Some((floc_old, rid)) = floc_old_opt { |
| 675 | let bots = res!(self.fbots()); |
| 676 | let (bot, b) = bots.choose_bot(&ChooseBot::ByFile(floc_old.file_number())); |
| 677 | if *b == self.wind().b() { |
| 678 | // This could be itself, and then the supersession passes the same guard as a |
| 679 | // `ScheduleOld` message. Applied directly to a file being collected, it went |
| 680 | // into the state that the collection's result then replaced: the record was |
| 681 | // carried into the rewritten file as current, its move entry was never cleared, |
| 682 | // and a file with a move entry is never collected again. With two file bots to a |
| 683 | // zone, half of all supersessions come this way; the online sweep lost 50 to 65 |
| 684 | // of its 224 to it (2026-09-23). |
| 685 | let msg = OzoneMsg::ScheduleOld(*floc_old, *rid, from.clone()); |
| 686 | if !self.gc_active(floc_old.file_number(), &msg, false) { |
| 687 | res!(self.schedule_deletion( |
| 688 | floc_old, |
| 689 | rid, |
| 690 | from, |
| 691 | )); |
| 692 | } |
| 693 | } else { |
| 694 | // Or another fbot. |
| 695 | hooks::forward_delay(); |
| 696 | res!(bot.send(OzoneMsg::ScheduleOld( |
| 697 | *floc_old, |
| 698 | *rid, |
| 699 | from.clone(), |
| 700 | ))); |
| 701 | } |
| 702 | } |
| 703 | |
| 704 | // A sealed file's last records land after its seal, and until they do it has not |
| 705 | // drained and cannot be collected. |
| 706 | self.maybe_collect(floc_new.file_number()) |
| 707 | } |
| 708 | |
| 709 | fn close_old_live_file_state( |
| 710 | &mut self, |
| 711 | fnum_old: FileNum, |
| 712 | fnum_new: FileNum, |
| 713 | new_dat_size: u64, |
| 714 | new_ind_size: u64, |
| 715 | resp: Responder<UIDL, UID, ENC, KH>, |
| 716 | ) |
| 717 | -> Outcome<()> |
| 718 | { |
| 719 | if fnum_old > 0 { |
| 720 | // [6] Update the state of the previous live file. |
| 721 | match self.states_mut().get_state_mut(fnum_old) { |
| 722 | Ok(fstat) => fstat.set_live(false), |
| 723 | Err(e) => return Err(err!(e, |
| 724 | "{}: Request to close old live file {} state.", self.ozid(), fnum_old; |
| 725 | Bug, Missing, Data)), |
| 726 | } |
| 727 | } |
| 728 | |
| 729 | // [7] Advise the appropriate fbot to add the new live file to the file state map. |
| 730 | let bots = res!(self.fbots()); |
| 731 | let (bot, b) = bots.choose_bot(&ChooseBot::ByFile(fnum_new)); |
| 732 | if *b == self.wind().b() { |
| 733 | // This could be itself... |
| 734 | res!(self.open_new_live_file_state( |
| 735 | fnum_new, |
| 736 | new_dat_size, |
| 737 | new_ind_size, |
| 738 | )); |
| 739 | self.respond(Ok(OzoneMsg::Ok), &resp); |
| 740 | } else { |
| 741 | // Or another fbot. |
| 742 | res!(bot.send(OzoneMsg::OpenNewLiveFileState{ |
| 743 | fnum_new, |
| 744 | new_dat_size, |
| 745 | new_ind_size, |
| 746 | resp, |
| 747 | })); |
| 748 | } |
| 749 | |
| 750 | // Supersessions that reached the file while it was live could not start a collection. |
| 751 | // The writer has its answer, so a failure here is only logged. |
| 752 | if fnum_old > 0 { |
| 753 | let result = self.maybe_collect(fnum_old); |
| 754 | self.result(&result); |
| 755 | } |
| 756 | |
| 757 | Ok(()) |
| 758 | } |
| 759 | |
| 760 | fn open_new_live_file_state( |
| 761 | &mut self, |
| 762 | fnum: FileNum, |
| 763 | dat_size: u64, |
| 764 | ind_size: u64, |
| 765 | ) |
| 766 | -> Outcome<OzoneMsg<UIDL, UID, ENC, KH>> |
| 767 | { |
| 768 | // [8] Update the state of the new live file. |
| 769 | self.states_mut().new_live_file( |
| 770 | fnum, |
| 771 | dat_size, |
| 772 | ind_size, |
| 773 | ); |
| 774 | Ok(OzoneMsg::Ok) |
| 775 | } |
| 776 | } |