| 522 | } |
| 523 | |
| 524 | pub fn from_versioned( |
| 525 | buffer_size: usize, |
| 526 | root_path: &Path, |
| 527 | mapping_fn: impl Fn(&Path) -> Option<(VersionNumber, u64)>, |
| 528 | latest_version: VersionNumber, |
| 529 | ) -> Result<Self, BufIoError> { |
| 530 | let regions = DashMap::new(); |
| 531 | let mut file_size = 0; |
| 532 | let mut region_versions_map = FxHashMap::<u64, (VersionNumber, PathBuf)>::default(); |
| 533 | |
| 534 | for entry in fs::read_dir(root_path)? { |
| 535 | let entry = entry?; |
| 536 | let path = entry.path(); |
| 537 | let Some((version, region_id)) = mapping_fn(&path) else { |
| 538 | continue; |
| 539 | }; |
| 540 | if *version > *latest_version { |
| 541 | continue; |
| 542 | } |
| 543 | if let Some((existing_version, _)) = region_versions_map.get(®ion_id) { |
| 544 | if **existing_version > *version { |
| 545 | continue; |
| 546 | } |
| 547 | } |
| 548 | region_versions_map.insert(region_id, (version, path)); |
| 549 | } |
| 550 | |
| 551 | for (region_id, (_, path)) in region_versions_map { |
| 552 | let offset = region_id * buffer_size as u64; |
| 553 | let mut file = OpenOptions::new().read(true).open(path)?; |
| 554 | let mut region = FilelessBufferRegion::new(offset, buffer_size); |
| 555 | let buffer = region.buffer.get_mut().map_err(|_| BufIoError::Locking)?; |
| 556 | let bytes_read = file.read(&mut buffer[..]).map_err(BufIoError::Io)?; |
| 557 | region.end.store(bytes_read, Ordering::SeqCst); |
| 558 | regions.insert(offset, Arc::new(region)); |
| 559 | let end = offset + bytes_read as u64; |
| 560 | file_size = file_size.max(end); |
| 561 | } |
| 562 | |
| 563 | Ok(Self { |
| 564 | regions, |
| 565 | cursors: RwLock::new(HashMap::new()), |
| 566 | next_cursor_id: AtomicU64::new(0), |
| 567 | file_size: RwLock::new(file_size), |
| 568 | buffer_size, |
| 569 | }) |
| 570 | } |
| 571 | |
| 572 | pub fn open_cursor(&self) -> Result<u64, BufIoError> { |
| 573 | let cursor_id = self.next_cursor_id.fetch_add(1, Ordering::SeqCst); |