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
4 changes: 2 additions & 2 deletions Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "cassadilia"
version = "0.4.3"
version = "0.4.4"
edition = "2024"
rust-version = "1.89.0"
authors = ["0xdeafbeef"]
Expand Down Expand Up @@ -87,6 +87,7 @@ mem_forget = "warn"
missing_enforced_import_renames = "warn"
mut_mut = "warn"
mutex_integer = "warn"
must_use_candidate = "warn"
needless_borrow = "warn"
needless_continue = "warn"
needless_for_each = "warn"
Expand All @@ -103,7 +104,6 @@ semicolon_if_nothing_returned = "warn"
string_add_assign = "warn"
string_add = "warn"
string_lit_as_bytes = "warn"
string_to_string = "warn"
todo = "warn"
trait_duplication_in_bounds = "warn"
unimplemented = "warn"
Expand Down
24 changes: 20 additions & 4 deletions src/index/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ use thiserror::Error;

pub use self::state::IndexStateItem;
use crate::paths::DbPaths;
use crate::types::{BlobHash, CheckpointReason, Config, KeyBytes, WalOp};
use crate::types::{BlobHash, CheckpointReason, Config, DbStats, KeyBytes, WalOp};
use crate::wal::{WalError, WalManager};

mod persistence;
Expand Down Expand Up @@ -149,15 +149,20 @@ where
.map_err(IndexError::InitCreateWalManager)?;

let persister = IndexStatePersister::new(&paths);
let state = persister.load()?;
let mut state = persister.load()?;

let index_file_size =
std::fs::metadata(paths.index_file_path()).map(|m| m.len()).unwrap_or(0);
state.recompute_stats(index_file_size);

let checkpoint_version = state.last_persisted_version;
let state = Arc::new(RwLock::new(state));

let mut replayed_count = 0u64;
wal_manager.replay_and_prepare(checkpoint_version, |op| {
replayed_count += 1;
let _ = state.write().apply_logical_op(&op).expect("Index is corrupted");
let mut guard = state.write();
let _ = guard.apply_logical_op(&op).expect("Index is corrupted");
})?;

let index = Self {
Expand Down Expand Up @@ -209,7 +214,8 @@ where
// 2. Set the version we're about to persist
snapshot.last_persisted_version = Some(target_version);

IndexStatePersister::new(&self.paths).save(snapshot)?;
let serialized_len = IndexStatePersister::new(&self.paths).save(snapshot)?;
snapshot.stats.index.serialized_size_bytes = serialized_len;
// 3. Prune segments up to the target
wal_guard
.commit_checkpoint(target_version, current_checkpoint)
Expand Down Expand Up @@ -412,29 +418,39 @@ where

/// Returns a snapshot of the current key map.
/// Calls `clone` inside.
#[must_use]
pub fn keys_snapshot(&self) -> BTreeMap<K, IndexStateItem> {
self.inner.key_to_hash.clone()
}

/// Returns the number of keys in the index.
#[must_use]
pub fn len(&self) -> usize {
self.inner.key_to_hash.len()
}

/// Returns `true` if the index contains no keys.
#[must_use]
pub fn is_empty(&self) -> bool {
self.inner.key_to_hash.is_empty()
}

/// Returns an iterator over known blob hashes and their reference counts.
#[must_use]
pub fn known_blobs(&self) -> std::collections::hash_map::Iter<'_, BlobHash, u32> {
self.inner.hash_to_ref_count.iter()
}

/// Returns true if the given blob hash is currently referenced.
#[must_use]
pub fn contains_blob_hash(&self, hash: &BlobHash) -> bool {
self.inner.hash_to_ref_count.contains_key(hash)
}

#[must_use]
pub fn stats(&self) -> DbStats {
self.inner.stats
}
}

#[cfg(test)]
Expand Down
5 changes: 3 additions & 2 deletions src/index/persistence.rs
Original file line number Diff line number Diff line change
Expand Up @@ -62,14 +62,15 @@ impl<'a> IndexStatePersister<'a> {
Ok(state)
}

pub fn save<K>(&self, state: &IndexState<K>) -> Result<(), PersisterError>
pub fn save<K>(&self, state: &IndexState<K>) -> Result<u64, PersisterError>
where
K: KeyBytes + Clone + Eq + Ord + Debug + Send + Sync + 'static,
{
let index_path = self.paths.index_file_path();
let index_tmp_path = self.paths.index_tmp_path();

let data_bytes = serialize_index_state(&state.key_to_hash, state.last_persisted_version);
let len = data_bytes.len() as u64;

atomically_write_file_bytes(index_path, index_tmp_path, &data_bytes)?;

Expand All @@ -78,7 +79,7 @@ impl<'a> IndexStatePersister<'a> {
state.key_to_hash.len(),
index_path.display()
);
Ok(())
Ok(len)
}
}

Expand Down
42 changes: 37 additions & 5 deletions src/index/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ use std::num::NonZeroU64;
use ahash::HashMap;
use thiserror::Error;

use crate::types::{BlobHash, WalOp};
use crate::types::{BlobHash, DbStats, WalOp};

#[derive(Error, Debug)]
pub enum IndexStateError {
Expand All @@ -19,6 +19,7 @@ pub(crate) struct IndexState<K> {
pub(crate) key_to_hash: BTreeMap<K, IndexStateItem>,
pub(crate) hash_to_ref_count: HashMap<BlobHash, u32>,
pub(crate) last_persisted_version: Option<NonZeroU64>,
pub(crate) stats: DbStats,
}

#[derive(Debug, Clone, Copy, PartialEq, Eq)]
Expand All @@ -36,9 +37,26 @@ where
key_to_hash: BTreeMap::new(),
hash_to_ref_count: HashMap::default(),
last_persisted_version: None,
stats: DbStats::default(),
}
}

pub fn recompute_stats(&mut self, index_file_size_bytes: u64) {
let mut unique =
HashMap::with_capacity_and_hasher(self.hash_to_ref_count.len(), Default::default());

for item in self.key_to_hash.values() {
unique.entry(item.blob_hash).or_insert(item.blob_size);
}

let unique_blobs = unique.len() as u64;
let total_bytes = unique.values().copied().sum::<u64>();

self.stats.cas.unique_blobs = unique_blobs;
self.stats.cas.total_bytes = total_bytes;
self.stats.index.serialized_size_bytes = index_file_size_bytes;
}

pub fn apply_logical_op(&mut self, op: &WalOp<K>) -> Result<Vec<BlobHash>, IndexStateError> {
let mut unreferenced_hashes = Vec::new();

Expand All @@ -49,16 +67,25 @@ where
match self.key_to_hash.insert(key.clone(), new_item) {
None => {
// New key → bump refcount of the new hash.
self.increment_ref(hash);
if self.increment_ref(hash) {
self.stats.cas.unique_blobs += 1;
self.stats.cas.total_bytes += *size;
}
}
Some(prev) if prev.blob_hash != *hash => {
// Repoint to a different blob:
// 1) decrement old, collect if it drops to zero
if let Some(h) = self.decrement_ref(&prev.blob_hash)? {
unreferenced_hashes.push(h);
self.stats.cas.unique_blobs -= 1;
self.stats.cas.total_bytes -= prev.blob_size;
}

// 2) increment new
self.increment_ref(hash);
if self.increment_ref(hash) {
self.stats.cas.unique_blobs += 1;
self.stats.cas.total_bytes += *size;
}
}
Some(prev) => {
// Same blob hash: refcounts unchanged.
Expand All @@ -80,6 +107,8 @@ where
&& let Some(h) = self.decrement_ref(&item.blob_hash)?
{
unreferenced_hashes.push(h);
self.stats.cas.unique_blobs -= 1;
self.stats.cas.total_bytes -= item.blob_size;
}
}
}
Expand All @@ -88,8 +117,11 @@ where
Ok(unreferenced_hashes)
}

pub(crate) fn increment_ref(&mut self, hash: &BlobHash) {
*self.hash_to_ref_count.entry(*hash).or_default() += 1;
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(
Expand Down
16 changes: 16 additions & 0 deletions src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -96,6 +96,7 @@ pub enum LibError {
},
}

#[must_use]
pub fn calculate_blob_hash(blob_data: &[u8]) -> BlobHash {
BlobHash(blake3::hash(blob_data).into())
}
Expand Down Expand Up @@ -192,6 +193,16 @@ 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()
}
}

pub struct CasInner<K> {
Expand Down Expand Up @@ -306,6 +317,11 @@ where
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
22 changes: 22 additions & 0 deletions src/tests/checkpoint.rs
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,28 @@ fn test_checkpoint_persists_overwrites_correctly() -> Result<()> {
Ok(())
}

#[test]
fn test_stats_index_size_on_checkpoint() -> Result<()> {
setup_tracing();
let harness = CasTestHarness::new(Config::default())?;
let db_path = harness.db_path().to_path_buf();

harness.run_session(|cas| {
let mut tx = cas.put("key".to_string())?;
tx.write(b"data")?;
tx.finish()?;

cas.checkpoint()?;

let stats = cas.stats();
let len = db_path.join("index").metadata()?.len();
assert_eq!(stats.index.serialized_size_bytes, len);
Ok(())
})?;

Ok(())
}

#[test]
fn test_checkpoint_prevents_double_replay() -> Result<()> {
setup_tracing();
Expand Down
Loading