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
| 1 | use 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 | |
| 32 | use oxedyne_fe2o3_iop_db::api::Meta; |
| 33 | use oxedyne_fe2o3_jdat::id::NumIdDat; |
| 34 | |
| 35 | use 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. |
| 56 | pub 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 | |
| 81 | impl< |
| 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 | |
| 94 | impl< |
| 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 | |
| 107 | impl< |
| 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 | |
| 188 | impl< |
| 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 | } |