use std::collections::HashMap;
use std::time::Duration;
use openkeyv::protocol::{AsyncEnumerateCollections, AsyncEnumerateKeys, AsyncKeyValue};
use openkeyv::store::memory::MemoryStore;
use openkeyv::wrapper::compression::CompressionWrapper;
use openkeyv::wrapper::default_value::DefaultValueWrapper;
use openkeyv::wrapper::encryption::EncryptionWrapper;
use openkeyv::wrapper::fallback::FallbackWrapper;
use openkeyv::wrapper::limit_size::LimitSizeWrapper;
use openkeyv::wrapper::logging::LoggingWrapper;
use openkeyv::wrapper::passthrough_cache::PassthroughCacheWrapper;
use openkeyv::wrapper::prefix_collections::PrefixCollectionsWrapper;
use openkeyv::wrapper::prefix_keys::PrefixKeysWrapper;
use openkeyv::wrapper::readonly::ReadOnlyWrapper;
use openkeyv::wrapper::retry::RetryWrapper;
use openkeyv::wrapper::routing::{CollectionRoutingWrapper, RoutingWrapper};
use openkeyv::wrapper::single_collection::SingleCollectionWrapper;
use openkeyv::wrapper::statistics::StatisticsWrapper;
use openkeyv::wrapper::timeout::TimeoutWrapper;
use openkeyv::wrapper::ttl_clamp::TtlClampWrapper;
use serde_json::Value;
fn sample_value() -> HashMap<String, Value> {
HashMap::from([(String::from("name"), Value::String(String::from("svc")))])
}
#[tokio::test]
async fn timeout_wrapper_preserves_keys_and_collections() {
let inner = MemoryStore::new();
inner
.put("svc", sample_value(), Some("services"), None)
.await
.unwrap();
let wrapper = TimeoutWrapper::new(inner, Duration::from_millis(50));
let keys = wrapper.keys(Some("services"), None).await.unwrap();
let collections = wrapper.collections(None).await.unwrap();
assert_eq!(keys, vec![String::from("svc")]);
assert!(collections.contains(&String::from("services")));
}
#[tokio::test]
async fn retry_wrapper_preserves_keys_and_collections() {
let inner = MemoryStore::new();
inner
.put("svc", sample_value(), Some("services"), None)
.await
.unwrap();
let wrapper = RetryWrapper::new(inner);
let keys = wrapper.keys(Some("services"), None).await.unwrap();
let collections = wrapper.collections(None).await.unwrap();
assert_eq!(keys, vec![String::from("svc")]);
assert!(collections.contains(&String::from("services")));
}
#[tokio::test]
async fn statistics_wrapper_preserves_keys_and_collections() {
let inner = MemoryStore::new();
inner
.put("svc", sample_value(), Some("services"), None)
.await
.unwrap();
let wrapper = StatisticsWrapper::new(inner);
let keys = wrapper.keys(Some("services"), None).await.unwrap();
let collections = wrapper.collections(None).await.unwrap();
assert_eq!(keys, vec![String::from("svc")]);
assert!(collections.contains(&String::from("services")));
}
#[tokio::test]
async fn prefix_keys_wrapper_strips_prefix_from_enumeration() {
let inner = MemoryStore::new();
let wrapper = PrefixKeysWrapper::new(inner, "tenant:");
wrapper
.put("svc", sample_value(), Some("services"), None)
.await
.unwrap();
let keys = wrapper.keys(Some("services"), None).await.unwrap();
let collections = wrapper.collections(None).await.unwrap();
assert_eq!(keys, vec![String::from("svc")]);
assert!(collections.contains(&String::from("services")));
}
#[tokio::test]
async fn prefix_collections_wrapper_strips_prefix_from_collections() {
let inner = MemoryStore::new();
let wrapper = PrefixCollectionsWrapper::new(inner, "tenant_");
wrapper
.put("svc", sample_value(), Some("services"), None)
.await
.unwrap();
let keys = wrapper.keys(Some("services"), None).await.unwrap();
let collections = wrapper.collections(None).await.unwrap();
assert_eq!(keys, vec![String::from("svc")]);
assert!(collections.contains(&String::from("services")));
assert!(!collections.iter().any(|name| name.starts_with("tenant_")));
}
#[tokio::test]
async fn single_collection_wrapper_reconstructs_collections_and_keys() {
let inner = MemoryStore::new();
let wrapper = SingleCollectionWrapper::new(inner, "all_data");
wrapper
.put("svc", sample_value(), Some("services"), None)
.await
.unwrap();
wrapper
.put("cfg", sample_value(), Some("configs"), None)
.await
.unwrap();
let service_keys = wrapper.keys(Some("services"), None).await.unwrap();
let collections = wrapper.collections(None).await.unwrap();
assert_eq!(service_keys, vec![String::from("svc")]);
assert!(collections.contains(&String::from("services")));
assert!(collections.contains(&String::from("configs")));
}
#[tokio::test]
async fn fallback_wrapper_unions_primary_and_fallback_enumeration() {
let primary = MemoryStore::new();
let fallback = MemoryStore::new();
primary
.put("primary", sample_value(), Some("services"), None)
.await
.unwrap();
fallback
.put("fallback", sample_value(), Some("services"), None)
.await
.unwrap();
fallback
.put("cfg", sample_value(), Some("configs"), None)
.await
.unwrap();
let wrapper = FallbackWrapper::new(primary, fallback);
let keys = wrapper.keys(Some("services"), None).await.unwrap();
let collections = wrapper.collections(None).await.unwrap();
assert!(keys.contains(&String::from("primary")));
assert!(keys.contains(&String::from("fallback")));
assert!(collections.contains(&String::from("services")));
assert!(collections.contains(&String::from("configs")));
}
#[tokio::test]
async fn passthrough_cache_wrapper_unions_cache_and_primary_enumeration() {
let cache = MemoryStore::new();
let primary = MemoryStore::new();
cache
.put("cached", sample_value(), Some("services"), None)
.await
.unwrap();
primary
.put("primary", sample_value(), Some("services"), None)
.await
.unwrap();
let wrapper = PassthroughCacheWrapper::new(cache, primary);
let keys = wrapper.keys(Some("services"), None).await.unwrap();
let collections = wrapper.collections(None).await.unwrap();
assert!(keys.contains(&String::from("cached")));
assert!(keys.contains(&String::from("primary")));
assert!(collections.contains(&String::from("services")));
}
#[tokio::test]
async fn routing_wrapper_unions_enumeration_across_shards() {
let left = MemoryStore::new();
let right = MemoryStore::new();
left.put("apple", sample_value(), Some("services"), None)
.await
.unwrap();
right
.put("banana", sample_value(), Some("services"), None)
.await
.unwrap();
right
.put("cfg", sample_value(), Some("configs"), None)
.await
.unwrap();
let wrapper = RoutingWrapper::new(vec![Box::new(left), Box::new(right)], |_, key| {
usize::from(key.bytes().next().unwrap_or_default() >= b'n')
});
let keys = wrapper.keys(Some("services"), None).await.unwrap();
let collections = wrapper.collections(None).await.unwrap();
assert!(keys.contains(&String::from("apple")));
assert!(keys.contains(&String::from("banana")));
assert!(collections.contains(&String::from("services")));
assert!(collections.contains(&String::from("configs")));
}
#[tokio::test]
async fn collection_routing_wrapper_unions_enumeration_across_routes() {
let users_store = MemoryStore::new();
let orders_store = MemoryStore::new();
let default_store = MemoryStore::new();
users_store
.put("u1", sample_value(), Some("users"), None)
.await
.unwrap();
orders_store
.put("o1", sample_value(), Some("orders"), None)
.await
.unwrap();
default_store
.put("misc1", sample_value(), Some("misc"), None)
.await
.unwrap();
let mut routes = std::collections::HashMap::new();
routes.insert(String::from("users"), Box::new(users_store) as _);
routes.insert(String::from("orders"), Box::new(orders_store) as _);
let wrapper = CollectionRoutingWrapper::new(routes, Box::new(default_store));
let user_keys = wrapper.keys(Some("users"), None).await.unwrap();
let collections = wrapper.collections(None).await.unwrap();
assert_eq!(user_keys, vec![String::from("u1")]);
assert!(collections.contains(&String::from("users")));
assert!(collections.contains(&String::from("orders")));
assert!(collections.contains(&String::from("misc")));
}
#[tokio::test]
async fn logging_wrapper_preserves_keys_and_collections() {
let inner = MemoryStore::new();
inner
.put("svc", sample_value(), Some("services"), None)
.await
.unwrap();
let wrapper = LoggingWrapper::new(inner);
let keys = wrapper.keys(Some("services"), None).await.unwrap();
let collections = wrapper.collections(None).await.unwrap();
assert_eq!(keys, vec![String::from("svc")]);
assert!(collections.contains(&String::from("services")));
}
#[tokio::test]
async fn readonly_wrapper_preserves_keys_and_collections() {
let inner = MemoryStore::new();
inner
.put("svc", sample_value(), Some("services"), None)
.await
.unwrap();
let wrapper = ReadOnlyWrapper::new(inner);
let keys = wrapper.keys(Some("services"), None).await.unwrap();
let collections = wrapper.collections(None).await.unwrap();
assert_eq!(keys, vec![String::from("svc")]);
assert!(collections.contains(&String::from("services")));
let err = wrapper
.put("new", sample_value(), Some("services"), None)
.await;
assert!(err.is_err());
}
#[tokio::test]
async fn default_value_wrapper_preserves_keys_and_collections() {
let inner = MemoryStore::new();
inner
.put("svc", sample_value(), Some("services"), None)
.await
.unwrap();
let default_value = HashMap::from([(String::from("default"), Value::Bool(true))]);
let wrapper = DefaultValueWrapper::new(inner, default_value);
let keys = wrapper.keys(Some("services"), None).await.unwrap();
let collections = wrapper.collections(None).await.unwrap();
assert_eq!(keys, vec![String::from("svc")]);
assert!(collections.contains(&String::from("services")));
let missing = wrapper.get("missing", Some("services")).await.unwrap();
assert_eq!(
missing,
Some(HashMap::from([(
String::from("default"),
Value::Bool(true)
)]))
);
}
#[tokio::test]
async fn ttl_clamp_wrapper_preserves_keys_and_collections() {
let inner = MemoryStore::new();
inner
.put("svc", sample_value(), Some("services"), None)
.await
.unwrap();
let wrapper = TtlClampWrapper::with_range(inner, 1.0, 3600.0);
let keys = wrapper.keys(Some("services"), None).await.unwrap();
let collections = wrapper.collections(None).await.unwrap();
assert_eq!(keys, vec![String::from("svc")]);
assert!(collections.contains(&String::from("services")));
}
#[tokio::test]
async fn limit_size_wrapper_preserves_keys_and_collections() {
let inner = MemoryStore::new();
inner
.put("svc", sample_value(), Some("services"), None)
.await
.unwrap();
let wrapper = LimitSizeWrapper::with_range(inner, 1, 1024 * 1024);
let keys = wrapper.keys(Some("services"), None).await.unwrap();
let collections = wrapper.collections(None).await.unwrap();
assert_eq!(keys, vec![String::from("svc")]);
assert!(collections.contains(&String::from("services")));
}
#[tokio::test]
async fn compression_wrapper_preserves_keys_and_collections() {
let inner = MemoryStore::new();
inner
.put("svc", sample_value(), Some("services"), None)
.await
.unwrap();
let wrapper = CompressionWrapper::new(inner);
let keys = wrapper.keys(Some("services"), None).await.unwrap();
let collections = wrapper.collections(None).await.unwrap();
assert_eq!(keys, vec![String::from("svc")]);
assert!(collections.contains(&String::from("services")));
let value = wrapper.get("svc", Some("services")).await.unwrap();
assert_eq!(value, Some(sample_value()));
}
#[tokio::test]
async fn encryption_wrapper_preserves_keys_and_collections() {
let inner = MemoryStore::new();
inner
.put("svc", sample_value(), Some("services"), None)
.await
.unwrap();
let key = b"testkey";
let encrypt = |data: &[u8]| {
Ok(data
.iter()
.enumerate()
.map(|(i, b)| b ^ key[i % key.len()])
.collect::<Vec<u8>>())
};
let decrypt = |data: &[u8], _version: u32| {
Ok(data
.iter()
.enumerate()
.map(|(i, b)| b ^ key[i % key.len()])
.collect::<Vec<u8>>())
};
let wrapper = EncryptionWrapper::new(inner, encrypt, decrypt, 1);
let keys = wrapper.keys(Some("services"), None).await.unwrap();
let collections = wrapper.collections(None).await.unwrap();
assert_eq!(keys, vec![String::from("svc")]);
assert!(collections.contains(&String::from("services")));
let value = wrapper.get("svc", Some("services")).await.unwrap();
assert_eq!(value, Some(sample_value()));
}