use super::*;
use crate::backend::{DashMapMemoryBackend, MokaMemoryBackend};
use crate::testing::MockBackend;
#[test]
fn test_chain_link_creation() {
let backend = MockBackend::new("test", 50, false);
let link = ChainLink::from_backend(backend);
assert_eq!(link.score(), 50);
assert!(!link.is_persistent());
assert_eq!(link.name(), "test");
}
#[test]
fn test_chain_cache_builder() {
let high = MockBackend::new("high", 100, false);
let low = MockBackend::new("low", 50, true);
let chain = ChainCache::builder()
.backend(low)
.backend(high)
.enable_backfill()
.build();
assert_eq!(chain.links().len(), 2);
assert_eq!(chain.links()[0].score(), 100);
assert_eq!(chain.links()[1].score(), 50);
}
#[test]
fn test_chain_cache_builder_empty_links_panics() {
let result = std::panic::catch_unwind(|| ChainCache::builder().build());
let err = match result {
Ok(_) => panic!("empty chain build must panic"),
Err(e) => e,
};
let msg = err
.downcast_ref::<String>()
.cloned()
.or_else(|| err.downcast_ref::<&str>().map(|s| s.to_string()))
.unwrap_or_default();
assert!(msg.contains("at least one link"), "unexpected panic: {msg}");
}
#[tokio::test]
async fn test_chain_cache_get_set() {
let high = MockBackend::new("high", 100, false);
let low = MockBackend::new("low", 50, true);
let chain = ChainCache::builder().backend(high).backend(low).build();
chain.set("key", b"value".to_vec(), None).await.unwrap();
let value = chain.get("key").await.unwrap();
assert_eq!(value, Some(b"value".to_vec()));
}
#[tokio::test]
async fn test_chain_cache_delete() {
let high = MockBackend::new("high", 100, false);
let low = MockBackend::new("low", 50, true);
let chain = ChainCache::builder().backend(high).backend(low).build();
chain.set("key", b"value".to_vec(), None).await.unwrap();
chain.delete("key").await.unwrap();
let exists = chain.exists("key").await.unwrap();
assert!(!exists);
}
#[tokio::test]
async fn test_chain_cache_backfill() {
let chain = ChainCache::builder()
.link(ChainLink::new(
MockBackend::new("high", 100, false),
100,
false,
"high",
))
.link(ChainLink::new(
MockBackend::new("low", 50, true),
50,
true,
"low",
))
.enable_backfill()
.build();
chain.set("key", b"value".to_vec(), None).await.unwrap();
let value = chain.get("key").await.unwrap();
assert_eq!(value, Some(b"value".to_vec()));
}
#[tokio::test]
async fn test_empty_chain() {
let chain = ChainCache::new(vec![]);
let value = chain.get("key").await.unwrap();
assert!(value.is_none());
let exists = chain.exists("key").await.unwrap();
assert!(!exists);
}
#[test]
fn test_chain_link_new_constructor() {
let backend = MokaMemoryBackend::new();
let link = ChainLink::new(backend, 75, true, "custom");
assert_eq!(link.score(), 75);
assert!(link.is_persistent());
assert_eq!(link.name(), "custom");
let _backend_ref = link.backend();
}
#[test]
fn test_chain_link_from_backend_moka() {
let backend = MokaMemoryBackend::new();
let link = ChainLink::from_backend(backend);
assert_eq!(link.score(), 100);
assert!(!link.is_persistent());
assert_eq!(link.name(), "moka");
}
#[test]
fn test_chain_link_debug() {
let backend = MokaMemoryBackend::new();
let link = ChainLink::new(backend, 80, true, "dbg");
let debug_str = format!("{:?}", link);
assert!(debug_str.contains("ChainLink"));
assert!(debug_str.contains("80"));
assert!(debug_str.contains("dbg"));
}
#[test]
fn test_chain_cache_new_constructor() {
let link = ChainLink::from_backend(MokaMemoryBackend::new());
let chain = ChainCache::new(vec![link]);
assert_eq!(chain.len(), 1);
assert!(!chain.is_empty());
}
#[test]
fn test_chain_cache_len_is_empty() {
let empty = ChainCache::new(vec![]);
assert!(empty.is_empty());
assert_eq!(empty.len(), 0);
let chain = ChainCache::builder()
.backend(MokaMemoryBackend::new())
.build();
assert!(!chain.is_empty());
assert_eq!(chain.len(), 1);
}
#[test]
fn test_chain_cache_get_by_score() {
let chain = ChainCache::builder()
.link(ChainLink::new(MokaMemoryBackend::new(), 100, false, "high"))
.link(ChainLink::new(MokaMemoryBackend::new(), 50, true, "low"))
.build();
assert!(chain.get_by_score(100).is_some());
assert!(chain.get_by_score(50).is_some());
assert!(chain.get_by_score(75).is_none());
}
#[test]
fn test_chain_cache_highest_lowest_backend() {
let chain = ChainCache::builder()
.link(ChainLink::new(MokaMemoryBackend::new(), 50, true, "low"))
.link(ChainLink::new(MokaMemoryBackend::new(), 100, false, "high"))
.build();
let highest = chain.highest_score_backend().unwrap();
assert_eq!(highest.score(), 100);
assert_eq!(highest.name(), "high");
let lowest = chain.lowest_score_backend().unwrap();
assert_eq!(lowest.score(), 50);
assert_eq!(lowest.name(), "low");
}
#[test]
fn test_chain_cache_highest_lowest_empty() {
let chain = ChainCache::new(vec![]);
assert!(chain.highest_score_backend().is_none());
assert!(chain.lowest_score_backend().is_none());
}
#[test]
fn test_chain_cache_persistent_filters() {
let chain = ChainCache::builder()
.link(ChainLink::new(MokaMemoryBackend::new(), 100, false, "high"))
.link(ChainLink::new(MokaMemoryBackend::new(), 50, true, "low"))
.build();
let persistent = chain.persistent_backends();
assert_eq!(persistent.len(), 1);
assert_eq!(persistent[0].name(), "low");
let non_persistent = chain.non_persistent_backends();
assert_eq!(non_persistent.len(), 1);
assert_eq!(non_persistent[0].name(), "high");
}
#[test]
fn test_chain_cache_links_accessor() {
let chain = ChainCache::builder()
.link(ChainLink::new(MokaMemoryBackend::new(), 100, false, "high"))
.build();
let links = chain.links();
assert_eq!(links.len(), 1);
assert_eq!(links[0].name(), "high");
}
#[test]
fn test_builder_link_method() {
let link = ChainLink::new(MokaMemoryBackend::new(), 100, false, "moka");
let chain = ChainCache::builder().link(link).build();
assert_eq!(chain.len(), 1);
}
#[test]
fn test_builder_links_method() {
let links = vec![
ChainLink::new(MokaMemoryBackend::new(), 100, false, "high"),
ChainLink::new(MokaMemoryBackend::new(), 50, true, "low"),
];
let chain = ChainCache::builder().links(links).build();
assert_eq!(chain.len(), 2);
assert_eq!(chain.links()[0].score(), 100);
assert_eq!(chain.links()[1].score(), 50);
}
#[tokio::test]
async fn test_builder_default_time_to_live() {
let chain = ChainCache::builder()
.backend(MokaMemoryBackend::new())
.default_time_to_live(Duration::from_secs(60))
.build();
chain.set("key", b"value".to_vec(), None).await.unwrap();
let value = chain.get("key").await.unwrap();
assert_eq!(value, Some(b"value".to_vec()));
}
#[test]
fn test_builder_disable_backfill() {
let chain = ChainCache::builder()
.backend(MokaMemoryBackend::new())
.enable_backfill()
.disable_backfill()
.build();
assert_eq!(chain.len(), 1);
}
#[tokio::test]
async fn test_chain_cache_get_bytes_set_bytes() {
use crate::UnifiedCache;
let chain = ChainCache::builder()
.backend(MokaMemoryBackend::new())
.build();
chain
.set_bytes("key", b"value".to_vec(), None)
.await
.unwrap();
let value = chain.get_bytes("key").await.unwrap();
assert_eq!(value, Some(b"value".to_vec()));
}
#[tokio::test]
async fn test_chain_cache_get_bytes_missing() {
use crate::UnifiedCache;
let chain = ChainCache::builder()
.backend(MokaMemoryBackend::new())
.build();
let value = chain.get_bytes("missing").await.unwrap();
assert!(value.is_none());
}
#[tokio::test]
async fn test_chain_cache_clear() {
let chain = ChainCache::builder()
.backend(MokaMemoryBackend::new())
.build();
chain.set("key", b"value".to_vec(), None).await.unwrap();
assert!(chain.exists("key").await.unwrap());
chain.clear().await.unwrap();
assert!(!chain.exists("key").await.unwrap());
}
#[tokio::test]
async fn test_chain_cache_clear_empty() {
let chain = ChainCache::new(vec![]);
assert!(chain.clear().await.is_ok());
}
#[tokio::test]
async fn test_chain_cache_expire() {
let chain = ChainCache::builder()
.backend(MokaMemoryBackend::new())
.build();
chain.set("key", b"value".to_vec(), None).await.unwrap();
let result = chain.expire("key", Duration::from_secs(60)).await.unwrap();
assert!(result);
}
#[tokio::test]
async fn test_chain_cache_expire_missing_key() {
let chain = ChainCache::builder()
.backend(MokaMemoryBackend::new())
.build();
let result = chain
.expire("missing", Duration::from_secs(60))
.await
.unwrap();
assert!(!result);
}
#[tokio::test]
async fn test_chain_cache_set_empty_chain_error() {
let chain = ChainCache::new(vec![]);
let result = chain.set("key", b"value".to_vec(), None).await;
assert!(result.is_err());
}
#[tokio::test]
async fn test_chain_cache_delete_empty_chain() {
let chain = ChainCache::new(vec![]);
assert!(chain.delete("key").await.is_ok());
}
#[tokio::test]
async fn test_chain_cache_set_with_explicit_ttl() {
let chain = ChainCache::builder()
.backend(MokaMemoryBackend::new())
.build();
chain
.set("key", b"value".to_vec(), Some(Duration::from_secs(60)))
.await
.unwrap();
let value = chain.get("key").await.unwrap();
assert_eq!(value, Some(b"value".to_vec()));
}
#[tokio::test]
async fn test_chain_cache_multi_backend_set_writes_all() {
let high = MokaMemoryBackend::new();
let low = MokaMemoryBackend::new();
let high_ref = high.clone();
let low_ref = low.clone();
let chain = ChainCache::builder()
.link(ChainLink::new(high, 100, false, "high"))
.link(ChainLink::new(low, 50, true, "low"))
.build();
chain.set("key", b"value".to_vec(), None).await.unwrap();
assert_eq!(high_ref.get("key").await.unwrap(), Some(b"value".to_vec()));
assert_eq!(low_ref.get("key").await.unwrap(), Some(b"value".to_vec()));
}
#[tokio::test]
async fn test_chain_cache_delete_removes_from_all() {
let high = MokaMemoryBackend::new();
let low = MokaMemoryBackend::new();
let high_ref = high.clone();
let low_ref = low.clone();
let chain = ChainCache::builder()
.link(ChainLink::new(high, 100, false, "high"))
.link(ChainLink::new(low, 50, true, "low"))
.build();
chain.set("key", b"value".to_vec(), None).await.unwrap();
chain.delete("key").await.unwrap();
assert!(high_ref.get("key").await.unwrap().is_none());
assert!(low_ref.get("key").await.unwrap().is_none());
}
#[tokio::test]
async fn test_chain_cache_backfill_populates_higher() {
let high = MokaMemoryBackend::new();
let low = MokaMemoryBackend::new();
let high_ref = high.clone();
let low_ref = low.clone();
let chain = ChainCache::builder()
.link(ChainLink::new(high, 100, false, "high"))
.link(ChainLink::new(low, 50, true, "low"))
.enable_backfill()
.build();
low_ref
.set(Arc::from("key"), Arc::new(b"low_value".to_vec()), None)
.await
.unwrap();
assert!(high_ref.get("key").await.unwrap().is_none());
let value = chain.get("key").await.unwrap();
assert_eq!(value, Some(b"low_value".to_vec()));
let mut backfilled = false;
for _ in 0..10 {
if high_ref.get("key").await.unwrap().is_some() {
backfilled = true;
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
let high_value = high_ref.get("key").await.unwrap();
assert_eq!(
high_value,
Some(b"low_value".to_vec()),
"backfill should populate high backend"
);
assert!(backfilled, "backfill should complete asynchronously");
}
#[tokio::test]
async fn test_chain_cache_no_backfill_when_disabled() {
let high = MokaMemoryBackend::new();
let low = MokaMemoryBackend::new();
let high_ref = high.clone();
let low_ref = low.clone();
let chain = ChainCache::builder()
.link(ChainLink::new(high, 100, false, "high"))
.link(ChainLink::new(low, 50, true, "low"))
.build();
low_ref
.set(Arc::from("key"), Arc::new(b"low_value".to_vec()), None)
.await
.unwrap();
let value = chain.get("key").await.unwrap();
assert_eq!(value, Some(b"low_value".to_vec()));
assert!(high_ref.get("key").await.unwrap().is_none());
}
#[tokio::test]
async fn test_chain_cache_ttl_len_capacity() {
let chain = ChainCache::builder()
.backend(MokaMemoryBackend::new())
.build();
chain.set("key", b"value".to_vec(), None).await.unwrap();
let ttl = chain.ttl("key").await.unwrap();
assert!(ttl.is_none());
let len = CacheReader::len(&chain).await.unwrap();
assert!(len <= 100, "len should be reasonable after single insert");
let capacity = chain.capacity().await.unwrap();
assert!(capacity > 0);
}
#[tokio::test]
async fn test_chain_cache_reader_empty() {
let chain = ChainCache::new(vec![]);
assert_eq!(CacheReader::len(&chain).await.unwrap(), 0);
assert!(CacheReader::is_empty(&chain).await.unwrap());
assert_eq!(chain.capacity().await.unwrap(), 0);
let ttl = chain.ttl("key").await.unwrap();
assert!(ttl.is_none());
}
#[tokio::test]
async fn test_chain_cache_stats() {
let chain = ChainCache::builder()
.link(ChainLink::new(MokaMemoryBackend::new(), 100, false, "high"))
.link(ChainLink::new(MokaMemoryBackend::new(), 50, true, "low"))
.build();
let stats = chain.stats().await.unwrap();
assert_eq!(stats.get("type"), Some(&"chain".to_string()));
assert_eq!(stats.get("backend_count"), Some(&"2".to_string()));
assert_eq!(stats.get("backend_0_name"), Some(&"high".to_string()));
assert_eq!(stats.get("backend_0_score"), Some(&"100".to_string()));
assert_eq!(stats.get("backend_1_name"), Some(&"low".to_string()));
assert_eq!(stats.get("backend_1_score"), Some(&"50".to_string()));
}
#[tokio::test]
async fn test_chain_cache_stats_empty() {
let chain = ChainCache::new(vec![]);
let stats = chain.stats().await.unwrap();
assert_eq!(stats.get("type"), Some(&"chain".to_string()));
assert_eq!(stats.get("backend_count"), Some(&"0".to_string()));
}
#[tokio::test]
async fn test_chain_cache_health_check() {
let chain = ChainCache::builder()
.backend(MokaMemoryBackend::new())
.build();
assert!(chain.health_check().await.is_ok());
}
#[tokio::test]
async fn test_chain_cache_health_check_empty() {
let chain = ChainCache::new(vec![]);
assert!(chain.health_check().await.is_ok());
}
#[tokio::test]
async fn test_chain_read_degrades_when_high_backend_fails() {
let high = MockBackend::new("high", 100, false).with_fail_get();
let low = MockBackend::new("low", 50, true);
let chain = ChainCache::builder()
.link(ChainLink::from_backend(high))
.link(ChainLink::from_backend(low))
.build();
chain.links()[1]
.backend()
.set(Arc::from("key"), Arc::new(b"low_value".to_vec()), None)
.await
.unwrap();
let value = chain.get("key").await.unwrap();
assert_eq!(
value,
Some(b"low_value".to_vec()),
"L1 get 失败时应降级到 L2 读取"
);
}
#[tokio::test]
async fn test_chain_read_returns_none_when_all_backends_fail() {
let high = MockBackend::new("high", 100, false).with_fail_get();
let low = MockBackend::new("low", 50, true).with_fail_get();
let chain = ChainCache::builder()
.link(ChainLink::from_backend(high))
.link(ChainLink::from_backend(low))
.build();
let result = chain.get("key").await;
assert!(result.is_err(), "所有后端 get 失败时应返回 Err");
}
#[tokio::test]
async fn test_chain_health_check_fails_when_backend_unhealthy() {
let healthy = MockBackend::new("ok", 100, false);
let unhealthy = MockBackend::new("down", 50, true).with_fail_health();
let chain = ChainCache::builder()
.link(ChainLink::from_backend(healthy))
.link(ChainLink::from_backend(unhealthy))
.build();
let result = chain.health_check().await;
assert!(result.is_err(), "含不健康后端时 health_check 应失败");
}
#[tokio::test]
async fn test_chain_write_succeeds_when_partial_backend_fails() {
let failing = MockBackend::new("failing", 100, false).with_fail_set();
let ok = MokaMemoryBackend::new();
let chain = ChainCache::builder()
.link(ChainLink::from_backend(failing))
.link(ChainLink::from_backend(ok))
.build();
chain.set("key", b"value".to_vec(), None).await.unwrap();
}
#[tokio::test]
async fn test_chain_cache_shutdown() {
let chain = ChainCache::builder()
.backend(MokaMemoryBackend::new())
.build();
chain.shutdown().await;
}
#[test]
fn test_chain_cache_backend_kind() {
let chain = ChainCache::builder()
.backend(MokaMemoryBackend::new())
.build();
assert_eq!(chain.backend_kind(), BackendKind::Chain);
}
#[tokio::test]
async fn test_chain_set_with_ttl_propagates_to_all_links() {
let moka = MokaMemoryBackend::new();
let dashmap = DashMapMemoryBackend::new();
let mock = MockBackend::new("mock", 30, false);
let moka_ref = moka.clone();
let dashmap_ref = dashmap.clone();
let chain = ChainCache::builder()
.link(ChainLink::from_backend(moka))
.link(ChainLink::from_backend(dashmap))
.link(ChainLink::new(mock, 30, false, "mock"))
.build();
chain
.set("k", b"v".to_vec(), Some(Duration::from_millis(50)))
.await
.unwrap();
assert_eq!(chain.get("k").await.unwrap(), Some(b"v".to_vec()));
tokio::time::sleep(Duration::from_millis(100)).await;
let mut moka_expired = false;
for _ in 0..10 {
if moka_ref.get("k").await.unwrap().is_none() {
moka_expired = true;
break;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
assert!(moka_expired, "moka link should expire after TTL");
assert_eq!(
dashmap_ref.get("k").await.unwrap(),
None,
"dashmap link should expire after TTL"
);
assert_eq!(
chain.get("k").await.unwrap(),
None,
"chain get should return None after all links expired"
);
}
#[tokio::test]
async fn test_chain_ttl_returns_highest_score_link_ttl() {
let moka = MokaMemoryBackend::new();
let dashmap = DashMapMemoryBackend::new();
let chain = ChainCache::builder()
.link(ChainLink::from_backend(moka))
.link(ChainLink::from_backend(dashmap))
.build();
chain
.set("k", b"v".to_vec(), Some(Duration::from_secs(60)))
.await
.unwrap();
let ttl = chain.ttl("k").await.unwrap();
assert!(
ttl.is_some(),
"chain ttl should return Some for highest-score link"
);
let ttl = ttl.unwrap();
assert!(
ttl > Duration::from_secs(58) && ttl <= Duration::from_secs(60),
"chain ttl={} should be in (58s, 60s]",
ttl.as_secs_f64()
);
}
#[tokio::test]
async fn test_chain_expire_any_link_success_returns_true() {
let moka = MokaMemoryBackend::new();
let dashmap = DashMapMemoryBackend::new();
let chain = ChainCache::builder()
.link(ChainLink::from_backend(moka))
.link(ChainLink::from_backend(dashmap))
.build();
chain
.set("k", b"v".to_vec(), Some(Duration::from_secs(60)))
.await
.unwrap();
let result = chain.expire("k", Duration::from_secs(120)).await.unwrap();
assert!(
result,
"chain expire should return true when any link succeeds"
);
}
#[tokio::test]
async fn test_chain_expire_all_missing_returns_false() {
let moka = MokaMemoryBackend::new();
let dashmap = DashMapMemoryBackend::new();
let chain = ChainCache::builder()
.link(ChainLink::from_backend(moka))
.link(ChainLink::from_backend(dashmap))
.build();
let result = chain
.expire("missing", Duration::from_secs(60))
.await
.unwrap();
assert!(
!result,
"chain expire should return false when all links miss"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn test_chain_sync_get_set() {
let moka = MokaMemoryBackend::new();
let dashmap = DashMapMemoryBackend::new();
let chain = ChainCache::builder()
.link(ChainLink::from_sync_backend(moka))
.link(ChainLink::from_sync_backend(dashmap))
.build();
chain.set_sync("k", b"v".to_vec(), None).unwrap();
let value = chain.get_sync("k").unwrap();
assert_eq!(value, Some(b"v".to_vec()));
}
#[tokio::test(flavor = "multi_thread")]
async fn test_chain_sync_get_returns_highest_score_hit() {
use crate::backend::SyncCacheWriter;
let moka = MokaMemoryBackend::new();
let dashmap = DashMapMemoryBackend::new();
SyncCacheWriter::set(&moka, Arc::from("k"), Arc::new(b"high".to_vec()), None).unwrap();
SyncCacheWriter::set(&dashmap, Arc::from("k"), Arc::new(b"low".to_vec()), None).unwrap();
let chain = ChainCache::builder()
.link(ChainLink::from_sync_backend(moka))
.link(ChainLink::from_sync_backend(dashmap))
.build();
let value = chain.get_sync("k").unwrap();
assert_eq!(
value,
Some(b"high".to_vec()),
"get_sync should return highest-score link's value (Moka=100 > DashMap=90)"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn test_chain_sync_with_unsupported_link_falls_back_to_err() {
use crate::error::OxCacheError;
let moka = MokaMemoryBackend::new();
let mock = MockBackend::new("mock", 30, false);
let chain = ChainCache::builder()
.link(ChainLink::from_sync_backend(moka))
.link(ChainLink::from_backend(mock)) .build();
let result = chain.get_sync("k");
assert!(
matches!(result, Err(OxCacheError::NotSupported(_))),
"get_sync should return NotSupported when chain has non-sync link, got {:?}",
result
);
let result = chain.set_sync("k", b"v".to_vec(), None);
assert!(
matches!(result, Err(OxCacheError::NotSupported(_))),
"set_sync should return NotSupported when chain has non-sync link, got {:?}",
result
);
}
#[tokio::test(flavor = "multi_thread")]
async fn test_chain_sync_set_propagates_ttl() {
let moka = MokaMemoryBackend::new();
let dashmap = DashMapMemoryBackend::new();
let chain = ChainCache::builder()
.link(ChainLink::from_sync_backend(moka))
.link(ChainLink::from_sync_backend(dashmap))
.build();
chain
.set_sync("k", b"v".to_vec(), Some(Duration::from_millis(50)))
.unwrap();
let value = chain.get_sync("k").unwrap();
assert_eq!(value, Some(b"v".to_vec()));
tokio::time::sleep(Duration::from_millis(100)).await;
let mut expired = false;
for _ in 0..10 {
if chain.get_sync("k").unwrap().is_none() {
expired = true;
break;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
assert!(
expired,
"chain get_sync should return None after TTL expires on all links"
);
}
#[tokio::test]
async fn test_chain_race_read_returns_earliest_hit() {
let high = MockBackend::new("high", 100, false);
let low = MockBackend::new("low", 50, true);
high.set(Arc::from("k"), Arc::new(b"v".to_vec()), None)
.await
.unwrap();
let chain = ChainCache::builder().backend(high).backend(low).build();
let value = chain.get("k").await.unwrap();
assert_eq!(value, Some(b"v".to_vec()));
let high = MockBackend::new("high", 100, false);
let low = MockBackend::new("low", 50, true);
low.set(Arc::from("k"), Arc::new(b"l2".to_vec()), None)
.await
.unwrap();
let chain = ChainCache::builder()
.backend(high)
.backend(low)
.enable_race_read()
.build();
let value = chain.get("k").await.unwrap();
assert_eq!(
value,
Some(b"l2".to_vec()),
"race read should return value from whichever backend has it"
);
}
#[tokio::test]
async fn test_chain_race_read_backs_off_on_backend_error() {
let failing = MockBackend::new("high", 100, false).with_fail_get();
let ok = MockBackend::new("low", 50, true);
ok.set(Arc::from("k"), Arc::new(b"l2".to_vec()), None)
.await
.unwrap();
let chain = ChainCache::builder()
.backend(failing)
.backend(ok)
.enable_race_read()
.build();
let value = chain.get("k").await.unwrap();
assert_eq!(
value,
Some(b"l2".to_vec()),
"race read should degrade past failing backend"
);
}
#[tokio::test]
async fn test_chain_race_read_all_backends_fail() {
use crate::error::OxCacheError;
let failing1 = MockBackend::new("high", 100, false).with_fail_get();
let failing2 = MockBackend::new("low", 50, true).with_fail_get();
let chain = ChainCache::builder()
.backend(failing1)
.backend(failing2)
.enable_race_read()
.build();
let result = chain.get("k").await;
assert!(
matches!(result, Err(OxCacheError::Operation(ref msg)) if msg.contains("All backends failed")),
"race read should error when all backends fail, got {:?}",
result
);
}
#[tokio::test]
async fn test_chain_race_read_miss_returns_none() {
let high = MockBackend::new("high", 100, false);
let low = MockBackend::new("low", 50, true);
let chain = ChainCache::builder()
.backend(high)
.backend(low)
.enable_race_read()
.build();
let value = chain.get("missing").await.unwrap();
assert_eq!(value, None, "race read with no hits should return None");
}
#[tokio::test]
async fn test_chain_atomic_incr() {
let mock = MockBackend::new("mock", 100, false);
let chain = ChainCache::builder().backend(mock).build();
let val = chain.incr("counter", 1, None).await.unwrap();
assert_eq!(val, 1);
let val = chain.incr("counter", 10, None).await.unwrap();
assert_eq!(val, 11);
let val = chain.incr("counter", -3, None).await.unwrap();
assert_eq!(val, 8);
}
#[tokio::test]
async fn test_chain_atomic_compare_and_swap() {
let mock = MockBackend::new("mock", 100, false);
let chain = ChainCache::builder().backend(mock).build();
let ok = chain
.compare_and_swap("cas_key", None, b"initial".to_vec(), None)
.await
.unwrap();
assert!(ok);
let ok = chain
.compare_and_swap("cas_key", Some(b"initial"), b"updated".to_vec(), None)
.await
.unwrap();
assert!(ok);
let ok = chain
.compare_and_swap("cas_key", Some(b"initial"), b"again".to_vec(), None)
.await
.unwrap();
assert!(!ok);
}
#[tokio::test]
async fn test_chain_atomic_set_if_absent() {
let mock = MockBackend::new("mock", 100, false);
let chain = ChainCache::builder().backend(mock).build();
let ok = chain
.set_if_absent("nx_key", b"first".to_vec(), None)
.await
.unwrap();
assert!(ok);
let ok = chain
.set_if_absent("nx_key", b"second".to_vec(), None)
.await
.unwrap();
assert!(!ok);
}
#[tokio::test]
async fn test_chain_atomic_no_atomic_backend_returns_not_supported() {
use crate::error::OxCacheError;
let dashmap = DashMapMemoryBackend::new();
let chain = ChainCache::builder().backend(dashmap).build();
let result = chain.incr("k", 1, None).await;
assert!(
matches!(result, Err(OxCacheError::NotSupported(_))),
"incr should return NotSupported when no link implements AtomicCacheWriter"
);
let result = chain.compare_and_swap("k", None, b"v".to_vec(), None).await;
assert!(matches!(result, Err(OxCacheError::NotSupported(_))));
let result = chain.set_if_absent("k", b"v".to_vec(), None).await;
assert!(matches!(result, Err(OxCacheError::NotSupported(_))));
}
#[tokio::test]
async fn test_chain_keys_merges_and_deduplicates() {
let high = MockBackend::new("high", 100, false);
let low = MockBackend::new("low", 50, true);
high.set(Arc::from("a"), Arc::new(b"1".to_vec()), None)
.await
.unwrap();
high.set(Arc::from("b"), Arc::new(b"2".to_vec()), None)
.await
.unwrap();
low.set(Arc::from("b"), Arc::new(b"2b".to_vec()), None)
.await
.unwrap();
low.set(Arc::from("c"), Arc::new(b"3".to_vec()), None)
.await
.unwrap();
let chain = ChainCache::builder().backend(high).backend(low).build();
let mut keys = chain.keys("*").await.unwrap();
keys.sort();
assert_eq!(
keys,
vec!["a", "b", "c"],
"keys should be merged and deduplicated"
);
}
#[tokio::test]
async fn test_chain_keys_empty_chain() {
let chain = ChainCache::new(vec![]);
let keys = chain.keys("*").await.unwrap();
assert!(keys.is_empty());
}
#[tokio::test(flavor = "multi_thread")]
async fn test_chain_sync_delete() {
let moka = MokaMemoryBackend::new();
let dashmap = DashMapMemoryBackend::new();
let chain = ChainCache::builder()
.link(ChainLink::from_sync_backend(moka))
.link(ChainLink::from_sync_backend(dashmap))
.build();
chain.set_sync("k", b"v".to_vec(), None).unwrap();
assert_eq!(chain.get_sync("k").unwrap(), Some(b"v".to_vec()));
chain.delete_sync("k").unwrap();
assert_eq!(chain.get_sync("k").unwrap(), None);
}
#[tokio::test(flavor = "multi_thread")]
async fn test_chain_sync_delete_with_unsupported_link() {
use crate::error::OxCacheError;
let moka = MokaMemoryBackend::new();
let mock = MockBackend::new("mock", 30, false);
let chain = ChainCache::builder()
.link(ChainLink::from_sync_backend(moka))
.link(ChainLink::from_backend(mock)) .build();
let result = chain.delete_sync("k");
assert!(
matches!(result, Err(OxCacheError::NotSupported(_))),
"delete_sync should return NotSupported when chain has non-sync link"
);
}
#[test]
fn test_chain_link_debug_format_detailed() {
let backend = MockBackend::new("debug_test", 42, true);
let link = ChainLink::from_backend(backend);
let debug_str = format!("{:?}", link);
assert!(debug_str.contains("ChainLink"));
assert!(debug_str.contains("42"));
assert!(debug_str.contains("debug_test"));
assert!(debug_str.contains("is_persistent"));
}
#[tokio::test]
async fn test_chain_backfill_preserves_ttl() {
let high = MockBackend::new("high", 100, false);
let low = MockBackend::new("low", 50, true);
low.set(
Arc::from("ttl_key"),
Arc::new(b"val".to_vec()),
Some(Duration::from_secs(300)),
)
.await
.unwrap();
let chain = ChainCache::builder()
.backend(high)
.backend(low)
.enable_backfill()
.build();
let value = chain.get("ttl_key").await.unwrap();
assert_eq!(value, Some(b"val".to_vec()));
}
#[tokio::test]
async fn test_chain_stats() {
let high = MockBackend::new("high", 100, false);
let low = MockBackend::new("low", 50, true);
let chain = ChainCache::builder().backend(high).backend(low).build();
let stats = chain.stats().await.unwrap();
assert_eq!(stats.get("type"), Some(&"chain".to_string()));
assert_eq!(stats.get("backend_count"), Some(&"2".to_string()));
assert_eq!(stats.get("backend_0_name"), Some(&"high".to_string()));
assert_eq!(stats.get("backend_1_score"), Some(&"50".to_string()));
}
#[test]
fn test_chain_persistent_and_non_persistent_backends() {
let high = MockBackend::new("high", 100, false);
let low = MockBackend::new("low", 50, true);
let chain = ChainCache::builder().backend(high).backend(low).build();
let persistent = chain.persistent_backends();
assert_eq!(persistent.len(), 1);
assert_eq!(persistent[0].name(), "low");
let non_persistent = chain.non_persistent_backends();
assert_eq!(non_persistent.len(), 1);
assert_eq!(non_persistent[0].name(), "high");
}
#[tokio::test]
async fn test_chain_len_is_empty_capacity() {
let empty_chain = ChainCache::new(vec![]);
assert!(empty_chain.is_empty());
assert_eq!(empty_chain.len(), 0);
let mock = MockBackend::new("mock", 100, false);
let chain = ChainCache::builder().backend(mock).build();
assert!(!chain.is_empty());
assert_eq!(chain.len(), 1);
assert!(CacheReader::is_empty(&empty_chain).await.unwrap()); assert!(CacheReader::is_empty(&chain).await.unwrap()); assert_eq!(CacheReader::len(&chain).await.unwrap(), 0); assert_eq!(CacheReader::capacity(&chain).await.unwrap(), 0);
chain.set("k", b"v".to_vec(), None).await.unwrap();
assert!(!CacheReader::is_empty(&chain).await.unwrap());
assert_eq!(CacheReader::len(&chain).await.unwrap(), 1);
}
#[tokio::test]
async fn test_chain_empty_chain_operations() {
let chain = ChainCache::new(vec![]);
assert_eq!(chain.get("k").await.unwrap(), None);
assert!(chain.set("k", b"v".to_vec(), None).await.is_err());
assert!(chain.delete("k").await.is_ok());
assert!(chain.clear().await.is_ok());
assert!(chain.health_check().await.is_ok());
assert!(!chain.expire("k", Duration::from_secs(1)).await.unwrap());
}
#[tokio::test]
async fn test_chain_exists_and_ttl() {
let high = MockBackend::new("high", 100, false);
let low = MockBackend::new("low", 50, true);
low.set(
Arc::from("k"),
Arc::new(b"v".to_vec()),
Some(Duration::from_secs(60)),
)
.await
.unwrap();
let chain = ChainCache::builder().backend(high).backend(low).build();
assert!(chain.exists("k").await.unwrap());
assert!(!chain.exists("missing").await.unwrap());
let ttl = chain.ttl("k").await.unwrap();
assert!(ttl.is_some());
assert!(ttl.unwrap() > Duration::from_secs(58));
assert!(chain.ttl("missing").await.unwrap().is_none());
}
#[tokio::test]
async fn test_chain_expire_propagates_to_all() {
let high = MockBackend::new("high", 100, false);
let low = MockBackend::new("low", 50, true);
high.set(Arc::from("k"), Arc::new(b"v".to_vec()), None)
.await
.unwrap();
low.set(Arc::from("k"), Arc::new(b"v".to_vec()), None)
.await
.unwrap();
let chain = ChainCache::builder().backend(high).backend(low).build();
let ok = chain.expire("k", Duration::from_secs(60)).await.unwrap();
assert!(
ok,
"expire should return true when at least one backend succeeds"
);
}
mod read_strategy_tests {
use super::*;
use crate::backend::CacheWriter;
use crate::cache::chain::ChainReadStrategy;
use std::sync::Arc;
use std::time::Duration;
async fn set_on(backend: &MockBackend, key: &str, value: &[u8], ttl: Option<Duration>) {
let key: Arc<str> = Arc::from(key);
let value = Arc::new(value.to_vec());
CacheWriter::set(backend, key, value, ttl).await.unwrap();
}
fn build_chain(
l1: MockBackend,
l2: MockBackend,
l3: MockBackend,
strategy: ChainReadStrategy,
) -> ChainCache {
ChainCache::builder()
.backend(l1)
.backend(l2)
.backend(l3)
.read_strategy(strategy)
.build()
}
#[tokio::test]
async fn parallel_freshest_picks_longest_remaining_ttl() {
let fast = MockBackend::new("fast", 100, false);
let mid = MockBackend::new("mid", 90, false);
let slow = MockBackend::new("slow", 50, true);
set_on(&fast, "k", b"from-fast", Some(Duration::from_secs(30))).await;
set_on(&mid, "k", b"from-mid", Some(Duration::from_secs(5))).await;
set_on(&slow, "k", b"from-slow", Some(Duration::from_secs(300))).await;
let chain = build_chain(fast, mid, slow, ChainReadStrategy::ParallelFreshest);
let v = chain.get("k").await.unwrap();
assert_eq!(
v,
Some(b"from-slow".to_vec()),
"freshest (longest remaining TTL) link must win, not the highest score"
);
}
#[tokio::test]
async fn parallel_freshest_falls_back_to_highest_score_when_no_ttl() {
let fast = MockBackend::new("fast", 100, false);
let slow = MockBackend::new("slow", 50, true);
set_on(&fast, "k", b"from-fast", None).await;
set_on(&slow, "k", b"from-slow", None).await;
let chain = build_chain(
fast,
MockBackend::new("empty", 90, false),
slow,
ChainReadStrategy::ParallelFreshest,
);
let v = chain.get("k").await.unwrap();
assert_eq!(v, Some(b"from-fast".to_vec()));
}
#[tokio::test]
async fn parallel_freshest_tolerates_partial_backend_failure() {
let fast = MockBackend::new("fast", 100, false).with_fail_get();
let slow = MockBackend::new("slow", 50, true);
set_on(&slow, "k", b"from-slow", Some(Duration::from_secs(300))).await;
let chain = build_chain(
fast,
MockBackend::new("empty", 90, false),
slow,
ChainReadStrategy::ParallelFreshest,
);
let v = chain.get("k").await.unwrap();
assert_eq!(
v,
Some(b"from-slow".to_vec()),
"single-link failure must not block"
);
}
#[tokio::test]
async fn race_alias_keeps_highest_score_semantics() {
let fast = MockBackend::new("fast", 100, false);
let slow = MockBackend::new("slow", 50, true);
set_on(&fast, "k", b"from-fast", Some(Duration::from_secs(5))).await;
set_on(&slow, "k", b"from-slow", Some(Duration::from_secs(300))).await;
let chain = build_chain(
fast,
MockBackend::new("empty", 90, false),
slow,
ChainReadStrategy::Race,
);
let v = chain.get("k").await.unwrap();
assert_eq!(
v,
Some(b"from-fast".to_vec()),
"race must pick the highest-score hit"
);
}
#[tokio::test]
async fn sequential_default_reads_highest_score_first() {
let fast = MockBackend::new("fast", 100, false);
let slow = MockBackend::new("slow", 50, true);
set_on(&fast, "k", b"from-fast", Some(Duration::from_secs(5))).await;
set_on(&slow, "k", b"from-slow", Some(Duration::from_secs(300))).await;
let chain = ChainCache::builder()
.backend(fast)
.backend(MockBackend::new("empty", 90, false))
.backend(slow)
.build();
let v = chain.get("k").await.unwrap();
assert_eq!(v, Some(b"from-fast".to_vec()));
}
}
use crate::core::CacheEvent;
struct BatchProbeBackend {
inner: MockBackend,
fail_batch: bool,
batch_calls: Arc<std::sync::Mutex<Vec<Vec<String>>>>,
}
impl BatchProbeBackend {
fn new(name: &'static str, score: u8, fail_batch: bool) -> Self {
Self {
inner: MockBackend::new(name, score, false),
fail_batch,
batch_calls: Arc::new(std::sync::Mutex::new(Vec::new())),
}
}
}
impl BackendScore for BatchProbeBackend {
fn score(&self) -> u8 {
self.inner.score()
}
fn is_persistent(&self) -> bool {
self.inner.is_persistent()
}
fn backend_name(&self) -> &'static str {
self.inner.backend_name()
}
}
#[async_trait]
impl CacheReader for BatchProbeBackend {
async fn get(&self, key: &str) -> OxCacheResult<Option<Vec<u8>>> {
self.inner.get(key).await
}
async fn exists(&self, key: &str) -> OxCacheResult<bool> {
self.inner.exists(key).await
}
async fn ttl(&self, key: &str) -> OxCacheResult<Option<Duration>> {
self.inner.ttl(key).await
}
async fn len(&self) -> OxCacheResult<u64> {
self.inner.len().await
}
async fn capacity(&self) -> OxCacheResult<u64> {
self.inner.capacity().await
}
async fn stats(&self) -> OxCacheResult<HashMap<String, String>> {
self.inner.stats().await
}
async fn keys(&self, pattern: &str) -> OxCacheResult<Vec<String>> {
self.inner.keys(pattern).await
}
async fn get_many(&self, keys: &[String]) -> OxCacheResult<Vec<Option<Vec<u8>>>> {
self.batch_calls.lock().unwrap().push(keys.to_vec());
if self.fail_batch {
return Err(OxCacheError::Operation("batch fault injected".to_string()));
}
let mut results = Vec::with_capacity(keys.len());
for key in keys {
results.push(self.inner.get(key).await?);
}
Ok(results)
}
}
#[async_trait]
impl CacheWriter for BatchProbeBackend {
async fn set(
&self,
key: Arc<str>,
value: Arc<Vec<u8>>,
ttl: Option<Duration>,
) -> OxCacheResult<()> {
self.inner.set(key, value, ttl).await
}
async fn delete(&self, key: &str) -> OxCacheResult<()> {
self.inner.delete(key).await
}
async fn clear(&self) -> OxCacheResult<()> {
self.inner.clear().await
}
async fn expire(&self, key: &str, ttl: Duration) -> OxCacheResult<bool> {
self.inner.expire(key, ttl).await
}
}
#[async_trait]
impl CacheConnector for BatchProbeBackend {
async fn health_check(&self) -> OxCacheResult<()> {
self.inner.health_check().await
}
async fn shutdown(&self) {
self.inner.shutdown().await;
}
fn backend_kind(&self) -> BackendKind {
self.inner.backend_kind()
}
}
#[derive(Default)]
struct CapturingPublisher {
errors: CapturedErrors,
}
type CapturedErrors = Arc<std::sync::Mutex<Vec<(Option<String>, String)>>>;
#[async_trait]
impl EventPublisher for CapturingPublisher {
async fn publish(&self, _event: CacheEvent) -> Result<(), OxCacheError> {
Ok(())
}
fn publish_error(&self, key: Option<String>, error: String) -> Result<(), OxCacheError> {
self.errors.lock().unwrap().push((key, error));
Ok(())
}
}
struct ExpireFailingBackend {
inner: MockBackend,
}
impl ExpireFailingBackend {
fn new(name: &'static str, score: u8) -> Self {
Self {
inner: MockBackend::new(name, score, false),
}
}
}
impl BackendScore for ExpireFailingBackend {
fn score(&self) -> u8 {
self.inner.score()
}
fn is_persistent(&self) -> bool {
self.inner.is_persistent()
}
fn backend_name(&self) -> &'static str {
self.inner.backend_name()
}
}
#[async_trait]
impl CacheReader for ExpireFailingBackend {
async fn get(&self, key: &str) -> OxCacheResult<Option<Vec<u8>>> {
self.inner.get(key).await
}
async fn exists(&self, key: &str) -> OxCacheResult<bool> {
self.inner.exists(key).await
}
async fn ttl(&self, key: &str) -> OxCacheResult<Option<Duration>> {
self.inner.ttl(key).await
}
async fn len(&self) -> OxCacheResult<u64> {
self.inner.len().await
}
async fn capacity(&self) -> OxCacheResult<u64> {
self.inner.capacity().await
}
async fn stats(&self) -> OxCacheResult<HashMap<String, String>> {
self.inner.stats().await
}
async fn keys(&self, pattern: &str) -> OxCacheResult<Vec<String>> {
self.inner.keys(pattern).await
}
async fn get_many(&self, keys: &[String]) -> OxCacheResult<Vec<Option<Vec<u8>>>> {
self.inner.get_many(keys).await
}
}
#[async_trait]
impl CacheWriter for ExpireFailingBackend {
async fn set(
&self,
key: Arc<str>,
value: Arc<Vec<u8>>,
ttl: Option<Duration>,
) -> OxCacheResult<()> {
self.inner.set(key, value, ttl).await
}
async fn delete(&self, key: &str) -> OxCacheResult<()> {
self.inner.delete(key).await
}
async fn clear(&self) -> OxCacheResult<()> {
self.inner.clear().await
}
async fn expire(&self, _key: &str, _ttl: Duration) -> OxCacheResult<bool> {
Err(OxCacheError::Operation("expire fault injected".to_string()))
}
}
#[async_trait]
impl CacheConnector for ExpireFailingBackend {
async fn health_check(&self) -> OxCacheResult<()> {
self.inner.health_check().await
}
async fn shutdown(&self) {
self.inner.shutdown().await;
}
fn backend_kind(&self) -> BackendKind {
self.inner.backend_kind()
}
}
#[tokio::test]
async fn chain_expire_backend_error_emits_event_and_keeps_partial_success() {
let publisher = Arc::new(CapturingPublisher::default());
let errors = publisher.errors.clone();
let good = MockBackend::new("good", 100, false);
let bad = ExpireFailingBackend::new("bad", 50);
let chain = ChainCache::builder()
.link(ChainLink::from_backend(good))
.link(ChainLink::from_backend(bad))
.event_publisher(publisher)
.build();
CacheWriter::set(&chain, Arc::from("k"), Arc::new(b"v".to_vec()), None)
.await
.unwrap();
let result = chain.expire("k", Duration::from_secs(60)).await.unwrap();
assert!(result, "partial success should still return true");
let captured = errors.lock().unwrap();
assert_eq!(captured.len(), 1, "one backend error event expected");
let (key, error) = &captured[0];
assert_eq!(key.as_deref(), Some("k"));
assert!(error.contains("backend bad:"), "error = {}", error);
}
#[tokio::test]
async fn chain_expire_all_backends_error_returns_false_with_events() {
let publisher = Arc::new(CapturingPublisher::default());
let errors = publisher.errors.clone();
let chain = ChainCache::builder()
.link(ChainLink::from_backend(ExpireFailingBackend::new(
"bad1", 100,
)))
.link(ChainLink::from_backend(ExpireFailingBackend::new(
"bad2", 50,
)))
.event_publisher(publisher)
.build();
let result = chain.expire("k", Duration::from_secs(60)).await.unwrap();
assert!(!result);
assert_eq!(errors.lock().unwrap().len(), 2);
}
#[tokio::test]
async fn iter_entries_resolves_hits_in_input_order_with_single_batch_per_layer() {
let l1 = BatchProbeBackend::new("probe-l1", 100, false);
let l2 = BatchProbeBackend::new("probe-l2", 50, false);
let l1_calls = l1.batch_calls.clone();
let l2_calls = l2.batch_calls.clone();
CacheWriter::set(&l1.inner, Arc::from("k1"), Arc::new(b"v1".to_vec()), None)
.await
.unwrap();
CacheWriter::set(&l2.inner, Arc::from("k2"), Arc::new(b"v2".to_vec()), None)
.await
.unwrap();
let chain = ChainCache::builder()
.link(ChainLink::from_backend(l1))
.link(ChainLink::from_backend(l2))
.build();
let result = chain.iter_entries(&["k1", "k2", "k3"]).await;
assert_eq!(
result,
vec![
("k1".to_string(), Some(b"v1".to_vec())),
("k2".to_string(), Some(b"v2".to_vec())),
("k3".to_string(), None),
]
);
assert_eq!(
*l1_calls.lock().unwrap(),
vec![vec!["k1".to_string(), "k2".to_string(), "k3".to_string()]]
);
assert_eq!(
*l2_calls.lock().unwrap(),
vec![vec!["k2".to_string(), "k3".to_string()]]
);
}
#[tokio::test]
async fn iter_entries_layer_batch_failure_splits_per_key_and_falls_through() {
let l1 = BatchProbeBackend::new("faulty-l1", 100, true);
let l2 = BatchProbeBackend::new("probe-l2", 50, false);
let l2_calls = l2.batch_calls.clone();
CacheWriter::set(&l2.inner, Arc::from("k1"), Arc::new(b"v1".to_vec()), None)
.await
.unwrap();
let publisher = Arc::new(CapturingPublisher::default());
let errors = publisher.errors.clone();
let chain = ChainCache::builder()
.link(ChainLink::from_backend(l1))
.link(ChainLink::from_backend(l2))
.event_publisher(publisher)
.build();
let result = chain.iter_entries(&["k1", "k2"]).await;
assert_eq!(
result,
vec![
("k1".to_string(), Some(b"v1".to_vec())),
("k2".to_string(), None),
]
);
let errors = errors.lock().unwrap();
assert_eq!(errors.len(), 2, "L1 整批失败应拆分为每键一条错误");
assert_eq!(errors[0].0.as_deref(), Some("k1"));
assert_eq!(errors[1].0.as_deref(), Some("k2"));
drop(errors);
assert_eq!(
*l2_calls.lock().unwrap(),
vec![vec!["k1".to_string(), "k2".to_string()]]
);
}
#[tokio::test]
async fn iter_entries_all_layers_failing_maps_every_key_to_none() {
let l1 = BatchProbeBackend::new("faulty-l1", 100, true);
let l2 = BatchProbeBackend::new("faulty-l2", 50, true);
let l1_calls = l1.batch_calls.clone();
let l2_calls = l2.batch_calls.clone();
let publisher = Arc::new(CapturingPublisher::default());
let errors = publisher.errors.clone();
let chain = ChainCache::builder()
.link(ChainLink::from_backend(l1))
.link(ChainLink::from_backend(l2))
.event_publisher(publisher)
.build();
let result = chain.iter_entries(&["k1", "k2"]).await;
assert_eq!(
result,
vec![("k1".to_string(), None), ("k2".to_string(), None)]
);
assert_eq!(errors.lock().unwrap().len(), 4);
assert_eq!(l1_calls.lock().unwrap().len(), 1);
assert_eq!(l2_calls.lock().unwrap().len(), 1);
}
#[tokio::test]
async fn iter_entries_empty_input_returns_empty_vec() {
let l1 = BatchProbeBackend::new("probe-l1", 100, false);
let chain = ChainCache::builder()
.link(ChainLink::from_backend(l1))
.build();
assert!(chain.iter_entries(&[]).await.is_empty());
}
#[cfg(feature = "invalidation")]
mod invalidation_integration {
use super::*;
use crate::cache::tiered_with_invalidation;
use crate::features::invalidation::{
DEFAULT_CHANNEL, InMemoryPubSubTransport, InvalidationBus, InvalidationConfig,
};
async fn setup_buses() -> (
Arc<InvalidationBus>,
Arc<InvalidationBus>,
Arc<dyn CacheBackend>,
crate::features::invalidation::ListenerHandle,
) {
let transport = Arc::new(InMemoryPubSubTransport::new());
let bus_a = Arc::new(InvalidationBus::new(
transport.clone(),
InvalidationConfig::new("instance-a").with_channel(DEFAULT_CHANNEL),
));
let bus_b = Arc::new(InvalidationBus::new(
transport.clone(),
InvalidationConfig::new("instance-b").with_channel(DEFAULT_CHANNEL),
));
let remote: Arc<dyn CacheBackend> = Arc::new(MockBackend::new("remote", 50, true));
let handle = bus_b.spawn_listener(remote.clone()).await.unwrap();
tokio::time::sleep(Duration::from_millis(50)).await;
(bus_a, bus_b, remote, handle)
}
async fn poll_until_gone(backend: &Arc<dyn CacheBackend>, key: &str) -> bool {
for _ in 0..50 {
if !backend.exists(key).await.unwrap() {
return true;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
!backend.exists(key).await.unwrap()
}
async fn poll_stays(backend: &Arc<dyn CacheBackend>, key: &str, window_ms: u64) {
let rounds = window_ms / 10;
for _ in 0..rounds {
assert!(
backend.exists(key).await.unwrap(),
"条目不应被广播失效: {key}"
);
tokio::time::sleep(Duration::from_millis(10)).await;
}
}
#[tokio::test]
async fn with_invalidation_wraps_persistent_links_and_broadcasts() {
let (bus_a, _bus_b, remote, _handle) = setup_buses().await;
let chain = ChainCache::builder()
.link(ChainLink::new(
MockBackend::new("l1", 100, false),
100,
false,
"l1",
))
.link(ChainLink::new(
MockBackend::new("l2", 50, true),
50,
true,
"l2",
))
.with_invalidation(bus_a)
.build();
remote
.set(Arc::from("user:1"), Arc::new(b"stale".to_vec()), None)
.await
.unwrap();
chain.set("user:1", b"fresh".to_vec(), None).await.unwrap();
assert_eq!(chain.get("user:1").await.unwrap(), Some(b"fresh".to_vec()));
assert!(
poll_until_gone(&remote, "user:1").await,
"持久层写入后远端旧值应被广播失效"
);
}
#[test]
fn with_invalidation_without_persistent_link_panics() {
let transport = Arc::new(InMemoryPubSubTransport::new());
let bus = Arc::new(InvalidationBus::new(
transport,
InvalidationConfig::new("instance-a").with_channel(DEFAULT_CHANNEL),
));
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
let _ = ChainCache::builder()
.link(ChainLink::new(
MockBackend::new("l1", 100, false),
100,
false,
"l1",
))
.with_invalidation(bus)
.build();
}));
let msg = result
.err()
.map(|e| {
e.downcast_ref::<String>()
.cloned()
.or_else(|| e.downcast_ref::<&str>().map(|s| s.to_string()))
.unwrap_or_default()
})
.unwrap_or_default();
assert!(
msg.contains("persistent link"),
"panic 应说明缺少持久层, got {msg}"
);
}
#[tokio::test]
async fn with_invalidation_expire_does_not_broadcast() {
let (bus_a, _bus_b, remote, _handle) = setup_buses().await;
let chain = ChainCache::builder()
.link(ChainLink::new(
MockBackend::new("l2", 50, true),
50,
true,
"l2",
))
.with_invalidation(bus_a)
.build();
remote
.set(Arc::from("user:2"), Arc::new(b"stale".to_vec()), None)
.await
.unwrap();
chain
.expire("user:2", Duration::from_secs(60))
.await
.unwrap();
poll_stays(&remote, "user:2", 150).await;
}
#[tokio::test]
async fn tiered_with_invalidation_assembles_wrapped_persistent_layer() {
use crate::cache::{L1Builder, L2Builder};
let (bus_a, _bus_b, remote, _handle) = setup_buses().await;
let l2_backend: Arc<dyn CacheBackend> = Arc::new(MockBackend::new("mock-l2", 50, true));
let chain = tiered_with_invalidation(
L1Builder::new().capacity(100),
L2Builder::new().custom(l2_backend).persistent(true),
bus_a,
)
.await
.unwrap();
remote
.set(Arc::from("user:3"), Arc::new(b"stale".to_vec()), None)
.await
.unwrap();
chain.set("user:3", b"fresh".to_vec(), None).await.unwrap();
assert_eq!(chain.get("user:3").await.unwrap(), Some(b"fresh".to_vec()));
assert_eq!(chain.len(), 2);
assert!(
poll_until_gone(&remote, "user:3").await,
"一站式装配的持久层写入应广播失效"
);
}
}