Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "cassadilia"
version = "0.4.6"
version = "0.4.7"
edition = "2024"
rust-version = "1.89.0"
authors = ["0xdeafbeef"]
Expand Down
104 changes: 59 additions & 45 deletions src/index/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.");

Expand Down Expand Up @@ -356,6 +352,12 @@ where
}
}

impl<K> Index<K> {
pub fn read_state(&self) -> IndexReadGuard<'_, K> {
IndexReadGuard { inner: self.state.read() }
}
}

impl<K: Debug> Debug for Index<K> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let state_guard = self.state.read();
Expand All @@ -381,46 +383,10 @@ pub struct IndexReadGuard<'a, K> {
inner: parking_lot::RwLockReadGuard<'a, IndexState<K>>,
}

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<IndexStateItem> {
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<IndexStateItem, IndexError> {
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<T, R>(&self, range: R) -> std::collections::btree_map::Range<'_, K, IndexStateItem>
where
T: ?Sized + Ord,
K: Borrow<T> + Ord,
R: RangeBounds<T>,
{
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<K, IndexStateItem> {
self.inner.key_to_hash.clone()
pub fn stats(&self) -> DbStats {
self.inner.stats
}

/// Returns the number of keys in the index.
Expand All @@ -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<T, R>(&self, range: R) -> std::collections::btree_map::Range<'_, K, IndexStateItem>
where
T: ?Sized + Ord,
K: Borrow<T> + Ord,
R: RangeBounds<T>,
{
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<IndexStateItem> {
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<IndexStateItem, IndexError> {
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<K, IndexStateItem> {
self.inner.key_to_hash.clone()
}
}

Expand Down
70 changes: 36 additions & 34 deletions src/index/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,10 +28,7 @@ pub struct IndexStateItem {
pub blob_size: u64,
}

impl<K> IndexState<K>
where
K: Clone + Eq + Ord,
{
impl<K> IndexState<K> {
pub fn new() -> Self {
IndexState {
key_to_hash: BTreeMap::new(),
Expand All @@ -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<Option<BlobHash>, 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<K> IndexState<K>
where
K: Clone + Ord,
{
pub fn apply_logical_op(&mut self, op: &WalOp<K>) -> Result<Vec<BlobHash>, IndexStateError> {
let mut unreferenced_hashes = Vec::new();

Expand Down Expand Up @@ -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<Option<BlobHash>, 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 }),
}
}
}
66 changes: 35 additions & 31 deletions src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -129,6 +129,23 @@ impl<K> Deref for Cas<K> {
}
}

impl<K> Cas<K> {
/// 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<CasInner<K>> {
&self.0
}
}

impl<K> Cas<K>
where
K: KeyBytes + Clone + Eq + Ord + std::hash::Hash + Debug + Send + Sync + 'static,
Expand Down Expand Up @@ -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<CasInner<K>> {
&self.0
}
}

pub struct CasInner<K> {
Expand Down Expand Up @@ -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<Transaction<'_, K>, LibError> {
Transaction::new(self, key).map_err(|e| match e {
Expand Down Expand Up @@ -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<T, F>(&self, key: &K, f: F) -> Result<Option<T>, LibError>
where
F: FnOnce(&IndexStateItem) -> Result<T, CasManagerError>,
Expand All @@ -471,6 +457,24 @@ where
}
}
}
}

impl<K> CasInner<K> {
/// 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 {
Expand Down
5 changes: 1 addition & 4 deletions src/orphan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -65,10 +65,7 @@ pub struct RecoveryResult {
pub errors: Vec<String>,
}

impl<K> OrphanStats<K>
where
K: KeyBytes + Clone + Eq + Ord + Hash + Debug + Send + Sync + 'static,
{
impl<K> OrphanStats<K> {
/// Delete orphaned blobs
pub fn delete_orphans(&self) -> Result<RecoveryResult, LibError> {
let mut result = RecoveryResult::default();
Expand Down
33 changes: 16 additions & 17 deletions src/transaction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<K>,
pub(crate) writer: BufWriter<File>,
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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(())
}
}