openkeyv 0.1.4

Async Key-Value Store - A pluggable interface for KV Stores
Documentation
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()));
}