diff --git a/Cargo.toml b/Cargo.toml index 6c35195..f5baf83 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "cassadilia" -version = "0.4.6" +version = "0.4.7" edition = "2024" rust-version = "1.89.0" authors = ["0xdeafbeef"] diff --git a/src/index/mod.rs b/src/index/mod.rs index 5717804..6cef9dc 100644 --- a/src/index/mod.rs +++ b/src/index/mod.rs @@ -180,10 +180,6 @@ where Ok(index) } - pub fn read_state(&self) -> IndexReadGuard<'_, K> { - IndexReadGuard { inner: self.state.read() } - } - pub fn checkpoint(&self, reason: CheckpointReason) -> Result<(), IndexError> { tracing::info!(?reason, "Starting checkpoint operation."); @@ -356,6 +352,12 @@ where } } +impl Index { + pub fn read_state(&self) -> IndexReadGuard<'_, K> { + IndexReadGuard { inner: self.state.read() } + } +} + impl Debug for Index { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { let state_guard = self.state.read(); @@ -381,46 +383,10 @@ pub struct IndexReadGuard<'a, K> { inner: parking_lot::RwLockReadGuard<'a, IndexState>, } -impl<'a, K> IndexReadGuard<'a, K> -where - K: Debug + Clone + Ord, -{ - /// Returns the item for the given key, if it exists. - pub fn get_item(&self, key: &K) -> Option { - self.inner.key_to_hash.get(key).copied() - } - - /// Returns the item for the given key, or an error if it is not present. - pub fn require_item(&self, key: &K) -> Result { - self.get_item(key).ok_or(IndexError::KeyNotFound { key: format!("{key:?}") }) - } - - /// Returns `true` if the index contains the specified key. - pub fn contains_key(&self, key: &K) -> bool { - self.inner.key_to_hash.contains_key(key) - } - - /// Returns an iterator over all entries in ascending key order. - pub fn iter(&self) -> std::collections::btree_map::Iter<'_, K, IndexStateItem> { - self.inner.key_to_hash.iter() - } - - /// Returns an iterator over entries within the specified key range, in ascending order. - /// The range may use a borrowed form of the key (e.g., `&str` for `String`). - pub fn range(&self, range: R) -> std::collections::btree_map::Range<'_, K, IndexStateItem> - where - T: ?Sized + Ord, - K: Borrow + Ord, - R: RangeBounds, - { - self.inner.key_to_hash.range(range) - } - - /// Returns a snapshot of the current key map. - /// Calls `clone` inside. +impl<'a, K> IndexReadGuard<'a, K> { #[must_use] - pub fn keys_snapshot(&self) -> BTreeMap { - self.inner.key_to_hash.clone() + pub fn stats(&self) -> DbStats { + self.inner.stats } /// Returns the number of keys in the index. @@ -447,9 +413,57 @@ where self.inner.hash_to_ref_count.contains_key(hash) } + /// Returns an iterator over entries within the specified key range, in ascending order. + /// The range may use a borrowed form of the key (e.g., `&str` for `String`). + pub fn range(&self, range: R) -> std::collections::btree_map::Range<'_, K, IndexStateItem> + where + T: ?Sized + Ord, + K: Borrow + Ord, + R: RangeBounds, + { + self.inner.key_to_hash.range(range) + } +} + +impl<'a, K> IndexReadGuard<'a, K> +where + K: Ord, +{ + /// Returns the item for the given key, if it exists. + pub fn get_item(&self, key: &K) -> Option { + self.inner.key_to_hash.get(key).copied() + } + + /// Returns `true` if the index contains the specified key. + pub fn contains_key(&self, key: &K) -> bool { + self.inner.key_to_hash.contains_key(key) + } + + /// Returns an iterator over all entries in ascending key order. + pub fn iter(&self) -> std::collections::btree_map::Iter<'_, K, IndexStateItem> { + self.inner.key_to_hash.iter() + } +} + +impl<'a, K> IndexReadGuard<'a, K> +where + K: Debug + Ord, +{ + /// Returns the item for the given key, or an error if it is not present. + pub fn require_item(&self, key: &K) -> Result { + self.get_item(key).ok_or(IndexError::KeyNotFound { key: format!("{key:?}") }) + } +} + +impl<'a, K> IndexReadGuard<'a, K> +where + K: Clone, +{ + /// Returns a snapshot of the current key map. + /// Calls `clone` inside. #[must_use] - pub fn stats(&self) -> DbStats { - self.inner.stats + pub fn keys_snapshot(&self) -> BTreeMap { + self.inner.key_to_hash.clone() } } diff --git a/src/index/state.rs b/src/index/state.rs index 5098156..01b3d14 100644 --- a/src/index/state.rs +++ b/src/index/state.rs @@ -28,10 +28,7 @@ pub struct IndexStateItem { pub blob_size: u64, } -impl IndexState -where - K: Clone + Eq + Ord, -{ +impl IndexState { pub fn new() -> Self { IndexState { key_to_hash: BTreeMap::new(), @@ -57,6 +54,41 @@ where self.stats.index.serialized_size_bytes = index_file_size_bytes; } + pub(crate) fn increment_ref(&mut self, hash: &BlobHash) -> bool { + let entry = self.hash_to_ref_count.entry(*hash).or_default(); + let was_zero = *entry == 0; + *entry += 1; + was_zero + } + + pub(crate) fn decrement_ref( + &mut self, + hash_to_decrement: &BlobHash, + ) -> Result, IndexStateError> { + match self.hash_to_ref_count.get_mut(hash_to_decrement) { + Some(count) => { + if *count == 0 { + return Err(IndexStateError::DecrementZeroRefCount { + hash: *hash_to_decrement, + }); + } + *count -= 1; + if *count == 0 { + self.hash_to_ref_count.remove(hash_to_decrement); + Ok(Some(*hash_to_decrement)) + } else { + Ok(None) + } + } + None => Err(IndexStateError::HashNotFoundForDecrement { hash: *hash_to_decrement }), + } + } +} + +impl IndexState +where + K: Clone + Ord, +{ pub fn apply_logical_op(&mut self, op: &WalOp) -> Result, IndexStateError> { let mut unreferenced_hashes = Vec::new(); @@ -116,34 +148,4 @@ where Ok(unreferenced_hashes) } - - pub(crate) fn increment_ref(&mut self, hash: &BlobHash) -> bool { - let entry = self.hash_to_ref_count.entry(*hash).or_default(); - let was_zero = *entry == 0; - *entry += 1; - was_zero - } - - pub(crate) fn decrement_ref( - &mut self, - hash_to_decrement: &BlobHash, - ) -> Result, IndexStateError> { - match self.hash_to_ref_count.get_mut(hash_to_decrement) { - Some(count) => { - if *count == 0 { - return Err(IndexStateError::DecrementZeroRefCount { - hash: *hash_to_decrement, - }); - } - *count -= 1; - if *count == 0 { - self.hash_to_ref_count.remove(hash_to_decrement); - Ok(Some(*hash_to_decrement)) - } else { - Ok(None) - } - } - None => Err(IndexStateError::HashNotFoundForDecrement { hash: *hash_to_decrement }), - } - } } diff --git a/src/lib.rs b/src/lib.rs index 55e2f70..96d15c5 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -129,6 +129,23 @@ impl Deref for Cas { } } +impl Cas { + /// Returns how many space the database occupies on disk. + /// # NOTE: + /// It does not include the size of dir in cas + /// So it can be around 256 * 256(l1) * 256(l2) bytes depending on the number of directories + /// created and fs used. + #[must_use] + pub fn stats(&self) -> DbStats { + self.0.stats() + } + + #[must_use] + pub fn as_arc(&self) -> &Arc> { + &self.0 + } +} + impl Cas where K: KeyBytes + Clone + Eq + Ord + std::hash::Hash + Debug + Send + Sync + 'static, @@ -193,21 +210,6 @@ where Ok((cas, orphan_stats)) } - - /// Returns how many space the database occupies on disk. - /// # NOTE: - /// It does not include the size of dir in cas - /// So it can be around 256 * 256(l1) * 256(l2) bytes depending on the number of directories - /// created and fs used. - #[must_use] - pub fn stats(&self) -> DbStats { - self.0.stats() - } - - #[must_use] - pub fn as_arc(&self) -> &Arc> { - &self.0 - } } pub struct CasInner { @@ -316,17 +318,6 @@ where Ok(Self { _lockfile: lockfile, paths, index, datasync_channel, cas_manager }) } - /// Returns a read-only view of the index state. - /// The returned guard holds a shared read lock until dropped. - pub fn read_index_state(&self) -> IndexReadGuard<'_, K> { - self.index.read_state() - } - - #[must_use] - pub fn stats(&self) -> DbStats { - self.index.read_state().stats() - } - /// Start a new transaction for the given key. pub fn put(&self, key: K) -> Result, LibError> { Transaction::new(self, key).map_err(|e| match e { @@ -442,11 +433,6 @@ where self.index.checkpoint(CheckpointReason::Explicit).map_err(LibError::Index) } - /// Returns root path of db provided at initialization - pub fn root_path(&self) -> &Path { - self.paths.db_root_path() - } - fn with_blob_item(&self, key: &K, f: F) -> Result, LibError> where F: FnOnce(&IndexStateItem) -> Result, @@ -471,6 +457,24 @@ where } } } +} + +impl CasInner { + /// Returns a read-only view of the index state. + /// The returned guard holds a shared read lock until dropped. + pub fn read_index_state(&self) -> IndexReadGuard<'_, K> { + self.index.read_state() + } + + /// Returns root path of db provided at initialization + pub fn root_path(&self) -> &Path { + self.paths.db_root_path() + } + + #[must_use] + pub fn stats(&self) -> DbStats { + self.index.read_state().stats() + } fn fdatasync(&self, file: File) -> Result<(), LibError> { match &self.datasync_channel { diff --git a/src/orphan.rs b/src/orphan.rs index b298198..9f2d7b8 100644 --- a/src/orphan.rs +++ b/src/orphan.rs @@ -65,10 +65,7 @@ pub struct RecoveryResult { pub errors: Vec, } -impl OrphanStats -where - K: KeyBytes + Clone + Eq + Ord + Hash + Debug + Send + Sync + 'static, -{ +impl OrphanStats { /// Delete orphaned blobs pub fn delete_orphans(&self) -> Result { let mut result = RecoveryResult::default(); diff --git a/src/transaction.rs b/src/transaction.rs index f798e18..81ca5c9 100644 --- a/src/transaction.rs +++ b/src/transaction.rs @@ -39,10 +39,7 @@ where } #[must_use = "Transaction must be completed by calling finish()"] -pub struct Transaction<'a, K> -where - K: Debug, -{ +pub struct Transaction<'a, K> { pub(crate) temp_file: NamedTempFile, pub(crate) cas_inner: &'a CasInner, pub(crate) writer: BufWriter, @@ -87,19 +84,6 @@ where }) } - /// Append `data` to the transaction. - /// Most common usage is to incrementally compress and write chunks. - pub fn write(&mut self, data: &[u8]) -> Result<(), TransactionError> { - self.size += data.len() as u64; - self.hasher.update(data); - self.writer.write_all(data).map_err(|e| TransactionError::StagingFileIo { - operation: StagingFileOp::Write, - path: self.temp_file.path().to_path_buf(), - source: e, - })?; - Ok(()) - } - pub fn finish(self) -> Result<(), crate::LibError> { tracing::debug!("Finishing transaction for key '{:?}'", self.key); self.commit() @@ -142,3 +126,18 @@ where Ok(()) } } + +impl<'a, K> Transaction<'a, K> { + /// Append `data` to the transaction. + /// Most common usage is to incrementally compress and write chunks. + pub fn write(&mut self, data: &[u8]) -> Result<(), TransactionError> { + self.size += data.len() as u64; + self.hasher.update(data); + self.writer.write_all(data).map_err(|e| TransactionError::StagingFileIo { + operation: StagingFileOp::Write, + path: self.temp_file.path().to_path_buf(), + source: e, + })?; + Ok(()) + } +}