From da823efbb9e2ecdaba074e19f98a44018f99b149 Mon Sep 17 00:00:00 2001 From: Roman Titov Date: Mon, 8 Jun 2026 22:43:41 +0200 Subject: [PATCH] Explicitly propagate `fs` into `Gridstore` (#9381) --- .../common/src/universal_io/disk_cache/mod.rs | 20 +++ .../common/src/universal_io/io_uring/mod.rs | 20 +++ .../common/src/universal_io/local_file_ops.rs | 27 ++++ .../common/src/universal_io/mmap/mod.rs | 20 +++ .../src/universal_io/simple_disk_cache/fs.rs | 20 +++ .../src/universal_io/traits/file_ops.rs | 27 +++- .../src/universal_io/wrappers/read_only.rs | 35 +++++- lib/common/io_bridge_object_store/src/file.rs | 20 +++ lib/common/io_bridge_object_store/src/fs.rs | 25 +++- lib/common/io_bridge_object_store/src/read.rs | 16 +++ .../io_bridge_object_store/src/source.rs | 105 +++++++++++++++- lib/gridstore/src/fixtures.rs | 5 +- lib/gridstore/src/gridstore/mod.rs | 118 +++++++++--------- lib/gridstore/src/gridstore/tests.rs | 33 ++--- .../mutable_text_index/lifecycle.rs | 5 +- .../geo_index/mutable_geo_index/lifecycle.rs | 5 +- .../map_index/mutable_map_index/lifecycle.rs | 5 +- .../mutable_numeric_index/lifecycle.rs | 6 +- .../payload_storage/payload_storage_impl.rs | 6 +- .../sparse/mmap_sparse_vector_storage.rs | 4 +- 20 files changed, 427 insertions(+), 95 deletions(-) diff --git a/lib/common/common/src/universal_io/disk_cache/mod.rs b/lib/common/common/src/universal_io/disk_cache/mod.rs index aca09a2443..6af34b9099 100644 --- a/lib/common/common/src/universal_io/disk_cache/mod.rs +++ b/lib/common/common/src/universal_io/disk_cache/mod.rs @@ -98,6 +98,26 @@ impl UniversalReadFileOps for BlockCacheFs { fn exists(&self, path: &Path) -> Result { fs::exists(path).map_err(UniversalIoError::from) } + + fn create(&self, path: &Path, expected_length: usize) -> Result<()> { + local_file_ops::local_create(path, expected_length) + } + + fn create_dir(&self, path: &Path) -> Result<()> { + local_file_ops::local_create_dir(path) + } + + fn remove(&self, path: &Path) -> Result<()> { + local_file_ops::local_remove(path) + } + + fn remove_dir(&self, path: &Path) -> Result<()> { + local_file_ops::local_remove_dir(path) + } + + fn atomic_save(&self, path: &Path, bytes: &[u8]) -> Result<()> { + local_file_ops::local_atomic_save(path, bytes) + } } impl UniversalReadFs for BlockCacheFs { diff --git a/lib/common/common/src/universal_io/io_uring/mod.rs b/lib/common/common/src/universal_io/io_uring/mod.rs index c649cbfea7..460b15d683 100644 --- a/lib/common/common/src/universal_io/io_uring/mod.rs +++ b/lib/common/common/src/universal_io/io_uring/mod.rs @@ -66,6 +66,26 @@ impl UniversalReadFileOps for IoUringFs { fn exists(&self, path: &Path) -> Result { fs::exists(path).map_err(UniversalIoError::from) } + + fn create(&self, path: &Path, expected_length: usize) -> Result<()> { + local_file_ops::local_create(path, expected_length) + } + + fn create_dir(&self, path: &Path) -> Result<()> { + local_file_ops::local_create_dir(path) + } + + fn remove(&self, path: &Path) -> Result<()> { + local_file_ops::local_remove(path) + } + + fn remove_dir(&self, path: &Path) -> Result<()> { + local_file_ops::local_remove_dir(path) + } + + fn atomic_save(&self, path: &Path, bytes: &[u8]) -> Result<()> { + local_file_ops::local_atomic_save(path, bytes) + } } /// Per-open backend extras for [`IoUringFs::open`]. diff --git a/lib/common/common/src/universal_io/local_file_ops.rs b/lib/common/common/src/universal_io/local_file_ops.rs index 6d20efa529..6368788b9f 100644 --- a/lib/common/common/src/universal_io/local_file_ops.rs +++ b/lib/common/common/src/universal_io/local_file_ops.rs @@ -1,7 +1,34 @@ +use std::io::Write as _; use std::path::{Path, PathBuf}; +use crate::fs::atomic_save; +use crate::mmap::create_and_ensure_length; use crate::universal_io::UniversalIoError; +pub fn local_create(path: &Path, expected_length: usize) -> crate::universal_io::Result<()> { + create_and_ensure_length(path, expected_length) + .map(drop) + .map_err(|err| UniversalIoError::extract_not_found(err, path)) +} + +pub fn local_create_dir(path: &Path) -> crate::universal_io::Result<()> { + fs_err::create_dir_all(path).map_err(|err| UniversalIoError::extract_not_found(err, path)) +} + +pub fn local_remove(path: &Path) -> crate::universal_io::Result<()> { + fs_err::remove_file(path).map_err(|err| UniversalIoError::extract_not_found(err, path)) +} + +pub fn local_remove_dir(path: &Path) -> crate::universal_io::Result<()> { + fs_err::remove_dir_all(path).map_err(|err| UniversalIoError::extract_not_found(err, path)) +} + +pub fn local_atomic_save(path: &Path, bytes: &[u8]) -> crate::universal_io::Result<()> { + atomic_save(path, |writer| { + writer.write_all(bytes).map_err(UniversalIoError::from) + }) +} + pub fn local_list_files(prefix_path: &Path) -> crate::universal_io::Result> { let dir = prefix_path.parent().unwrap_or(Path::new(".")); let file_prefix = prefix_path diff --git a/lib/common/common/src/universal_io/mmap/mod.rs b/lib/common/common/src/universal_io/mmap/mod.rs index 572144182e..eb32752010 100644 --- a/lib/common/common/src/universal_io/mmap/mod.rs +++ b/lib/common/common/src/universal_io/mmap/mod.rs @@ -35,6 +35,26 @@ impl UniversalReadFileOps for MmapFs { fn exists(&self, path: &Path) -> Result { fs_err::exists(path).map_err(UniversalIoError::from) } + + fn create(&self, path: &Path, expected_length: usize) -> Result<()> { + local_file_ops::local_create(path, expected_length) + } + + fn create_dir(&self, path: &Path) -> Result<()> { + local_file_ops::local_create_dir(path) + } + + fn remove(&self, path: &Path) -> Result<()> { + local_file_ops::local_remove(path) + } + + fn remove_dir(&self, path: &Path) -> Result<()> { + local_file_ops::local_remove_dir(path) + } + + fn atomic_save(&self, path: &Path, bytes: &[u8]) -> Result<()> { + local_file_ops::local_atomic_save(path, bytes) + } } impl UniversalReadFs for MmapFs { diff --git a/lib/common/common/src/universal_io/simple_disk_cache/fs.rs b/lib/common/common/src/universal_io/simple_disk_cache/fs.rs index 9db789ef62..4d603a487f 100644 --- a/lib/common/common/src/universal_io/simple_disk_cache/fs.rs +++ b/lib/common/common/src/universal_io/simple_disk_cache/fs.rs @@ -84,6 +84,26 @@ where fn exists(&self, path: &Path) -> Result { self.remote_fs.exists(path) } + + fn create(&self, path: &Path, expected_length: usize) -> Result<()> { + self.remote_fs.create(path, expected_length) + } + + fn create_dir(&self, path: &Path) -> Result<()> { + self.remote_fs.create_dir(path) + } + + fn remove(&self, path: &Path) -> Result<()> { + self.remote_fs.remove(path) + } + + fn remove_dir(&self, path: &Path) -> Result<()> { + self.remote_fs.remove_dir(path) + } + + fn atomic_save(&self, path: &Path, bytes: &[u8]) -> Result<()> { + self.remote_fs.atomic_save(path, bytes) + } } impl UniversalReadFs for DiskCacheFs diff --git a/lib/common/common/src/universal_io/traits/file_ops.rs b/lib/common/common/src/universal_io/traits/file_ops.rs index ae662f161d..fdb0beef81 100644 --- a/lib/common/common/src/universal_io/traits/file_ops.rs +++ b/lib/common/common/src/universal_io/traits/file_ops.rs @@ -16,7 +16,7 @@ use crate::universal_io::{OpenOptions, Result}; /// a backend can implement this trait to expose metadata-style operations /// without ever opening file handles. The "open files" capability lives on /// the [`UniversalReadFs`] subtrait. -pub trait UniversalReadFileOps: Sized + Debug { +pub trait UniversalReadFileOps: Clone + Debug + Sized { /// Implementation-specific construction config. Backends are free to /// require explicit construction; callers that want to opt into the /// `::default()` pattern must constrain @@ -37,6 +37,31 @@ pub trait UniversalReadFileOps: Sized + Debug { /// Check whether a file exists at the given path. fn exists(&self, path: &Path) -> Result; + /// Create or truncate a file at the given path. + /// + /// Local backends use `expected_length` to pre-size the file. Backends + /// without fixed-size file objects may ignore it. + fn create(&self, path: &Path, expected_length: usize) -> Result<()>; + + /// Create a directory at the given path. + /// + /// Backends without materialized directories may treat this as a no-op. + fn create_dir(&self, path: &Path) -> Result<()>; + + /// Remove a file at the given path. + fn remove(&self, path: &Path) -> Result<()>; + + /// Remove a directory at the given path. + /// + /// Backends without materialized directories may treat this as a no-op. + fn remove_dir(&self, path: &Path) -> Result<()>; + + /// Atomically save bytes at the given path. + /// + /// Local backends should use an atomic file replacement. Object-store + /// backends may overwrite the full object. + fn atomic_save(&self, path: &Path, bytes: &[u8]) -> Result<()>; + // When adding provided methods, don't forget to update impls in // `crate::universal_io::wrappers::*`. } diff --git a/lib/common/common/src/universal_io/wrappers/read_only.rs b/lib/common/common/src/universal_io/wrappers/read_only.rs index 72775e866e..e6b4d3baf7 100644 --- a/lib/common/common/src/universal_io/wrappers/read_only.rs +++ b/lib/common/common/src/universal_io/wrappers/read_only.rs @@ -10,8 +10,8 @@ use crate::ext::aligned_vec::ACow; use crate::generic_consts::AccessPattern; use crate::universal_io::traits::UniversalReadFileOps; use crate::universal_io::{ - Item, OpenOptions, ReadBytesItem, ReadRange, Result, UniversalKind, UniversalRead, - UniversalReadFs, UserData, + Item, OpenOptions, ReadBytesItem, ReadRange, Result, UniversalIoError, UniversalKind, + UniversalRead, UniversalReadFs, UserData, }; #[derive(Debug, TransparentWrapper)] @@ -26,6 +26,7 @@ pub struct ReadOnly(S); /// and asserts read-only semantics on `open`. In practice this Fs is /// rarely instantiated; callers use [`ReadOnly::open`] with the /// underlying `&S::Fs` directly. +#[derive(Clone)] pub struct ReadOnlyFs(F); impl fmt::Debug for ReadOnlyFs { @@ -48,6 +49,36 @@ impl UniversalReadFileOps for ReadOnlyFs { fn exists(&self, path: &Path) -> Result { self.0.exists(path) } + + fn create(&self, _path: &Path, _expected_length: usize) -> Result<()> { + Err(UniversalIoError::uninitialized( + "ReadOnlyFs does not support creating files", + )) + } + + fn create_dir(&self, _path: &Path) -> Result<()> { + Err(UniversalIoError::uninitialized( + "ReadOnlyFs does not support creating directories", + )) + } + + fn remove(&self, _path: &Path) -> Result<()> { + Err(UniversalIoError::uninitialized( + "ReadOnlyFs does not support removing files", + )) + } + + fn remove_dir(&self, _path: &Path) -> Result<()> { + Err(UniversalIoError::uninitialized( + "ReadOnlyFs does not support removing directories", + )) + } + + fn atomic_save(&self, _path: &Path, _bytes: &[u8]) -> Result<()> { + Err(UniversalIoError::uninitialized( + "ReadOnlyFs does not support atomic saves", + )) + } } impl UniversalReadFs for ReadOnlyFs { diff --git a/lib/common/io_bridge_object_store/src/file.rs b/lib/common/io_bridge_object_store/src/file.rs index 5158852de6..58f28a1338 100644 --- a/lib/common/io_bridge_object_store/src/file.rs +++ b/lib/common/io_bridge_object_store/src/file.rs @@ -153,6 +153,26 @@ mod tests { std::future::ready(Ok(true)) } + fn create(&self, _path: &Path) -> impl Future> + Send + 'static { + std::future::ready(Ok(())) + } + + fn remove(&self, _path: &Path) -> impl Future> + Send + 'static { + std::future::ready(Ok(())) + } + + fn remove_dir(&self, _path: &Path) -> impl Future> + Send + 'static { + std::future::ready(Ok(())) + } + + fn atomic_save( + &self, + _path: &Path, + _bytes: Bytes, + ) -> impl Future> + Send + 'static { + std::future::ready(Ok(())) + } + fn read_range( &self, _path: &Path, diff --git a/lib/common/io_bridge_object_store/src/fs.rs b/lib/common/io_bridge_object_store/src/fs.rs index c61cf9ba51..9e7ccb5219 100644 --- a/lib/common/io_bridge_object_store/src/fs.rs +++ b/lib/common/io_bridge_object_store/src/fs.rs @@ -1,5 +1,6 @@ use std::path::{Path, PathBuf}; +use bytes::Bytes; use common::universal_io::{OpenOptions, Result, UniversalReadFileOps, UniversalReadFs}; use crate::{AsyncRead, BlobFile, BridgeRuntime}; @@ -8,6 +9,7 @@ use crate::{AsyncRead, BlobFile, BridgeRuntime}; /// the [`BridgeRuntime`] used to drive its async operations. Opens per-object /// [`BlobFile`] handles via [`UniversalReadFs::open`] and answers metadata /// queries (`list_files`, `exists`) by blocking on the backend. +#[derive(Clone)] pub struct BlobFs { inner: A, runtime: BridgeRuntime, @@ -27,7 +29,7 @@ impl BlobFs { } } -impl UniversalReadFileOps for BlobFs { +impl UniversalReadFileOps for BlobFs { type ContextConfig = A::Config; fn from_context(config: Self::ContextConfig) -> Result { @@ -43,6 +45,27 @@ impl UniversalReadFileOps for BlobFs { fn exists(&self, path: &Path) -> Result { self.runtime.block_on(self.inner.exists(path)) } + + fn create(&self, path: &Path, _expected_length: usize) -> Result<()> { + self.runtime.block_on(self.inner.create(path)) + } + + fn create_dir(&self, _path: &Path) -> Result<()> { + Ok(()) + } + + fn remove(&self, path: &Path) -> Result<()> { + self.runtime.block_on(self.inner.remove(path)) + } + + fn remove_dir(&self, path: &Path) -> Result<()> { + self.runtime.block_on(self.inner.remove_dir(path)) + } + + fn atomic_save(&self, path: &Path, bytes: &[u8]) -> Result<()> { + self.runtime + .block_on(self.inner.atomic_save(path, Bytes::copy_from_slice(bytes))) + } } impl UniversalReadFs for BlobFs { diff --git a/lib/common/io_bridge_object_store/src/read.rs b/lib/common/io_bridge_object_store/src/read.rs index eaf6f5f6ca..9405f534ca 100644 --- a/lib/common/io_bridge_object_store/src/read.rs +++ b/lib/common/io_bridge_object_store/src/read.rs @@ -36,6 +36,22 @@ pub trait AsyncRead: Send + Sync + Sized + 'static { fn exists(&self, path: &Path) -> impl Future> + Send + 'static; + /// Create or truncate an empty object at `path`. + fn create(&self, path: &Path) -> impl Future> + Send + 'static; + + /// Remove the object at `path`. + fn remove(&self, path: &Path) -> impl Future> + Send + 'static; + + /// Remove all objects matching the directory prefix at `path`. + fn remove_dir(&self, path: &Path) -> impl Future> + Send + 'static; + + /// Save bytes by overwriting the full object at `path`. + fn atomic_save( + &self, + path: &Path, + bytes: Bytes, + ) -> impl Future> + Send + 'static; + /// Fetch `range` from `path` as a stream of byte chunks. /// /// The returned future resolves once the request has been initiated and the diff --git a/lib/common/io_bridge_object_store/src/source.rs b/lib/common/io_bridge_object_store/src/source.rs index d503b89bd9..8f6faeaff8 100644 --- a/lib/common/io_bridge_object_store/src/source.rs +++ b/lib/common/io_bridge_object_store/src/source.rs @@ -69,6 +69,73 @@ impl AsyncRead for Arc { } } + fn create(&self, path: &Path) -> impl Future> + Send + 'static { + let store = self.clone(); + let key = build_key(path); + + async move { + store + .put(&key, Bytes::new().into()) + .await + .map(drop) + .map_err(UniversalIoError::s3) + } + } + + fn remove(&self, path: &Path) -> impl Future> + Send + 'static { + let store = self.clone(); + let key = build_key(path); + + async move { + store.delete(&key).await.map_err(|err| match err { + object_store::Error::NotFound { .. } => UniversalIoError::NotFound { + path: PathBuf::from(key.to_string()), + }, + err => UniversalIoError::s3(err), + }) + } + } + + fn remove_dir(&self, path: &Path) -> impl Future> + Send + 'static { + let store = self.clone(); + let prefix_path = path.to_path_buf(); + let prefix = build_dir_prefix(path); + + async move { + let mut objects = store.list(Some(&prefix)); + while let Some(meta) = objects.try_next().await.map_err(|err| match err { + object_store::Error::NotFound { .. } => UniversalIoError::NotFound { + path: prefix_path.clone(), + }, + other => UniversalIoError::s3(other), + })? { + match store.delete(&meta.location).await { + Ok(()) | Err(object_store::Error::NotFound { .. }) => {} + Err(other) => return Err(UniversalIoError::s3(other)), + } + } + + Ok(()) + } + } + + fn atomic_save( + &self, + path: &Path, + bytes: Bytes, + ) -> impl Future> + Send + 'static { + let store = self.clone(); + let key = build_key(path); + + async move { + store + .put(&key, bytes.into()) + .await + .map(drop) + .map_err(UniversalIoError::s3) + } + } + fn read_range( &self, path: &Path, @@ -117,10 +184,20 @@ fn build_key(path: &Path) -> object_store::path::Path { object_store::path::Path::from(path.to_string_lossy().as_ref()) } +fn build_dir_prefix(path: &Path) -> object_store::path::Path { + let path = path.to_string_lossy(); + let path = path.trim_end_matches('/'); + if path.is_empty() { + object_store::path::Path::from("") + } else { + object_store::path::Path::from(format!("{path}/")) + } +} + #[cfg(test)] mod tests { use bytes::Bytes; - use common::universal_io::{ReadRange, UniversalRead}; + use common::universal_io::{ReadRange, UniversalRead, UniversalReadFileOps}; use object_store::ObjectStoreExt as _; use object_store::memory::InMemory; @@ -237,6 +314,32 @@ mod tests { assert_eq!(len, 2); } + #[test] + fn fs_create_writes_empty_object_and_create_dir_is_noop() { + let runtime = BridgeRuntime::global(); + let store = Arc::new(InMemory::new()); + let fs = crate::BlobFs::new(store.clone(), runtime.clone()); + + fs.create_dir(Path::new("prefix")).unwrap(); + fs.create(Path::new("prefix/empty"), 1024).unwrap(); + fs.create(Path::new("prefix/nested/empty"), 1024).unwrap(); + fs.create(Path::new("prefix-sibling/empty"), 1024).unwrap(); + fs.atomic_save(Path::new("prefix/saved"), b"saved").unwrap(); + + assert!(fs.exists(Path::new("prefix/empty")).unwrap()); + let len = runtime.block_on(store.head(&object_store::path::Path::from("prefix/empty"))); + assert_eq!(len.unwrap().size, 0); + let len = runtime.block_on(store.head(&object_store::path::Path::from("prefix/saved"))); + assert_eq!(len.unwrap().size, 5); + + fs.remove(Path::new("prefix/empty")).unwrap(); + fs.remove_dir(Path::new("prefix")).unwrap(); + assert!(!fs.exists(Path::new("prefix/empty")).unwrap()); + assert!(!fs.exists(Path::new("prefix/nested/empty")).unwrap()); + assert!(!fs.exists(Path::new("prefix/saved")).unwrap()); + assert!(fs.exists(Path::new("prefix-sibling/empty")).unwrap()); + } + #[test] fn read_only_wrapper_compiles_with_blob_file() { use common::universal_io::ReadOnly; diff --git a/lib/gridstore/src/fixtures.rs b/lib/gridstore/src/fixtures.rs index d5df309bed..cf4c3ecf70 100644 --- a/lib/gridstore/src/fixtures.rs +++ b/lib/gridstore/src/fixtures.rs @@ -1,3 +1,4 @@ +use common::universal_io::MmapFs; use rand::distr::{Distribution, Uniform}; use rand::{Rng, RngExt}; use serde::{Deserialize, Serialize}; @@ -29,7 +30,7 @@ impl Blob for Payload { /// Create an empty storage with the default configuration pub fn empty_storage() -> (TempDir, Gridstore) { let dir = Builder::new().prefix("test-storage").tempdir().unwrap(); - let storage = Gridstore::new(dir.path().to_path_buf(), Default::default()).unwrap(); + let storage = Gridstore::new(MmapFs, dir.path().to_path_buf(), Default::default()).unwrap(); (dir, storage) } @@ -44,7 +45,7 @@ pub fn empty_storage_sized( compression: Some(compression), ..Default::default() }; - let storage = Gridstore::new(dir.path().to_path_buf(), options).unwrap(); + let storage = Gridstore::new(MmapFs, dir.path().to_path_buf(), options).unwrap(); (dir, storage) } diff --git a/lib/gridstore/src/gridstore/mod.rs b/lib/gridstore/src/gridstore/mod.rs index b6cd957619..918d2e8383 100644 --- a/lib/gridstore/src/gridstore/mod.rs +++ b/lib/gridstore/src/gridstore/mod.rs @@ -11,12 +11,9 @@ use ahash::AHashMap; use common::counter::counter_cell::CounterCell; use common::counter::hardware_counter::HardwareCounterCell; use common::counter::referenced_counter::HwMetricRefCounter; -use common::fs::atomic_save_json; use common::generic_consts::{AccessPattern, Random}; use common::is_alive_lock::IsAliveLock; -use common::mmap::create_and_ensure_length; -use common::universal_io::{MmapFile, UniversalReadFs, UniversalWrite}; -use fs_err as fs; +use common::universal_io::{MmapFile, UniversalReadFileOps, UniversalWrite}; use itertools::Itertools; use parking_lot::RwLock; use reader::CONFIG_FILENAME; @@ -24,7 +21,7 @@ pub use reader::GridstoreReader; pub use view::GridstoreView; use crate::Result; -use crate::bitmask::MmapBitmask; +use crate::bitmask::Bitmask; use crate::blob::Blob; use crate::config::{StorageConfig, StorageOptions}; use crate::error::GridstoreError; @@ -38,14 +35,18 @@ pub type Flusher = Box std::result::Result<(), GridstoreError> + /// Uses `Arc>` for pages and tracker to support concurrent flushing. /// Assumes sequential IDs to the values (0, 1, 2, 3, ...) #[derive(Debug)] -pub struct Gridstore { +pub struct Gridstore +where + S: UniversalWrite + 'static, +{ + pub(super) fs: S::Fs, pub(super) config: StorageConfig, pub(super) tracker: Arc>>, pub(super) pages: Arc>>, /// MmapBitmask to represent which "blocks" of data in the pages are used and which are free. /// /// 0 is free, 1 is used. - pub(super) bitmask: Arc>, + pub(super) bitmask: Arc>>, pub(super) base_path: PathBuf, pub(super) _value_type: std::marker::PhantomData, /// Lock to prevent concurrent flushes and used for waiting for ongoing flushes to finish. @@ -56,7 +57,6 @@ impl Gridstore where V: Blob, S: UniversalWrite + 'static, - S::Fs: Default, { /// Create a [`GridstoreView`] by locking pages and tracker, then call `f` with the view. fn with_view(&self, f: impl FnOnce(GridstoreView<'_, V, S>) -> R) -> R { @@ -93,17 +93,17 @@ where /// Depends on the existence of the config file at the `base_path`. /// /// In case of opening, it ignores the `create_options` parameter. - pub fn open_or_create(base_path: PathBuf, create_options: StorageOptions) -> Result { + pub fn open_or_create( + fs: S::Fs, + base_path: PathBuf, + create_options: StorageOptions, + ) -> Result { let config_path = base_path.join(CONFIG_FILENAME); if config_path.exists() { - Self::open(base_path) + Self::open(fs, base_path) } else { - fs::create_dir_all(&base_path).map_err(|err| { - GridstoreError::service_error(format!( - "Failed to create gridstore storage directory: {err}" - )) - })?; - Self::new(base_path, create_options) + fs.create_dir(&base_path)?; + Self::new(fs, base_path, create_options) } } @@ -111,17 +111,18 @@ where /// /// `base_path` is the directory where the storage files will be stored. /// It should exist already. - pub fn new(base_path: PathBuf, options: StorageOptions) -> Result { + pub fn new(fs: S::Fs, base_path: PathBuf, options: StorageOptions) -> Result { let config = StorageConfig::try_from(options).map_err(GridstoreError::service_error)?; let config_path = base_path.join(CONFIG_FILENAME); - let bitmask = MmapBitmask::create(&default_fs(), &base_path, config.clone())?; + let bitmask = Bitmask::create(&fs, &base_path, config.clone())?; let storage = Self { - tracker: Arc::new(RwLock::new(Tracker::new(&default_fs(), &base_path, None)?)), + tracker: Arc::new(RwLock::new(Tracker::new(&fs, &base_path, None)?)), pages: Arc::new(RwLock::new(Pages::new(base_path.clone(), true))), base_path, config, + fs, _value_type: std::marker::PhantomData, bitmask: Arc::new(RwLock::new(bitmask)), is_alive_flush_lock: IsAliveLock::new(), @@ -129,11 +130,11 @@ where let new_page_id = storage.next_page_id(); let path = page_path(&storage.base_path, new_page_id); - create_and_ensure_length(&path, storage.config.page_size_bytes)?; - storage.pages.write().attach_page(&default_fs(), &path)?; + storage.create_page_file(&path)?; + storage.pages.write().attach_page(&storage.fs, &path)?; - atomic_save_json(&config_path, &storage.config) - .map_err(|err| GridstoreError::service_error(err.to_string()))?; + let config_bytes = serde_json::to_vec(&storage.config)?; + storage.fs.atomic_save(&config_path, &config_bytes)?; Ok(storage) } @@ -141,13 +142,13 @@ where /// Open an existing storage at the given path. /// /// Uses the bitmask to infer page count for consistency with the write path. - pub fn open(base_path: PathBuf) -> Result { + pub fn open(fs: S::Fs, base_path: PathBuf) -> Result { // Writable store: open pages and tracker writable so it can append. - let (config, tracker) = reader::read_config_and_tracker(&default_fs(), &base_path, true)?; - let bitmask = MmapBitmask::open(&default_fs(), &base_path, config.clone())?; + let (config, tracker) = reader::read_config_and_tracker(&fs, &base_path, true)?; + let bitmask = Bitmask::open(&fs, &base_path, config.clone())?; let num_pages = bitmask.infer_num_pages(); - let pages = Pages::open(&default_fs(), &base_path, true)?; + let pages = Pages::open(&fs, &base_path, true)?; let loaded_pages = pages.num_pages(); if loaded_pages != num_pages { @@ -157,6 +158,7 @@ where } Ok(Self { + fs, config, tracker: Arc::new(RwLock::new(tracker)), pages: Arc::new(RwLock::new(pages)), @@ -172,14 +174,19 @@ where fn create_new_page(&mut self) -> Result { let new_page_id = self.next_page_id(); let path = page_path(&self.base_path, new_page_id); - create_and_ensure_length(&path, self.config.page_size_bytes)?; - self.pages.write().attach_page(&default_fs(), &path)?; + self.create_page_file(&path)?; + self.pages.write().attach_page(&self.fs, &path)?; self.bitmask.write().cover_new_page()?; Ok(new_page_id) } + fn create_page_file(&self, path: &std::path::Path) -> Result<()> { + self.fs.create(path, self.config.page_size_bytes)?; + Ok(()) + } + fn find_or_create_available_blocks( &mut self, num_blocks: u32, @@ -329,24 +336,17 @@ where /// /// Completely wipes the storage, and recreates it with a single empty page. pub fn clear(&mut self) -> Result<()> { - let create_options = StorageOptions::from(&self.config); - let base_path = self.base_path.clone(); - self.is_alive_flush_lock.blocking_mark_dead(); - self.pages.write().clear(); - fs::remove_dir_all(&base_path).map_err(|err| { - GridstoreError::service_error(format!( - "Failed to remove gridstore storage directory: {err}" - )) - })?; - fs::create_dir_all(&base_path).map_err(|err| { - GridstoreError::service_error(format!( - "Failed to create gridstore storage directory: {err}" - )) - })?; - *self = Self::new(base_path, create_options)?; + self.fs.remove_dir(&self.base_path)?; + self.fs.create_dir(&self.base_path)?; + + *self = Self::new( + self.fs.clone(), + self.base_path.clone(), + StorageOptions::from(&self.config), + )?; Ok(()) } @@ -356,17 +356,22 @@ where /// Takes ownership because this function leaves Gridstore in an inconsistent state which does /// not allow further usage. Use [`clear`](Self::clear) instead to clear and reuse the storage. pub fn wipe(self) -> Result<()> { - let base_path = self.base_path.clone(); + let Self { + fs, + tracker, + pages, + bitmask, + base_path, + config: _, + _value_type, + is_alive_flush_lock, + } = self; - self.is_alive_flush_lock.blocking_mark_dead(); + is_alive_flush_lock.blocking_mark_dead(); + drop((tracker, pages, bitmask)); - drop(self); - - fs::remove_dir_all(base_path).map_err(|err| { - GridstoreError::service_error(format!( - "Failed to remove gridstore storage directory: {err}" - )) - }) + fs.remove_dir(&base_path)?; + Ok(()) } /// Return the storage size in bytes (precise, based on bitmask occupancy). @@ -503,7 +508,7 @@ impl Gridstore { /// Update all free blocks in the bitmask for old pointers and flush it. fn flush_free_blocks( - bitmask: &Arc>, + bitmask: &Arc>>, old_pointers: Vec, block_size_bytes: usize, ) -> crate::Result<()> { @@ -537,6 +542,7 @@ impl Gridstore { /// Drop disk cache. pub fn clear_cache(&self) -> crate::Result<()> { let Self { + fs: _, config: _, tracker: _, pages, @@ -550,7 +556,3 @@ impl Gridstore { Ok(()) } } - -fn default_fs() -> T { - T::default() -} diff --git a/lib/gridstore/src/gridstore/tests.rs b/lib/gridstore/src/gridstore/tests.rs index a872b69983..e2d5daae4d 100644 --- a/lib/gridstore/src/gridstore/tests.rs +++ b/lib/gridstore/src/gridstore/tests.rs @@ -5,6 +5,7 @@ use std::time::Duration; use common::counter::hardware_counter::HardwareCounterCell; use common::generic_consts::{Random, Sequential}; use common::universal_io::MmapFs; +use fs_err as fs; use fs_err::File; use itertools::Itertools; use rand::distr::Uniform; @@ -542,7 +543,7 @@ fn test_behave_like_hashmap( drop(storage); // reopen storage - let storage = Gridstore::::open(dir.path().to_path_buf()).unwrap(); + let storage = Gridstore::::open(MmapFs, dir.path().to_path_buf()).unwrap(); // assert same size assert_eq!(storage.get_storage_size_bytes().unwrap(), before_size); // assert same length @@ -640,7 +641,7 @@ fn test_storage_persistence_basic() { let hw_counter = HardwareCounterCell::new(); let hw_counter_ref = hw_counter.ref_payload_io_write_counter(); { - let mut storage = Gridstore::<_>::new(path.clone(), Default::default()).unwrap(); + let mut storage = Gridstore::<_>::new(MmapFs, path.clone(), Default::default()).unwrap(); storage.put_value(0, &payload, hw_counter_ref).unwrap(); assert_eq!(storage.pages.read().num_pages(), 1); @@ -657,7 +658,7 @@ fn test_storage_persistence_basic() { } // reopen storage - let storage = Gridstore::::open(path).unwrap(); + let storage = Gridstore::::open(MmapFs, path).unwrap(); assert_eq!(storage.pages.read().num_pages(), 1); let stored_payload = storage.get_value::(0, &hw_counter).unwrap(); @@ -750,7 +751,7 @@ fn test_with_real_hm_data() { storage.flusher()().unwrap(); drop(storage); - let mut storage = Gridstore::open(dir.path().to_path_buf()).unwrap(); + let mut storage = Gridstore::open(MmapFs, dir.path().to_path_buf()).unwrap(); assert_eq!(point_offset, EXPECTED_LEN as u32 * 2); assert_eq!(storage.pages.read().num_pages(), 4); assert_eq!( @@ -810,7 +811,7 @@ fn test_different_block_sizes(#[case] block_size_bytes: usize) { block_size_bytes: Some(block_size_bytes), ..Default::default() }; - let mut storage = Gridstore::<_>::new(dir.path().to_path_buf(), options).unwrap(); + let mut storage = Gridstore::<_>::new(MmapFs, dir.path().to_path_buf(), options).unwrap(); let hw_counter = HardwareCounterCell::new(); let hw_counter_ref = hw_counter.ref_payload_io_write_counter(); @@ -907,7 +908,7 @@ fn test_deferred_flush() { // Reopen gridstore drop(storage); - let mut storage = Gridstore::::open(path).unwrap(); + let mut storage = Gridstore::::open(MmapFs, path).unwrap(); assert_eq!(storage.pages.read().num_pages(), 1); // On reopen, we expect to read the data at the time the flusher was created @@ -989,7 +990,7 @@ fn test_deferred_flush_with_delete() { // Reopen gridstore drop(storage); - let mut storage = Gridstore::::open(path.clone()).unwrap(); + let mut storage = Gridstore::::open(MmapFs, path.clone()).unwrap(); assert_eq!(storage.pages.read().num_pages(), 1); let flusher = storage.flusher(); @@ -1008,7 +1009,7 @@ fn test_deferred_flush_with_delete() { // Reopen gridstore drop(storage); - let mut storage = Gridstore::::open(path.clone()).unwrap(); + let mut storage = Gridstore::::open(MmapFs, path.clone()).unwrap(); assert_eq!(storage.pages.read().num_pages(), 1); // On reopen, delete was flushed this time, expect point to be missing @@ -1033,7 +1034,7 @@ fn test_deferred_flush_with_delete() { // Reopen gridstore drop(storage); - let mut storage = Gridstore::::open(path.clone()).unwrap(); + let mut storage = Gridstore::::open(MmapFs, path.clone()).unwrap(); assert_eq!(storage.pages.read().num_pages(), 1); // On reopen, value 4 was flushed, expect to read it @@ -1059,28 +1060,28 @@ fn test_deferred_flush_with_delete() { // Not flushed, still expect to read value 4 { - let tmp_storage = Gridstore::::open(path.clone()).unwrap(); + let tmp_storage = Gridstore::::open(MmapFs, path.clone()).unwrap(); assert_eq!(get_payload(&tmp_storage).unwrap(), "value 4"); } // First flusher flushed, expect to read value 5 if we load from disk flusher_1_value_5().unwrap(); { - let tmp_storage = Gridstore::::open(path.clone()).unwrap(); + let tmp_storage = Gridstore::::open(MmapFs, path.clone()).unwrap(); assert_eq!(get_payload(&tmp_storage).unwrap(), "value 5"); } // Second flusher flushed, expect point to be missing if we load from disk flusher_2_delete().unwrap(); { - let tmp_storage = Gridstore::::open(path.clone()).unwrap(); + let tmp_storage = Gridstore::::open(MmapFs, path.clone()).unwrap(); assert!(get_payload(&tmp_storage).is_none()); } // Third flusher flushed, expect to read value 6 if we load from disk flusher_3_value_6().unwrap(); { - let tmp_storage = Gridstore::::open(path).unwrap(); + let tmp_storage = Gridstore::::open(MmapFs, path).unwrap(); assert_eq!(get_payload(&tmp_storage).unwrap(), "value 6"); } @@ -1120,7 +1121,7 @@ fn test_live_reload() { }; // Step 1: Write initial data and flush - let mut storage = Gridstore::<_>::new(path.clone(), Default::default()).unwrap(); + let mut storage = Gridstore::<_>::new(MmapFs, path.clone(), Default::default()).unwrap(); let payload_0 = make_payload("key", "value_0"); let payload_1 = make_payload("key", "value_1"); @@ -1190,7 +1191,7 @@ fn test_live_reload_across_pages() { page_size_bytes: Some(page_size), ..Default::default() }; - let mut storage = Gridstore::<_>::new(path.clone(), options).unwrap(); + let mut storage = Gridstore::<_>::new(MmapFs, path.clone(), options).unwrap(); let payload = minimal_payload(); @@ -1309,7 +1310,7 @@ fn test_skip_deferred_flush_after_clear() { // If we reopen the storage it must still be empty drop(storage); - let storage = Gridstore::::open(path.clone()).unwrap(); + let storage = Gridstore::::open(MmapFs, path.clone()).unwrap(); assert_eq!(storage.pages.read().num_pages(), 1); assert!(storage.get_pointer(0).is_none(), "point must not exist"); assert_eq!(storage.max_point_offset(), 0, "must have zero points"); diff --git a/lib/segment/src/index/field_index/full_text_index/mutable_text_index/lifecycle.rs b/lib/segment/src/index/field_index/full_text_index/mutable_text_index/lifecycle.rs index ffcd894322..4d4a4f69c3 100644 --- a/lib/segment/src/index/field_index/full_text_index/mutable_text_index/lifecycle.rs +++ b/lib/segment/src/index/field_index/full_text_index/mutable_text_index/lifecycle.rs @@ -3,6 +3,7 @@ use std::path::PathBuf; use common::counter::hardware_counter::HardwareCounterCell; use common::types::PointOffsetType; +use common::universal_io::MmapFs; use gridstore::Gridstore; use itertools::Itertools; @@ -29,13 +30,13 @@ impl MutableFullTextIndex { create_if_missing: bool, ) -> OperationResult> { let store = if create_if_missing { - Gridstore::open_or_create(path, GRIDSTORE_OPTIONS).map_err(|err| { + Gridstore::open_or_create(MmapFs, path, GRIDSTORE_OPTIONS).map_err(|err| { OperationError::service_error(format!( "failed to open mutable full text index on gridstore: {err}" )) })? } else if path.exists() { - Gridstore::open(path).map_err(|err| { + Gridstore::open(MmapFs, path).map_err(|err| { OperationError::service_error(format!( "failed to open mutable full text index on gridstore: {err}" )) diff --git a/lib/segment/src/index/field_index/geo_index/mutable_geo_index/lifecycle.rs b/lib/segment/src/index/field_index/geo_index/mutable_geo_index/lifecycle.rs index 92906a07a8..a8418c826f 100644 --- a/lib/segment/src/index/field_index/geo_index/mutable_geo_index/lifecycle.rs +++ b/lib/segment/src/index/field_index/geo_index/mutable_geo_index/lifecycle.rs @@ -2,6 +2,7 @@ use std::path::PathBuf; use common::counter::hardware_counter::HardwareCounterCell; use common::types::PointOffsetType; +use common::universal_io::MmapFs; use gridstore::Gridstore; use gridstore::config::StorageOptions; @@ -30,13 +31,13 @@ impl MutableGeoIndex { /// loaded. pub fn open(path: PathBuf, create_if_missing: bool) -> OperationResult> { let store = if create_if_missing { - Gridstore::open_or_create(path, GRIDSTORE_OPTIONS).map_err(|err| { + Gridstore::open_or_create(MmapFs, path, GRIDSTORE_OPTIONS).map_err(|err| { OperationError::service_error(format!( "failed to open mutable geo index on gridstore: {err}" )) })? } else if path.exists() { - Gridstore::open(path).map_err(|err| { + Gridstore::open(MmapFs, path).map_err(|err| { OperationError::service_error(format!( "failed to open mutable geo index on gridstore: {err}" )) diff --git a/lib/segment/src/index/field_index/map_index/mutable_map_index/lifecycle.rs b/lib/segment/src/index/field_index/map_index/mutable_map_index/lifecycle.rs index 8785955540..5fdc9d5ed7 100644 --- a/lib/segment/src/index/field_index/map_index/mutable_map_index/lifecycle.rs +++ b/lib/segment/src/index/field_index/map_index/mutable_map_index/lifecycle.rs @@ -2,6 +2,7 @@ use std::path::PathBuf; use common::counter::hardware_counter::HardwareCounterCell; use common::types::PointOffsetType; +use common::universal_io::MmapFs; use gridstore::config::StorageOptions; use gridstore::error::GridstoreError; use gridstore::{Blob, Gridstore}; @@ -35,13 +36,13 @@ where pub fn open_gridstore(path: PathBuf, create_if_missing: bool) -> OperationResult> { let store = if create_if_missing { let options = default_gridstore_options(N::gridstore_block_size()); - Gridstore::open_or_create(path, options).map_err(|err| { + Gridstore::open_or_create(MmapFs, path, options).map_err(|err| { OperationError::service_error(format!( "failed to open mutable map index on gridstore: {err}" )) })? } else if path.exists() { - Gridstore::open(path).map_err(|err| { + Gridstore::open(MmapFs, path).map_err(|err| { OperationError::service_error(format!( "failed to open mutable map index on gridstore: {err}" )) diff --git a/lib/segment/src/index/field_index/numeric_index/mutable_numeric_index/lifecycle.rs b/lib/segment/src/index/field_index/numeric_index/mutable_numeric_index/lifecycle.rs index e78a0b3b86..e14d4bf739 100644 --- a/lib/segment/src/index/field_index/numeric_index/mutable_numeric_index/lifecycle.rs +++ b/lib/segment/src/index/field_index/numeric_index/mutable_numeric_index/lifecycle.rs @@ -4,7 +4,7 @@ use std::path::PathBuf; use common::counter::hardware_counter::HardwareCounterCell; use common::types::PointOffsetType; -use common::universal_io::UniversalRead; +use common::universal_io::{MmapFs, UniversalRead}; use gridstore::error::GridstoreError; use gridstore::{Blob, Gridstore}; @@ -157,13 +157,13 @@ where pub fn open_gridstore(path: PathBuf, create_if_missing: bool) -> OperationResult> { let store = if create_if_missing { let options = default_gridstore_options::(); - Gridstore::open_or_create(path, options).map_err(|err| { + Gridstore::open_or_create(MmapFs, path, options).map_err(|err| { OperationError::service_error(format!( "failed to open mutable numeric index on gridstore: {err}" )) })? } else if path.exists() { - Gridstore::open(path).map_err(|err| { + Gridstore::open(MmapFs, path).map_err(|err| { OperationError::service_error(format!( "failed to open mutable numeric index on gridstore: {err}" )) diff --git a/lib/segment/src/payload_storage/payload_storage_impl.rs b/lib/segment/src/payload_storage/payload_storage_impl.rs index bd41043dd5..aa1b13d94f 100644 --- a/lib/segment/src/payload_storage/payload_storage_impl.rs +++ b/lib/segment/src/payload_storage/payload_storage_impl.rs @@ -28,7 +28,7 @@ impl Blob for Payload { } #[derive(Debug)] -pub struct PayloadStorageImpl { +pub struct PayloadStorageImpl { storage: Gridstore, populate: bool, } @@ -52,7 +52,7 @@ where } fn open(path: PathBuf, populate: bool) -> OperationResult { - let storage = Gridstore::open(path).map_err(|err| { + let storage = Gridstore::open(S::Fs::default(), path).map_err(|err| { OperationError::service_error(format!("Failed to open mmap payload storage: {err}")) })?; @@ -64,7 +64,7 @@ where } fn new(path: PathBuf, populate: bool) -> OperationResult { - let storage = Gridstore::new(path, StorageOptions::default())?; + let storage = Gridstore::new(S::Fs::default(), path, StorageOptions::default())?; if populate { storage.populate()?; diff --git a/lib/segment/src/vector_storage/sparse/mmap_sparse_vector_storage.rs b/lib/segment/src/vector_storage/sparse/mmap_sparse_vector_storage.rs index 5de8239d81..fe15ab43cd 100644 --- a/lib/segment/src/vector_storage/sparse/mmap_sparse_vector_storage.rs +++ b/lib/segment/src/vector_storage/sparse/mmap_sparse_vector_storage.rs @@ -55,7 +55,7 @@ impl MmapSparseVectorStorage { fn open(path: &Path) -> OperationResult { // Storage let storage_dir = path.join(STORAGE_DIRNAME); - let storage = Gridstore::open(storage_dir).map_err(|err| { + let storage = Gridstore::open(MmapFs, storage_dir).map_err(|err| { OperationError::service_error(format!( "Failed to open mmap sparse vector storage: {err}" )) @@ -99,7 +99,7 @@ impl MmapSparseVectorStorage { ..Default::default() }; - let storage = Gridstore::new(storage_dir, storage_config).map_err(|err| { + let storage = Gridstore::new(MmapFs, storage_dir, storage_config).map_err(|err| { OperationError::service_error(format!( "Failed to create storage for mmap sparse vectors: {err}" ))