From 690d92e751de7b44bbcf453ff4e147a04421fab1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Luis=20Coss=C3=ADo?= Date: Tue, 8 Sep 2026 16:56:34 -0300 Subject: [PATCH] [updater] genericize fs to use `UniversalAppendFs` (#10451) * introduce UniversalAppendFs helper * AI: migrate to UniversalAppendFs bound Co-Authored-By: Claude Fable 5 * don't use it in Gridstore --------- Co-authored-by: Claude Fable 5 --- .../src/blobstore/gridstore/bitmask/gaps.rs | 10 ++- .../src/blobstore/gridstore/bitmask/mod.rs | 8 ++- lib/blobstore/src/blobstore/logstore/mod.rs | 47 +++++++++----- lib/blobstore/src/blobstore/logstore/page.rs | 19 ++++-- lib/blobstore/src/tracker/append_only.rs | 10 ++- lib/blobstore/src/tracker/mod.rs | 6 +- .../common/src/universal_io/cached_fs/mod.rs | 63 ++++++++++++++++--- lib/common/common/src/universal_io/mod.rs | 6 +- .../common/src/universal_io/traits/append.rs | 15 ++++- .../common/src/universal_io/traits/mod.rs | 2 +- lib/edge/src/update_only/apply.rs | 10 +-- lib/edge/src/update_only/lifecycle.rs | 62 +++++++++--------- lib/edge/src/update_only/mod.rs | 15 +++-- lib/edge/src/update_only/preview.rs | 4 +- lib/edge/src/update_only/tests.rs | 32 +++++----- lib/edge/tools/shard_update/src/main.rs | 22 +++---- .../src/common/update_only_blobstore.rs | 6 +- .../id_tracker/disk_id_tracker/update_only.rs | 17 ++--- .../delete_only_tracker_enum.rs | 10 +-- .../immutable_id_tracker/deleted_storage.rs | 8 +-- .../immutable_id_tracker/update_only.rs | 17 ++--- .../mutable_id_tracker/update_only/heal.rs | 6 +- .../mutable_id_tracker/update_only/mod.rs | 19 +++--- .../field_index_base/update_only/mod.rs | 12 ++-- .../field_index_base/update_only/writer.rs | 10 ++- .../struct_payload_index/update_only/mod.rs | 9 ++- .../src/payload_storage/update_only/mod.rs | 9 ++- .../src/segment/update_only/appendable/mod.rs | 30 ++++----- .../segment/update_only/delete_only/mod.rs | 14 ++--- .../segment/update_only/lookup/lifecycle.rs | 4 +- .../segment/update_only/lookup/live_reload.rs | 4 +- .../src/segment/update_only/segment_enum.rs | 18 +++--- .../update_only.rs | 29 +++++---- .../quantized/update_only/mod.rs | 18 +++--- .../quantized/update_only/tests.rs | 8 +-- .../vector_storage/sparse/update_only/mod.rs | 6 +- .../src/vector_storage/update_only/mod.rs | 15 +++-- 37 files changed, 351 insertions(+), 249 deletions(-) diff --git a/lib/blobstore/src/blobstore/gridstore/bitmask/gaps.rs b/lib/blobstore/src/blobstore/gridstore/bitmask/gaps.rs index 3d18b9d7c2..bc5763beda 100644 --- a/lib/blobstore/src/blobstore/gridstore/bitmask/gaps.rs +++ b/lib/blobstore/src/blobstore/gridstore/bitmask/gaps.rs @@ -3,7 +3,9 @@ use std::ops::Range; use std::path::{Path, PathBuf}; use common::mmap::{Advice, AdviceSetting, create_and_ensure_length}; -use common::universal_io::{Flusher, OpenOptions, Populate, TypedStorage, UniversalWrite}; +use common::universal_io::{ + Flusher, OpenOptions, Populate, TypedStorage, UniversalReadFs, UniversalWrite, +}; use itertools::Itertools; use super::{GridstoreConfig, RegionId}; @@ -122,7 +124,11 @@ impl BitmaskGaps { }) } - pub fn open(fs: &S::Fs, dir: &Path, config: GridstoreConfig) -> Result { + pub fn open( + fs: &impl UniversalReadFs, + dir: &Path, + config: GridstoreConfig, + ) -> Result { let path = gaps_file_path(dir); let options = OpenOptions { writeable: true, diff --git a/lib/blobstore/src/blobstore/gridstore/bitmask/mod.rs b/lib/blobstore/src/blobstore/gridstore/bitmask/mod.rs index 60672d5c4a..c9b555b03a 100644 --- a/lib/blobstore/src/blobstore/gridstore/bitmask/mod.rs +++ b/lib/blobstore/src/blobstore/gridstore/bitmask/mod.rs @@ -8,7 +8,7 @@ use common::bitvec::BitSlice; use common::mmap::{Advice, AdviceSetting, create_and_ensure_length}; use common::stored_bitslice::StoredBitSlice; use common::universal_io::{ - MmapFile, OpenOptions, Populate, UniversalWrite, UniversalWriteFileOps, + MmapFile, OpenOptions, Populate, UniversalReadFs, UniversalWrite, UniversalWriteFileOps, }; use gaps::{BitmaskGaps, RegionGaps}; use itertools::Itertools; @@ -115,7 +115,11 @@ impl Bitmask { }) } - pub(crate) fn open(fs: &S::Fs, dir: &Path, config: GridstoreConfig) -> Result { + pub(crate) fn open( + fs: &impl UniversalReadFs, + dir: &Path, + config: GridstoreConfig, + ) -> Result { debug_assert!( config .page_size_bytes diff --git a/lib/blobstore/src/blobstore/logstore/mod.rs b/lib/blobstore/src/blobstore/logstore/mod.rs index 543c4f1d52..b082208bd9 100644 --- a/lib/blobstore/src/blobstore/logstore/mod.rs +++ b/lib/blobstore/src/blobstore/logstore/mod.rs @@ -15,7 +15,8 @@ use common::counter::referenced_counter::HwMetricRefCounter; use common::generic_consts::{AccessPattern, Sequential}; use common::is_alive_lock::IsAliveLock; use common::universal_io::{ - OkNotFound as _, Populate, UniversalAppend, UniversalRead, UniversalWriteFileOps, UserData, + OkNotFound as _, Populate, UniversalAppend, UniversalAppendFs, UniversalRead, UniversalReadFs, + UniversalWriteFileOps, UserData, }; use page::AppendOnlyPages; use parking_lot::RwLock; @@ -121,7 +122,10 @@ where /// /// `base_path` is the directory where the storage files will be stored. /// It should exist already. - pub(super) fn new(fs: &S::Fs, base_path: PathBuf, config: LogstoreConfig) -> Result { + pub(super) fn new(fs: &Fs, base_path: PathBuf, config: LogstoreConfig) -> Result + where + Fs: UniversalWriteFileOps + UniversalReadFs, + { let tracker = AppendOnlyTracker::new(fs, &base_path)?; let pages = AppendOnlyPages::new(fs, &base_path)?; @@ -146,7 +150,7 @@ where /// In case of opening, it ignores the `config_if_create` parameter. A storage created in /// mutable mode is rejected, the two modes persist incompatible file formats. pub fn open_or_create( - fs: &S::Fs, + fs: &impl UniversalAppendFs, base_path: PathBuf, config_if_create: LogstoreConfig, populate: Populate, @@ -165,12 +169,15 @@ where } /// Open an existing storage at the given path, with the already read config. - pub(super) fn open( - fs: &S::Fs, + pub(super) fn open( + fs: &Fs, base_path: PathBuf, config: LogstoreConfig, populate: Populate, - ) -> Result { + ) -> Result + where + Fs: UniversalWriteFileOps + UniversalReadFs, + { let tracker = AppendOnlyTracker::open_writable(fs, &base_path, populate)?; let pages = AppendOnlyPages::open(fs, &base_path, true, populate)?; validate_consistency(&tracker, &pages)?; @@ -204,13 +211,16 @@ where /// rejected, the storage is append-only. /// /// Always returns false on success, as values can never be updated. - pub fn put_value( + pub fn put_value( &mut self, - fs: &S::Fs, + fs: &Fs, point_offset: PointOffset, value: &V, hw_counter: HwMetricRefCounter, - ) -> Result { + ) -> Result + where + Fs: UniversalWriteFileOps + UniversalReadFs, + { self.put_value_bytes(fs, point_offset, value.to_bytes(), hw_counter) } @@ -222,13 +232,16 @@ where /// See [`put_value`](Self::put_value) for buffering and append-only semantics. // Takes &mut self for signature parity with the mutable variant #[allow(clippy::needless_pass_by_ref_mut)] - pub(super) fn put_value_bytes( + pub(super) fn put_value_bytes( &mut self, - fs: &S::Fs, + fs: &Fs, point_offset: PointOffset, value_bytes: Vec, hw_counter: HwMetricRefCounter, - ) -> Result { + ) -> Result + where + Fs: UniversalWriteFileOps + UniversalReadFs, + { // Validate before buffering anything, a rejected put must not leave data behind let next = self.tracker.read().pointer_count(); if point_offset < next { @@ -270,7 +283,10 @@ where /// Clear the storage, going back to the initial state. /// /// Completely wipes the storage, and recreates it in append-only mode. - pub(super) fn clear(&mut self, fs: &S::Fs) -> Result<()> { + pub(super) fn clear(&mut self, fs: &Fs) -> Result<()> + where + Fs: UniversalWriteFileOps + UniversalReadFs, + { self.is_alive_flush_lock.blocking_mark_dead(); fs.remove_dir(&self.base_path)?; @@ -286,7 +302,10 @@ where /// Takes ownership because this function leaves the storage in an inconsistent state which /// does not allow further usage. Use [`clear`](Self::clear) instead to clear and reuse the /// storage. - pub(super) fn wipe(self, fs: &S::Fs) -> Result<()> { + pub(super) fn wipe(self, fs: &Fs) -> Result<()> + where + Fs: UniversalWriteFileOps, + { let Self { config: _, tracker, diff --git a/lib/blobstore/src/blobstore/logstore/page.rs b/lib/blobstore/src/blobstore/logstore/page.rs index e73f2fee5c..36ddde3732 100644 --- a/lib/blobstore/src/blobstore/logstore/page.rs +++ b/lib/blobstore/src/blobstore/logstore/page.rs @@ -305,7 +305,10 @@ impl AppendOnlyPages { /// Create a new empty page 0 in the given directory, truncating it if it already exists. /// /// The directory must exist already. - pub(super) fn new(fs: &S::Fs, dir: &Path) -> Result { + pub(super) fn new(fs: &Fs, dir: &Path) -> Result + where + Fs: UniversalWriteFileOps + UniversalReadFs, + { let page = AppendOnlyPage::new(fs, page_file_name(dir, 0))?; Ok(Self { dir: dir.to_path_buf(), @@ -322,12 +325,15 @@ impl AppendOnlyPages { /// /// The value is buffered in memory until the next flush; only a page rollover touches disk /// by creating the new, empty page file. - pub(super) fn append_value( + pub(super) fn append_value( &mut self, - fs: &S::Fs, + fs: &Fs, value: &[u8], page_capacity_bytes: u64, - ) -> Result<(PageId, BlockOffset)> { + ) -> Result<(PageId, BlockOffset)> + where + Fs: UniversalWriteFileOps + UniversalReadFs, + { let last = self .pages .last() @@ -526,7 +532,10 @@ impl AppendOnlyPage { /// Create a new empty page file at the given path, truncating it if it already exists. /// /// The directory must exist already. - fn new(fs: &S::Fs, path: PathBuf) -> Result { + fn new(fs: &Fs, path: PathBuf) -> Result + where + Fs: UniversalWriteFileOps + UniversalReadFs, + { fs.create(&path, 0)?; let file = fs.open( &path, diff --git a/lib/blobstore/src/tracker/append_only.rs b/lib/blobstore/src/tracker/append_only.rs index afec7e7cdc..54275c1d1f 100644 --- a/lib/blobstore/src/tracker/append_only.rs +++ b/lib/blobstore/src/tracker/append_only.rs @@ -309,7 +309,10 @@ impl AppendOnlyTracker { /// exists. /// /// The directory must exist already. - pub fn new(fs: &S::Fs, dir: &Path) -> Result { + pub fn new(fs: &Fs, dir: &Path) -> Result + where + Fs: UniversalWriteFileOps + UniversalReadFs, + { let path = Self::tracker_file_name(dir); fs.create(&path, 0)?; let file = fs.open( @@ -331,7 +334,10 @@ impl AppendOnlyTracker { /// /// A trailing partial entry due to a torn write is truncated away, so that appends always /// start at a whole entry offset. - pub fn open_writable(fs: &S::Fs, dir: &Path, populate: Populate) -> Result { + pub fn open_writable(fs: &Fs, dir: &Path, populate: Populate) -> Result + where + Fs: UniversalWriteFileOps + UniversalReadFs, + { let path = Self::tracker_file_name(dir); let mut file = Self::open_file(fs, &path, populate, true)?; diff --git a/lib/blobstore/src/tracker/mod.rs b/lib/blobstore/src/tracker/mod.rs index f69d341cc4..fe6f06465c 100644 --- a/lib/blobstore/src/tracker/mod.rs +++ b/lib/blobstore/src/tracker/mod.rs @@ -455,7 +455,11 @@ where /// Create a new PageTracker at the given dir path /// The file is created with the default size if no size hint is given - pub fn new(fs: &S::Fs, path: &Path, size_hint: Option) -> Result { + pub fn new( + fs: &impl UniversalReadFs, + path: &Path, + size_hint: Option, + ) -> Result { let path = Self::tracker_file_name(path); let size = size_hint.unwrap_or(Self::DEFAULT_SIZE).next_power_of_two(); assert!( diff --git a/lib/common/common/src/universal_io/cached_fs/mod.rs b/lib/common/common/src/universal_io/cached_fs/mod.rs index a6f943f0b7..eeb494ba0d 100644 --- a/lib/common/common/src/universal_io/cached_fs/mod.rs +++ b/lib/common/common/src/universal_io/cached_fs/mod.rs @@ -12,10 +12,9 @@ use parking_lot::Mutex; mod async_io; use crate::mmap::AdviceSetting; -use crate::universal_io::traits::CachedReadFs; use crate::universal_io::{ - ListedFile, OpenExtra, OpenOptions, Populate, UioResult, UniversalIoError, - UniversalReadFileOps, UniversalReadFs, UniversalReadFsAsync, + CachedReadFs, ListedFile, OpenExtra, OpenOptions, Populate, UioResult, UniversalIoError, + UniversalReadFileOps, UniversalReadFs, UniversalReadFsAsync, UniversalWriteFileOps, }; #[derive(Clone, Debug)] @@ -53,10 +52,15 @@ impl FileInfo { } } -/// Read-only filesystem wrapper that snapshots the file listing and serves +/// Filesystem wrapper that snapshots the file listing and serves read-only /// opens from explicitly prefetched handles. The only [`CachedReadFs`] /// implementation. /// +/// The write side is a passthrough: writable opens and +/// [`UniversalWriteFileOps`] forward to the wrapped filesystem and do NOT +/// update the snapshot — reads that must observe a post-snapshot mutation +/// need a fresh [`CachedFs::cache_file_info`]. +/// /// Opens produce the *wrapped* backend's file type (`Fs::File`), so /// components generic over `impl UniversalReadFs` accept a raw /// backend and this wrapper interchangeably, and stored handle types never @@ -397,6 +401,47 @@ impl Debug for CachedFs { } } +impl UniversalWriteFileOps for CachedFs +where + Fs: UniversalReadFs + UniversalWriteFileOps, +{ + type AppendFile = Fs::AppendFile; + + fn create(&self, path: &Path, expected_length: usize) -> UioResult<()> { + self.fs.create(path, expected_length) + } + + fn create_dir(&self, path: &Path) -> UioResult<()> { + self.fs.create_dir(path) + } + + fn remove(&self, path: &Path) -> UioResult<()> { + // Drop any prefetched handle so a pooled open cannot resurrect the + // removed file. + self.files_prefetched.lock().remove(path); + self.fs.remove(path) + } + + fn remove_dir(&self, path: &Path) -> UioResult<()> { + self.files_prefetched + .lock() + .retain(|pooled, _| !pooled.starts_with(path)); + self.fs.remove_dir(path) + } + + fn atomic_save(&self, path: &Path, bytes: &[u8]) -> UioResult<()> { + self.fs.atomic_save(path, bytes) + } + + fn open_append( + &self, + path: impl AsRef, + options: OpenOptions, + ) -> UioResult { + self.fs.open_append(path, options) + } +} + impl UniversalReadFs for CachedFs { /// The *wrapped* backend's file type: opening through the cache hands /// out the very handles the inner filesystem produced (prefetched or @@ -412,12 +457,12 @@ impl UniversalReadFs for CachedFs { ) -> UioResult { let path = path.as_ref(); + // Writable opens bypass the cache machinery entirely: prefetched + // handles are read-only, the snapshot's size may lag appends (a + // writer must observe the true length), and the snapshot cannot know + // files this writer created after the LIST. if options.writeable { - return Err(UniversalIoError::Uninitialized { - description: - "CachedReadFs only supports read-only files, writeable option is not allowed" - .to_string(), - }); + return self.fs.open(path, options, extra); } if let Some(file) = self.files_prefetched.lock().remove(path) { diff --git a/lib/common/common/src/universal_io/mod.rs b/lib/common/common/src/universal_io/mod.rs index d82121a4fb..92afae5dac 100644 --- a/lib/common/common/src/universal_io/mod.rs +++ b/lib/common/common/src/universal_io/mod.rs @@ -33,9 +33,9 @@ pub use self::simple_disk_cache::{ }; pub use self::sorted_block_index::SortedBlockIndex; pub use self::traits::{ - CachedReadFs, Item, OpenExtra, OwnedPipeline, ReadPipeline, UniversalAppend, UniversalFlush, - UniversalRead, UniversalReadAsync, UniversalReadFileOps, UniversalReadFs, UniversalReadFsAsync, - UniversalWrite, UniversalWriteFileOps, UserData, + CachedReadFs, Item, OpenExtra, OwnedPipeline, ReadPipeline, UniversalAppend, UniversalAppendFs, + UniversalFlush, UniversalRead, UniversalReadAsync, UniversalReadFileOps, UniversalReadFs, + UniversalReadFsAsync, UniversalWrite, UniversalWriteFileOps, UserData, }; pub use self::types::{ ByteOffset, FileIndex, Flusher, ListedFile, OpenOptions, Populate, ReadBytesItem, ReadRange, diff --git a/lib/common/common/src/universal_io/traits/append.rs b/lib/common/common/src/universal_io/traits/append.rs index f912cfaa53..453aed891d 100644 --- a/lib/common/common/src/universal_io/traits/append.rs +++ b/lib/common/common/src/universal_io/traits/append.rs @@ -1,5 +1,5 @@ use super::{UniversalFlush, UniversalRead, UniversalWriteFileOps}; -use crate::universal_io::{ByteOffset, UioResult}; +use crate::universal_io::{ByteOffset, UioResult, UniversalReadFsAsync}; /// A file handle that supports atomic appends at a caller-provided offset. /// @@ -95,3 +95,16 @@ pub trait UniversalAppend: // When adding provided methods, don't forget to update impls in // crate::universal_io::wrappers::*. } + +pub trait UniversalAppendFs: + UniversalReadFsAsync::AppendFile> + + UniversalWriteFileOps +{ +} + +impl UniversalAppendFs for Fs +where + Fs: UniversalReadFsAsync + UniversalWriteFileOps, + Self::File: UniversalAppend, +{ +} diff --git a/lib/common/common/src/universal_io/traits/mod.rs b/lib/common/common/src/universal_io/traits/mod.rs index c1085f7d92..187ede268c 100644 --- a/lib/common/common/src/universal_io/traits/mod.rs +++ b/lib/common/common/src/universal_io/traits/mod.rs @@ -8,7 +8,7 @@ mod write; use std::fmt; -pub use append::UniversalAppend; +pub use append::{UniversalAppend, UniversalAppendFs}; pub use async_io::{UniversalReadAsync, UniversalReadFsAsync}; pub use file_ops::{CachedReadFs, UniversalReadFileOps, UniversalReadFs, UniversalWriteFileOps}; pub use open_extra::OpenExtra; diff --git a/lib/edge/src/update_only/apply.rs b/lib/edge/src/update_only/apply.rs index 1d0d60f19c..8f0b2e2ab2 100644 --- a/lib/edge/src/update_only/apply.rs +++ b/lib/edge/src/update_only/apply.rs @@ -7,7 +7,7 @@ use std::collections::HashMap; use ahash::AHashMap; use common::counter::hardware_counter::HardwareCounterCell; use common::types::PointOffsetType; -use common::universal_io::{UniversalAppend, UniversalRead}; +use common::universal_io::{UniversalAppendFs, UniversalRead}; use rayon::ThreadPool; use rayon::prelude::*; use segment::common::operation_error::{OperationError, OperationResult}; @@ -123,7 +123,7 @@ pub(super) struct PointLocations { pub(super) slots: Vec<(Uuid, PointOffsetType)>, } -impl UpdateOnlyEdgeShard { +impl UpdateOnlyEdgeShard { /// Apply a batch of update operations, each paired with the operation /// number to record as its version. Operations are expected in ascending /// operation-number order; see [`UpdateBatchPlan::build`] for what is @@ -306,10 +306,10 @@ impl UpdateOnlyEdgeShard { /// The held writer for segment `uuid`; an error when there is none, which can /// only mean the shard's inventory changed under a batch in flight. -fn get_writer( - writers: &mut HashMap>, +fn get_writer( + writers: &mut HashMap>, uuid: Uuid, -) -> OperationResult<&mut UpdateOnlySegmentEnum> { +) -> OperationResult<&mut UpdateOnlySegmentEnum> { writers .get_mut(&uuid) .ok_or_else(|| OperationError::service_error(format!("No writer open for segment {uuid}"))) diff --git a/lib/edge/src/update_only/lifecycle.rs b/lib/edge/src/update_only/lifecycle.rs index 480b6e5f3c..353d927943 100644 --- a/lib/edge/src/update_only/lifecycle.rs +++ b/lib/edge/src/update_only/lifecycle.rs @@ -4,8 +4,8 @@ use std::path::Path; use common::counter::hardware_counter::HardwareCounterCell; use common::mmap::AdviceSetting; use common::universal_io::{ - MmapFile, MmapFs, OpenOptions, Populate, UniversalAppend, UniversalFlush as _, UniversalReadFs, - UniversalReadFsAsync, UniversalWriteFileOps, + MmapFs, OpenOptions, Populate, UniversalAppend, UniversalAppendFs, UniversalFlush as _, + UniversalWriteFileOps, }; use parking_lot::RwLock; use rayon::prelude::*; @@ -23,7 +23,7 @@ use crate::read_view::build_segment_pool; use crate::update_only::UpdateOnlyEdgeShard; use crate::update_only::holder::LookupSegmentHolder; -impl UpdateOnlyEdgeShard { +impl UpdateOnlyEdgeShard { /// Open a writer over local memory-mapped files, discovering segments by /// scanning the `segments/` directory — the writer owns the directory it /// writes to, so there is no manifest to agree with. @@ -32,7 +32,7 @@ impl UpdateOnlyEdgeShard { } } -impl UpdateOnlyEdgeShard { +impl UpdateOnlyEdgeShard { /// Open a writer over the shard directory at `path`, using `fs` as the /// backend and `enumerator` to discover the segments. /// @@ -46,13 +46,10 @@ impl UpdateOnlyEdgeShard { /// is an error, not a skip — a writer that misses a segment would resolve /// a point against a stale copy of itself, or duplicate it. pub fn open( - fs: S::Fs, + fs: Fs, path: &Path, enumerator: impl SegmentEnumerator + 'static, - ) -> OperationResult - where - S::Fs: UniversalReadFs + UniversalReadFsAsync, - { + ) -> OperationResult { // Sized like the search pools: over-provisioned relative to the CPU // count, since on a remote backend the threads mostly wait on IO. let pool = build_segment_pool( @@ -63,26 +60,25 @@ impl UpdateOnlyEdgeShard { let segments: Vec<(Uuid, ListedSegment)> = enumerator.list_segments()?.into_iter().collect(); - let opened: Vec<(Uuid, LookupSegment, UpdateOnlySegmentEnum, bool)> = - pool.install(|| { - segments - .into_par_iter() - .map(|(uuid, listing)| { - let ListedSegment { path, writable } = listing; - // No deferred threshold yet: it belongs to the coordination - // with an external rebuilder, which does not exist in this - // iteration. - let segment = LookupSegment::::open(&fs, &path, None)?; - let writer = UpdateOnlySegmentEnum::open( - &fs, - &path, - &segment.segment_config, - segment.writer_state(), - )?; - Ok((uuid, segment, writer, writable)) - }) - .collect::>>() - })?; + let opened: Vec<_> = pool.install(|| { + segments + .into_par_iter() + .map(|(uuid, listing)| { + let ListedSegment { path, writable } = listing; + // No deferred threshold yet: it belongs to the coordination + // with an external rebuilder, which does not exist in this + // iteration. + let segment = LookupSegment::open(&fs, &path, None)?; + let writer = UpdateOnlySegmentEnum::open( + fs.clone(), + &path, + &segment.segment_config, + segment.writer_state(), + )?; + Ok((uuid, segment, writer, writable)) + }) + .collect::>>() + })?; let mut holder = LookupSegmentHolder::default(); let mut writers = HashMap::new(); @@ -101,9 +97,9 @@ impl UpdateOnlyEdgeShard { } } -impl UpdateOnlyEdgeShard +impl UpdateOnlyEdgeShard where - S::Fs: UniversalReadFs + UniversalReadFsAsync, + Fs: UniversalAppendFs, { /// [`create_appendable`](Self::create_appendable) with `source`'s config. #[cfg(test)] @@ -155,9 +151,9 @@ where let remote = self.path.join(SEGMENTS_PATH).join(uuid.to_string()); copy_dir_via(&self.fs, &local, &remote)?; - let lookup = LookupSegment::::open(&self.fs, &remote, None)?; + let lookup = LookupSegment::::open(&self.fs, &remote, None)?; let writer = UpdateOnlySegmentEnum::open( - &self.fs, + self.fs.clone(), &remote, &lookup.segment_config, lookup.writer_state(), diff --git a/lib/edge/src/update_only/mod.rs b/lib/edge/src/update_only/mod.rs index 94045cae25..03b321ba1e 100644 --- a/lib/edge/src/update_only/mod.rs +++ b/lib/edge/src/update_only/mod.rs @@ -36,7 +36,7 @@ use std::collections::HashMap; use std::path::{Path, PathBuf}; use std::sync::Arc; -use common::universal_io::UniversalAppend; +use common::universal_io::UniversalAppendFs; use parking_lot::RwLock; use rayon::ThreadPool; use segment::segment::update_only::UpdateOnlySegmentEnum; @@ -48,23 +48,22 @@ pub use self::batch::{PointUpdates, UpdateBatchPlan}; use self::holder::LookupSegmentHolder; pub use self::preview::{PointAction, PointCopy, PointPreview, UpdateBatchPreview}; -/// A batch writer over the segments of one shard directory, generic over the -/// backend `S`. +/// A batch writer over the segments of one shard directory. /// /// Compared to [`EdgeShard`](crate::EdgeShard), there is no WAL, no /// optimizers, and no `EdgeConfig` — the write target's own segment config is /// the only configuration a write needs. -pub struct UpdateOnlyEdgeShard { +pub struct UpdateOnlyEdgeShard { path: PathBuf, /// Backend the segments were opened on; live-reloads their lookup /// halves after a batch writes to them. - fs: S::Fs, - segments: RwLock>, + fs: Fs, + segments: RwLock>, /// One writer per segment, opened at shard open from the state its /// [`LookupSegment`](segment::segment::update_only::LookupSegment) /// observed — which also decided whether the segment accepts appends or /// deletes only. - writers: HashMap>, + writers: HashMap>, /// Thread pool the per-segment work of a batch runs on: on a remote /// backend each segment's reads block on the network, so segments are /// visited in parallel. @@ -81,7 +80,7 @@ pub struct SegmentConfigInfo { pub config: SegmentConfig, } -impl UpdateOnlyEdgeShard { +impl UpdateOnlyEdgeShard { pub fn path(&self) -> &Path { &self.path } diff --git a/lib/edge/src/update_only/preview.rs b/lib/edge/src/update_only/preview.rs index 21ce003aba..2ce25af63c 100644 --- a/lib/edge/src/update_only/preview.rs +++ b/lib/edge/src/update_only/preview.rs @@ -8,7 +8,7 @@ //! [`apply_batch`]: UpdateOnlyEdgeShard::apply_batch use common::types::PointOffsetType; -use common::universal_io::{UniversalAppend, UniversalRead}; +use common::universal_io::{UniversalAppendFs, UniversalRead}; use rayon::ThreadPool; use segment::common::operation_error::OperationResult; use segment::data_types::fully_qualified_point::FullyQualifiedPoint; @@ -128,7 +128,7 @@ pub(super) fn resolve_batch( Ok(points) } -impl UpdateOnlyEdgeShard { +impl UpdateOnlyEdgeShard { /// Resolve a batch without writing anything: what /// [`apply_batch`](Self::apply_batch) would do, reported per point. /// diff --git a/lib/edge/src/update_only/tests.rs b/lib/edge/src/update_only/tests.rs index ba29c07cbf..0094657b97 100644 --- a/lib/edge/src/update_only/tests.rs +++ b/lib/edge/src/update_only/tests.rs @@ -6,7 +6,7 @@ //! saw, and what it appends is verified through an ordinary follower opened //! afterwards. -use common::universal_io::MmapFile; +use common::universal_io::MmapFs; use segment::types::{ExtendedPointId, SeqNumberType}; use shard::operations::CollectionUpdateOperations; use shard::operations::CollectionUpdateOperations::PointOperation; @@ -48,7 +48,7 @@ fn delete_batch( fn delete_batch_retires_points_and_leaves_the_rest() { let dir = leader_with_ten_points("edge-update-delete"); - let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); + let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); let (_writer, outcome) = writer.apply_batch(delete_batch([3, 7])).unwrap(); assert_eq!(outcome.deleted, 2); @@ -91,7 +91,7 @@ fn delete_batch_retires_points_and_leaves_the_rest() { fn replayed_delete_batch_is_a_no_op() { let dir = leader_with_ten_points("edge-update-delete-replay"); - let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); + let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); let (writer, outcome) = writer.apply_batch(delete_batch([3])).unwrap(); assert_eq!(outcome.deleted, 1); @@ -99,7 +99,7 @@ fn replayed_delete_batch_is_a_no_op() { assert_eq!(replayed.deleted, 0); assert_eq!(replayed.missing, 1); - let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); + let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); let (_writer, replayed) = writer.apply_batch(delete_batch([3])).unwrap(); assert_eq!(replayed.deleted, 0); assert_eq!(replayed.missing, 1); @@ -144,7 +144,7 @@ fn vacuumed_leader(prefix: &str) -> TempDir { fn delete_batch_tombstones_points_in_immutable_segments() { let dir = vacuumed_leader("edge-update-delete-immutable"); - let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); + let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); // The segment holding point 500 must not be the write target, or the // test dodges the delete-only path. @@ -171,7 +171,7 @@ fn delete_batch_tombstones_points_in_immutable_segments() { assert_eq!(replayed.deleted, 0); assert_eq!(replayed.missing, 1); - let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); + let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); let (_writer, replayed) = writer.apply_batch(delete_batch([500])).unwrap(); assert_eq!(replayed.deleted, 0); assert_eq!(replayed.missing, 1); @@ -273,7 +273,7 @@ mod store { ..point(5) }; - let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); + let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); let (_writer, outcome) = writer .apply_batch(store_batch(100, vec![new_point, rewritten])) .unwrap(); @@ -340,7 +340,7 @@ mod store { let dir = leader_with_ten_points("edge-update-store-replay"); recreate_payload_storages_append_only(dir.path()); - let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); + let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); let batch = || store_batch(100, vec![point(11)]); let (writer, outcome) = writer.apply_batch(batch()).unwrap(); assert_eq!(outcome.stored, 1); @@ -349,7 +349,7 @@ mod store { assert_eq!(replayed.stored, 0); assert_eq!(replayed.skipped, 1); - let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); + let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); let (_writer, replayed) = writer.apply_batch(batch()).unwrap(); assert_eq!(replayed.stored, 0); assert_eq!(replayed.skipped, 1); @@ -367,13 +367,13 @@ mod store { let dir = leader_with_ten_points("edge-update-store-resume"); recreate_payload_storages_append_only(dir.path()); - let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); + let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); let (_writer, outcome) = writer .apply_batch(store_batch(100, vec![point(11)])) .unwrap(); assert_eq!(outcome.stored, 1); - let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); + let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); let (_writer, second) = writer .apply_batch([ ( @@ -410,7 +410,7 @@ mod store { let dir = leader_with_ten_points("edge-update-store-sequential"); recreate_payload_storages_append_only(dir.path()); - let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); + let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); let (writer, outcome) = writer .apply_batch(store_batch(100, vec![point(11)])) .unwrap(); @@ -465,7 +465,7 @@ mod store { ..point(11) }; - let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); + let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); let (writer, outcome) = writer .apply_batch(conditional_batch( 100, @@ -530,7 +530,7 @@ mod store { ..point(5) }; - let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); + let writer = UpdateOnlyEdgeShard::::open_mmap(dir.path()).unwrap(); let (_writer, outcome) = writer .apply_batch(conditional_batch( 100, @@ -594,7 +594,7 @@ fn optimizing_target_gets_a_created_appendable() { ); fs_err::write(&manifest_path, serde_json::to_vec(&manifest).unwrap()).unwrap(); - let writer = UpdateOnlyEdgeShard::::open( + let writer = UpdateOnlyEdgeShard::open( MmapFs, dir.path(), ManifestSegmentEnumerator::new(MmapFs, dir.path()), @@ -647,7 +647,7 @@ fn empty_manifest_shard_bootstraps_an_appendable() { .unwrap(); fs_err::write(segment_manifest_path(dir.path()), "{}").unwrap(); - let writer = UpdateOnlyEdgeShard::::open( + let writer = UpdateOnlyEdgeShard::open( MmapFs, dir.path(), ManifestSegmentEnumerator::new(MmapFs, dir.path()), diff --git a/lib/edge/tools/shard_update/src/main.rs b/lib/edge/tools/shard_update/src/main.rs index 59878bf499..d21f83ba12 100644 --- a/lib/edge/tools/shard_update/src/main.rs +++ b/lib/edge/tools/shard_update/src/main.rs @@ -63,7 +63,7 @@ use std::sync::Arc; use anyhow::{Context, Result, anyhow}; use clap::{Args as ClapArgs, Parser, ValueEnum}; use common::universal_io::{ - DiskCacheConfig, MmapFile, OkNotFound as _, UniversalAppend, UniversalReadFileOps, + DiskCacheConfig, MmapFs, OkNotFound as _, UniversalAppendFs, UniversalReadFileOps, UniversalReadFs, read_json_via, }; use edge::external::uuid::Uuid; @@ -76,9 +76,7 @@ use edge::{ }; use io_bridge_object_store::backends::aws::{AwsConfig, AwsCredentials}; use io_bridge_object_store::backends::gcp::{GcsConfig, GcsCredentials}; -use io_bridge_object_store::{ - AsyncAppend, CachedBlobFile, CachedBlobFs, CachedBlobFsContext, ObjectStoreSource, -}; +use io_bridge_object_store::{AsyncAppend, CachedBlobFs, CachedBlobFsContext, ObjectStoreSource}; use object_store::aws::AmazonS3; use object_store::gcp::GoogleCloudStorage; use rand::rngs::StdRng; @@ -217,8 +215,8 @@ struct ShardSchema { /// parsed during the open — no extra reads) plus the payload-index schema, /// which is the one file the writer deliberately never opens, so the tool /// reads it through `fs` itself. -fn read_schema( - shard: &UpdateOnlyEdgeShard, +fn read_schema( + shard: &UpdateOnlyEdgeShard, fs: &F, shard_path: &Path, ) -> Result { @@ -514,8 +512,8 @@ fn generate_batch(schema: &ShardSchema, ids: &[PointId], seed: u64) -> UpdateOpe /// Generate the random batch, resolve it against the open shard, and log what /// it would do. The backend is behind `S`, so this is the whole dry run for /// local and object-storage shards alike. -fn dry_run( - shard: &UpdateOnlyEdgeShard, +fn dry_run( + shard: &UpdateOnlyEdgeShard, schema: &ShardSchema, ids: &[PointId], op_num: u64, @@ -568,8 +566,8 @@ fn dry_run( /// the backend once this returns `Ok`. With `interactive`, keeps prompting /// for the next round's ids and applies them through the writer the previous /// `apply_batch` handed back, one op-num (and seed) higher per round. -fn apply_run( - mut shard: UpdateOnlyEdgeShard, +fn apply_run( + mut shard: UpdateOnlyEdgeShard, schema: &ShardSchema, ids: &[PointId], mut op_num: u64, @@ -764,7 +762,7 @@ fn run_local(cli: &Cli, ids: &[PointId]) -> Result<()> { .clone() .ok_or_else(|| anyhow!("--path is required for the local backend"))?; - let shard = UpdateOnlyEdgeShard::::open_mmap(&path) + let shard = UpdateOnlyEdgeShard::::open_mmap(&path) .context("failed to open update-only edge shard")?; log::info!( "opened update-only shard with {} segment(s)", @@ -799,7 +797,7 @@ where let enumerator = ManifestSegmentEnumerator::new(cached_fs.clone(), &prefix); let shard = - UpdateOnlyEdgeShard::>::open(cached_fs.clone(), &prefix, enumerator) + UpdateOnlyEdgeShard::>::open(cached_fs.clone(), &prefix, enumerator) .context("failed to open update-only edge shard over object storage")?; log::info!( "opened update-only shard with {} segment(s)", diff --git a/lib/segment/src/common/update_only_blobstore.rs b/lib/segment/src/common/update_only_blobstore.rs index 421079b615..785eecee94 100644 --- a/lib/segment/src/common/update_only_blobstore.rs +++ b/lib/segment/src/common/update_only_blobstore.rs @@ -6,7 +6,7 @@ use blobstore::config::LogstoreConfig; use blobstore::{Blob, Logstore}; use common::counter::referenced_counter::HwMetricRefCounter; use common::types::PointOffsetType; -use common::universal_io::{Populate, UniversalAppend}; +use common::universal_io::{Populate, UniversalAppend, UniversalAppendFs}; use crate::common::operation_error::OperationResult; @@ -33,7 +33,7 @@ impl UpdateOnlyBlobstore { /// not there yet. An existing storage keeps the config it was created with, /// so `config_if_create` only decides the layout of a brand new one. pub fn open( - fs: &S::Fs, + fs: &impl UniversalAppendFs, path: &Path, config_if_create: LogstoreConfig, ) -> OperationResult { @@ -57,7 +57,7 @@ impl UpdateOnlyBlobstore { /// rollover. pub fn put( &mut self, - fs: &S::Fs, + fs: &impl UniversalAppendFs, slot: PointOffsetType, value: &V, hw_counter: HwMetricRefCounter, diff --git a/lib/segment/src/id_tracker/disk_id_tracker/update_only.rs b/lib/segment/src/id_tracker/disk_id_tracker/update_only.rs index cbff9320c8..e3c091d4cf 100644 --- a/lib/segment/src/id_tracker/disk_id_tracker/update_only.rs +++ b/lib/segment/src/id_tracker/disk_id_tracker/update_only.rs @@ -4,7 +4,7 @@ use std::path::{Path, PathBuf}; use common::bitvec::BitVec; use common::types::PointOffsetType; -use common::universal_io::{UniversalRead, UniversalWriteFileOps}; +use common::universal_io::UniversalAppendFs; use crate::common::operation_error::OperationResult; use crate::id_tracker::immutable_id_tracker::tombstone_points_in_stored_mask; @@ -16,18 +16,18 @@ use crate::types::PointIdType; /// [`atomic_save`] from the backend, so object stores qualify. /// /// [`atomic_save`]: UniversalWriteFileOps::atomic_save -pub struct UpdateOnlyDiskIdTracker + 'static> { - fs: S::Fs, +pub struct UpdateOnlyDiskIdTracker { + fs: Fs, segment_path: PathBuf, /// Consumed by the first [`tombstone_points`](Self::tombstone_points) in /// place of reading the mask file. deleted: Option, } -impl + 'static> UpdateOnlyDiskIdTracker { +impl UpdateOnlyDiskIdTracker { /// `deleted` is the mask as the read phase held it in memory, when it /// did; nothing is read here. - pub fn new(fs: S::Fs, segment_path: &Path, deleted: Option) -> Self { + pub fn new(fs: Fs, segment_path: &Path, deleted: Option) -> Self { Self { fs, segment_path: segment_path.to_path_buf(), @@ -44,11 +44,6 @@ impl + 'static> UpdateOnlyDiskIdTrac &mut self, points: &[(PointIdType, PointOffsetType)], ) -> OperationResult<()> { - tombstone_points_in_stored_mask::( - &self.fs, - &self.segment_path, - &mut self.deleted, - points, - ) + tombstone_points_in_stored_mask(&self.fs, &self.segment_path, &mut self.deleted, points) } } diff --git a/lib/segment/src/id_tracker/id_tracker_base/delete_only_tracker_enum.rs b/lib/segment/src/id_tracker/id_tracker_base/delete_only_tracker_enum.rs index 185fcc3505..83cff4310d 100644 --- a/lib/segment/src/id_tracker/id_tracker_base/delete_only_tracker_enum.rs +++ b/lib/segment/src/id_tracker/id_tracker_base/delete_only_tracker_enum.rs @@ -1,5 +1,5 @@ use common::types::PointOffsetType; -use common::universal_io::{UniversalRead, UniversalWriteFileOps}; +use common::universal_io::UniversalAppendFs; use crate::common::operation_error::OperationResult; use crate::id_tracker::disk_id_tracker::update_only::UpdateOnlyDiskIdTracker; @@ -8,12 +8,12 @@ use crate::types::PointIdType; /// The update-only tracker of whichever immutable id-tracker format a segment /// holds. Each variant decides where its tombstones go. -pub enum DeleteOnlyIdTrackerEnum + 'static> { - Immutable(UpdateOnlyImmutableIdTracker), - DiskResident(UpdateOnlyDiskIdTracker), +pub enum DeleteOnlyIdTrackerEnum { + Immutable(UpdateOnlyImmutableIdTracker), + DiskResident(UpdateOnlyDiskIdTracker), } -impl + 'static> DeleteOnlyIdTrackerEnum { +impl DeleteOnlyIdTrackerEnum { /// Retire the given points by marking the slots they occupy in the stored /// deleted mask — the only thing written, the data on those slots stays. pub fn tombstone_points( diff --git a/lib/segment/src/id_tracker/immutable_id_tracker/deleted_storage.rs b/lib/segment/src/id_tracker/immutable_id_tracker/deleted_storage.rs index 3bcca2e90d..a3431d789d 100644 --- a/lib/segment/src/id_tracker/immutable_id_tracker/deleted_storage.rs +++ b/lib/segment/src/id_tracker/immutable_id_tracker/deleted_storage.rs @@ -4,7 +4,7 @@ use common::bitvec::BitVec; use common::mmap::AdviceSetting; use common::stored_bitslice::StoredBitSlice; use common::types::PointOffsetType; -use common::universal_io::{OpenOptions, Populate, UniversalRead, UniversalWriteFileOps}; +use common::universal_io::{OpenOptions, Populate, UniversalReadFs, UniversalWriteFileOps}; use crate::common::operation_error::{OperationError, OperationResult}; use crate::types::PointIdType; @@ -20,8 +20,8 @@ pub(crate) fn deleted_path(base: &Path) -> PathBuf { /// [`StoredBitSlice::atomic_update`]; `seed` is consumed in place of reading /// the mask file when the caller held it in memory — and kept for a later /// call when `points` is empty and nothing is written. -pub(crate) fn tombstone_points_in_stored_mask>( - fs: &S::Fs, +pub(crate) fn tombstone_points_in_stored_mask( + fs: &Fs, segment_path: &Path, seed: &mut Option, points: &[(PointIdType, PointOffsetType)], @@ -29,7 +29,7 @@ pub(crate) fn tombstone_points_in_stored_mask::atomic_update( + StoredBitSlice::atomic_update( fs, deleted_path(segment_path), OpenOptions { diff --git a/lib/segment/src/id_tracker/immutable_id_tracker/update_only.rs b/lib/segment/src/id_tracker/immutable_id_tracker/update_only.rs index 7682157f73..ace523c9f1 100644 --- a/lib/segment/src/id_tracker/immutable_id_tracker/update_only.rs +++ b/lib/segment/src/id_tracker/immutable_id_tracker/update_only.rs @@ -5,7 +5,7 @@ use std::path::{Path, PathBuf}; use common::bitvec::BitVec; use common::types::PointOffsetType; -use common::universal_io::{UniversalRead, UniversalWriteFileOps}; +use common::universal_io::UniversalAppendFs; use super::deleted_storage::tombstone_points_in_stored_mask; use crate::common::operation_error::OperationResult; @@ -16,18 +16,18 @@ use crate::types::PointIdType; /// plus [`atomic_save`] from the backend, so object stores qualify. /// /// [`atomic_save`]: UniversalWriteFileOps::atomic_save -pub struct UpdateOnlyImmutableIdTracker + 'static> { - fs: S::Fs, +pub struct UpdateOnlyImmutableIdTracker { + fs: Fs, segment_path: PathBuf, /// Consumed by the first [`tombstone_points`](Self::tombstone_points) in /// place of reading the mask file. deleted: Option, } -impl + 'static> UpdateOnlyImmutableIdTracker { +impl UpdateOnlyImmutableIdTracker { /// `deleted` is the mask as the read phase held it in memory, when it /// did; nothing is read here. - pub fn new(fs: S::Fs, segment_path: &Path, deleted: Option) -> Self { + pub fn new(fs: Fs, segment_path: &Path, deleted: Option) -> Self { Self { fs, segment_path: segment_path.to_path_buf(), @@ -44,11 +44,6 @@ impl + 'static> UpdateOnlyImmutableI &mut self, points: &[(PointIdType, PointOffsetType)], ) -> OperationResult<()> { - tombstone_points_in_stored_mask::( - &self.fs, - &self.segment_path, - &mut self.deleted, - points, - ) + tombstone_points_in_stored_mask(&self.fs, &self.segment_path, &mut self.deleted, points) } } diff --git a/lib/segment/src/id_tracker/mutable_id_tracker/update_only/heal.rs b/lib/segment/src/id_tracker/mutable_id_tracker/update_only/heal.rs index 1bc7178c80..e1ea996112 100644 --- a/lib/segment/src/id_tracker/mutable_id_tracker/update_only/heal.rs +++ b/lib/segment/src/id_tracker/mutable_id_tracker/update_only/heal.rs @@ -2,7 +2,7 @@ use std::cmp::Ordering; use std::path::Path; use common::generic_consts::Sequential; -use common::universal_io::{UniversalRead as _, UniversalWriteFileOps}; +use common::universal_io::{UniversalAppendFs, UniversalRead as _}; use super::UpdateOnlyAppendableIdTracker; use crate::common::operation_error::{OperationError, OperationResult}; @@ -16,7 +16,7 @@ impl UpdateOnlyAppendableIdTracker { /// /// Takes ownership of `file` and drops it before the rewrite, Windows cannot replace a path /// that still has an mmap open. Callers open a fresh handle afterwards. - pub(super) fn heal_versions( + pub(super) fn heal_versions( fs: &Fs, path: &Path, file: Fs::AppendFile, @@ -48,7 +48,7 @@ impl UpdateOnlyAppendableIdTracker { /// The caller must drop any open handle on `path` first, Windows cannot replace a path that /// still has an mmap open. Opens its own handle, drops it before the rewrite, and leaves the /// caller to open a fresh one. - pub(super) fn heal_mappings( + pub(super) fn heal_mappings( &self, fs: &Fs, path: &Path, diff --git a/lib/segment/src/id_tracker/mutable_id_tracker/update_only/mod.rs b/lib/segment/src/id_tracker/mutable_id_tracker/update_only/mod.rs index 9420c5cea0..0f8052ec7e 100644 --- a/lib/segment/src/id_tracker/mutable_id_tracker/update_only/mod.rs +++ b/lib/segment/src/id_tracker/mutable_id_tracker/update_only/mod.rs @@ -8,8 +8,8 @@ use std::path::{Path, PathBuf}; use common::mmap::{Advice, AdviceSetting}; use common::types::PointOffsetType; use common::universal_io::{ - IsNotFound as _, OkNotFound as _, OpenOptions, Populate, UniversalAppend, UniversalFlush as _, - UniversalWriteFileOps, + IsNotFound as _, OkNotFound as _, OpenOptions, Populate, UniversalAppend, UniversalAppendFs, + UniversalFlush as _, }; use super::change::{MappingChange, write_entry}; @@ -85,7 +85,7 @@ impl UpdateOnlyAppendableIdTracker { /// `mappings_end` is not a hint: the first append cuts the file back to it (see /// [`heal_mappings`](Self::heal_mappings)), which drops a torn entry, or good data if it lags /// for any other reason. - pub fn new( + pub fn new( fs: &Fs, segment_path: impl Into, max_claimed_internal_id: Option, @@ -117,7 +117,7 @@ impl UpdateOnlyAppendableIdTracker { /// opened. /// /// [`DELETED_POINT_VERSION`]: crate::id_tracker::DELETED_POINT_VERSION - pub fn set_internal_versions( + pub fn set_internal_versions( &mut self, fs: &Fs, internal_ids: &[PointOffsetType], @@ -208,7 +208,7 @@ impl UpdateOnlyAppendableIdTracker { /// /// Slots are consecutive above every slot the log has claimed, and stay invisible to readers /// until [`set_internal_versions`](Self::set_internal_versions) covers them. - pub fn insert_operations( + pub fn insert_operations( &mut self, fs: &Fs, operations: &[MappingOperation], @@ -296,7 +296,7 @@ impl UpdateOnlyAppendableIdTracker { /// /// Runs at construction rather than lazily on the first write, so that no later write path can /// be added that forgets it and publishes one of these points. - fn retire_pending_inserts( + fn retire_pending_inserts( &mut self, fs: &Fs, pending_inserts: impl IntoIterator, @@ -306,7 +306,7 @@ impl UpdateOnlyAppendableIdTracker { /// Retire `point_ids`: each stops resolving, and the slot it held keeps its data and is never /// handed out again. Tombstoning that data is the caller's business. - pub fn delete_points( + pub fn delete_points( &mut self, fs: &Fs, point_ids: impl IntoIterator, @@ -334,10 +334,7 @@ impl UpdateOnlyAppendableIdTracker { } /// Open the append handle for `path`, creating the file if it is not there yet. - fn open_append( - fs: &Fs, - path: &Path, - ) -> OperationResult { + fn open_append(fs: &Fs, path: &Path) -> OperationResult { match fs.open_append(path, Self::open_options()) { Ok(file) => Ok(file), Err(err) if err.is_not_found() => { diff --git a/lib/segment/src/index/field_index/field_index_base/update_only/mod.rs b/lib/segment/src/index/field_index/field_index_base/update_only/mod.rs index 67b06913c4..1cebb2b613 100644 --- a/lib/segment/src/index/field_index/field_index_base/update_only/mod.rs +++ b/lib/segment/src/index/field_index/field_index_base/update_only/mod.rs @@ -18,7 +18,7 @@ use std::path::Path; use common::counter::hardware_counter::HardwareCounterCell; use common::types::PointOffsetType; -use common::universal_io::UniversalAppend; +use common::universal_io::{UniversalAppend, UniversalAppendFs}; use serde_json::Value; pub use self::writer::{UpdateOnlyIndexKind, UpdateOnlyValueIndex}; @@ -64,7 +64,7 @@ impl UpdateOnlyFieldIndex { /// Open the writer for `field`'s index of type `index_type`, under the /// payload index root `dir`, creating its storage if it is not there yet. pub fn open( - fs: &S::Fs, + fs: &impl UniversalAppendFs, dir: &Path, field: &JsonPath, schema: &PayloadFieldSchema, @@ -116,7 +116,7 @@ impl UpdateOnlyFieldIndex { /// [`flush`](Self::flush). pub fn add_point( &mut self, - fs: &S::Fs, + fs: &impl UniversalAppendFs, slot: PointOffsetType, values: &[&Value], hw_counter: &HardwareCounterCell, @@ -139,7 +139,11 @@ impl UpdateOnlyFieldIndex { /// /// `hw_counter` is charged only by the bitmask-backed indexes, which do /// their writing here; the rest charge each value as it is put. - pub fn flush(&mut self, fs: &S::Fs, hw_counter: &HardwareCounterCell) -> OperationResult<()> { + pub fn flush( + &mut self, + fs: &impl UniversalAppendFs, + hw_counter: &HardwareCounterCell, + ) -> OperationResult<()> { match self { Self::IntIndex(index) => index.flush(), Self::DatetimeIndex(index) => index.flush(), diff --git a/lib/segment/src/index/field_index/field_index_base/update_only/writer.rs b/lib/segment/src/index/field_index/field_index_base/update_only/writer.rs index 9d36cb67d4..f53bc8dae3 100644 --- a/lib/segment/src/index/field_index/field_index_base/update_only/writer.rs +++ b/lib/segment/src/index/field_index/field_index_base/update_only/writer.rs @@ -7,7 +7,7 @@ use blobstore::Blob; use blobstore::config::{Compression, DEFAULT_PAGE_SIZE_BYTES, LogstoreConfig}; use common::counter::hardware_counter::HardwareCounterCell; use common::types::PointOffsetType; -use common::universal_io::UniversalAppend; +use common::universal_io::{UniversalAppend, UniversalAppendFs}; use serde_json::Value; use crate::common::operation_error::OperationResult; @@ -48,7 +48,11 @@ pub struct UpdateOnlyValueIndex UpdateOnlyValueIndex { /// Open the index storage directory at `dir` for appending, creating it if /// the field has no index there yet. - pub fn open(fs: &S::Fs, dir: &Path, kind: K) -> OperationResult { + pub fn open( + fs: &impl UniversalAppendFs, + dir: &Path, + kind: K, + ) -> OperationResult { let storage = UpdateOnlyBlobstore::open(fs, dir, INDEX_LOGSTORE_CONFIG)?; Ok(Self { kind, storage }) } @@ -64,7 +68,7 @@ impl UpdateOnlyValueIndex< /// empty slot reads back as a point this index does not cover. pub fn add_point( &mut self, - fs: &S::Fs, + fs: &impl UniversalAppendFs, slot: PointOffsetType, values: &[&Value], hw_counter: &HardwareCounterCell, diff --git a/lib/segment/src/index/struct_payload_index/update_only/mod.rs b/lib/segment/src/index/struct_payload_index/update_only/mod.rs index 4b57793697..bfd666358c 100644 --- a/lib/segment/src/index/struct_payload_index/update_only/mod.rs +++ b/lib/segment/src/index/struct_payload_index/update_only/mod.rs @@ -8,7 +8,7 @@ use std::path::Path; use ahash::AHashMap; use common::counter::hardware_counter::HardwareCounterCell; use common::types::PointOffsetType; -use common::universal_io::UniversalAppend; +use common::universal_io::{UniversalAppend, UniversalAppendFs}; use crate::common::operation_error::{OperationError, OperationResult}; use crate::index::field_index::UpdateOnlyFieldIndex; @@ -40,7 +40,10 @@ impl UpdateOnlyStructPayloadIndex { /// is a config from before those types were recorded, which only the /// writable index can repair, by deriving them from the schema on its next /// open. - pub fn open(fs: &S::Fs, segment_path: &Path) -> OperationResult { + pub fn open( + fs: &impl UniversalAppendFs, + segment_path: &Path, + ) -> OperationResult { let path = get_payload_index_path(segment_path); let config = PayloadConfig::load_universal(fs, &PayloadConfig::get_config_path(&path))? .unwrap_or_default(); @@ -81,7 +84,7 @@ impl UpdateOnlyStructPayloadIndex { /// value there. pub fn append_many<'a>( &mut self, - fs: &S::Fs, + fs: &impl UniversalAppendFs, points: impl IntoIterator, hw_counter: &HardwareCounterCell, ) -> OperationResult<()> { diff --git a/lib/segment/src/payload_storage/update_only/mod.rs b/lib/segment/src/payload_storage/update_only/mod.rs index 2b5efd50a7..6122cdf14d 100644 --- a/lib/segment/src/payload_storage/update_only/mod.rs +++ b/lib/segment/src/payload_storage/update_only/mod.rs @@ -6,7 +6,7 @@ use std::path::Path; use blobstore::config::LogstoreConfig; use common::counter::hardware_counter::HardwareCounterCell; use common::types::PointOffsetType; -use common::universal_io::UniversalAppend; +use common::universal_io::{UniversalAppend, UniversalAppendFs}; use crate::common::operation_error::OperationResult; use crate::common::update_only_blobstore::UpdateOnlyBlobstore; @@ -35,7 +35,10 @@ impl UpdateOnlyPayloadStorage { /// its file format is not one this writer can append to. /// /// [1]: crate::payload_storage::payload_storage_impl::PayloadStorageImpl::open_or_create - pub fn open(fs: &S::Fs, segment_path: &Path) -> OperationResult { + pub fn open( + fs: &impl UniversalAppendFs, + segment_path: &Path, + ) -> OperationResult { let storage = UpdateOnlyBlobstore::open(fs, &storage_dir(segment_path), LogstoreConfig::DEFAULT)?; @@ -54,7 +57,7 @@ impl UpdateOnlyPayloadStorage { /// an unwritten slot already reads back as an empty payload. pub fn append_many<'a>( &mut self, - fs: &S::Fs, + fs: &impl UniversalAppendFs, payloads: impl IntoIterator, hw_counter: &HardwareCounterCell, ) -> OperationResult<()> { diff --git a/lib/segment/src/segment/update_only/appendable/mod.rs b/lib/segment/src/segment/update_only/appendable/mod.rs index dc30e0f070..751d2575b2 100644 --- a/lib/segment/src/segment/update_only/appendable/mod.rs +++ b/lib/segment/src/segment/update_only/appendable/mod.rs @@ -6,7 +6,7 @@ use std::path::{Path, PathBuf}; use common::counter::hardware_counter::HardwareCounterCell; use common::types::PointOffsetType; -use common::universal_io::UniversalAppend; +use common::universal_io::UniversalAppendFs; use super::AppendableIdTrackerState; use crate::common::operation_error::OperationResult; @@ -23,40 +23,40 @@ use crate::vector_storage::update_only::{UpdateOnlyVectorStorage, VectorToStore} /// A segment open for appends: the write target. Every point a batch stores /// lands here, in a fresh slot — nothing is ever rewritten in place. -pub struct AppendableSegment { +pub struct AppendableSegment { id_tracker: UpdateOnlyAppendableIdTracker, /// What the id tracker appends through and /// [`store_components`](Self::store_components) opens with. - fs: S::Fs, + fs: Fs, segment_path: PathBuf, config: SegmentConfig, /// The components a stored point's data goes into, opened on the first /// [`store_points`](Self::store_points). A batch that only deletes writes /// nothing but the mappings log, so it never pays for these opens. - store: Option>, + store: Option>, } /// Everywhere a stored point's data goes. The mappings log that publishes the /// point is not here: it belongs to the segment itself, since deletes need it /// too. -struct StoreComponents { +struct StoreComponents { /// What the components write through, passed down per call. - fs: S::Fs, - payload_storage: UpdateOnlyPayloadStorage, - payload_indexes: UpdateOnlyStructPayloadIndex, + fs: Fs, + payload_storage: UpdateOnlyPayloadStorage, + payload_indexes: UpdateOnlyStructPayloadIndex, /// One writer per named vector, dense and sparse alike. - vector_storages: Vec<(VectorNameBuf, UpdateOnlyVectorStorage)>, + vector_storages: Vec<(VectorNameBuf, UpdateOnlyVectorStorage)>, /// Quantized overlays, keyed by vector name — only for dense, non-multivector vectors /// whose quantization method supports incremental appends (Binary/Turbo — see /// `QuantizationConfig::supports_appendable`). A vector name with no entry here simply /// has no live quantization (never configured, an unsupported method, or a multivector — /// that support is a follow-up, needing its own append-only offsets storage): it stays /// searchable exactly, through the raw storage alone, same as before this existed. - quantized_vectors: HashMap>, + quantized_vectors: HashMap>, } -impl StoreComponents { - fn open(fs: S::Fs, segment_path: &Path, config: &SegmentConfig) -> OperationResult { +impl StoreComponents { + fn open(fs: Fs, segment_path: &Path, config: &SegmentConfig) -> OperationResult { let payload_storage = UpdateOnlyPayloadStorage::open(&fs, segment_path)?; let payload_indexes = UpdateOnlyStructPayloadIndex::open(&fs, segment_path)?; @@ -90,7 +90,7 @@ impl StoreComponents { } } -impl AppendableSegment { +impl AppendableSegment { /// Resume the segment directory at `segment_path` from the mappings-log /// state the read phase observed. /// @@ -98,7 +98,7 @@ impl AppendableSegment { /// versions were never committed are retired here, since which components /// got to write their data is unknowable. pub fn open( - fs: S::Fs, + fs: Fs, segment_path: &Path, config: &SegmentConfig, state: AppendableIdTrackerState, @@ -126,7 +126,7 @@ impl AppendableSegment { }) } - fn store_components(&mut self) -> OperationResult<&mut StoreComponents> { + fn store_components(&mut self) -> OperationResult<&mut StoreComponents> { if self.store.is_none() { self.store = Some(StoreComponents::open( self.fs.clone(), diff --git a/lib/segment/src/segment/update_only/delete_only/mod.rs b/lib/segment/src/segment/update_only/delete_only/mod.rs index e0f30ee1df..ae64cbeb76 100644 --- a/lib/segment/src/segment/update_only/delete_only/mod.rs +++ b/lib/segment/src/segment/update_only/delete_only/mod.rs @@ -4,7 +4,7 @@ use std::path::Path; use common::types::PointOffsetType; -use common::universal_io::{UniversalRead, UniversalWriteFileOps}; +use common::universal_io::UniversalAppendFs; use super::DeleteOnlyIdTrackerState; use crate::common::operation_error::OperationResult; @@ -16,19 +16,15 @@ use crate::types::PointIdType; /// A segment open for deletes: nothing in it can grow, so the only thing a /// batch can do here is retire points that are already there. Where their /// tombstones go is the id tracker's decision. -pub struct DeleteOnlySegment + 'static> { - id_tracker: DeleteOnlyIdTrackerEnum, +pub struct DeleteOnlySegment { + id_tracker: DeleteOnlyIdTrackerEnum, } -impl + 'static> DeleteOnlySegment { +impl DeleteOnlySegment { /// Open the segment directory at `segment_path` for deletes, resuming the /// tracker kind its [`DeleteOnlyIdTrackerState`] variant names; nothing is /// read. - pub fn open( - fs: S::Fs, - segment_path: &Path, - id_tracker_state: DeleteOnlyIdTrackerState, - ) -> Self { + pub fn open(fs: Fs, segment_path: &Path, id_tracker_state: DeleteOnlyIdTrackerState) -> Self { let id_tracker = match id_tracker_state { DeleteOnlyIdTrackerState::Immutable(deleted) => DeleteOnlyIdTrackerEnum::Immutable( UpdateOnlyImmutableIdTracker::new(fs, segment_path, deleted), diff --git a/lib/segment/src/segment/update_only/lookup/lifecycle.rs b/lib/segment/src/segment/update_only/lookup/lifecycle.rs index 6bff80a9b7..e522aae021 100644 --- a/lib/segment/src/segment/update_only/lookup/lifecycle.rs +++ b/lib/segment/src/segment/update_only/lookup/lifecycle.rs @@ -54,7 +54,7 @@ fn build_cached_fs( Ok(cached_fs) } -impl + 'static> LookupSegment { +impl LookupSegment { /// Open the segment over a per-segment [`CachedFs`]: every file the /// components will read is prefetched concurrently /// ([`preopen`](Self::preopen)) before the component opens consume it, so @@ -67,7 +67,7 @@ impl + 'static> LookupSegment { /// `deferred_internal_id` is the cutoff agreed with an external rebuilder /// working the same directory — see [`open_via`](Self::open_via). pub fn open( - fs: &S::Fs, + fs: &impl UniversalReadFsAsync, segment_path: &Path, deferred_internal_id: Option, ) -> OperationResult { diff --git a/lib/segment/src/segment/update_only/lookup/live_reload.rs b/lib/segment/src/segment/update_only/lookup/live_reload.rs index 410d6d1293..d11a3499c2 100644 --- a/lib/segment/src/segment/update_only/lookup/live_reload.rs +++ b/lib/segment/src/segment/update_only/lookup/live_reload.rs @@ -1,6 +1,6 @@ use common::counter::hardware_counter::HardwareCounterCell; use common::sorted_slice::SortedSlice; -use common::universal_io::UniversalRead; +use common::universal_io::{UniversalRead, UniversalReadFs}; use super::LookupSegment; use crate::common::live_reload::LiveReload as _; @@ -20,7 +20,7 @@ impl LookupSegment { /// [`ReadOnlySegment::live_reload`]: crate::segment::read_only::ReadOnlySegment::live_reload pub fn live_reload( &mut self, - fs: &S::Fs, + fs: &impl UniversalReadFs, hw_counter: &HardwareCounterCell, ) -> OperationResult<()> { let Self { diff --git a/lib/segment/src/segment/update_only/segment_enum.rs b/lib/segment/src/segment/update_only/segment_enum.rs index 0c50d3b021..5dd1742eb0 100644 --- a/lib/segment/src/segment/update_only/segment_enum.rs +++ b/lib/segment/src/segment/update_only/segment_enum.rs @@ -4,40 +4,40 @@ use std::path::Path; use common::types::PointOffsetType; -use common::universal_io::UniversalAppend; +use common::universal_io::UniversalAppendFs; use super::{AppendableSegment, DeleteOnlySegment, WriterIdTrackerState}; use crate::common::operation_error::OperationResult; use crate::types::{PointIdType, SegmentConfig}; /// A segment opened for writing: appendable, or accepting deletes only. -pub enum UpdateOnlySegmentEnum { - DeleteOnly(DeleteOnlySegment), - Appendable(Box>), +pub enum UpdateOnlySegmentEnum { + DeleteOnly(DeleteOnlySegment), + Appendable(Box>), } -impl UpdateOnlySegmentEnum { +impl UpdateOnlySegmentEnum { /// Open a writer over the segment directory at `segment_path`, of the /// kind the read phase's `id_tracker_state` dictates. pub fn open( - fs: &S::Fs, + fs: Fs, segment_path: &Path, config: &SegmentConfig, id_tracker_state: WriterIdTrackerState, ) -> OperationResult { Ok(match id_tracker_state { WriterIdTrackerState::Appendable(state) => Self::Appendable(Box::new( - AppendableSegment::open(fs.clone(), segment_path, config, state)?, + AppendableSegment::open(fs, segment_path, config, state)?, )), WriterIdTrackerState::DeleteOnly(state) => { - Self::DeleteOnly(DeleteOnlySegment::open(fs.clone(), segment_path, state)) + Self::DeleteOnly(DeleteOnlySegment::open(fs, segment_path, state)) } }) } /// The appendable writer, when this segment is one; `None` when it accepts /// deletes only. Storing points is the one operation the two do not share. - pub fn as_appendable_mut(&mut self) -> Option<&mut AppendableSegment> { + pub fn as_appendable_mut(&mut self) -> Option<&mut AppendableSegment> { match self { Self::Appendable(segment) => Some(segment), Self::DeleteOnly(_) => None, diff --git a/lib/segment/src/vector_storage/quantized/quantized_chunked_mmap_storage/update_only.rs b/lib/segment/src/vector_storage/quantized/quantized_chunked_mmap_storage/update_only.rs index e2548eba9b..29f6d95d51 100644 --- a/lib/segment/src/vector_storage/quantized/quantized_chunked_mmap_storage/update_only.rs +++ b/lib/segment/src/vector_storage/quantized/quantized_chunked_mmap_storage/update_only.rs @@ -7,7 +7,7 @@ use std::path::Path; use common::counter::hardware_counter::HardwareCounterCell; use common::mmap::Flusher; use common::types::PointOffsetType; -use common::universal_io::UniversalAppend; +use common::universal_io::UniversalAppendFs; use quantization::{EncodedStorageBuilder, EncodedStorageWrite}; use crate::common::operation_error::OperationResult; @@ -20,7 +20,8 @@ use crate::vector_storage::chunked_vectors::update_only::UpdateOnlyChunkedVector /// A quantized vector persists through [`UpdateOnlyChunkedVectors`] — the same append-only, /// fresh-slot-per-upsert storage [`UpdateOnlyDenseVectorStorage`] uses for raw vectors — instead /// of the positional writes `ChunkedVectors::insert` needs, so this only requires -/// `S: UniversalAppend`, not `UniversalWrite`. It writes files in the exact layout +/// append-capable files +/// (`Fs::File: UniversalAppend`), not `UniversalWrite`. It writes files in the exact layout /// `QuantizedChunkedStorage` reads, so no new reading code is needed once a segment is promoted. /// /// Implements [`EncodedStorageWrite`], not the full [`EncodedStorage`](quantization::EncodedStorage): @@ -31,15 +32,15 @@ use crate::vector_storage::chunked_vectors::update_only::UpdateOnlyChunkedVector /// exist. /// /// [`UpdateOnlyDenseVectorStorage`]: crate::vector_storage::dense::update_only::UpdateOnlyDenseVectorStorage -pub struct UpdateOnlyQuantizedChunkedStorage { +pub struct UpdateOnlyQuantizedChunkedStorage { vectors: UpdateOnlyChunkedVectors, /// Owned rather than passed down: [`EncodedStorageWrite`] cannot carry a /// filesystem per call. - fs: S::Fs, + fs: Fs, } -impl UpdateOnlyQuantizedChunkedStorage { - pub fn open(fs: S::Fs, path: &Path, quantized_vector_size: usize) -> OperationResult { +impl UpdateOnlyQuantizedChunkedStorage { + pub fn open(fs: Fs, path: &Path, quantized_vector_size: usize) -> OperationResult { Ok(Self { vectors: UpdateOnlyChunkedVectors::open(&fs, path, quantized_vector_size)?, fs, @@ -47,7 +48,7 @@ impl UpdateOnlyQuantizedChunkedStorage { } } -impl EncodedStorageWrite for UpdateOnlyQuantizedChunkedStorage { +impl EncodedStorageWrite for UpdateOnlyQuantizedChunkedStorage { fn is_in_ram_or_mmap() -> bool { false } @@ -101,22 +102,20 @@ impl EncodedStorageWrite for UpdateOnlyQuantizedCh /// Builder counterpart, used once at segment creation to open an empty overlay — an update-only /// segment always starts with zero quantized vectors, so [`build`](Self::build) never has /// [`push_vector_data`](Self::push_vector_data) called on it. -pub struct UpdateOnlyQuantizedChunkedStorageBuilder { - storage: UpdateOnlyQuantizedChunkedStorage, +pub struct UpdateOnlyQuantizedChunkedStorageBuilder { + storage: UpdateOnlyQuantizedChunkedStorage, } -impl UpdateOnlyQuantizedChunkedStorageBuilder { - pub fn new(fs: S::Fs, path: &Path, quantized_vector_size: usize) -> OperationResult { +impl UpdateOnlyQuantizedChunkedStorageBuilder { + pub fn new(fs: Fs, path: &Path, quantized_vector_size: usize) -> OperationResult { Ok(Self { storage: UpdateOnlyQuantizedChunkedStorage::open(fs, path, quantized_vector_size)?, }) } } -impl EncodedStorageBuilder - for UpdateOnlyQuantizedChunkedStorageBuilder -{ - type Storage = UpdateOnlyQuantizedChunkedStorage; +impl EncodedStorageBuilder for UpdateOnlyQuantizedChunkedStorageBuilder { + type Storage = UpdateOnlyQuantizedChunkedStorage; type Error = std::io::Error; fn build(self) -> std::io::Result { diff --git a/lib/segment/src/vector_storage/quantized/update_only/mod.rs b/lib/segment/src/vector_storage/quantized/update_only/mod.rs index ed642a4ebc..1667833d33 100644 --- a/lib/segment/src/vector_storage/quantized/update_only/mod.rs +++ b/lib/segment/src/vector_storage/quantized/update_only/mod.rs @@ -40,7 +40,7 @@ use std::path::Path; use common::counter::hardware_counter::HardwareCounterCell; use common::types::PointOffsetType; -use common::universal_io::{UniversalAppend, UniversalReadFileOps as _, read_json_via}; +use common::universal_io::{UniversalAppendFs, read_json_via}; use quantization::encoded_vectors_binary::EncodedVectorsBin; use quantization::encoded_vectors_tq::EncodedVectorsTQ; @@ -54,9 +54,9 @@ use crate::vector_storage::quantized::quantized_vectors::{ }; use crate::vector_storage::update_only::VectorToStore; -enum UpdateOnlyQuantizedVectorStorage { - Binary(Box>>), - Turbo(Box>>), +enum UpdateOnlyQuantizedVectorStorage { + Binary(Box>>), + Turbo(Box>>), } /// The write half of a dense vector's quantized overlay, for one update-only appendable @@ -64,15 +64,15 @@ enum UpdateOnlyQuantizedVectorStorage { /// segment's quantization config supports it. /// /// [`UpdateOnlyDenseVectorStorage`]: crate::vector_storage::dense::update_only::UpdateOnlyDenseVectorStorage -pub struct UpdateOnlyQuantizedVectors { - storage: UpdateOnlyQuantizedVectorStorage, +pub struct UpdateOnlyQuantizedVectors { + storage: UpdateOnlyQuantizedVectorStorage, config: QuantizedVectorsConfig, /// Raw-storage properties, needed to decode a [`VectorToStore::Raw`]. distance: Distance, datatype: VectorStorageDatatype, } -impl UpdateOnlyQuantizedVectors { +impl UpdateOnlyQuantizedVectors { /// Reopen the quantized overlay persisted at `path`, if one is there. /// /// This never creates anything: whether a vector gets a quantized overlay is a decision made @@ -81,7 +81,7 @@ impl UpdateOnlyQuantizedVectors { /// for this vector, or the configured method didn't support incremental appends (Scalar, /// Product — see [`QuantizationConfig::supports_appendable`]) at creation time. pub fn open( - fs: S::Fs, + fs: Fs, path: &Path, vector_config: &VectorDataConfig, ) -> OperationResult> { @@ -108,7 +108,7 @@ impl UpdateOnlyQuantizedVectors { /// every vector it writes is sized from this same metadata, so the invariant `load`'s check /// protects holds by construction here, not by verification. fn open_existing( - fs: S::Fs, + fs: Fs, config: QuantizedVectorsConfig, path: &Path, distance: Distance, diff --git a/lib/segment/src/vector_storage/quantized/update_only/tests.rs b/lib/segment/src/vector_storage/quantized/update_only/tests.rs index f3b391956b..1f6abd6764 100644 --- a/lib/segment/src/vector_storage/quantized/update_only/tests.rs +++ b/lib/segment/src/vector_storage/quantized/update_only/tests.rs @@ -96,7 +96,7 @@ fn some_vectors(n: usize) -> Vec> { fn create_empty_overlay( config: &QuantizationConfig, path: &std::path::Path, -) -> UpdateOnlyQuantizedVectors { +) -> UpdateOnlyQuantizedVectors { let storage_type = QuantizedVectorsStorageType::Mutable; let vector_parameters = QuantizedVectors::construct_vector_parameters(config, Distance::Dot, DIM, 0, storage_type); @@ -205,7 +205,7 @@ fn write_all(config: &QuantizationConfig, path: &std::path::Path, vectors: &[Vec drop(writer); let mut writer = - UpdateOnlyQuantizedVectors::::open(MmapFs, path, &dense_vector_config()) + UpdateOnlyQuantizedVectors::::open(MmapFs, path, &dense_vector_config()) .unwrap() .expect("overlay was already created by the first writer"); writer @@ -371,7 +371,7 @@ fn turbo_bytes_match_the_standard_batch_encode_path() { fn open_returns_none_when_nothing_persisted() { let dir = TempDir::with_prefix("update_only_quantized_no_config").unwrap(); let overlay = - UpdateOnlyQuantizedVectors::::open(MmapFs, dir.path(), &dense_vector_config()) + UpdateOnlyQuantizedVectors::::open(MmapFs, dir.path(), &dense_vector_config()) .unwrap(); assert!(overlay.is_none()); } @@ -398,7 +398,7 @@ fn reopening_a_nonempty_overlay_works() { drop(writer); let reopened = - UpdateOnlyQuantizedVectors::::open(MmapFs, dir.path(), &dense_vector_config()) + UpdateOnlyQuantizedVectors::::open(MmapFs, dir.path(), &dense_vector_config()) .unwrap(); assert!(reopened.is_some()); } diff --git a/lib/segment/src/vector_storage/sparse/update_only/mod.rs b/lib/segment/src/vector_storage/sparse/update_only/mod.rs index d1c9639c8c..e60c9c9ea3 100644 --- a/lib/segment/src/vector_storage/sparse/update_only/mod.rs +++ b/lib/segment/src/vector_storage/sparse/update_only/mod.rs @@ -6,7 +6,7 @@ use std::path::Path; use blobstore::config::{Compression, DEFAULT_PAGE_SIZE_BYTES, LogstoreConfig}; use common::counter::hardware_counter::HardwareCounterCell; use common::types::PointOffsetType; -use common::universal_io::UniversalAppend; +use common::universal_io::{UniversalAppend, UniversalAppendFs}; use sparse::common::sparse_vector::SparseVector; use super::mmap_sparse_vector_storage::{DELETED_DIRNAME, STORAGE_DIRNAME}; @@ -37,7 +37,7 @@ pub struct UpdateOnlySparseVectorStorage { impl UpdateOnlySparseVectorStorage { /// Open the storage at `path` for appending, creating it if it is not there /// yet. - pub fn open(fs: &S::Fs, path: &Path) -> OperationResult { + pub fn open(fs: &impl UniversalAppendFs, path: &Path) -> OperationResult { Ok(Self { storage: UpdateOnlyBlobstore::open(fs, &path.join(STORAGE_DIRNAME), STORAGE_CONFIG)?, deleted: UpdateOnlyStoredFlags::open(fs, &path.join(DELETED_DIRNAME))?, @@ -52,7 +52,7 @@ impl UpdateOnlySparseVectorStorage { /// deleted, which is what the writable side records for it. pub fn append_many<'a>( &mut self, - fs: &S::Fs, + fs: &impl UniversalAppendFs, start_slot: PointOffsetType, vectors: impl IntoIterator>, hw_counter: &HardwareCounterCell, diff --git a/lib/segment/src/vector_storage/update_only/mod.rs b/lib/segment/src/vector_storage/update_only/mod.rs index ffa6cd7de6..ab7da8df9d 100644 --- a/lib/segment/src/vector_storage/update_only/mod.rs +++ b/lib/segment/src/vector_storage/update_only/mod.rs @@ -14,7 +14,7 @@ use std::path::Path; use common::counter::hardware_counter::HardwareCounterCell; use common::types::PointOffsetType; -use common::universal_io::UniversalAppend; +use common::universal_io::{UniversalAppend, UniversalAppendFs}; use crate::common::operation_error::{OperationError, OperationResult}; use crate::data_types::vectors::{ @@ -80,7 +80,11 @@ impl UpdateOnlyVectorStorage { /// /// Fails for a storage type an update-only segment cannot have: the mmap /// ones are immutable, built whole rather than appended to. - pub fn open(fs: &S::Fs, path: &Path, config: &VectorDataConfig) -> OperationResult { + pub fn open( + fs: &impl UniversalAppendFs, + path: &Path, + config: &VectorDataConfig, + ) -> OperationResult { match config.storage_type { VectorStorageType::ChunkedMmap | VectorStorageType::InRamChunkedMmap => {} storage_type @ (VectorStorageType::Mmap @@ -128,7 +132,10 @@ impl UpdateOnlyVectorStorage { /// Open the writer for a sparse vector storage at `path`. Sparse vectors /// are configured separately from dense ones, so they do not go through /// [`open`](Self::open). - pub fn open_sparse(fs: &S::Fs, path: &Path) -> OperationResult { + pub fn open_sparse( + fs: &impl UniversalAppendFs, + path: &Path, + ) -> OperationResult { Ok(Self::Sparse(Box::new(UpdateOnlySparseVectorStorage::open( fs, path, )?))) @@ -138,7 +145,7 @@ impl UpdateOnlyVectorStorage { /// persist them. pub fn append_many<'a>( &mut self, - fs: &S::Fs, + fs: &impl UniversalAppendFs, start_slot: PointOffsetType, vectors: impl IntoIterator>, hw_counter: &HardwareCounterCell,