Oregami
Repositories/oxedyne/fe2o3

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

57.6 KiB, 208 runs

created by r1870400018:755, 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 choose::ChooseCache,
10 },
11 file::{
12 core::FileAccess,
13 floc::{
14 DataLocation,
15 FileNum,
16 StoredFileLocation,
17 },
18 state::{
19 DataState,
20 FileState,
21 },
22 stored::{
23 RecordDigest,
24 StoredIndex,
25 StoredKey,
26 StoredValue,
27 },
28 },
29 test::hooks,
30};
31
32use oxedyne_fe2o3_iop_db::api::Meta;
33use oxedyne_fe2o3_jdat::id::NumIdDat;
34
35use std::{
36 fs::{
37 self,
38 File,
39 OpenOptions,
40 },
41 io::{
42 BufReader,
43 BufWriter,
44 Seek,
45 SeekFrom,
46 Read,
47 Write,
48 },
49 sync::Arc,
50};
51
52/// `InitGarbageBot`s have two functions:
53/// 1. Initialisation where they are asked to read files and fill the caches.
54/// 2. Garbage collection where they subsequently regularly rewrite data files to remove stale
55/// data.
56pub struct InitGarbageBot<
57 const UIDL: usize,
58 UID: NumIdDat<UIDL>,
59 ENC: Encrypter,
60 KH: Hasher,
61 PR: Hasher,
62 CS: Checksummer,
63>{
64 // Identity
65 wind: WorkerInd,
66 wtyp: WorkerType,
67 // Bot
68 sem: Semaphore,
69 errc: Arc<Mutex<usize>>,
70 log_stream_id: String,
71 // Config
72 zdir: ZoneDir,
73 // Comms
74 chan_in: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>,
75 // API
76 api: OzoneApi<UIDL, UID, ENC, KH, PR, CS>,
77 // State
78 inited: bool,
79}
80
81impl<
82 const UIDL: usize,
83 UID: NumIdDat<UIDL> + 'static,
84 ENC: Encrypter + 'static,
85 KH: Hasher + 'static,
86 PR: Hasher,
87 CS: Checksummer,
88>
89 WorkerBot<UIDL, UID, ENC, KH, PR, CS> for InitGarbageBot<UIDL, UID, ENC, KH, PR, CS>
90{
91 workerbot_methods!();
92}
93
94impl<
95 const UIDL: usize,
96 UID: NumIdDat<UIDL> + 'static,
97 ENC: Encrypter + 'static,
98 KH: Hasher + 'static,
99 PR: Hasher,
100 CS: Checksummer,
101>
102 OzoneBot<UIDL, UID, ENC, KH, PR, CS> for InitGarbageBot<UIDL, UID, ENC, KH, PR, CS>
103{
104 ozonebot_methods!();
105}
106
107impl<
108 const UIDL: usize,
109 UID: NumIdDat<UIDL> + 'static,
110 ENC: Encrypter + 'static,
111 KH: Hasher + 'static,
112 PR: Hasher,
113 CS: Checksummer,
114>
115 Bot<{ BID_LEN }, Bid, OzoneMsg<UIDL, UID, ENC, KH>> for InitGarbageBot<UIDL, UID, ENC, KH, PR, CS>
116{
117 bot_methods!();
118
119 fn go(&mut self) {
120
121 sync_log::set_stream(self.log_stream_id());
122
123 if self.no_init() { return; }
124 self.now_listening();
125 loop {
126 if self.listen().must_end() { break; }
127 }
128 }
129
130 fn listen(&mut self) -> LoopBreak {
131 match self.chan_in().recv() {
132 Err(e) => self.err_cannot_receive(err!(e,
133 "{}: Waiting for message.", self.ozid();
134 IO, Channel)),
135 Ok(msg) => {
136 if let Some(msg) = self.listen_worker(msg) {
137 match msg {
138 // Init
139 OzoneMsg::CacheDataFile {
140 fnum,
141 dat_file_size,
142 resp,
143 } => {
144 let result = self.cache_file(
145 fnum,
146 &FileType::Data,
147 dat_file_size,
148 0,
149 );
150 self.respond(result, &resp);
151 },
152 OzoneMsg::CacheIndexFile {
153 fnum,
154 dat_file_size,
155 ind_file_size,
156 resp,
157 } => {
158 let result = self.cache_file(
159 fnum,
160 &FileType::Index,
161 dat_file_size,
162 ind_file_size,
163 );
164 self.respond(result, &resp);
165 },
166 // Garbage collection
167 OzoneMsg::CollectGarbage {
168 fnum,
169 fstat,
170 fbot_index,
171 } => {
172 let result = self.collect_garbage(
173 fnum,
174 fstat,
175 fbot_index,
176 );
177 self.result(&result);
178 },
179 _ => return self.listen_more(msg),
180 }
181 }
182 },
183 }
184 LoopBreak(false)
185 }
186}
187
188impl<
189 const UIDL: usize,
190 UID: NumIdDat<UIDL> + 'static,
191 ENC: Encrypter + 'static,
192 KH: Hasher + 'static,
193 PR: Hasher,
194 CS: Checksummer,
195>
196 InitGarbageBot<UIDL, UID, ENC, KH, PR, CS>
197{
198 pub fn new(
199 args: ZoneWorkerInitArgs<UIDL, UID, ENC, KH, PR, CS>,
200 )
201 -> Self
202 {
203 Self {
204 // Identity
205 wind: args.wind,
206 wtyp: args.wtyp,
207 // Bot
208 sem: args.sem,
209 errc: Arc::new(Mutex::new(0)),
210 log_stream_id: args.log_stream_id,
211 // Config
212 zdir: ZoneDir::default(),
213 // Comms
214 chan_in: args.chan_in,
215 // API
216 api: args.api,
217 // State
218 inited: false,
219 }
220 }
221
222 fn cache_file(
223 &mut self,
224 fnum: FileNum,
225 typ: &FileType,
226 dat_size: usize,
227 ind_size: usize,
228 )
229 -> Outcome<OzoneMsg<UIDL, UID, ENC, KH>>
230 {
231 let (_, file) = res!(self.zdir().open_ozone_file(
232 fnum,
233 typ,
234 &FileAccess::Reading,
235 ));
236 let meta = res!(file.metadata());
237 let reader = BufReader::new(file);
238 match typ {
239 FileType::Index => {
240 if meta.len() == 0 {
241 // Nothing to cache from, so read the data file and write the index back
242 // from what it holds. Until that happens every scan of this zone is
243 // short by everything this file holds, and says so to nobody.
244 warn!(sync_log::stream(),
245 "{}: Index file {} is empty beside a data file of {} bytes; \
246 rebuilding the index by reading the data file. Until this \
247 completes, a scan of this zone under-reports what it holds.",
248 self.ozid(), fnum, dat_size);
249 let (_, file) = res!(self.zdir().open_ozone_file(
250 fnum,
251 &FileType::Data,
252 &FileAccess::Reading,
253 ));
254 let reader = BufReader::new(file);
255 res!(self.init_cache_data_file(
256 reader,
257 fnum,
258 dat_size,
259 ));
260 } else {
261 match self.init_cache_index_file(
262 reader,
263 fnum,
264 dat_size,
265 ind_size,
266 ) {
267 Err(e) => {
268 // The index file cannot be read, so it is of no use to a scan
269 // either. Read the data file and write the index back from it.
270 warn!(sync_log::stream(),
271 "{}: Index file {} could not be read, so it is rebuilt by \
272 reading data file {}; until this completes, a scan of this \
273 zone under-reports what it holds. Caused by {}.",
274 self.ozid(), fnum, fnum, e);
275 let (_, file) = res!(self.zdir().open_ozone_file(
276 fnum,
277 &FileType::Data,
278 &FileAccess::Reading,
279 ));
280 let reader = BufReader::new(file);
281 res!(self.init_cache_data_file(
282 reader,
283 fnum,
284 dat_size,
285 ));
286 },
287 Ok(()) => (),
288 }
289 }
290 },
291 FileType::Data => res!(self.init_cache_data_file(
292 reader,
293 fnum,
294 dat_size,
295 )),
296 }
297 Ok(OzoneMsg::Ok)
298 }
299
300 /// Use the index file to update the cache with data file value locations, which should be
301 /// quicker than scanning the data file itself.
302 #[allow(unused_assignments, unused_variables)]
303 fn init_cache_index_file(
304 &mut self,
305 mut reader: BufReader<File>,
306 fnum: FileNum,
307 dat_size1: usize,
308 ind_size: usize,
309 )
310 -> Outcome<()>
311 {
312 let mut pos = 0;
313 let mut count = 0;
314 let mut dat_size2: u64 = 0;
315 let typ = FileType::Index;
316
317 loop {
318 // 1. Load the key Daticle bytes and while we're at it, compare the checksum.
319 let (key, meta, chash) = match StoredKey::load(
320 &mut reader,
321 self.api().schemes().checksummer().clone(),
322 ) {
323 Err(e) => return Err(err!(e,
324 "{}: While reading from position {} in {:?} file {}.",
325 self.ozid(), pos, typ, fnum;
326 IO, File, Read)),
327 Ok(None) => break, // We're done.
328 Ok(Some((skey, _, n))) => {
329 count += 1;
330 pos += n;
331 let meta = skey.meta().clone();
332 let chash = skey.ref_chash().clone();
333 (skey.into_key(), meta, chash)
334 },
335 };
336 // 2. Read the StoredIndex.
337 match StoredIndex::read(
338 &mut reader,
339 fnum,
340 self.api().schemes().checksummer().clone(),
341 ) {
342 Err(e) => return Err(err!(e,
343 "{}: While reading from position {} in {:?} file {}.",
344 self.ozid(), pos, typ, fnum;
345 IO, File, Read)),
346 Ok((None, _)) => return Err(err!(
347 "{}: Missing StoredIndex at end of {:?} file {}.",
348 self.ozid(), typ, fnum;
349 Missing)),
350 Ok((Some(sindex), n)) => {
351 count += 1;
352 pos += n;
353 dat_size2 += sindex.keyval_len();
354 // 3. Insert the key and location into the bot cache, informing an fbot about
355 // new data and old data that can be scheduled for garbage collection. The
356 // bot we advise actually performs any garbage collection, so instead of
357 // choosing randomly, we allocate each bot to an exclusive fraction of files
358 // based on their number.
359 let cind = key.index();
360 let kbyts = key.into_bytes();
361 let chash = res!(<alias::ChooseHash>::try_from(
362 &chash[..constant::CACHE_HASH_BYTES]));
363 let cbwind = ChooseCache::<PR>::choose_cbot_select(
364 alias::ChooseHashUint::from_be_bytes(chash),
365 self.cfg().num_zones,
366 self.cfg().num_cbots_per_zone,
367 );
368 let cbots = res!(self.cbots());
369 let bot = res!(cbots.get_bot(**cbwind.bpind()));
370 res!(bot.send(
371 OzoneMsg::Insert(
372 kbyts,
373 None,
374 cind,
375 sindex.ref_file_location().clone(),
376 sindex.ref_stored_file_location().buf.len(),
377 meta,
378 Responder::none(Some(self.ozid())),
379 None,
380 )
381 ));
382 },
383 }
384 }
385
386 // 8. Do size check.
387 if pos != ind_size {
388 return Err(err!(
389 "{}: After initial caching of data file {} using the index file, the \
390 index file data count came to {} bytes, but the originally surveyed file \
391 size was {}.", self.ozid(), fnum, pos, ind_size;
392 Mismatch, Data));
393 }
394 if dat_size1 != res!(usize::try_from(dat_size2)) {
395 return Err(err!(
396 "{}: After initial caching of data file {} using the index file, the \
397 file data count came to {} bytes, but the originally surveyed file \
398 size was {}.", self.ozid(), fnum, dat_size2, dat_size1;
399 Mismatch, Data));
400 }
401
402 Ok(())
403 }
404
405 /// When a valid index file is not available, scan the data file directly and update the cache
406 /// with value locations, and write the index file back from what the scan found. We do not
407 /// read and decode the values, so only the locations go into the cache.
408 ///
409 /// The rebuilt index is written **into the existing index file, in place**, and not renamed
410 /// over it from a temporary the way garbage collection does. The difference matters, and it
411 /// is the writer bot. `ZoneBot::survey_files` hands each wbot its live `(data, index)` pair
412 /// at step 8, which is before `ZoneBot::init_caches` asks for any of this caching, so by the
413 /// time this method runs a wbot is already holding an open append handle on this very index
414 /// file. A rename replaces the inode underneath it: the handle goes on referring to the old
415 /// one, now unlinked, and every index record the wbot writes for the rest of the process goes
416 /// into a file nothing will ever open again. The data file is never renamed and so keeps
417 /// every record, which is why `get()` by key stays correct throughout while `scan`, which
418 /// walks index files, reports the zone as holding nothing.
419 ///
420 /// That is the whole of the symptom: an operator changes a limit, is answered `200`, sees the
421 /// new value in the console, and the gateway goes on using the old one for the life of the
422 /// process, because the read that would have found the new one goes through a scan and the
423 /// scan is blind. It recurs on every restart, since a restart is what runs this method.
424 ///
425 /// Rewriting the same inode keeps the wbot's handle valid. That handle is `O_APPEND`, so its
426 /// writes go to the true end of the file whatever length it believes the file to have, and
427 /// they land after the records written here.
428 #[allow(unused_assignments, unused_variables)]
429 pub fn init_cache_data_file(
430 &mut self,
431 mut reader: BufReader<File>,
432 fnum: FileNum,
433 dat_size: usize,
434 )
435 -> Outcome<()>
436 {
437 let typ = FileType::Data;
438 let mut index_file_buffer = Vec::new();
439
440 // 1. Name the index file. The records are gathered in memory and written at the end,
441 // after the size check below has confirmed that the data file read whole.
442 let mut ind_path = self.zdir().dir.clone();
443 ind_path.push(ZoneDir::relative_file_path(&FileType::Index, fnum));
444 // The data file's own path, needed if the walk hits a torn final record
445 // and has to truncate it away (step 4a).
446 let mut dat_path = self.zdir().dir.clone();
447 dat_path.push(ZoneDir::relative_file_path(&FileType::Data, fnum));
448
449 let csummer = self.api().schemes().checksummer().clone();
450
451 let mut pos = 0;
452 let mut kpos = 0;
453 let mut klen = 0;
454 let mut count = 0;
455 // Byte offset just past the last record that decoded whole, the number
456 // of such records, and whether a torn tail was cut. An append-only log
457 // can only be interrupted at its end, so a decode failure with nothing
458 // valid after it is the torn tail: the file is truncated here and the
459 // rebuild carries on, so one interrupted append costs one record rather
460 // than the whole file. A good record decoding after a bad one is not a
461 // crash tail but mid-file corruption, and is surfaced instead (see
462 // `try_truncate_torn_tail`).
463 let mut last_good_pos: u64 = 0;
464 let mut good_count = 0usize;
465 let mut tail_truncated = false;
466
467 let csum_len = res!(csummer.len());
468
469 loop {
470 // 3. Load the key Daticle bytes and while we're at it, compare the checksum.
471 let (key, meta, chash, mut pending_index) = match StoredKey::load(
472 &mut reader,
473 csummer.clone(),
474 ) {
475 Err(e) => {
476 // A decode failure at the key. If this is the append-only
477 // crash tail, truncate and finish the rebuild from what was
478 // recovered; otherwise surface it.
479 if res!(self.try_truncate_torn_tail(
480 &dat_path, last_good_pos, dat_size, fnum, csummer.clone(),
481 )) {
482 tail_truncated = true;
483 break;
484 }
485 return Err(err!(e,
486 "{}: Decode failure reading a key at position {} in {:?} file {} \
487 after {} good records, and it is not an append-only crash tail \
488 (a valid record decodes further on, or a writer is extending the \
489 file): this is mid-file corruption, and truncating here would \
490 discard live data.",
491 self.ozid(), last_good_pos, typ, fnum, good_count;
492 IO, File, Data, Mismatch));
493 },
494 Ok(None) => break,
495 Ok(Some((skey, skbyts, n))) => {
496 count += 1;
497 kpos = pos;
498 klen = n;
499 pos += n;
500 let chash = skey.ref_chash().clone();
501 // Build the index record (cache hash followed by the stored
502 // key bytes) but hold it back: it is appended to the index
503 // buffer only once the value beside it has also decoded, so
504 // a torn tail that left a key without its value does not
505 // leave a dangling key in the rebuilt index.
506 let mut pending = skey.ref_chash().to_vec();
507 pending.extend_from_slice(&skbyts);
508 let meta = skey.meta().clone();
509 (skey.into_key(), meta, chash, pending)
510 },
511 };
512 // 4. Count the value Daticle bytes. Dat::count_bytes also moves the
513 // reader cursor.
514 match StoredValue::count(
515 &mut reader,
516 csum_len,
517 ) {
518 Err(e) => {
519 // 4a. The key decoded but its value did not. Treat a torn
520 // tail the same way as a torn key above.
521 if res!(self.try_truncate_torn_tail(
522 &dat_path, last_good_pos, dat_size, fnum, csummer.clone(),
523 )) {
524 tail_truncated = true;
525 break;
526 }
527 return Err(err!(e,
528 "{}: Decode failure reading a value at position {} in {:?} file {} \
529 after {} good records, and it is not an append-only crash tail: \
530 this is mid-file corruption, and truncating here would discard \
531 live data.",
532 self.ozid(), last_good_pos, typ, fnum, good_count;
533 IO, File, Data, Mismatch));
534 },
535 Ok(0) => {
536 // A key with no value is the classic interrupted append: the
537 // key reached disk and the crash fell before its value. The
538 // reader is at EOF, so nothing follows -- a torn tail.
539 if res!(self.try_truncate_torn_tail(
540 &dat_path, last_good_pos, dat_size, fnum, csummer.clone(),
541 )) {
542 tail_truncated = true;
543 break;
544 }
545 return Err(err!(
546 "{}: Missing value at position {} in {:?} file {} after {} good \
547 records, and it is not an append-only crash tail.",
548 self.ozid(), last_good_pos, typ, fnum, good_count;
549 IO, File, Data, Missing));
550 },
551 Ok(n) => {
552 count += 1;
553 pos += n;
554
555 // 4b. `StoredValue::count` walks the value by seeking over
556 // its declared length rather than reading it, so a value
557 // whose length header survived but whose body was cut off
558 // by a crash is not caught above -- it seeks past the end
559 // of the file and reports the full length. Catch it here:
560 // a record that claims to end past the surveyed file size
561 // is a torn final value. Treat it as the torn tail (drop
562 // it, truncate, continue); if a good record still decodes
563 // after it, that is mid-file corruption and is surfaced.
564 if pos as u64 > dat_size as u64 {
565 if res!(self.try_truncate_torn_tail(
566 &dat_path, last_good_pos, dat_size, fnum, csummer.clone(),
567 )) {
568 tail_truncated = true;
569 break;
570 }
571 return Err(err!(
572 "{}: Value at position {} in {:?} file {} declares a length \
573 that runs {} bytes past the surveyed file size of {}, after \
574 {} good records, and it is not an append-only crash tail: \
575 this is mid-file corruption.",
576 self.ozid(), last_good_pos, typ, fnum, pos - dat_size, dat_size,
577 good_count;
578 IO, File, Data, Mismatch));
579 }
580
581 // 5. Create the FileLocation.
582 let sfloc = res!(StoredFileLocation::new(
583 fnum,
584 kpos as u64,
585 klen as u64,
586 n as u64,
587 csummer.clone(),
588 ));
589
590 // 6. Insert the key and location into the bot cache, informing a gbot about
591 // new data and old data that can be scheduled for garbage collection. The
592 // bot we advise actually performs any garbage collection, so instead of
593 // choosing randomly, we allocate each bot to an exclusive fraction of files
594 // based on their number.
595 let cind = key.index();
596 let kbyts = key.into_bytes();
597 let chash = res!(<alias::ChooseHash>::try_from(
598 &chash[..constant::CACHE_HASH_BYTES]));
599 let cbwind = ChooseCache::<PR>::choose_cbot_select(
600 alias::ChooseHashUint::from_be_bytes(chash),
601 self.cfg().num_zones,
602 self.cfg().num_cbots_per_zone,
603 );
604 let cbots = res!(self.cbots());
605 let bot = res!(cbots.get_bot(**cbwind.bpind()));
606 let ibuf = &sfloc.buf;
607 res!(bot.send(
608 OzoneMsg::Insert(
609 kbyts,
610 None,
611 cind,
612 sfloc.ref_file_location().clone(),
613 ibuf.len(),
614 meta,
615 Responder::none(Some(self.ozid())),
616 None,
617 )
618 ));
619
620 // 7. The record is whole: commit its held-back key bytes and
621 // its index entry to the buffer, and mark this as the last
622 // good boundary.
623 index_file_buffer.append(&mut pending_index);
624 index_file_buffer.extend_from_slice(ibuf);
625 last_good_pos = pos as u64;
626 good_count += 1;
627 },
628 }
629 }
630
631 // 8. Do size check. The walk ended without a decode failure but read a
632 // different number of bytes than the survey measured. `pos > dat_size`
633 // is caught inside the loop (step 4b) and cannot reach here. A
634 // deliberate tail truncation leaves `pos < dat_size` by design and is
635 // exempt. What remains is `pos < dat_size`: trailing bytes after the
636 // last decoded record that were too few to read as another record (a
637 // sub-header fragment of an interrupted append, which `StoredKey::load`
638 // reports as a clean end). That is the torn tail too, so truncate to
639 // the last good record rather than abandoning the whole file -- unless
640 // a valid record still decodes in those trailing bytes (corruption) or
641 // a writer has grown the file past the survey (a live append), both of
642 // which `try_truncate_torn_tail` refuses and which are surfaced here.
643 if !tail_truncated && pos != dat_size {
644 if pos < dat_size
645 && res!(self.try_truncate_torn_tail(
646 &dat_path, last_good_pos, dat_size, fnum, csummer.clone(),
647 ))
648 {
649 tail_truncated = true;
650 } else {
651 return Err(err!(
652 "{}: After initial caching of data file {}, the total data count \
653 came to {} bytes, but the originally surveyed file size was {}.",
654 self.ozid(), fnum, pos, dat_size;
655 Mismatch, Data));
656 }
657 }
658
659 // 9. Write the rebuilt index into the existing index file, keeping its inode so that
660 // the wbot holding this file open for append goes on writing to it. See the note on
661 // this method for what renaming a fresh file over it costs.
662 let mut ind_file = match OpenOptions::new()
663 .create(true)
664 .write(true)
665 .truncate(true)
666 .open(&ind_path)
667 {
668 Err(e) => return Err(err!(e,
669 "{}: While opening index file {:?} to rebuild it from data file {}.",
670 self.ozid(), ind_path, fnum;
671 IO, File, Write, Create)),
672 Ok(file) => file,
673 };
674 res!(ind_file.write_all(&index_file_buffer));
675 res!(ind_file.flush());
676 // The rebuild is the only copy of this index. Losing it to a crash before it reaches
677 // the disk puts the zone straight back into the state this method exists to repair, and
678 // the repair is not free -- it reads the whole data file.
679 res!(ind_file.sync_data());
680
681 Ok(())
682 }
683
684 /// Performs garbage collection on the given file. Assumes that the move map for the file is
685 /// empty. The basic idea is to transcribe (re-write) the data file, skipping sections
686 /// scheduled for deletion. While this process is going on, deletion messages can continue to
687 /// be queue up on the fbot write channel. It is therefore necessary to create a "move map" in
688 /// the file state which maps the old data locations to their new locations in the file, later
689 /// allowing those queued deletions to correctly apply to the new locations. After
690 /// transcription, the data in the new file is cached using a similar approach to a cbot.
691 ///```ignore
692 /// Garbage collection of file fg
693 /// file to live involves transcription of
694 /// be gc'd file current key-value pairs
695 /// with starting locations s01
696 /// f1 fg fL and s02 to new starting
697 /// +------+ +------+ +------+ locations s11 and s12. Old
698 /// | | | | | | values in fg are not copied
699 /// | | | | | | and thereby deleted.
700 /// original | | k1 |\\\\\\| s01 | |
701 /// data | | | | | | However, while this occurs we
702 /// files | | ... | | ... | | want the writer bot(s) to
703 /// | | | | | | continue appending data to
704 /// | | k2 |//////| s02 | | live files, which could include
705 /// | | | | | | new values for k1 and k2.
706 /// +------+ +------+ +------+
707 ///
708 /// new fg
709 /// new data
710 /// files s11 |\\\\\\|
711 /// (post-gc) s12 |//////|
712 /// cache changes
713 ///
714 /// | scheduled for | third party | gc updates
715 /// scenarios | deletion in fg | changes | required
716 /// | after gc started | (examples) |
717 /// --------------------+------------------------+-----------------------+------------------------
718 /// 1. None | <none> | | k1:(fg,s01)->(fg,s11)
719 /// | | |
720 /// 2. New value(s) | s11 <- s01 <- | k1:(fg,s01)->(fL,s21) |
721 /// | | k1:(fL,s21)->(fL,s31) |
722 /// | | |
723 /// | | |
724 /// | | |
725 ///
726 ///```
727 /// Returns whether the file state can be eliminated because the data file has been completely
728 /// deleted, or the new data file size.
729 fn collect_garbage(
730 &mut self,
731 fnum: FileNum,
732 mut fstat: FileState,
733 fbot_index: usize,
734 )
735 -> Outcome<()>
736 {
737 // [19] Perform transcription from data_reader to data_writer.
738
739 hooks::collect_delay();
740 trace!(sync_log::stream(), "{}: Performing garbage collection on file {}...", self.ozid(), fnum);
741 let typ = FileType::Data;
742 // 1. Open the data file for reading.
743 let (data_path, file) = res!(self.zdir().open_ozone_file(
744 fnum,
745 &typ,
746 &FileAccess::Reading,
747 ));
748 let data_file_len = res!(file.metadata()).len();
749 let old_size = data_file_len as usize;
750 let mut data_reader = BufReader::new(file);
751
752 // 2. Create new, temporary data file for writing.
753 let mut tmp_data_path = self.zdir().dir.clone();
754 tmp_data_path.push(ZoneDir::relative_gc_temp_path(&typ, fnum));
755 let mut new_start: u64 = 0;
756 let old_sum = try_into!(usize, fstat.get_old_sum());
757
758 {
759 let file = res!(ZoneDir::open_file(
760 &tmp_data_path,
761 &FileAccess::Writing,
762 ));
763 // Rust and/or linux seems to require that this BufWriter on a write-only file (
764 // creation sets it to write-only) be closed before we can open a BufReader to the
765 // same file, so we create this special scope for data_writer.
766 let mut data_writer = BufWriter::new(file);
767
768 // 3. Transcribe the existing data file to the temporary file, skipping old key-value pairs.
769 let mut old_start1: u64 = 0;
770 let mut dstat1 = None;
771 let mut first = true;
772 for old_start2 in res!(fstat.get_data_start_positions()) {
773 if !first {
774 let dloc = DataLocation {
775 start: old_start1,
776 len: old_start2 - old_start1,
777 };
778 match dstat1 {
779 Some(DataState::Cur) => {
780 let mut buf = vec![0u8; dloc.len as usize];
781 res!(data_reader.seek(SeekFrom::Start(dloc.start)));
782 match data_reader.read_exact(&mut buf) {
783 Err(e) => return Err(err!(e,
784 "{}: While trying to read exactly {} bytes from position \
785 {} in file {} of {} bytes length, {:?}. The file state is {:?}.",
786 self.ozid(), dloc.len, dloc.start, fnum,
787 data_file_len, data_path, fstat;
788 IO, File, Read)),
789 Ok(()) => (),
790 }
791 // The move is kept against the record it carries, named by the key
792 // and stamp at the head of its bytes.
793 let rid = match res!(StoredKey::<UIDL, UID>::load(
794 &mut &buf[..],
795 self.api().schemes().checksummer().clone(),
796 )) {
797 Some((skey, _, _)) => res!(RecordDigest::new(
798 skey.key().as_bytes(),
799 skey.meta(),
800 )),
801 None => return Err(err!(
802 "{}: No key at position {} in file {}, where the file state \
803 has a current record.", self.ozid(), dloc.start, fnum;
804 Bug, Missing, Data)),
805 };
806 res!(data_writer.write_all(&mut buf));
807 fstat.update_moved(&dloc, new_start, rid);
808 new_start += dloc.len;
809 },
810 Some(DataState::Old) => {
811 res!(fstat.retire_old(&dloc));
812 },
813 None => break,
814 }
815 } else {
816 first = false;
817 }
818 old_start1 = old_start2;
819 dstat1 = fstat.get_data_state(old_start2).cloned();
820 }
821
822 // Durability barrier before the rename below: force the transcribed
823 // temporary data file to stable storage. Dropping the BufWriter
824 // flushes the buffer into the page cache, but the rename at step 7
825 // is a directory operation that can reach disk before the file's
826 // contents do. A power loss in that window would leave the rename
827 // durable and the file torn -- a renamed, torn file replacing a
828 // previously good one. Syncing the contents first closes that hole.
829 res!(data_writer.flush());
830 if let Err(e) = data_writer.get_ref().sync_data() {
831 return Err(err!(e,
832 "{}: sync_data on the transcribed temporary data file {:?} \
833 failed before renaming it over data file {}.",
834 self.ozid(), tmp_data_path, fnum;
835 IO, File, Write));
836 }
837 }
838
839 let new_size = new_start as usize;
840
841 // 4. Do some checks.
842 if new_size > old_size {
843 return Err(err!(
844 "{}: The file {} has grown in size from {} to {} after garbage \
845 collection, this should not occur in Ozone.",
846 self.ozid(), fnum, old_size, new_size;
847 Bug, Missing, Data));
848 }
849 if old_sum != old_size - new_size {
850 return Err(err!(
851 "{}: The file {} was scheduled to remove {} bytes, but instead \
852 removed {} bytes, going from {} to {} bytes.",
853 self.ozid(), fnum, old_sum, old_size - new_size, old_size, new_size;
854 Bug, Mismatch, Data));
855 }
856 if !fstat.data_map_empty() {
857 return Err(err!(
858 "{}: Garbage collection for file {} should have cleared out the \
859 data map, instead it still contains entries, {:?}.",
860 self.ozid(), fnum, fstat.data_map();
861 Bug, Mismatch, Data));
862 }
863
864 if new_size == 0 {
865 return Err(err!(
866 "{}: Garbage collection for file {} has deleted the entire data file, \
867 however this should have been done by the fbot.",
868 self.ozid(), fnum;
869 Bug, Mismatch, Data));
870 }
871
872 let mut dat_ind_file_size_decrease = old_sum;
873 fstat.set_data_file_size(new_size);
874 let old_ind_size = fstat.get_index_file_size();
875
876 // 6. Re-create the index file by scanning the new data file. Any values remaining in the
877 // data file are unique and we must handle a few scenarios.
878 let file = match OpenOptions::new().read(true).open(&tmp_data_path) {
879 Err(e) => return Err(err!(e, "While opening file {:?}", tmp_data_path; IO, File, Read)),
880 Ok(f) => f,
881 };
882 let new_data_reader = BufReader::new(file);
883 fstat = res!(self.cache_data_file(
884 new_data_reader,
885 fnum,
886 fstat,
887 ));
888
889 if fstat.get_index_file_size() > old_ind_size {
890 return Err(err!(
891 "{}: Indexing of the new data file {} after garbage collection \
892 has resulted in unexpected growth of the index file from {} to \
893 {} bytes.",
894 self.ozid(), fnum, old_ind_size, fstat.get_index_file_size();
895 Bug, Mismatch, Data));
896 }
897 dat_ind_file_size_decrease += old_ind_size - fstat.get_index_file_size();
898
899 // 7. Replace the old data file with the new temporary file.
900 res!(fs::rename(tmp_data_path, data_path));
901
902 // Persist the directory entries changed by the renames above (this data
903 // file here, and the index file inside `cache_data_file`). A rename is a
904 // directory metadata operation; fsyncing the file contents does not
905 // persist the rename itself, so without this a power loss could leave
906 // the directory pointing at a file that is not yet on disk. Both files
907 // live directly in the zone directory, so one fsync of it covers both.
908 res!(Self::sync_dir(&self.zdir().dir));
909
910 // Both renames put a new inode behind an unchanged path and unlinked the old one, but
911 // an rbot that read this file earlier still holds the old inode open in its file cache
912 // for `constant::FILE_CACHE_EXPIRY_SECS`, and nothing about a rename reaches that cache.
913 // It would go on seeking to the NEW offsets in the OLD inode, which are not record
914 // boundaries there, so every read of a carried record would fail its checksum until the
915 // entry expired a quarter of an hour later. Hence this notice, and hence its position:
916 // after both renames, because an rbot told beforehand would simply reopen the path and
917 // cache the old inode again. Telling an rbot that has no entry, or one opened since the
918 // rename, costs it a reopen and nothing else.
919 res!(self.notify_file_replaced(fnum, &[FileType::Data, FileType::Index]));
920
921 // 8. Reset FileState.
922 fstat.reset_old_accounting();
923
924 // [22] Send updated file state back to the fbot.
925 let bots = res!(self.fbots());
926 let bot = res!(bots.get_bot(fbot_index));
927 if let Err(e) = bot.send(
928 OzoneMsg::GcCompleted(
929 fnum,
930 fstat,
931 dat_ind_file_size_decrease,
932 )
933 ) {
934 return Err(err!(e,
935 "{}: Cannot send updated file state for file number {} to fbot {}",
936 self.ozid(), fnum, fbot_index;
937 Channel, Write));
938 }
939 debug!(sync_log::stream(), "{}: Reduced file {} size by {:.1}% from {} to {} bytes.",
940 self.ozid(),
941 fnum,
942 100.0 * ((old_size - new_size) as f32) / (old_size as f32),
943 old_size,
944 new_size,
945 );
946 Ok(())
947 }
948
949 /// Tells every rbot to drop its cached handle on the given files, because a collection has
950 /// just renamed new ones over them. Every pool in every zone is told: a file number is only
951 /// unique within a zone, so a same-numbered file elsewhere is invalidated needlessly, but that
952 /// costs one reopen and keeps the notice independent of which zone a reader serves.
953 fn notify_file_replaced(
954 &self,
955 fnum: FileNum,
956 typs: &[FileType],
957 )
958 -> Outcome<()>
959 {
960 for pool in self.chans().get_all_workers_of_type(&WorkerType::Reader) {
961 for i in 0..pool.len() {
962 let bot = res!(pool.get_bot(i));
963 for typ in typs {
964 if let Err(e) = bot.send(OzoneMsg::FileReplaced(fnum, typ.clone())) {
965 return Err(err!(e,
966 "{}: Cannot tell rbot {} that {:?} file {} has been replaced.",
967 self.ozid(), i, typ, fnum;
968 Channel, Write));
969 }
970 }
971 }
972 }
973 Ok(())
974 }
975
976 /// This method is similar to the initialisation method `init_cache_data_file` in that the
977 /// data file is scanned and a new index file is created, however we must deal with the cache
978 /// and deletion scheduling differently. Initialisation can rely on the chronological order of
979 /// data and work its way sequentially through files, but here we want to allow values to be
980 /// added to live files while we collect the garbage and therefore must update the cache and
981 /// file state depending on live file changes.
982 pub fn cache_data_file(
983 &mut self,
984 mut reader: BufReader<File>,
985 fnum: FileNum,
986 mut fstat: FileState,
987 )
988 -> Outcome<FileState>
989 {
990 let typ = FileType::Data;
991 // 1. Name the index file and the temporary the rebuild is written to. The rebuild
992 // goes to the temporary and is renamed over the index at the end, so a reader
993 // walking index files alongside this collection sees either the whole old index
994 // or the whole new one, never a gap. Writing in place would leave the file
995 // absent, then partial, for the length of the rebuild, and a scan arriving in
996 // that window would drop every key whose current value lives in this file.
997 let mut ind_path = self.zdir().dir.clone();
998 ind_path.push(ZoneDir::relative_file_path(&FileType::Index, fnum));
999 let mut tmp_ind_path = self.zdir().dir.clone();
1000 tmp_ind_path.push(ZoneDir::relative_gc_temp_path(&FileType::Index, fnum));
1001 // An abandoned rebuild would otherwise be appended to, because ozone opens files
1002 // for writing in append mode.
1003 if tmp_ind_path.is_file() {
1004 res!(fs::remove_file(&tmp_ind_path));
1005 }
1006 fstat.reset_index_file_size();
1007
1008 let mut pos = 0;
1009 let mut count = 0;
1010
1011 // Prepare to buffer request to update caches.
1012 let resp_g1 = Responder::new(Some(self.ozid()));
1013 let nc = self.cfg().num_cbots_per_zone();
1014 let mut buffers = Vec::new();
1015 for _ in 0..nc {
1016 buffers.push(Vec::new());
1017 }
1018
1019 {
1020 // 2. Create the new index file.
1021 let file = res!(ZoneDir::open_file(
1022 &tmp_ind_path,
1023 &FileAccess::Writing,
1024 ));
1025 let mut writer = BufWriter::new(file);
1026
1027 loop {
1028 // 3. Load the key Daticle bytes and while we're at it, compare the checksum.
1029 let (key, klen, kpos, meta, mut index_entry, chash) =
1030 match StoredKey::load(
1031 &mut reader,
1032 self.api().schemes().checksummer().clone(),
1033 ) {
1034 Err(e) => return Err(err!(e,
1035 "{}: While reading from position {} in {:?} file {}, \
1036 having read {} items.",
1037 self.ozid(), pos, typ, fnum, count;
1038 IO, File, Read)),
1039 Ok(None) => break,
1040 Ok(Some((skey, skbyts, n))) => {
1041 count += 1;
1042 let kpos = pos;
1043 pos += n;
1044 let klen = n;
1045 let meta: Meta<UIDL, UID> = skey.meta().clone();
1046 let chash = skey.ref_chash().clone();
1047 let key = skey.into_key();
1048 // An index record starts with the cache hash, exactly as the
1049 // wbot writes it and as `StoredKey::load` expects to read it.
1050 // `load` returns the key bytes with that hash already drained
1051 // off the front, so it has to be put back; a rebuild that
1052 // omitted it left every record in the file misaligned by four
1053 // bytes and the file undecodable from its first byte.
1054 let mut index_entry = chash.to_vec();
1055 index_entry.extend_from_slice(&skbyts);
1056 (key, klen, kpos, meta, index_entry, chash)
1057 },
1058 };
1059 // 4. Count the value Daticle bytes. Dat::count_bytes also moves the
1060 // reader cursor.
1061 match StoredValue::count(
1062 &mut reader,
1063 res!(self.api().schemes().checksummer().len()),
1064 ) {
1065 Err(e) => return Err(err!(e,
1066 "{}: While reading from position {} in {:?} file {}, \
1067 having read {} items.",
1068 self.ozid(), pos, typ, fnum, count;
1069 IO, File, Read)),
1070 Ok(0) => return Err(err!(
1071 "{}: Missing value at end of {:?} file {}.",
1072 self.ozid(), typ, fnum;
1073 Missing, IO, File)),
1074 Ok(n) => {
1075 let new_start = kpos as u64;
1076 // 5. Create the FileLocation.
1077 let sfloc = res!(StoredFileLocation::new( // do this before incrementing pos
1078 fnum,
1079 new_start,
1080 klen as u64,
1081 n as u64,
1082 self.api().schemes().checksummer().clone(),
1083 ));
1084 count += 1;
1085 pos += n;
1086
1087 // 6. Update existing cache reference to this key. If the cached location
1088 // still points to this file, update the starting position.
1089 let chash = res!(<alias::ChooseHash>::try_from(
1090 &chash[..constant::CACHE_HASH_BYTES]));
1091 let cbwind = ChooseCache::<PR>::choose_cbot_select(
1092 alias::ChooseHashUint::from_be_bytes(chash),
1093 self.cfg().num_zones,
1094 self.cfg().num_cbots_per_zone,
1095 );
1096 let bpind = cbwind.bpind();
1097 buffers[**bpind].push((
1098 key.into_bytes(),
1099 sfloc.ref_file_location().clone(),
1100 meta,
1101 ));
1102
1103 // 8. Append to the index file and update the effect of the file size increase
1104 // on the file state and the directory size.
1105 //let ibuf = StoredIndex::as_bytes(&mut floc); // update floc encoded length
1106 let ibuf = sfloc.buf;
1107 index_entry.extend_from_slice(&ibuf);
1108 res!(writer.write(&index_entry));
1109 res!(fstat.inc_index_file_size(index_entry.len()));
1110 },
1111 }
1112 }
1113
1114 // Flush explicitly: dropping a BufWriter flushes too, but discards any error
1115 // doing so, and a rebuild that silently lost its tail would be renamed over a
1116 // good index.
1117 res!(writer.flush());
1118 // Durability barrier before the rename below: force the rebuilt
1119 // index contents to stable storage, for the same reason the data
1120 // transcription is synced before its rename -- the rename can reach
1121 // disk before the contents, and a power loss there would rename a
1122 // torn index over a good one.
1123 if let Err(e) = writer.get_ref().sync_data() {
1124 return Err(err!(e,
1125 "{}: sync_data on the rebuilt temporary index file {:?} \
1126 failed before renaming it over index file {}.",
1127 self.ozid(), tmp_ind_path, fnum;
1128 IO, File, Write));
1129 }
1130 } // close out that writer
1131
1132 // Put the rebuilt index in place in one step.
1133 res!(fs::rename(&tmp_ind_path, &ind_path));
1134
1135 // Send cache update request batches to each cbot.
1136 let bots = res!(self.cbots());
1137 for i in (0..nc).rev() {
1138 let bot = res!(bots.get_bot(i));
1139 if let Some(buf) = buffers.pop() { // Transfer ownership to OzoneMsg.
1140 if let Err(e) = bot.send(
1141 OzoneMsg::GcCacheUpdateRequest(
1142 buf,
1143 resp_g1.clone(),
1144 )
1145 ) {
1146 return Err(err!(e,
1147 "Cannot send gc cache update requests to cbot {}", i;
1148 Channel, Write));
1149 }
1150 } else {
1151 return Err(err!(
1152 "The number of cache bots is {} should match the number of buffers, {}.",
1153 nc, nc - i;
1154 Bug, Mismatch, Size));
1155 }
1156 }
1157 // Wait for each response.
1158 for _ in 0..nc {
1159 match resp_g1.recv_timeout(constant::BOT_REQUEST_TIMEOUT) {
1160 Err(e) => return Err(err!(e,
1161 "While collecting gc cache update response.";
1162 IO, Channel, Read)),
1163 Ok(OzoneMsg::GcCacheUpdateResponse(old_flocs)) => {
1164 for (old_floc, rid) in old_flocs {
1165 // The cache re-anchored this record, so it is still current and nothing
1166 // will ask for it at its old start again but a read already on its way,
1167 // which the reader confirms and retries. Its move entry is spent, and
1168 // taken by record: the entry at that offset may be another's.
1169 fstat.map_and_remove(&old_floc.keyval(), &rid);
1170 }
1171 },
1172 Ok(msg) => return Err(err!(
1173 "Unrecognised cache initialisation request response: {:?}", msg;
1174 Channel, Unknown)),
1175 }
1176 }
1177 Ok(fstat)
1178 }
1179
1180 /// Is there a complete, checksum-valid key/value record beginning at any
1181 /// offset in `[from, end)` of the data file? Called only after the rebuild
1182 /// walk hits a decode failure, to tell an append-only crash tail (nothing
1183 /// valid follows the torn record) from mid-file corruption (a good record
1184 /// decodes after the bad one). A crash leaves only the one partially
1185 /// written record after the last good one, so this scans a short region and
1186 /// answers no; corruption of an otherwise whole file finds the next real
1187 /// record quickly and answers yes. The region is read once into memory and
1188 /// every offset is tried, so the cost is bounded by the data file size and
1189 /// paid only on the rare torn-rebuild path. A false positive would need a
1190 /// run of bytes that parses as a whole key and value and passes the key's
1191 /// checksum by chance, which is negligible.
1192 fn valid_record_after<C: Checksummer>(
1193 dat_path: &std::path::Path,
1194 csummer: C,
1195 from: u64,
1196 end: u64,
1197 )
1198 -> Outcome<bool>
1199 {
1200 if from >= end {
1201 return Ok(false);
1202 }
1203 let csum_len = res!(csummer.len());
1204 let mut file = match OpenOptions::new().read(true).open(dat_path) {
1205 Err(e) => return Err(err!(e,
1206 "While opening data file {:?} to scan for a valid record past a \
1207 torn one.", dat_path;
1208 IO, File, Read)),
1209 Ok(f) => f,
1210 };
1211 res!(file.seek(SeekFrom::Start(from)));
1212 let mut tail = Vec::new();
1213 res!(file.read_to_end(&mut tail));
1214 // Offset 0 is the record that already failed to decode in the caller, so
1215 // start one byte in: the question is whether a GOOD record follows it.
1216 for off in 1..tail.len() {
1217 let mut cur = std::io::Cursor::new(&tail[off..]);
1218 if let Ok(Some((_, _, _))) = StoredKey::<UIDL, UID>::load(&mut cur, csummer.clone()) {
1219 if let Ok(n) = StoredValue::count(&mut cur, csum_len) {
1220 if n > 0 {
1221 return Ok(true);
1222 }
1223 }
1224 }
1225 }
1226 Ok(false)
1227 }
1228
1229 /// Decides whether a decode failure during the data-file rebuild is the
1230 /// append-only crash tail and, if so, truncates the data file to the last
1231 /// good record and returns `true` so the caller can finish the rebuild from
1232 /// what was recovered. Returns `false` when a good record decodes after the
1233 /// failed one (mid-file corruption) or the file has grown past the survey (a
1234 /// writer is appending to it): in both cases the tail is not ours to cut and
1235 /// the caller surfaces the original error.
1236 fn try_truncate_torn_tail<C: Checksummer>(
1237 &self,
1238 dat_path: &std::path::Path,
1239 last_good_pos: u64,
1240 dat_size: usize,
1241 fnum: FileNum,
1242 csummer: C,
1243 )
1244 -> Outcome<bool>
1245 {
1246 if res!(Self::valid_record_after(dat_path, csummer, last_good_pos, dat_size as u64)) {
1247 return Ok(false);
1248 }
1249 // A file longer than the survey means a writer appended under the
1250 // rebuild; the tail is not a crash artefact and must not be cut.
1251 let phys_len = res!(fs::metadata(dat_path)).len();
1252 if phys_len != dat_size as u64 {
1253 return Ok(false);
1254 }
1255 let tf = match OpenOptions::new().write(true).open(dat_path) {
1256 Err(e) => return Err(err!(e,
1257 "{}: Opening data file {:?} ({}) to truncate a torn tail to {}.",
1258 self.ozid(), dat_path, fnum, last_good_pos;
1259 IO, File, Write)),
1260 Ok(f) => f,
1261 };
1262 if let Err(e) = tf.set_len(last_good_pos) {
1263 return Err(err!(e,
1264 "{}: Truncating data file {:?} ({}) to {} to drop a torn tail.",
1265 self.ozid(), dat_path, fnum, last_good_pos;
1266 IO, File, Write));
1267 }
1268 if let Err(e) = tf.sync_all() {
1269 return Err(err!(e,
1270 "{}: Syncing data file {:?} ({}) after truncating a torn tail.",
1271 self.ozid(), dat_path, fnum;
1272 IO, File, Write));
1273 }
1274 warn!(sync_log::stream(),
1275 "{}: Data file {} had a torn final record at position {}; truncated to \
1276 the last good record and rebuilding the index from it. One interrupted \
1277 append costs one record.",
1278 self.ozid(), fnum, last_good_pos);
1279 Ok(true)
1280 }
1281
1282 /// Forces a directory's entries to stable storage. A `rename` is a
1283 /// directory metadata operation, so fsyncing a renamed file's contents does
1284 /// not persist the rename itself; this is called after the GC renames so a
1285 /// power loss cannot leave the directory pointing at a file that never
1286 /// reached disk.
1287 fn sync_dir(dir: &std::path::Path) -> Outcome<()> {
1288 match File::open(dir) {
1289 Err(e) => Err(err!(e,
1290 "While opening directory {:?} to fsync it.", dir;
1291 IO, File, Read)),
1292 Ok(d) => match d.sync_all() {
1293 Err(e) => Err(err!(e,
1294 "While fsyncing directory {:?}.", dir;
1295 IO, File, Write)),
1296 Ok(()) => Ok(()),
1297 },
1298 }
1299 }
1300
1301}