Oregami
Repositories/oxedyne/fe2o3

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

1use 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
21use oxedyne_fe2o3_core::channels::Recv;
22use oxedyne_fe2o3_jdat::id::NumIdDat;
23
24use std::{
25 collections::BTreeMap,
26 fs::self,
27 sync::Arc,
28 time::Instant,
29};
30
31#[derive(Clone, Debug)]
32pub 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)]
39pub 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
70impl<
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
83impl<
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
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 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
159impl<
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}