use std::time::Duration;
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub enum BackendType {
#[default]
Memory,
Redis,
Chain,
}
#[derive(Debug, Clone)]
pub struct RedisConfig {
pub connection_string: String,
pub pool_size: usize,
pub connection_timeout: Duration,
pub retry_count: u32,
pub retry_delay: Duration,
pub circuit_breaker_threshold: u32,
pub circuit_breaker_reset_timeout: Duration,
}
impl Default for RedisConfig {
fn default() -> Self {
Self {
connection_string: String::new(),
pool_size: 8,
connection_timeout: Duration::from_secs(2),
retry_count: 3,
retry_delay: Duration::from_millis(100),
circuit_breaker_threshold: 5,
circuit_breaker_reset_timeout: Duration::from_secs(30),
}
}
}
#[derive(Debug, Clone)]
pub struct ChainLinkConfig {
pub backend: BackendType,
pub score: u8,
pub redis: Option<RedisConfig>,
}
#[derive(Debug, Clone)]
pub struct OxcacheConfig {
pub backend: BackendType,
pub capacity: u64,
pub ttl: Option<Duration>,
pub tti: Option<Duration>,
pub redis: Option<RedisConfig>,
pub chain: Vec<ChainLinkConfig>,
}
impl Default for OxcacheConfig {
fn default() -> Self {
Self {
backend: BackendType::Memory,
capacity: 10_000,
ttl: None,
tti: None,
redis: None,
chain: Vec::new(),
}
}
}
use std::any::TypeId;
use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use trait_kit::prelude::*;
use crate::backend::{CacheBackend, MokaMemoryBackend};
use crate::error::OxCacheError;
#[cfg(feature = "redis")]
use crate::backend::RedisBackend;
#[cfg(feature = "redis")]
use crate::cache::{ChainCacheBuilder, ChainLink};
pub struct OxcacheModule;
impl ModuleMeta for OxcacheModule {
const NAME: &'static str = "oxcache";
fn dependencies() -> &'static [(&'static str, TypeId)] {
&[]
}
}
impl AsyncAutoBuilder for OxcacheModule {
type Capability = Arc<dyn CacheBackend + Send + Sync>;
type Error = OxCacheError;
fn build<'a>(
kit: &'a AsyncKit,
) -> Pin<Box<dyn Future<Output = Result<Self::Capability, Self::Error>> + Send + 'a>> {
Box::pin(async move {
let config: OxcacheConfig = kit
.config()
.map_err(|e| OxCacheError::Internal(format!("OxcacheModule: read config: {e}")))?;
match config.backend {
BackendType::Memory => build_memory_backend(&config),
BackendType::Redis => build_redis_backend(&config).await,
BackendType::Chain => build_chain_cache(&config).await,
}
})
}
}
fn build_memory_backend(
config: &OxcacheConfig,
) -> Result<Arc<dyn CacheBackend + Send + Sync>, OxCacheError> {
let mut builder = MokaMemoryBackend::builder().capacity(config.capacity);
if let Some(ttl) = config.ttl {
builder = builder.ttl(ttl);
}
if let Some(tti) = config.tti {
builder = builder.time_to_idle(tti);
}
let backend = builder.build();
Ok(Arc::new(backend) as Arc<dyn CacheBackend + Send + Sync>)
}
#[cfg(feature = "redis")]
async fn build_redis_backend(
config: &OxcacheConfig,
) -> Result<Arc<dyn CacheBackend + Send + Sync>, OxCacheError> {
let rc = config.redis.as_ref().ok_or_else(|| {
OxCacheError::InvalidInput("OxcacheModule: backend=Redis requires `redis` config".into())
})?;
let backend = apply_redis_config(RedisBackend::builder(), rc)
.build()
.await?;
Ok(Arc::new(backend) as Arc<dyn CacheBackend + Send + Sync>)
}
#[cfg(not(feature = "redis"))]
async fn build_redis_backend(
_config: &OxcacheConfig,
) -> Result<Arc<dyn CacheBackend + Send + Sync>, OxCacheError> {
Err(OxCacheError::InvalidInput(
"OxcacheModule: backend=Redis requires the `redis` cargo feature to be enabled".into(),
))
}
#[cfg(feature = "redis")]
async fn build_chain_cache(
config: &OxcacheConfig,
) -> Result<Arc<dyn CacheBackend + Send + Sync>, OxCacheError> {
let mut chain_builder = ChainCacheBuilder::default();
if !config.chain.is_empty() {
for link_cfg in &config.chain {
let link = build_chain_link(link_cfg, config).await?;
chain_builder = chain_builder.link(link);
}
} else {
let l1 = {
let mut b = MokaMemoryBackend::builder().capacity(config.capacity);
if let Some(ttl) = config.ttl {
b = b.ttl(ttl);
}
if let Some(tti) = config.tti {
b = b.time_to_idle(tti);
}
b.build()
};
chain_builder = chain_builder.link(ChainLink::from_backend(l1));
let rc = config.redis.as_ref().ok_or_else(|| {
OxCacheError::InvalidInput(
"OxcacheModule: backend=Chain (default mode) requires `redis` config".into(),
)
})?;
let l2 = apply_redis_config(RedisBackend::builder(), rc)
.build()
.await?;
chain_builder = chain_builder.link(ChainLink::from_backend(l2));
}
let chain = chain_builder.build();
Ok(Arc::new(chain) as Arc<dyn CacheBackend + Send + Sync>)
}
#[cfg(not(feature = "redis"))]
async fn build_chain_cache(
_config: &OxcacheConfig,
) -> Result<Arc<dyn CacheBackend + Send + Sync>, OxCacheError> {
Err(OxCacheError::InvalidInput(
"OxcacheModule: backend=Chain requires the `redis` cargo feature to be enabled".into(),
))
}
#[cfg(feature = "redis")]
async fn build_chain_link(
link_cfg: &ChainLinkConfig,
config: &OxcacheConfig,
) -> Result<ChainLink, OxCacheError> {
match link_cfg.backend {
BackendType::Memory => {
let mut b = MokaMemoryBackend::builder().capacity(config.capacity);
if let Some(ttl) = config.ttl {
b = b.ttl(ttl);
}
if let Some(tti) = config.tti {
b = b.time_to_idle(tti);
}
let backend = b.build();
Ok(ChainLink::new(backend, link_cfg.score, false, "memory"))
}
BackendType::Redis => {
let rc = link_cfg
.redis
.as_ref()
.or(config.redis.as_ref())
.ok_or_else(|| {
OxCacheError::InvalidInput(
"OxcacheModule: chain link backend=Redis requires `redis` config".into(),
)
})?;
let backend = apply_redis_config(RedisBackend::builder(), rc)
.build()
.await?;
Ok(ChainLink::new(backend, link_cfg.score, true, "redis"))
}
BackendType::Chain => Err(OxCacheError::InvalidInput(
"OxcacheModule: nested Chain inside Chain is not supported".into(),
)),
}
}
#[cfg(feature = "redis")]
fn apply_redis_config(
mut builder: crate::backend::RedisBackendBuilder,
rc: &RedisConfig,
) -> crate::backend::RedisBackendBuilder {
if !rc.connection_string.is_empty() {
builder = builder.connection_string(&rc.connection_string);
}
builder = builder
.pool_size(rc.pool_size)
.connection_timeout(rc.connection_timeout)
.retry_count(rc.retry_count)
.retry_delay(rc.retry_delay)
.circuit_breaker_threshold(rc.circuit_breaker_threshold)
.circuit_breaker_reset_timeout(rc.circuit_breaker_reset_timeout);
builder
}
impl trait_kit::core::health::AsyncHealthCheck for OxcacheModule {
fn check(cap: &Self::Capability) -> trait_kit::core::health::HealthStatus {
match futures::executor::block_on(cap.health_check()) {
Ok(()) => trait_kit::core::health::HealthStatus::Healthy,
Err(e) => trait_kit::core::health::HealthStatus::unhealthy(format!(
"cache backend health check failed: {e}"
)),
}
}
}
impl trait_kit::core::lifecycle::AsyncLifecycle for OxcacheModule {
fn on_shutdown<'a>(cap: &'a Self::Capability) -> Pin<Box<dyn Future<Output = ()> + Send + 'a>> {
Box::pin(async move {
cap.shutdown().await;
})
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn oxcache_module_meta_name() {
assert_eq!(OxcacheModule::NAME, "oxcache");
}
#[test]
fn oxcache_module_meta_dependencies_empty() {
assert_eq!(
OxcacheModule::dependencies(),
&[] as &[(&'static str, TypeId)]
);
}
#[tokio::test]
async fn oxcache_module_build_returns_cache_capability() {
let mut kit = AsyncKit::new();
kit.set_config(OxcacheConfig::default());
kit.register::<OxcacheModule>()
.expect("register OxcacheModule");
let kit = kit.build().await.expect("AsyncKit::build");
let cache: Arc<dyn CacheBackend + Send + Sync> = kit
.require::<OxcacheModule>()
.expect("require OxcacheModule");
cache
.set(Arc::from("k"), Arc::new(b"v".to_vec()), None)
.await
.expect("set");
let got = cache.get("k").await.expect("get");
assert_eq!(got, Some(b"v".to_vec()));
}
#[tokio::test]
async fn oxcache_module_build_reads_config_from_kit() {
let mut kit = AsyncKit::new();
kit.set_config(OxcacheConfig {
capacity: 5,
..OxcacheConfig::default()
});
kit.register::<OxcacheModule>()
.expect("register OxcacheModule");
let kit = kit.build().await.expect("AsyncKit::build");
let cache: Arc<dyn CacheBackend + Send + Sync> = kit
.require::<OxcacheModule>()
.expect("require OxcacheModule");
for i in 0..6u8 {
cache
.set(Arc::from(format!("k{i}")), Arc::new(vec![i]), None)
.await
.expect("set");
}
cache.health_check().await.expect("health_check");
}
#[tokio::test]
async fn oxcache_module_build_is_async() {
let kit = AsyncKit::new();
kit.set_config(OxcacheConfig::default());
let fut = <OxcacheModule as AsyncAutoBuilder>::build(&kit);
let cache: Arc<dyn CacheBackend + Send + Sync> = fut.await.expect("build future resolves");
cache.clear().await.expect("clear");
fn assert_send_sync<T: Send + Sync>() {}
assert_send_sync::<Arc<dyn CacheBackend + Send + Sync>>();
}
#[tokio::test]
async fn oxcache_module_health_check_returns_healthy() {
let mut kit = AsyncKit::new();
kit.set_config(OxcacheConfig::default());
kit.register::<OxcacheModule>()
.expect("register OxcacheModule");
kit.register_health_check::<OxcacheModule>();
let kit = kit.build().await.expect("AsyncKit::build");
let status = kit.health_check::<OxcacheModule>().expect("health_check");
assert!(
status.is_healthy(),
"expected Healthy status, got {status:?}"
);
}
#[tokio::test]
async fn oxcache_module_health_check_direct_call() {
let kit = AsyncKit::new();
kit.set_config(OxcacheConfig::default());
let fut = <OxcacheModule as AsyncAutoBuilder>::build(&kit);
let cache = fut.await.expect("build");
let status = <OxcacheModule as trait_kit::core::health::AsyncHealthCheck>::check(&cache);
assert!(status.is_healthy(), "expected Healthy, got {status:?}");
}
#[tokio::test]
async fn oxcache_module_lifecycle_on_shutdown() {
let kit = AsyncKit::new();
kit.set_config(OxcacheConfig::default());
let fut = <OxcacheModule as AsyncAutoBuilder>::build(&kit);
let cache = fut.await.expect("build");
<OxcacheModule as trait_kit::core::lifecycle::AsyncLifecycle>::on_shutdown(&cache).await;
}
#[tokio::test]
async fn oxcache_module_lifecycle_full_kit_integration() {
let mut kit = AsyncKit::new();
kit.set_config(OxcacheConfig::default());
kit.register::<OxcacheModule>()
.expect("register OxcacheModule");
kit.register_lifecycle::<OxcacheModule>();
let kit = kit.build().await.expect("AsyncKit::build");
kit.shutdown_async().await;
}
#[test]
fn backend_type_default_is_memory() {
assert_eq!(BackendType::default(), BackendType::Memory);
}
#[test]
fn oxcache_config_default_selects_memory() {
let cfg = OxcacheConfig::default();
assert_eq!(cfg.backend, BackendType::Memory);
assert!(cfg.redis.is_none());
assert!(cfg.chain.is_empty());
}
#[tokio::test]
async fn explicit_memory_backend_builds() {
let cfg = OxcacheConfig {
backend: BackendType::Memory,
capacity: 100,
ttl: Some(Duration::from_secs(60)),
..OxcacheConfig::default()
};
let cache = build_memory_backend(&cfg).expect("build memory backend");
cache
.set(Arc::from("t021"), Arc::new(b"ok".to_vec()), None)
.await
.expect("set");
let got = cache.get("t021").await.expect("get");
assert_eq!(got, Some(b"ok".to_vec()));
}
#[tokio::test]
async fn redis_backend_without_config_errors() {
let cfg = OxcacheConfig {
backend: BackendType::Redis,
..OxcacheConfig::default()
};
let result = build_redis_backend(&cfg).await;
assert!(result.is_err(), "expected error for Redis without config");
let err_msg = format!("{}", result.err().unwrap());
assert!(
err_msg.contains("redis") || err_msg.contains("feature"),
"error should mention redis: {err_msg}"
);
}
#[tokio::test]
async fn chain_backend_without_config_errors() {
let cfg = OxcacheConfig {
backend: BackendType::Chain,
..OxcacheConfig::default()
};
let result = build_chain_cache(&cfg).await;
assert!(result.is_err(), "expected error for Chain without config");
}
#[cfg(feature = "redis")]
#[tokio::test]
async fn chain_nested_chain_errors() {
let link = ChainLinkConfig {
backend: BackendType::Chain,
score: 50,
redis: None,
};
let cfg = OxcacheConfig {
backend: BackendType::Chain,
chain: vec![link],
..OxcacheConfig::default()
};
let result = build_chain_cache(&cfg).await;
assert!(result.is_err());
assert!(format!("{}", result.err().unwrap()).contains("nested"));
}
#[test]
fn redis_config_defaults() {
let rc = RedisConfig::default();
assert_eq!(rc.pool_size, 8);
assert_eq!(rc.retry_count, 3);
assert_eq!(rc.connection_timeout, Duration::from_secs(2));
assert_eq!(rc.retry_delay, Duration::from_millis(100));
assert_eq!(rc.circuit_breaker_threshold, 5);
assert_eq!(rc.circuit_breaker_reset_timeout, Duration::from_secs(30));
}
}