use crate::address::Address;
use crate::address::GuardianDBAddress;
use crate::log::identity::{Identity, Signatures};
use crate::p2p::EventBus;
use crate::p2p::network::client::IrohClient;
use crate::p2p::network::config::ClientConfig;
use crate::stores::kv_store::GuardianDBKeyValue;
use crate::traits::NewStoreOptions;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use tempfile::TempDir;
static TEST_COUNTER: AtomicU64 = AtomicU64::new(0);
fn test_identity() -> Arc<Identity> {
Arc::new(Identity::new(
"test-user",
"test-public-key",
Signatures::new("test-id-sig", "test-pub-sig"),
))
}
async fn test_address() -> Arc<dyn Address + Send + Sync> {
use blake3;
use std::time::{SystemTime, UNIX_EPOCH};
let timestamp = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_nanos();
let counter = TEST_COUNTER.fetch_add(1, Ordering::SeqCst);
let unique_data = format!("test-kvstore-{}-{}", timestamp, counter);
let hash_bytes: [u8; 32] = blake3::hash(unique_data.as_bytes()).into();
let hash = iroh_blobs::Hash::from(hash_bytes);
Arc::new(GuardianDBAddress::new(
hash,
format!("test-kvstore-{}-{}", timestamp, counter),
))
}
async fn create_test_store() -> Result<(GuardianDBKeyValue, TempDir), Box<dyn std::error::Error>> {
let temp_dir = TempDir::new()?;
let client_config = ClientConfig {
data_store_path: Some(temp_dir.path().to_path_buf()),
..ClientConfig::development()
};
let client = Arc::new(IrohClient::new(client_config).await?);
let identity = test_identity();
let address = test_address().await;
let event_bus = EventBus::new();
let options = NewStoreOptions {
event_bus: Some(event_bus),
directory: temp_dir.path().join("cache").to_string_lossy().to_string(),
..Default::default()
};
let store = GuardianDBKeyValue::new(client, identity, address, Some(options)).await?;
Ok((store, temp_dir))
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn test_kv_store_creation() {
let result = create_test_store().await;
assert!(result.is_ok(), "Should create KeyValueStore successfully");
let (store, _temp_dir) = result.unwrap();
assert_eq!(store.get_type(), "keyvalue");
assert!(store.is_empty());
assert_eq!(store.len(), 0);
}
#[tokio::test]
async fn test_put_and_get() {
let (store, _temp_dir) = create_test_store()
.await
.expect("Failed to create test store");
let key = "test_key";
let value = b"test_value".to_vec();
let result = store.put_impl(key, value.clone()).await;
assert!(
result.is_ok(),
"Should put value successfully: {:?}",
result
);
let retrieved = store.get_impl(key).await.expect("Failed to get value");
assert!(retrieved.is_some(), "Should find the key");
assert_eq!(retrieved.unwrap(), value);
}
#[tokio::test]
async fn test_put_updates_existing_key() {
let (store, _temp_dir) = create_test_store()
.await
.expect("Failed to create test store");
let key = "update_key";
store
.put_impl(key, b"initial".to_vec())
.await
.expect("Failed to put initial value");
store
.put_impl(key, b"updated".to_vec())
.await
.expect("Failed to update value");
let retrieved = store.get_impl(key).await.expect("Failed to get value");
assert_eq!(retrieved.unwrap(), b"updated");
}
#[tokio::test]
async fn test_delete() {
let (store, _temp_dir) = create_test_store()
.await
.expect("Failed to create test store");
let key = "delete_key";
let value = b"delete_value".to_vec();
store
.put_impl(key, value)
.await
.expect("Failed to put value");
assert!(store.contains_key(key));
let result = store.delete_impl(key).await;
assert!(result.is_ok(), "Should delete successfully: {:?}", result);
assert!(!store.contains_key(key));
let retrieved = store.get_impl(key).await.expect("Failed to get value");
assert!(retrieved.is_none(), "Key should not exist after deletion");
}
#[tokio::test]
async fn test_delete_nonexistent_key() {
let (store, _temp_dir) = create_test_store()
.await
.expect("Failed to create test store");
let result = store.delete_impl("nonexistent").await;
assert!(
result.is_err(),
"Should error when deleting nonexistent key"
);
}
#[tokio::test]
async fn test_empty_key_validation() {
let (store, _temp_dir) = create_test_store()
.await
.expect("Failed to create test store");
let result = store.put_impl("", b"value".to_vec()).await;
assert!(result.is_err(), "Should error on empty key for PUT");
let result = store.get_impl("").await;
assert!(result.is_err(), "Should error on empty key for GET");
let result = store.delete_impl("").await;
assert!(result.is_err(), "Should error on empty key for DELETE");
}
#[tokio::test]
async fn test_empty_value_validation() {
let (store, _temp_dir) = create_test_store()
.await
.expect("Failed to create test store");
let result = store.put_impl("key", Vec::new()).await;
assert!(result.is_err(), "Should error on empty value");
}
#[tokio::test]
async fn test_contains_key() {
let (store, _temp_dir) = create_test_store()
.await
.expect("Failed to create test store");
let key = "contains_key";
assert!(!store.contains_key(key));
store
.put_impl(key, b"value".to_vec())
.await
.expect("Failed to put value");
assert!(store.contains_key(key));
}
#[tokio::test]
async fn test_keys() {
let (store, _temp_dir) = create_test_store()
.await
.expect("Failed to create test store");
assert!(store.keys().is_empty());
store
.put_impl("key1", b"value1".to_vec())
.await
.expect("Failed to put key1");
store
.put_impl("key2", b"value2".to_vec())
.await
.expect("Failed to put key2");
store
.put_impl("key3", b"value3".to_vec())
.await
.expect("Failed to put key3");
let keys = store.keys();
assert_eq!(keys.len(), 3);
assert!(keys.contains(&"key1".to_string()));
assert!(keys.contains(&"key2".to_string()));
assert!(keys.contains(&"key3".to_string()));
}
#[tokio::test]
async fn test_all() {
let (store, _temp_dir) = create_test_store()
.await
.expect("Failed to create test store");
store
.put_impl("key1", b"value1".to_vec())
.await
.expect("Failed to put key1");
store
.put_impl("key2", b"value2".to_vec())
.await
.expect("Failed to put key2");
store
.put_impl("key3", b"value3".to_vec())
.await
.expect("Failed to put key3");
let all = store.all();
assert_eq!(all.len(), 3);
assert_eq!(all.get("key1").unwrap(), b"value1");
assert_eq!(all.get("key2").unwrap(), b"value2");
assert_eq!(all.get("key3").unwrap(), b"value3");
}
#[tokio::test]
async fn test_len_and_is_empty() {
let (store, _temp_dir) = create_test_store()
.await
.expect("Failed to create test store");
assert!(store.is_empty());
assert_eq!(store.len(), 0);
store
.put_impl("key1", b"value1".to_vec())
.await
.expect("Failed to put key1");
assert!(!store.is_empty());
assert_eq!(store.len(), 1);
store
.put_impl("key2", b"value2".to_vec())
.await
.expect("Failed to put key2");
store
.put_impl("key3", b"value3".to_vec())
.await
.expect("Failed to put key3");
assert_eq!(store.len(), 3);
store
.delete_impl("key2")
.await
.expect("Failed to delete key2");
assert_eq!(store.len(), 2);
}
#[tokio::test]
async fn test_multiple_operations_sequence() {
let (store, _temp_dir) = create_test_store()
.await
.expect("Failed to create test store");
store
.put_impl("seq_key", b"value1".to_vec())
.await
.expect("Failed to put");
assert_eq!(store.get_impl("seq_key").await.unwrap().unwrap(), b"value1");
store
.put_impl("seq_key", b"value2".to_vec())
.await
.expect("Failed to update");
assert_eq!(store.get_impl("seq_key").await.unwrap().unwrap(), b"value2");
store
.delete_impl("seq_key")
.await
.expect("Failed to delete");
assert!(store.get_impl("seq_key").await.unwrap().is_none());
store
.put_impl("seq_key", b"value3".to_vec())
.await
.expect("Failed to put again");
assert_eq!(store.get_impl("seq_key").await.unwrap().unwrap(), b"value3");
}
#[tokio::test]
async fn test_binary_values() {
let (store, _temp_dir) = create_test_store()
.await
.expect("Failed to create test store");
let binary_data = vec![0x00, 0xFF, 0x42, 0x00, 0xAB, 0xCD, 0xEF];
store
.put_impl("binary_key", binary_data.clone())
.await
.expect("Failed to put binary data");
let retrieved = store
.get_impl("binary_key")
.await
.expect("Failed to get binary data");
assert_eq!(retrieved.unwrap(), binary_data);
}
#[tokio::test]
async fn test_large_value() {
let (store, _temp_dir) = create_test_store()
.await
.expect("Failed to create test store");
let large_value = vec![0x42u8; 1024 * 1024];
store
.put_impl("large_key", large_value.clone())
.await
.expect("Failed to put large value");
let retrieved = store
.get_impl("large_key")
.await
.expect("Failed to get large value");
assert_eq!(retrieved.unwrap(), large_value);
}
#[tokio::test]
async fn test_special_characters_in_keys() {
let (store, _temp_dir) = create_test_store()
.await
.expect("Failed to create test store");
let special_keys = vec![
"key-with-dash",
"key_with_underscore",
"key.with.dots",
"key/with/slashes",
"key:with:colons",
"key with spaces",
"key@with#special$chars",
];
for key in &special_keys {
let value = format!("value for {}", key).into_bytes();
store
.put_impl(key, value.clone())
.await
.unwrap_or_else(|_| panic!("Failed to put key: {}", key));
let retrieved = store
.get_impl(key)
.await
.unwrap_or_else(|_| panic!("Failed to get key: {}", key));
assert_eq!(retrieved.unwrap(), value);
}
}
#[tokio::test]
async fn test_integration_workflow() {
let (store, _temp_dir) = create_test_store()
.await
.expect("Failed to create test store");
store
.put_impl("session:user1", b"token123".to_vec())
.await
.expect("Failed to create session1");
store
.put_impl("session:user2", b"token456".to_vec())
.await
.expect("Failed to create session2");
store
.put_impl("session:user3", b"token789".to_vec())
.await
.expect("Failed to create session3");
assert_eq!(store.len(), 3);
let all_sessions = store.all();
assert_eq!(all_sessions.len(), 3);
store
.put_impl("session:user2", b"new_token456".to_vec())
.await
.expect("Failed to update session");
assert_eq!(
store.get_impl("session:user2").await.unwrap().unwrap(),
b"new_token456"
);
store
.delete_impl("session:user1")
.await
.expect("Failed to delete session");
assert_eq!(store.len(), 2);
let remaining_keys = store.keys();
assert_eq!(remaining_keys.len(), 2);
assert!(remaining_keys.contains(&"session:user2".to_string()));
assert!(remaining_keys.contains(&"session:user3".to_string()));
assert!(!remaining_keys.contains(&"session:user1".to_string()));
}
#[tokio::test]
async fn test_concurrent_operations() {
let (store, _temp_dir) = create_test_store()
.await
.expect("Failed to create test store");
for i in 0..10 {
let key = format!("concurrent_key_{}", i);
let value = format!("value_{}", i).into_bytes();
store
.put_impl(&key, value)
.await
.unwrap_or_else(|_| panic!("Failed to put key: {}", i));
}
assert_eq!(store.len(), 10);
for i in 0..10 {
let key = format!("concurrent_key_{}", i);
let expected_value = format!("value_{}", i).into_bytes();
let retrieved = store
.get_impl(&key)
.await
.unwrap_or_else(|_| panic!("Failed to get key: {}", key));
assert_eq!(retrieved.unwrap(), expected_value);
}
}
#[tokio::test]
async fn test_utf8_values() {
let (store, _temp_dir) = create_test_store()
.await
.expect("Failed to create test store");
let utf8_values = [
("english", "Hello World"),
("portuguese", "Olá Mundo"),
("chinese", "你好世界"),
("arabic", "مرحبا بالعالم"),
("emoji", "🚀🎉✨"),
];
for (key, value) in utf8_values {
store
.put_impl(key, value.as_bytes().to_vec())
.await
.unwrap_or_else(|_| panic!("Failed to put {}", key));
let retrieved = store
.get_impl(key)
.await
.unwrap_or_else(|_| panic!("Failed to get {}", key));
let retrieved_str = String::from_utf8(retrieved.unwrap()).expect("Invalid UTF-8");
assert_eq!(retrieved_str, value);
}
}
#[tokio::test]
async fn test_read_only_store_cannot_create_namespace() {
let temp_dir = TempDir::new().expect("temp dir");
let client_config = ClientConfig {
data_store_path: Some(temp_dir.path().to_path_buf()),
..ClientConfig::development()
};
let client = Arc::new(IrohClient::new(client_config).await.expect("client"));
let identity = test_identity();
let address = test_address().await;
let options = NewStoreOptions {
event_bus: Some(EventBus::new()),
directory: temp_dir.path().join("cache").to_string_lossy().to_string(),
read_only: Some(true),
..Default::default()
};
let result = GuardianDBKeyValue::new(client, identity, address, Some(options)).await;
assert!(
result.is_err(),
"Read-only store with no namespace to import must fail, not create one"
);
}
#[tokio::test]
async fn test_reopened_read_only_replica_rejects_writes() {
use crate::traits::Store;
let temp_dir = TempDir::new().expect("temp dir");
let cache_dir = temp_dir.path().join("cache").to_string_lossy().to_string();
let client_config = ClientConfig {
data_store_path: Some(temp_dir.path().to_path_buf()),
..ClientConfig::development()
};
let client = Arc::new(IrohClient::new(client_config).await.expect("client"));
let identity = test_identity();
let address = test_address().await;
{
let opts = NewStoreOptions {
event_bus: Some(EventBus::new()),
directory: cache_dir.clone(),
..Default::default()
};
let store = GuardianDBKeyValue::new(
client.clone(),
identity.clone(),
address.clone(),
Some(opts),
)
.await
.expect("create writable store");
assert!(
store.is_writable(),
"freshly created store must be writable"
);
store
.put_impl("k", b"v".to_vec())
.await
.expect("write should succeed on writable store");
store.close().await.expect("close writable store");
}
let opts_ro = NewStoreOptions {
event_bus: Some(EventBus::new()),
directory: cache_dir.clone(),
read_only: Some(true),
..Default::default()
};
let ro = GuardianDBKeyValue::new(client, identity, address, Some(opts_ro))
.await
.expect("reopen read-only");
assert!(
!ro.is_writable(),
"reopened read-only replica must not be writable"
);
assert!(
ro.put_impl("k2", b"v2".to_vec()).await.is_err(),
"put must be rejected on a read-only replica"
);
assert!(
ro.delete_impl("k").await.is_err(),
"delete must be rejected on a read-only replica"
);
}
#[tokio::test]
async fn test_rotation_copies_all_key_value_state() {
use crate::traits::KeyValueStore;
let (src, _src_dir) = create_test_store()
.await
.expect("Failed to create source store");
src.put_impl("a", b"1".to_vec()).await.expect("put a");
src.put_impl("b", b"2".to_vec()).await.expect("put b");
src.put_impl("c", b"3".to_vec()).await.expect("put c");
let (dst, _dst_dir) = create_test_store()
.await
.expect("Failed to create destination store");
assert!(dst.is_empty(), "fresh destination should start empty");
let copied = crate::rotation::copy_key_value_state(
&src as &dyn KeyValueStore<Error = crate::guardian::error::GuardianError>,
&dst as &dyn KeyValueStore<Error = crate::guardian::error::GuardianError>,
)
.await
.expect("rotation copy should succeed");
assert_eq!(copied, 3, "all three keys should be copied");
let dst_state = dst.all();
assert_eq!(dst_state.len(), 3);
assert_eq!(
dst_state.get("a").map(|v| v.as_slice()),
Some(b"1".as_ref())
);
assert_eq!(
dst_state.get("b").map(|v| v.as_slice()),
Some(b"2".as_ref())
);
assert_eq!(
dst_state.get("c").map(|v| v.as_slice()),
Some(b"3".as_ref())
);
}
#[tokio::test]
async fn test_rotation_into_read_only_destination_fails() {
use crate::traits::KeyValueStore;
let temp_dir = TempDir::new().expect("temp dir");
let cache_dir = temp_dir.path().join("cache").to_string_lossy().to_string();
let client_config = ClientConfig {
data_store_path: Some(temp_dir.path().to_path_buf()),
..ClientConfig::development()
};
let client = Arc::new(IrohClient::new(client_config).await.expect("client"));
let identity = test_identity();
let address = test_address().await;
{
let opts = NewStoreOptions {
event_bus: Some(EventBus::new()),
directory: cache_dir.clone(),
..Default::default()
};
let w = GuardianDBKeyValue::new(
client.clone(),
identity.clone(),
address.clone(),
Some(opts),
)
.await
.expect("create writable");
w.put_impl("seed", b"x".to_vec()).await.expect("seed write");
crate::traits::Store::close(&w).await.expect("close");
}
let ro_opts = NewStoreOptions {
event_bus: Some(EventBus::new()),
directory: cache_dir.clone(),
read_only: Some(true),
..Default::default()
};
let dst_ro = GuardianDBKeyValue::new(client, identity, address, Some(ro_opts))
.await
.expect("reopen read-only");
let (src, _src_dir) = create_test_store().await.expect("source");
src.put_impl("a", b"1".to_vec()).await.expect("put a");
let result = crate::rotation::copy_key_value_state(
&src as &dyn KeyValueStore<Error = crate::guardian::error::GuardianError>,
&dst_ro as &dyn KeyValueStore<Error = crate::guardian::error::GuardianError>,
)
.await;
assert!(
result.is_err(),
"rotation into a read-only destination must fail"
);
}
}