diff --git a/Cargo.toml b/Cargo.toml index c593fdf..a56f827 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "cassadilia" -version = "0.4.3" +version = "0.4.4" edition = "2024" rust-version = "1.89.0" authors = ["0xdeafbeef"] @@ -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" @@ -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" diff --git a/src/index/mod.rs b/src/index/mod.rs index fc86dce..5717804 100644 --- a/src/index/mod.rs +++ b/src/index/mod.rs @@ -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; @@ -149,7 +149,11 @@ 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)); @@ -157,7 +161,8 @@ where 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 { @@ -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) @@ -412,29 +418,39 @@ where /// Returns a snapshot of the current key map. /// Calls `clone` inside. + #[must_use] pub fn keys_snapshot(&self) -> BTreeMap { 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)] diff --git a/src/index/persistence.rs b/src/index/persistence.rs index bb209ae..95a659c 100644 --- a/src/index/persistence.rs +++ b/src/index/persistence.rs @@ -62,7 +62,7 @@ impl<'a> IndexStatePersister<'a> { Ok(state) } - pub fn save(&self, state: &IndexState) -> Result<(), PersisterError> + pub fn save(&self, state: &IndexState) -> Result where K: KeyBytes + Clone + Eq + Ord + Debug + Send + Sync + 'static, { @@ -70,6 +70,7 @@ impl<'a> IndexStatePersister<'a> { 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)?; @@ -78,7 +79,7 @@ impl<'a> IndexStatePersister<'a> { state.key_to_hash.len(), index_path.display() ); - Ok(()) + Ok(len) } } diff --git a/src/index/state.rs b/src/index/state.rs index cfa9ca5..5098156 100644 --- a/src/index/state.rs +++ b/src/index/state.rs @@ -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 { @@ -19,6 +19,7 @@ pub(crate) struct IndexState { pub(crate) key_to_hash: BTreeMap, pub(crate) hash_to_ref_count: HashMap, pub(crate) last_persisted_version: Option, + pub(crate) stats: DbStats, } #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -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::(); + + 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) -> Result, IndexStateError> { let mut unreferenced_hashes = Vec::new(); @@ -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. @@ -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; } } } @@ -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( diff --git a/src/lib.rs b/src/lib.rs index 3eb3f60..155616e 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -96,6 +96,7 @@ pub enum LibError { }, } +#[must_use] pub fn calculate_blob_hash(blob_data: &[u8]) -> BlobHash { BlobHash(blake3::hash(blob_data).into()) } @@ -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 { @@ -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, LibError> { Transaction::new(self, key).map_err(|e| match e { diff --git a/src/tests/checkpoint.rs b/src/tests/checkpoint.rs index db161d9..4b4beba 100644 --- a/src/tests/checkpoint.rs +++ b/src/tests/checkpoint.rs @@ -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(); diff --git a/src/tests/mod.rs b/src/tests/mod.rs index 8c4f5ee..dbf881a 100644 --- a/src/tests/mod.rs +++ b/src/tests/mod.rs @@ -123,6 +123,108 @@ fn test_overwrite_persists_across_reopen() -> Result<()> { Ok(()) } +#[test] +fn test_stats_unique_and_bytes_shared_blob() -> Result<()> { + setup_tracing(); + let dir = tempdir()?; + let cas = Cas::open(dir.path(), Config::default())?; + + let data = b"abc"; + let size = data.len() as u64; + let key1 = "key1".to_string(); + let key2 = "key2".to_string(); + + let mut tx = cas.put(key1.clone())?; + tx.write(data)?; + tx.finish()?; + + let stats = cas.stats(); + assert_eq!(stats.cas.unique_blobs, 1); + assert_eq!(stats.cas.total_bytes, size); + + let mut tx = cas.put(key2.clone())?; + tx.write(data)?; + tx.finish()?; + + let stats = cas.stats(); + assert_eq!(stats.cas.unique_blobs, 1); + assert_eq!(stats.cas.total_bytes, size); + + assert!(cas.remove(&key1)?); + let stats = cas.stats(); + assert_eq!(stats.cas.unique_blobs, 1); + assert_eq!(stats.cas.total_bytes, size); + + assert!(cas.remove(&key2)?); + let stats = cas.stats(); + assert_eq!(stats.cas.unique_blobs, 0); + assert_eq!(stats.cas.total_bytes, 0); + + Ok(()) +} + +#[test] +fn test_stats_repoint_overwrite_updates_bytes() -> Result<()> { + setup_tracing(); + let dir = tempdir()?; + let cas = Cas::open(dir.path(), Config::default())?; + + let key = "key".to_string(); + let data_a = b"a"; + let data_b = b"bbbbb"; + + let mut tx = cas.put(key.clone())?; + tx.write(data_a)?; + tx.finish()?; + + let stats = cas.stats(); + assert_eq!(stats.cas.unique_blobs, 1); + assert_eq!(stats.cas.total_bytes, data_a.len() as u64); + + let mut tx = cas.put(key.clone())?; + tx.write(data_b)?; + tx.finish()?; + + let stats = cas.stats(); + assert_eq!(stats.cas.unique_blobs, 1); + assert_eq!(stats.cas.total_bytes, data_b.len() as u64); + + Ok(()) +} + +#[test] +fn test_stats_recomputed_after_reopen() -> Result<()> { + setup_tracing(); + let dir = tempdir()?; + let db_path = dir.path(); + + { + let cas = Cas::open(db_path, Config::default())?; + let data = b"abc"; + + let mut tx = cas.put("k1".to_string())?; + tx.write(data)?; + tx.finish()?; + + let mut tx = cas.put("k2".to_string())?; + tx.write(data)?; + tx.finish()?; + + let stats = cas.stats(); + assert_eq!(stats.cas.unique_blobs, 1); + assert_eq!(stats.cas.total_bytes, data.len() as u64); + } + + { + let cas = Cas::::open(db_path, Config::default())?; + let stats = cas.stats(); + assert_eq!(stats.cas.unique_blobs, 1); + assert_eq!(stats.cas.total_bytes, 3); + } + + Ok(()) +} + #[test] fn test_remove_persists_across_reopen() -> Result<()> { setup_tracing(); diff --git a/src/types.rs b/src/types.rs index efa1165..a3a6337 100644 --- a/src/types.rs +++ b/src/types.rs @@ -10,14 +10,33 @@ use crate::cas_manager::CasManagerError; pub const HASH_SIZE: usize = blake3::OUT_LEN; +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct CasStats { + pub unique_blobs: u64, + pub total_bytes: u64, +} + +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct IndexStats { + pub serialized_size_bytes: u64, +} + +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct DbStats { + pub cas: CasStats, + pub index: IndexStats, +} + #[derive(Clone, Copy, Eq, Ord, PartialOrd)] pub struct BlobHash(pub [u8; HASH_SIZE]); impl BlobHash { + #[must_use] pub fn as_bytes(&self) -> &[u8; HASH_SIZE] { &self.0 } + #[must_use] pub fn to_hex(&self) -> String { hex::encode(self.0) } @@ -30,6 +49,7 @@ impl BlobHash { /// Returns the relative path for this blob hash, suitable for storing in a directory structure. /// The path will have 2 levels /h/a/sh + #[must_use] pub fn relative_path(&self) -> PathBuf { let hex = self.to_hex(); let top_level_dir = &hex[0..2]; @@ -61,6 +81,7 @@ impl BlobHash { Ok(Self(res)) } + #[must_use] pub fn from_bytes(bytes: [u8; HASH_SIZE]) -> Self { Self(bytes) }