oxedyne/fe2o3/fe2o3_o3db_sync/src/file/state.rs
19.3 KiB, 107 runs
created by r1870400018:799, 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 | file::{ |
| 4 | floc::{ |
| 5 | DataLocation, |
| 6 | FileLocation, |
| 7 | FileNum, |
| 8 | }, |
| 9 | stored::RecordDigest, |
| 10 | }, |
| 11 | }; |
| 12 | |
| 13 | use std::{ |
| 14 | collections::BTreeMap, |
| 15 | }; |
| 16 | |
| 17 | #[derive(Clone, Debug)] |
| 18 | pub enum Present { |
| 19 | Solo(FileType), |
| 20 | Pair, |
| 21 | } |
| 22 | |
| 23 | impl Default for Present { |
| 24 | fn default() -> Self { |
| 25 | Self::Pair |
| 26 | } |
| 27 | } |
| 28 | |
| 29 | #[derive(Clone, Debug, Eq, PartialEq)] |
| 30 | pub enum DataState { |
| 31 | Cur, // Current version of value for this key. |
| 32 | Old, // Value flagged for garbage collection. |
| 33 | } |
| 34 | |
| 35 | impl std::fmt::Display for DataState { |
| 36 | fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result { |
| 37 | match self { |
| 38 | Self::Cur => write!(f, "cur"), |
| 39 | Self::Old => write!(f, "old"), |
| 40 | } |
| 41 | } |
| 42 | } |
| 43 | |
| 44 | /// Where a collection carried a record, and which record it was. |
| 45 | #[derive(Clone, Copy, Debug)] |
| 46 | pub struct Move { |
| 47 | to: u64, |
| 48 | rid: RecordDigest, |
| 49 | } |
| 50 | |
| 51 | #[derive(Clone, Debug, Default)] |
| 52 | pub struct FileState { |
| 53 | present: Present, |
| 54 | dat_size: usize, |
| 55 | ind_size: usize, |
| 56 | live: bool, |
| 57 | oldsum: u64, |
| 58 | oldcnt: usize, |
| 59 | dmap: BTreeMap<u64, DataState>, // Map of key-value pair starting positions in data file. |
| 60 | mmap: BTreeMap<u64, Move>, // Ephemeral map of the movement of starting positions due to gc. |
| 61 | pending_old: BTreeMap<u64, u64>, // Supersessions that arrived before the record's insert; start -> registered length. |
| 62 | gc_active: bool, |
| 63 | readers: usize, |
| 64 | } |
| 65 | |
| 66 | impl FileState { |
| 67 | |
| 68 | // Getters. |
| 69 | pub fn present(&self) -> &Present { &self.present } |
| 70 | pub fn get_data_file_size(&self) -> usize { self.dat_size } |
| 71 | pub fn get_index_file_size(&self) -> usize { self.ind_size } |
| 72 | pub fn is_live(&self) -> bool { self.live } |
| 73 | pub fn get_old_sum(&self) -> u64 { self.oldsum } |
| 74 | pub fn get_old_count(&self) -> usize { self.oldcnt } |
| 75 | pub fn data_map(&self) -> &BTreeMap<u64, DataState> { &self.dmap } |
| 76 | pub fn data_map_mut(&mut self) -> &mut BTreeMap<u64, DataState> { &mut self.dmap } |
| 77 | pub fn gc_active(&self) -> bool { self.gc_active } |
| 78 | pub fn readers(&self) -> usize { self.readers } |
| 79 | pub fn no_readers(&self) -> bool { self.readers == 0 } |
| 80 | |
| 81 | pub fn get_data_state(&self, start: u64) -> Option<&DataState> { |
| 82 | self.dmap.get(&start) |
| 83 | } |
| 84 | /// Provides mutable access to a data state entry. Used during initialisation and state |
| 85 | /// updates to modify entry states. |
| 86 | pub fn get_data_state_mut(&mut self, start: u64) -> Option<&mut DataState> { |
| 87 | self.dmap.get_mut(&start) |
| 88 | } |
| 89 | pub fn get_data_start_positions(&self) -> Outcome<Vec<u64>> { |
| 90 | let mut starts = Vec::new(); |
| 91 | for start in self.dmap.keys() { |
| 92 | starts.push(*start); |
| 93 | } |
| 94 | starts.push(res!(self.dat_size.try_into())); |
| 95 | Ok(starts) |
| 96 | } |
| 97 | |
| 98 | pub fn is_all_old(&self) -> bool { |
| 99 | for (_, dstat) in &self.dmap { |
| 100 | if *dstat == DataState::Cur { |
| 101 | return false; |
| 102 | } |
| 103 | } |
| 104 | true |
| 105 | } |
| 106 | |
| 107 | // Setters. |
| 108 | pub fn set_data_map(&mut self, dmap: BTreeMap<u64, DataState>) { |
| 109 | self.dmap = dmap; |
| 110 | } |
| 111 | pub fn set_data_file_size(&mut self, size: usize) { |
| 112 | self.dat_size = size; |
| 113 | } |
| 114 | |
| 115 | /// Updates the index file size directly. Separate from data file size for accurate space |
| 116 | /// tracking. |
| 117 | pub fn set_index_file_size(&mut self, size: usize) { |
| 118 | self.ind_size = size; |
| 119 | } |
| 120 | pub fn set_live(&mut self, live: bool) { |
| 121 | self.live = live; |
| 122 | } |
| 123 | pub fn set_present(&mut self, present: Present) { |
| 124 | self.present = present; |
| 125 | } |
| 126 | pub fn set_gc(&mut self, active: bool) { |
| 127 | self.gc_active = active; |
| 128 | } |
| 129 | pub fn inc_readers(&mut self) -> Outcome<()> { |
| 130 | let (new, oflow) = self.readers().overflowing_add(1); |
| 131 | if oflow { |
| 132 | Err(err!( |
| 133 | "Attempt to increment the number of readers for the file state. \ |
| 134 | The number, {}, is already at maximum.", self.readers(); |
| 135 | Bug, Overflow, Integer)) |
| 136 | } else { |
| 137 | self.readers = new; |
| 138 | Ok(()) |
| 139 | } |
| 140 | } |
| 141 | pub fn dec_readers(&mut self) -> Outcome<()> { |
| 142 | let (new, uflow) = self.readers().overflowing_sub(1); |
| 143 | if uflow { |
| 144 | Err(err!( |
| 145 | "Attempt to decrement the number of readers for the file state. \ |
| 146 | The number, {}, is already at minimum.", self.readers(); |
| 147 | Bug, Underflow, Integer)) |
| 148 | } else { |
| 149 | self.readers = new; |
| 150 | Ok(()) |
| 151 | } |
| 152 | } |
| 153 | |
| 154 | pub fn reset(&mut self) { |
| 155 | self.dat_size = 0; |
| 156 | self.ind_size = 0; |
| 157 | self.oldsum = 0; |
| 158 | self.oldcnt = 0; |
| 159 | self.dmap = BTreeMap::new(); |
| 160 | self.mmap = BTreeMap::new(); |
| 161 | } |
| 162 | |
| 163 | pub fn reset_data_file_size(&mut self) -> usize { |
| 164 | let dat_size = self.dat_size; |
| 165 | self.dat_size = 0; |
| 166 | dat_size |
| 167 | } |
| 168 | pub fn reset_index_file_size(&mut self) -> usize { |
| 169 | let ind_size = self.ind_size; |
| 170 | self.ind_size = 0; |
| 171 | ind_size |
| 172 | } |
| 173 | pub fn reset_old_accounting(&mut self) { |
| 174 | self.oldsum = 0; |
| 175 | self.oldcnt = 0; |
| 176 | } |
| 177 | |
| 178 | // Queries. |
| 179 | pub fn is_all_data_old(&self) -> bool { |
| 180 | self.oldcnt == self.dmap.len() |
| 181 | } |
| 182 | |
| 183 | pub fn no_pending_moves(&self) -> bool { |
| 184 | self.mmap.len() == 0 |
| 185 | } |
| 186 | |
| 187 | /// Are there no supersessions still waiting for their record's insertion to land? |
| 188 | pub fn pending_old_empty(&self) -> bool { |
| 189 | self.pending_old.is_empty() |
| 190 | } |
| 191 | pub fn pending_old(&self) -> &BTreeMap<u64, u64> { &self.pending_old } |
| 192 | |
| 193 | pub fn data_map_empty(&self) -> bool { |
| 194 | self.dmap.len() == 0 |
| 195 | } |
| 196 | |
| 197 | pub fn data_map_len(&self) -> usize { |
| 198 | self.dmap.len() |
| 199 | } |
| 200 | |
| 201 | pub fn move_map_len(&self) -> usize { |
| 202 | self.mmap.len() |
| 203 | } |
| 204 | |
| 205 | // Data map mutation - this is where the size of old data and the file are modified. |
| 206 | |
| 207 | pub fn insert_new( |
| 208 | &mut self, |
| 209 | floc: &FileLocation, |
| 210 | ilen: usize, // encoded index length |
| 211 | ) |
| 212 | -> Outcome<usize> |
| 213 | { |
| 214 | self.dmap.insert(floc.start, DataState::Cur); |
| 215 | // If a supersession of this record arrived before the record itself (see `register_old` |
| 216 | // case (c)), apply the deferred flag now that the record is present. The parked length |
| 217 | // must match the record actually inserted here; a mismatch means the parked supersession |
| 218 | // referred to a different record at this position -- a genuine inconsistency, not a race. |
| 219 | if let Some(plen) = self.pending_old.remove(&floc.start) { |
| 220 | let rec_len = floc.klen + floc.vlen; |
| 221 | if plen != rec_len { |
| 222 | return Err(err!( |
| 223 | "A supersession parked for position {} expected a record of length {}, but \ |
| 224 | the record inserted there has length {}.", floc.start, plen, rec_len; |
| 225 | Bug, Mismatch, Data)); |
| 226 | } |
| 227 | match self.dmap.get_mut(&floc.start) { |
| 228 | Some(dstat) => *dstat = DataState::Old, |
| 229 | None => return Err(err!( |
| 230 | "The record just inserted at position {} vanished before its parked \ |
| 231 | supersession could be applied.", floc.start; |
| 232 | Bug, Missing, Data)), |
| 233 | } |
| 234 | match self.oldsum.checked_add(rec_len) { |
| 235 | Some(sum) => self.oldsum = sum, |
| 236 | None => { |
| 237 | self.oldsum = u64::MAX; |
| 238 | return Err(err!( |
| 239 | "Applying a parked supersession at position {} overflowed oldsum.", |
| 240 | floc.start; Bug, Overflow, Integer)); |
| 241 | }, |
| 242 | } |
| 243 | match self.oldcnt.checked_add(1) { |
| 244 | Some(sum) => self.oldcnt = sum, |
| 245 | None => { |
| 246 | self.oldcnt = usize::MAX; |
| 247 | return Err(err!( |
| 248 | "Applying a parked supersession at position {} overflowed oldcnt.", |
| 249 | floc.start; Bug, Overflow, Integer)); |
| 250 | }, |
| 251 | } |
| 252 | } |
| 253 | let dat_len = try_into!(usize, floc.klen + floc.vlen); |
| 254 | match self.dat_size.checked_add(dat_len) { |
| 255 | Some(sum) => self.dat_size = sum, |
| 256 | None => { |
| 257 | let old_dat_size = self.dat_size; |
| 258 | self.dat_size = usize::MAX; |
| 259 | return Err(err!( |
| 260 | "When inserting the new {:?}, the data file size, {}, overflowed. \ |
| 261 | It has been set to the maximum {}.", floc, old_dat_size, usize::MAX; |
| 262 | Bug, Overflow, Integer)); |
| 263 | }, |
| 264 | } |
| 265 | let ind_len = try_into!(usize, floc.klen) + ilen; |
| 266 | match self.ind_size.checked_add(ind_len) { |
| 267 | Some(sum) => self.ind_size = sum, |
| 268 | None => { |
| 269 | let old_ind_size = self.ind_size; |
| 270 | self.ind_size = usize::MAX; |
| 271 | return Err(err!( |
| 272 | "When inserting the new {:?}, the index file size, {}, overflowed. \ |
| 273 | It has been set to the maximum {}.", floc, old_ind_size, usize::MAX; |
| 274 | Bug, Overflow, Integer)); |
| 275 | }, |
| 276 | } |
| 277 | match dat_len.checked_add(ind_len) { |
| 278 | Some(sum) => Ok(sum), |
| 279 | None => Err(err!( |
| 280 | "When inserting the new {:?}, the sum of the data file size, {}, \ |
| 281 | and the index file size, {}, overflowed.", |
| 282 | floc, self.dat_size, self.ind_size; |
| 283 | Bug, Overflow, Integer)), |
| 284 | } |
| 285 | } |
| 286 | |
| 287 | pub fn inc_index_file_size( |
| 288 | &mut self, |
| 289 | len: usize, |
| 290 | ) |
| 291 | -> Outcome<usize> |
| 292 | { |
| 293 | match self.ind_size.checked_add(len) { |
| 294 | Some(sum) => self.ind_size = sum, |
| 295 | None => { |
| 296 | let old_ind_size = self.ind_size; |
| 297 | self.ind_size = usize::MAX; |
| 298 | return Err(err!( |
| 299 | "When incrementing the size of the index file, {}, by {}, \ |
| 300 | an overflow occurred. It has been set to the maximum {}.", |
| 301 | old_ind_size, len, usize::MAX; |
| 302 | Bug, Overflow, Integer)); |
| 303 | }, |
| 304 | } |
| 305 | Ok(len) |
| 306 | } |
| 307 | |
| 308 | /// Records that a collection carried the record `rid` names from `dloc` to `new_start`. |
| 309 | pub fn update_moved( |
| 310 | &mut self, |
| 311 | dloc: &DataLocation, |
| 312 | new_start: u64, |
| 313 | rid: RecordDigest, |
| 314 | ) { |
| 315 | self.mmap.insert(dloc.start, Move { to: new_start, rid }); |
| 316 | self.dmap.remove(&dloc.start); |
| 317 | } |
| 318 | |
| 319 | pub fn register_old( |
| 320 | &mut self, |
| 321 | dloc: &DataLocation, |
| 322 | ) |
| 323 | -> Outcome<()> |
| 324 | { |
| 325 | // A write's bytes reach the data file (in `WriterBot::write`) before its accounting |
| 326 | // does: the cbot acknowledges the caller and only then forwards the `UpdateData` that |
| 327 | // drives `insert_new`, so a supersession of a key can reach the fbot before the very |
| 328 | // record it supersedes has been inserted into this map. That is a general property of |
| 329 | // the write path, not a chunk peculiarity, but deterministic chunk keys make it routine: |
| 330 | // one overwrite supersedes a whole value's worth of same-keyed records at once, racing |
| 331 | // their sibling insertions. So a lookup miss here is not proof of a fault -- it may be a |
| 332 | // supersession that has merely overtaken its record. Three cases, kept distinct so a |
| 333 | // real accounting fault cannot hide behind a tolerant one: |
| 334 | // (a) the record is present and current -- flag it old, the ordinary path; |
| 335 | // (b) the record is present and already old -- the same supersession seen twice (a |
| 336 | // start is unique within a file generation, so an old entry here is provably the |
| 337 | // same record), absorbed without double counting; |
| 338 | // (c) the record is absent -- park the supersession and let `insert_new` apply it when |
| 339 | // the record lands. A park that never reconciles is caught as a hard error once |
| 340 | // the file has fully drained (see `schedule_deletion`), so a genuinely missing |
| 341 | // record -- the fault class that masked the 2026-07-28 rollover bug -- still fails |
| 342 | // loudly rather than being swallowed. |
| 343 | match self.dmap.get_mut(&dloc.start) { |
| 344 | Some(dstat @ DataState::Cur) => *dstat = DataState::Old, |
| 345 | // (b) Provable duplicate: same location, already accounted old. Nothing to do. |
| 346 | Some(DataState::Old) => return Ok(()), |
| 347 | // (c) The record has not been inserted yet; park until it is. |
| 348 | None => { |
| 349 | match self.pending_old.get(&dloc.start) { |
| 350 | Some(len) if *len == dloc.len => (), // already parked, same record |
| 351 | Some(len) => return Err(err!( |
| 352 | "Two different supersessions were parked for position {}: lengths {} \ |
| 353 | and {}. A start is unique within a file generation, so this is a genuine \ |
| 354 | accounting inconsistency, not a race.", dloc.start, len, dloc.len; |
| 355 | Bug, Mismatch, Data)), |
| 356 | None => { self.pending_old.insert(dloc.start, dloc.len); }, |
| 357 | } |
| 358 | return Ok(()); |
| 359 | }, |
| 360 | } |
| 361 | match self.oldsum.checked_add(dloc.len) { |
| 362 | Some(sum) => self.oldsum = sum, |
| 363 | None => { |
| 364 | let old_oldsum = self.oldsum; |
| 365 | self.oldsum = u64::MAX; |
| 366 | return Err(err!( |
| 367 | "When registering the old {:?}, the sum of old data sizes, \ |
| 368 | {} overflowed. It has been set to the maximum {}.", |
| 369 | dloc, old_oldsum, u64::MAX; |
| 370 | Bug, Overflow, Integer)); |
| 371 | }, |
| 372 | } |
| 373 | match self.oldcnt.checked_add(1) { |
| 374 | Some(sum) => self.oldcnt = sum, |
| 375 | None => { |
| 376 | let old_oldcnt = self.oldcnt; |
| 377 | self.oldcnt = usize::MAX; |
| 378 | return Err(err!( |
| 379 | "When registering the old {:?}, the count of old data entries, \ |
| 380 | {} overflowed. It has been set to the maximum {}.", |
| 381 | dloc, old_oldcnt, usize::MAX; |
| 382 | Bug, Overflow, Integer)); |
| 383 | }, |
| 384 | } |
| 385 | Ok(()) |
| 386 | } |
| 387 | |
| 388 | pub fn retire_old( |
| 389 | &mut self, |
| 390 | dloc: &DataLocation, |
| 391 | ) |
| 392 | -> Outcome<usize> |
| 393 | { |
| 394 | self.dmap.remove(&dloc.start); |
| 395 | let dat_len = try_into!(usize, dloc.len); |
| 396 | if self.dat_size >= dat_len { |
| 397 | self.dat_size -= dat_len; |
| 398 | } else { |
| 399 | return Err(err!( |
| 400 | "While retiring {:?} from {:?}, the data file size, {}, will become negative.", |
| 401 | dloc, self, self.dat_size; |
| 402 | Bug, Underflow, Integer)); |
| 403 | } |
| 404 | if self.oldsum >= dloc.len { |
| 405 | self.oldsum -= dloc.len; |
| 406 | } else { |
| 407 | return Err(err!( |
| 408 | "While retiring {:?} from {:?}, oldsum, {}, will become negative.", |
| 409 | dloc, self, self.oldsum; |
| 410 | Bug, Underflow, Integer)); |
| 411 | } |
| 412 | if self.oldcnt > 0 { |
| 413 | self.oldcnt -= 1; |
| 414 | } else { |
| 415 | return Err(err!( |
| 416 | "While retiring {:?} from {:?}, oldcnt, {}, will become negative.", |
| 417 | dloc, self, self.oldcnt; |
| 418 | Bug, Underflow, Integer)); |
| 419 | } |
| 420 | Ok(dat_len) |
| 421 | } |
| 422 | |
| 423 | /// Where the record `rid` names, last seen at `dloc` before a collection, now starts, if the |
| 424 | /// collection moved it and a supersession has still to find it there. A move entry is |
| 425 | /// matched by the record as well as the offset: with records of one size, a new offset can |
| 426 | /// equal an old one whose move is still pending, and matched by offset alone a read of the |
| 427 | /// record now at that offset was handed the other record (2026-09-23). A read only looks: |
| 428 | /// the entry stays for the supersession it is kept for. |
| 429 | pub fn moved_to( |
| 430 | &self, |
| 431 | dloc: &DataLocation, |
| 432 | rid: &RecordDigest, |
| 433 | ) |
| 434 | -> Option<u64> |
| 435 | { |
| 436 | match self.mmap.get(&dloc.start) { |
| 437 | Some(mv) if mv.rid == *rid => Some(mv.to), |
| 438 | _ => None, |
| 439 | } |
| 440 | } |
| 441 | |
| 442 | /// As `moved_to`, and the entry is then spent: its record is current at its new start until |
| 443 | /// the caller says otherwise. For a supersession, and for the collector re-anchoring a record |
| 444 | /// in its cache. An entry of another record at the same offset is left alone. |
| 445 | pub fn map_and_remove( |
| 446 | &mut self, |
| 447 | dloc: &DataLocation, |
| 448 | rid: &RecordDigest, |
| 449 | ) |
| 450 | -> Option<u64> |
| 451 | { |
| 452 | let to = match self.moved_to(dloc, rid) { |
| 453 | Some(to) => to, |
| 454 | None => return None, |
| 455 | }; |
| 456 | self.mmap.remove(&dloc.start); |
| 457 | self.dmap.insert(to, DataState::Cur); |
| 458 | Some(to) |
| 459 | } |
| 460 | } |
| 461 | |
| 462 | /// A portion, or shard of the FileState data for a zone. |
| 463 | #[derive(Clone, Debug, Default)] |
| 464 | pub struct FileStateMap { |
| 465 | map: BTreeMap<FileNum, FileState>, |
| 466 | size: usize, // Sum of all data and index file sizes for this shard. |
| 467 | } |
| 468 | |
| 469 | impl FileStateMap { |
| 470 | pub fn map(&self) -> &BTreeMap<FileNum, FileState> { &self.map } |
| 471 | pub fn map_mut(&mut self) -> &mut BTreeMap<FileNum, FileState> { &mut self.map } |
| 472 | |
| 473 | pub fn get_state(&self, fnum: FileNum) -> Outcome<&FileState> { |
| 474 | match self.map.get(&fnum) { |
| 475 | Some(fstat) => Ok(fstat), |
| 476 | None => Err(err!( |
| 477 | "No state entry for file number {}.", fnum; |
| 478 | Bug, Missing, Data)), |
| 479 | } |
| 480 | } |
| 481 | |
| 482 | pub fn get_state_mut(&mut self, fnum: FileNum) -> Outcome<&mut FileState> { |
| 483 | match self.map.get_mut(&fnum) { |
| 484 | Some(fstat) => Ok(fstat), |
| 485 | None => Err(err!( |
| 486 | "No state entry for file number {}.", fnum; |
| 487 | Bug, Missing, Data)), |
| 488 | } |
| 489 | } |
| 490 | |
| 491 | pub fn set_size(&mut self, size: usize) { self.size = size; } |
| 492 | pub fn inc_size(&mut self, len: usize) -> Outcome<()> { |
| 493 | self.size = try_add!(&self.size, len); |
| 494 | Ok(()) |
| 495 | } |
| 496 | pub fn dec_size(&mut self, len: usize) -> Outcome<()> { |
| 497 | self.size = try_sub!(&self.size, len); |
| 498 | Ok(()) |
| 499 | } |
| 500 | pub fn get_size(&self) -> usize { self.size } |
| 501 | |
| 502 | #[inline] |
| 503 | pub fn shard_index(fnum: FileNum, nf: usize) -> usize { (fnum as usize) % nf } |
| 504 | |
| 505 | pub fn new_file_state(&mut self, fnum: FileNum, fs: FileState) { |
| 506 | self.map.insert(fnum, fs); |
| 507 | } |
| 508 | |
| 509 | pub fn new_live_file( |
| 510 | &mut self, |
| 511 | num: FileNum, |
| 512 | dat_size: u64, |
| 513 | ind_size: u64, |
| 514 | ) { |
| 515 | match self.map.get_mut(&num) { |
| 516 | Some(fstat) => fstat.set_live(true), |
| 517 | None => self.new_file_state(num, FileState { |
| 518 | dat_size: dat_size as usize, |
| 519 | ind_size: ind_size as usize, |
| 520 | live: true, |
| 521 | ..Default::default() |
| 522 | }), |
| 523 | } |
| 524 | } |
| 525 | |
| 526 | pub fn insert_new( |
| 527 | &mut self, |
| 528 | floc: &FileLocation, |
| 529 | ilen: usize, |
| 530 | ) |
| 531 | -> Outcome<()> |
| 532 | { |
| 533 | if self.map.get(&floc.file_number()).is_none() { |
| 534 | self.new_file_state(floc.file_number(), FileState::default()); |
| 535 | } |
| 536 | let fstat = ok!(self.get_state_mut(floc.file_number())); |
| 537 | let len = res!(fstat.insert_new(&floc, ilen)); |
| 538 | self.inc_size(len) |
| 539 | } |
| 540 | } |