MCPcopy Create free account
hub / github.com/cosdata/cosdata / from_versioned

Method from_versioned

src/models/buffered_io.rs:524–570  ·  view source on GitHub ↗
(
        buffer_size: usize,
        root_path: &Path,
        mapping_fn: impl Fn(&Path) -> Option<(VersionNumber, u64)>,
        latest_version: VersionNumber,
    )

Source from the content-addressed store, hash-verified

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(&region_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);

Callers

nothing calls this directly

Calls 4

maxMethod · 0.80
getMethod · 0.45
insertMethod · 0.45
readMethod · 0.45

Tested by

no test coverage detected