pub(crate) mod utils;
use utils::setup_tracing;
mod checkpoint;
use std::fs;
use std::num::NonZeroU64;
use std::os::unix::fs::PermissionsExt;
use anyhow::Result;
use tempfile::tempdir;
use crate::{Cas, Config, LibError, LibIoOperation, SyncMode};
#[test]
fn test_put_get_remove_string_key() -> Result<()> {
setup_tracing();
let dir = tempdir()?;
let cas = Cas::open(dir.path(), Config::default())?;
let key = "my_first_blob".to_string();
let data = b"Hello, Bubs!";
let mut tx = cas.put(key.clone())?;
tx.write(data)?;
tx.finish()?;
let retrieved = cas.get(&key)?.expect("blob should exist");
assert_eq!(retrieved.as_ref(), data);
assert!(cas.get(&"does_not_exist".to_string())?.is_none());
assert!(cas.remove(&key)?);
assert!(cas.get(&key)?.is_none());
assert!(!cas.remove(&key)?);
Ok(())
}
#[test]
fn test_put_get_remove_bytes_key() -> Result<()> {
setup_tracing();
let dir = tempdir()?;
let cas = Cas::open(dir.path(), Config::default())?;
let key = b"my_bytes_blob".to_vec();
let data = b"rofls";
let mut tx = cas.put(key.clone())?;
tx.write(data)?;
tx.finish()?;
let retrieved = cas.get(&key)?.expect("blob should exist");
assert_eq!(retrieved.as_ref(), data);
assert!(cas.remove(&key)?);
assert!(cas.get(&key)?.is_none());
Ok(())
}
#[test]
fn test_get_range() -> Result<()> {
setup_tracing();
let dir = tempdir()?;
let cas = Cas::open(dir.path(), Config::default())?;
let key = "range_blob".to_string();
let data = b"0123456789abcdef";
let mut tx = cas.put(key.clone())?;
tx.write(data)?;
tx.finish()?;
assert_eq!(cas.get_range(&key, 0, 16)?.unwrap().as_ref(), b"0123456789abcdef");
assert_eq!(cas.get_range(&key, 4, 8)?.unwrap().as_ref(), b"4567");
assert_eq!(cas.get_range(&key, 10, 16)?.unwrap().as_ref(), b"abcdef");
assert_eq!(cas.get_range(&key, 12, 100)?.unwrap().as_ref(), b"cdef");
assert!(cas.get_range(&key, 100, 200)?.unwrap().is_empty());
assert!(cas.get_range(&key, 5, 5)?.unwrap().is_empty());
assert!(cas.get_range(&key, 8, 4).is_err());
Ok(())
}
#[test]
fn test_overwrite_persists_across_reopen() -> Result<()> {
setup_tracing();
let dir = tempdir()?;
let db_path = dir.path();
let key = "overwrite_test".to_string();
{
let cas = Cas::open(db_path, Config::default())?;
let mut tx = cas.put(key.clone())?;
tx.write(b"Version 1")?;
tx.finish()?;
assert_eq!(cas.get(&key)?.unwrap().as_ref(), b"Version 1");
}
{
let cas = Cas::open(db_path, Config::default())?;
assert_eq!(cas.get(&key)?.unwrap().as_ref(), b"Version 1");
let mut tx = cas.put(key.clone())?;
tx.write(b"Version 2 is better")?;
tx.finish()?;
assert_eq!(cas.get(&key)?.unwrap().as_ref(), b"Version 2 is better");
}
{
let cas = Cas::open(db_path, Config::default())?;
let retrieved = cas.get(&key)?.unwrap();
assert_eq!(retrieved.as_ref(), b"Version 2 is better");
}
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::<String>::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();
let dir = tempdir()?;
let db_path = dir.path();
let key = "remove_persist_test".to_string();
{
let cas = Cas::open(db_path, Config::default())?;
let mut tx = cas.put(key.clone())?;
tx.write(b"Data to be removed")?;
tx.finish()?;
assert!(cas.get(&key)?.is_some());
assert!(cas.remove(&key)?);
assert!(cas.get(&key)?.is_none());
}
{
let cas = Cas::open(db_path, Config::default())?;
assert!(cas.get(&key)?.is_none(), "data should still be removed after reopen");
}
Ok(())
}
#[test]
fn test_transaction_drop_cleans_up_staging_file() -> Result<()> {
setup_tracing();
let dir = tempdir()?;
let cas = Cas::open(dir.path(), Config::default())?;
let key = "dropped_tx_key".to_string();
let staging_path;
{
let mut tx = cas.put(key.clone())?;
tx.write(b"This data won't be saved")?;
staging_path = tx.temp_file.path().to_path_buf();
assert!(staging_path.exists(), "staging file should exist during tx");
}
assert!(!staging_path.exists(), "staging file should be removed on drop");
assert!(cas.get(&key)?.is_none(), "data should not be present if tx was dropped");
Ok(())
}
#[test]
fn test_checkpoint_persists_index() -> Result<()> {
setup_tracing();
let dir = tempdir()?;
let db_path = dir.path();
let index_file_path = db_path.join("index");
{
let cas = Cas::open(db_path, Config::default())?;
let mut tx = cas.put("key1".to_string())?;
tx.write(b"data1")?;
tx.finish()?;
cas.checkpoint()?;
assert!(index_file_path.exists(), "no index");
assert!(index_file_path.metadata()?.len() > 0, "empty index");
assert_eq!(cas.0.index.state.read().last_persisted_version, NonZeroU64::new(1));
}
{
let cas = Cas::open(db_path, Config::default())?;
assert_eq!(cas.get(&"key1".to_string())?.unwrap().as_ref(), b"data1");
}
Ok(())
}
#[test]
fn test_wal_rollover_and_cleanup() -> Result<()> {
setup_tracing();
let dir = tempdir()?;
let db_path = dir.path();
let config = Config {
sync_mode: SyncMode::Sync,
num_ops_per_wal: NonZeroU64::new(2).unwrap(),
pre_create_cas_dirs: false,
..Default::default()
};
let wal0_path = db_path.join("0_index.wal");
let wal1_path = db_path.join("1_index.wal");
{
let cas = Cas::open(db_path, config.clone())?;
let mut tx = cas.put("key1".to_string())?;
tx.write(b"d1")?;
tx.finish()?;
let mut tx = cas.put("key2".to_string())?;
tx.write(b"d2")?;
tx.finish()?;
assert!(wal0_path.exists());
assert!(!wal1_path.exists());
let mut tx = cas.put("key3".to_string())?;
tx.write(b"d3")?;
tx.finish()?;
assert!(wal1_path.exists());
cas.checkpoint()?;
assert!(!wal0_path.exists(), "stale wal (0) should be cleaned up by checkpoint");
assert!(wal1_path.exists(), "active wal (1) should remain");
}
{
let cas = Cas::open(db_path, config)?;
assert_eq!(cas.get(&"key1".to_string())?.unwrap().as_ref(), b"d1");
assert_eq!(cas.get(&"key2".to_string())?.unwrap().as_ref(), b"d2");
assert_eq!(cas.get(&"key3".to_string())?.unwrap().as_ref(), b"d3");
}
Ok(())
}
#[test]
fn test_remove_range_persists() -> Result<()> {
setup_tracing();
let dir = tempdir()?;
let db_path = dir.path();
{
let cas = Cas::open(db_path, Config::default())?;
for i in 1..=4 {
let mut tx = cas.put(format!("key_{i}"))?;
tx.write(format!("data_{i}",).as_bytes())?;
tx.finish()?;
}
let count = cas.remove_range("key_1".to_string().."key_3".to_string())?;
assert_eq!(count, 2);
assert!(cas.get(&"key_1".to_string())?.is_none());
assert!(cas.get(&"key_2".to_string())?.is_none());
assert!(cas.get(&"key_3".to_string())?.is_some());
}
{
let cas = Cas::open(db_path, Config::default())?;
assert!(cas.get(&"key_1".to_string())?.is_none());
assert!(cas.get(&"key_2".to_string())?.is_none());
assert!(cas.get(&"key_3".to_string())?.is_some());
assert!(cas.get(&"key_4".to_string())?.is_some());
assert_eq!(cas.index.read_state().known_blobs().count(), 2);
}
Ok(())
}
#[test]
fn test_api_on_nonexistent_key() -> Result<()> {
setup_tracing();
let dir = tempdir()?;
let cas = Cas::open(dir.path(), Config::default())?;
let key = "nonexistent_key".to_string();
let result = cas.get_reader(&key).unwrap();
match result {
None => {}
_ => panic!("Expected KeyNotFound error, got: {result:?}"),
}
assert!(cas.get(&key)?.is_none());
assert!(cas.get_size(&key)?.is_none());
assert!(cas.get_range(&key, 0, 10)?.is_none());
Ok(())
}
#[test]
fn test_io_error_on_staging_file_creation() -> anyhow::Result<()> {
setup_tracing();
let dir = tempdir()?;
let cas = Cas::open(dir.path(), Config::default())?;
let staging_dir = dir.path().join("staging");
fs::set_permissions(&staging_dir, fs::Permissions::from_mode(0o555))?;
let result = cas.put("test_key".to_string());
assert!(result.is_err());
let err = result.unwrap_err();
let LibError::Io { operation, path, .. } = err else {
panic!("Expected a specific IO error, but got: {err}");
};
assert!(matches!(operation, LibIoOperation::CreateStagingFile));
let path_str = path.unwrap().to_string_lossy().to_string();
assert!(path_str.contains("staging"));
Ok(())
}
#[test]
fn test_orphan_detection_and_cleanup() -> Result<()> {
setup_tracing();
let dir = tempdir()?;
let cas = Cas::open(dir.path(), Config::default())?;
let key1 = "valid_key1".to_string();
let data1 = b"valid data 1";
let mut tx = cas.put(key1.clone())?;
tx.write(data1)?;
tx.finish()?;
let _valid_hash = cas.index.read_state().get_item(&key1).unwrap();
let orphan_data = b"orphaned data";
let orphan_hash = crate::calculate_blob_hash(orphan_data);
let orphan_path = dir.path().join("cas").join(orphan_hash.relative_path());
if let Some(parent) = orphan_path.parent() {
fs::create_dir_all(parent)?;
}
fs::write(&orphan_path, orphan_data)?;
let invalid_file_path = dir.path().join("cas").join(".DS_Store");
fs::write(&invalid_file_path, b"invalid")?;
let staging_file = dir.path().join("staging").join("old_file.tmp");
fs::write(&staging_file, b"old staging data")?;
drop(cas);
let config = Config { scan_orphans_on_startup: true, ..Default::default() };
let (cas, orphan_stats) = Cas::open_with_recover(dir.path(), config)?;
let stats = orphan_stats.expect("Should have orphan stats");
let [only] = stats.orphaned_blobs.as_slice() else {
panic!("Expected one orphaned blob");
};
assert_eq!(*only, orphan_hash);
assert_eq!(stats.invalid_files.len(), 1);
let retrieved = cas.get(&key1)?.expect("valid blob should still exist");
assert_eq!(retrieved.as_ref(), data1);
assert!(orphan_path.exists(), "orphaned blob should still exist before cleanup");
let result = stats.delete_orphans()?;
assert_eq!(result.orphans_deleted, 1);
assert_eq!(result.invalid_files_removed, 1);
assert!(!orphan_path.exists(), "orphaned blob should be removed after cleanup");
assert!(!invalid_file_path.exists(), "invalid file should be removed after cleanup");
Ok(())
}
#[test]
fn test_orphan_detection_with_integrity_check() -> Result<()> {
setup_tracing();
let dir = tempdir()?;
let cas = Cas::open(dir.path(), Config::default())?;
let key1 = "valid_key1".to_string();
let data1 = b"valid data 1";
let mut tx = cas.put(key1.clone())?;
tx.write(data1)?;
tx.finish()?;
let item = cas.index.read_state().get_item(&key1).unwrap();
let valid_path = dir.path().join("cas").join(item.blob_hash.relative_path());
fs::write(&valid_path, b"corrupted data")?;
drop(cas);
let config = Config {
scan_orphans_on_startup: true,
verify_blob_integrity: true,
fail_on_integrity_errors: true,
..Default::default()
};
let result = Cas::<String>::open(dir.path(), config);
match result {
Err(LibError::IntegrityCheckFailed { corrupted_blobs, .. }) => {
let [only] = corrupted_blobs.as_slice() else {
panic!("Expected one corrupted blob");
};
assert_eq!(*only, item.blob_hash);
}
_ => panic!("Expected IntegrityCheckFailed error"),
}
Ok(())
}
#[test]
fn test_orphan_detection_with_missing_blobs() -> Result<()> {
setup_tracing();
let dir = tempdir()?;
let cas = Cas::open(dir.path(), Config::default())?;
let key1 = "key1".to_string();
let data1 = b"data 1";
let mut tx = cas.put(key1.clone())?;
tx.write(data1)?;
tx.finish()?;
let item1 = cas.index.read_state().get_item(&key1).unwrap();
let blob_path = dir.path().join("cas").join(item1.blob_hash.relative_path());
fs::remove_file(&blob_path)?;
drop(cas);
let config = Config {
scan_orphans_on_startup: true,
fail_on_integrity_errors: true,
..Default::default()
};
let result = Cas::<String>::open(dir.path(), config);
match result {
Err(LibError::IntegrityCheckFailed { missing_blobs, .. }) => {
let [only] = missing_blobs.as_slice() else {
panic!("Expected one missing blob");
};
assert_eq!(*only, item1.blob_hash);
}
_ => panic!("Expected IntegrityCheckFailed error"),
}
Ok(())
}
#[test]
fn test_orphan_quarantine() -> Result<()> {
setup_tracing();
let dir = tempdir()?;
let quarantine_dir = dir.path().join("quarantine");
let cas = Cas::open(dir.path(), Config::default())?;
let key1 = "key1".to_string();
let data1 = b"data 1";
let mut tx = cas.put(key1.clone())?;
tx.write(data1)?;
tx.finish()?;
let orphan_data = b"orphan";
let orphan_hash = crate::calculate_blob_hash(orphan_data);
let orphan_path = dir.path().join("cas").join(orphan_hash.relative_path());
if let Some(parent) = orphan_path.parent() {
fs::create_dir_all(parent)?;
}
fs::write(&orphan_path, orphan_data)?;
drop(cas);
let config = Config { scan_orphans_on_startup: true, ..Default::default() };
let (_cas, orphan_stats) = Cas::<String>::open_with_recover(dir.path(), config)?;
let stats = orphan_stats.expect("Should have orphan stats");
let result = stats.quarantine_orphans(&quarantine_dir)?;
assert_eq!(result.orphans_quarantined, 1);
assert!(!orphan_path.exists());
assert!(quarantine_dir.join(orphan_hash.to_string()).exists());
Ok(())
}
#[test]
fn test_orphan_stats_holds_lock() -> Result<()> {
setup_tracing();
let dir = tempdir()?;
let cas = Cas::<String>::open(dir.path(), Config::default())?;
let orphan_data = b"orphan";
let orphan_hash = crate::calculate_blob_hash(orphan_data);
let orphan_path = dir.path().join("cas").join(orphan_hash.relative_path());
if let Some(parent) = orphan_path.parent() {
fs::create_dir_all(parent)?;
}
fs::write(&orphan_path, orphan_data)?;
drop(cas);
let config = Config { scan_orphans_on_startup: true, ..Default::default() };
let (_cas, orphan_stats) = Cas::<String>::open_with_recover(dir.path(), config)?;
let stats = orphan_stats.expect("Should have orphan stats");
assert_eq!(stats.orphaned_blobs.len(), 1);
let result = stats.delete_orphans()?;
assert_eq!(result.orphans_deleted, 1);
assert!(!orphan_path.exists());
drop(stats);
Ok(())
}
#[test]
fn test_cleanup_disabled() -> Result<()> {
setup_tracing();
let dir = tempdir()?;
let cas = Cas::<String>::open(dir.path(), Config::default())?;
let orphan_data = b"orphan";
let orphan_hash = crate::calculate_blob_hash(orphan_data);
let orphan_path = dir.path().join("cas").join(orphan_hash.relative_path());
if let Some(parent) = orphan_path.parent() {
fs::create_dir_all(parent)?;
}
fs::write(&orphan_path, orphan_data)?;
drop(cas);
let config = Config { scan_orphans_on_startup: false, ..Default::default() };
let (_cas, orphan_stats) = Cas::<String>::open_with_recover(dir.path(), config)?;
assert!(orphan_stats.is_none());
assert!(orphan_path.exists(), "orphan should not be removed when cleanup is disabled");
Ok(())
}
#[test]
fn regression_put_same_content_should_not_delete_blob() {
use tempfile::tempdir;
use crate::{Cas, Config};
let dir = tempdir().unwrap();
let cas = Cas::open(dir.path(), Config::default()).unwrap();
let key = b"same key";
let data = b"same content";
{
let mut tx = cas.put(*key).unwrap();
tx.write(data).unwrap();
tx.finish().unwrap();
}
{
let mut tx = cas.put(*key).unwrap();
tx.write(data).unwrap();
tx.finish().unwrap();
}
let got = cas.get(key);
assert!(got.is_ok(), "unexpected error: {:?}", got);
assert_eq!(got.unwrap().unwrap(), bytes::Bytes::from_static(data));
}