use crate::error::{OxCacheError, OxCacheResult};
use crate::i18n::messages::{
MSG_DETAIL_CONFERS_EXPECTS_BOOL, MSG_DETAIL_CONFERS_EXPECTS_F64,
MSG_DETAIL_CONFERS_EXPECTS_STRING, MSG_DETAIL_CONFERS_EXPECTS_U64,
MSG_DETAIL_CONFERS_READ_FAILED, MSG_DETAIL_CONFERS_VALUE_EXCEEDS_RANGE, t,
};
use serde::{Deserialize, Serialize};
use std::sync::{Arc, RwLock};
use std::time::Duration;
pub const CACHE_KEY_PREFIX: &str = "cache";
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(default)]
pub struct CircuitBreakerSettings {
pub failure_threshold: u32,
pub recovery_timeout_ms: u64,
}
impl Default for CircuitBreakerSettings {
fn default() -> Self {
Self {
failure_threshold: 5,
recovery_timeout_ms: 30_000,
}
}
}
#[derive(Clone, PartialEq, Serialize, Deserialize)]
#[serde(default)]
pub struct OxcacheConfig {
pub capacity: u64,
pub default_ttl_ms: u64,
pub tti_ms: Option<u64>,
pub null_cache_ttl_ms: Option<u64>,
pub ttl_jitter_factor: Option<f64>,
pub sync_mode: Option<bool>,
pub backend: Option<String>,
pub metrics_enabled: Option<bool>,
pub redis_url: Option<String>,
pub disk_path: Option<String>,
pub serialization_format: Option<String>,
pub connection_pool_size: Option<usize>,
pub service_name: Option<String>,
pub circuit_breaker: CircuitBreakerSettings,
}
impl std::fmt::Debug for OxcacheConfig {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("OxcacheConfig")
.field("capacity", &self.capacity)
.field("default_ttl_ms", &self.default_ttl_ms)
.field("tti_ms", &self.tti_ms)
.field("null_cache_ttl_ms", &self.null_cache_ttl_ms)
.field("ttl_jitter_factor", &self.ttl_jitter_factor)
.field("sync_mode", &self.sync_mode)
.field("backend", &self.backend)
.field("metrics_enabled", &self.metrics_enabled)
.field(
"redis_url",
&self.redis_url.as_deref().map(|_| "***" as &str),
)
.field("disk_path", &self.disk_path)
.field("serialization_format", &self.serialization_format)
.field("connection_pool_size", &self.connection_pool_size)
.field("service_name", &self.service_name)
.field("circuit_breaker", &self.circuit_breaker)
.finish()
}
}
impl Default for OxcacheConfig {
fn default() -> Self {
Self {
capacity: 10_000,
default_ttl_ms: 60_000,
tti_ms: None,
null_cache_ttl_ms: None,
ttl_jitter_factor: None,
sync_mode: None,
backend: None,
metrics_enabled: None,
redis_url: None,
disk_path: None,
serialization_format: None,
connection_pool_size: None,
service_name: None,
circuit_breaker: CircuitBreakerSettings::default(),
}
}
}
impl OxcacheConfig {
pub async fn load_from<C>(connector: &C) -> OxCacheResult<Self>
where
C: confers::ConfigConnector,
{
let mut cfg = Self::default();
if let Some(v) = get_u64(connector, "cache.capacity").await? {
cfg.capacity = v;
}
if let Some(v) = get_u64(connector, "cache.default_ttl_ms").await? {
cfg.default_ttl_ms = v;
}
if let Some(v) = get_u64(connector, "cache.circuit_breaker.failure_threshold").await? {
cfg.circuit_breaker.failure_threshold = u32::try_from(v).map_err(|_| {
out_of_range_u64("cache.circuit_breaker.failure_threshold", v, "u32")
})?;
}
if let Some(v) = get_u64(connector, "cache.circuit_breaker.recovery_timeout_ms").await? {
cfg.circuit_breaker.recovery_timeout_ms = v;
}
if let Some(v) = get_u64(connector, "cache.tti_ms").await? {
cfg.tti_ms = Some(v);
}
if let Some(v) = get_u64(connector, "cache.null_cache_ttl_ms").await? {
cfg.null_cache_ttl_ms = Some(v);
}
if let Some(v) = get_f64(connector, "cache.ttl_jitter_factor").await? {
cfg.ttl_jitter_factor = Some(v);
}
if let Some(v) = get_bool(connector, "cache.sync_mode").await? {
cfg.sync_mode = Some(v);
}
if let Some(v) = get_string(connector, "cache.backend").await? {
cfg.backend = Some(v);
}
if let Some(v) = get_bool(connector, "cache.metrics_enabled").await? {
cfg.metrics_enabled = Some(v);
}
if let Some(v) = get_string(connector, "cache.redis_url").await? {
cfg.redis_url = Some(v);
}
if let Some(v) = get_string(connector, "cache.disk_path").await? {
cfg.disk_path = Some(v);
}
if let Some(v) = get_u64(connector, "cache.connection_pool_size").await? {
cfg.connection_pool_size = Some(
usize::try_from(v)
.map_err(|_| out_of_range_u64("cache.connection_pool_size", v, "usize"))?,
);
}
if let Some(v) = get_string(connector, "cache.serialization_format").await? {
cfg.serialization_format = Some(v);
}
if let Some(v) = get_string(connector, "cache.service_name").await? {
cfg.service_name = Some(v);
}
Ok(cfg)
}
pub fn default_ttl(&self) -> Duration {
Duration::from_millis(self.default_ttl_ms)
}
}
impl From<&OxcacheConfig> for crate::cache::L1Builder {
fn from(config: &OxcacheConfig) -> Self {
crate::cache::L1Builder::new()
.capacity(config.capacity)
.ttl(config.default_ttl())
}
}
async fn get_u64<C>(connector: &C, key: &str) -> OxCacheResult<Option<u64>>
where
C: confers::ConfigConnector,
{
match connector.get_raw(key).await.map_err(|e| {
OxCacheError::Operation(t(
MSG_DETAIL_CONFERS_READ_FAILED,
&[("key", key.to_string()), ("err", e.to_string())],
))
})? {
None => Ok(None),
Some(v) => v.as_u64().map(Some).ok_or_else(|| {
OxCacheError::InvalidInput(t(
MSG_DETAIL_CONFERS_EXPECTS_U64,
&[("key", key.to_string()), ("value", format!("{v:?}"))],
))
}),
}
}
fn out_of_range_u64(key: &str, value: u64, target: &str) -> OxCacheError {
OxCacheError::InvalidInput(t(
MSG_DETAIL_CONFERS_VALUE_EXCEEDS_RANGE,
&[
("key", key.to_string()),
("value", value.to_string()),
("target", target.to_string()),
],
))
}
async fn get_bool<C>(connector: &C, key: &str) -> OxCacheResult<Option<bool>>
where
C: confers::ConfigConnector,
{
match connector.get_raw(key).await.map_err(|e| {
OxCacheError::Operation(t(
MSG_DETAIL_CONFERS_READ_FAILED,
&[("key", key.to_string()), ("err", e.to_string())],
))
})? {
None => Ok(None),
Some(v) => v.as_bool().map(Some).ok_or_else(|| {
OxCacheError::InvalidInput(t(
MSG_DETAIL_CONFERS_EXPECTS_BOOL,
&[("key", key.to_string()), ("value", format!("{v:?}"))],
))
}),
}
}
async fn get_f64<C>(connector: &C, key: &str) -> OxCacheResult<Option<f64>>
where
C: confers::ConfigConnector,
{
match connector.get_raw(key).await.map_err(|e| {
OxCacheError::Operation(t(
MSG_DETAIL_CONFERS_READ_FAILED,
&[("key", key.to_string()), ("err", e.to_string())],
))
})? {
None => Ok(None),
Some(v) => v.as_f64().map(Some).ok_or_else(|| {
OxCacheError::InvalidInput(t(
MSG_DETAIL_CONFERS_EXPECTS_F64,
&[("key", key.to_string()), ("value", format!("{v:?}"))],
))
}),
}
}
async fn get_string<C>(connector: &C, key: &str) -> OxCacheResult<Option<String>>
where
C: confers::ConfigConnector,
{
match connector.get_raw(key).await.map_err(|e| {
OxCacheError::Operation(t(
MSG_DETAIL_CONFERS_READ_FAILED,
&[("key", key.to_string()), ("err", e.to_string())],
))
})? {
None => Ok(None),
Some(v) => v.as_str().map(|s| Some(s.to_string())).ok_or_else(|| {
OxCacheError::InvalidInput(t(
MSG_DETAIL_CONFERS_EXPECTS_STRING,
&[("key", key.to_string()), ("value", format!("{v:?}"))],
))
}),
}
}
impl crate::config::CacheConfig {
pub fn try_from_confers(config: &OxcacheConfig) -> OxCacheResult<Self> {
Ok(Self {
capacity: Some(config.capacity),
ttl: Some(config.default_ttl()),
tti: config.tti_ms.map(Duration::from_millis),
null_cache_ttl: config.null_cache_ttl_ms.map(Duration::from_millis),
ttl_jitter_factor: config.ttl_jitter_factor,
sync_mode: config.sync_mode,
backend: config.backend.clone(),
metrics_enabled: config.metrics_enabled,
serialization_format: config.serialization_format.clone(),
redis_url: config.redis_url.clone(),
disk_path: config.disk_path.clone(),
connection_pool_size: config.connection_pool_size,
circuit_breaker_failure_threshold: Some(config.circuit_breaker.failure_threshold),
circuit_breaker_reset_timeout: Some(Duration::from_millis(
config.circuit_breaker.recovery_timeout_ms,
)),
service_name: config.service_name.clone(),
})
}
}
#[async_trait::async_trait]
pub trait CacheConfigSource: Send + Sync {
async fn load(&self) -> OxCacheResult<OxcacheConfig>;
}
pub struct ConfersConfigSource<C> {
connector: Arc<C>,
}
impl<C> ConfersConfigSource<C> {
pub fn new(connector: Arc<C>) -> Self {
Self { connector }
}
}
#[async_trait::async_trait]
impl<C> CacheConfigSource for ConfersConfigSource<C>
where
C: confers::ConfigConnector + 'static,
{
async fn load(&self) -> OxCacheResult<OxcacheConfig> {
OxcacheConfig::load_from(self.connector.as_ref()).await
}
}
#[derive(Clone)]
pub struct ConfigSnapshot {
inner: Arc<RwLock<Arc<OxcacheConfig>>>,
}
impl ConfigSnapshot {
pub fn new(config: OxcacheConfig) -> Self {
Self {
inner: Arc::new(RwLock::new(Arc::new(config))),
}
}
pub fn get(&self) -> Arc<OxcacheConfig> {
self.inner
.read()
.map(|guard| guard.clone())
.unwrap_or_else(|_| Arc::new(OxcacheConfig::default()))
}
fn swap(&self, config: OxcacheConfig) -> Arc<OxcacheConfig> {
match self.inner.write() {
Ok(mut guard) => std::mem::replace(&mut guard, Arc::new(config)),
Err(_) => Arc::new(OxcacheConfig::default()),
}
}
}
pub trait ConfigChangeListener: Send + Sync {
fn on_change(&self, old: &OxcacheConfig, new: &OxcacheConfig);
}
pub struct ConfersConfigWatcher {
snapshot: ConfigSnapshot,
listeners: Vec<Arc<dyn ConfigChangeListener>>,
}
impl ConfersConfigWatcher {
pub fn new(initial: OxcacheConfig) -> Self {
Self {
snapshot: ConfigSnapshot::new(initial),
listeners: Vec::new(),
}
}
pub fn on_change(mut self, listener: Arc<dyn ConfigChangeListener>) -> Self {
self.listeners.push(listener);
self
}
pub fn snapshot(&self) -> ConfigSnapshot {
self.snapshot.clone()
}
pub fn watch<B, S>(self: &Arc<Self>, bus: Arc<B>, source: Arc<S>) -> tokio::task::JoinHandle<()>
where
B: confers::ConfigBus + 'static,
S: CacheConfigSource + 'static,
{
let watcher = Arc::clone(self);
tokio::spawn(async move {
let Ok(mut stream) = bus.subscribe().await else {
return;
};
use futures::stream::StreamExt;
while let Some(event) = stream.next().await {
let relevant = event.changed_keys.is_empty()
|| event
.changed_keys
.iter()
.any(|k| k.starts_with(CACHE_KEY_PREFIX));
if !relevant {
continue;
}
match source.load().await {
Ok(new_cfg) => {
let old_arc = watcher.snapshot.swap(new_cfg);
let new_arc = watcher.snapshot.get();
for listener in &watcher.listeners {
listener.on_change(&old_arc, &new_arc);
}
}
Err(err) => {
oxcache_report_config_reload_rejected(&err);
continue;
}
}
}
})
}
}
#[cfg(feature = "metrics")]
#[inline]
fn record_reload_rejected_counter() {
crate::infra::metrics::unified::GLOBAL_UNIFIED_METRICS
.increment_counter("oxcache_config_reload_rejected_total", 1);
}
#[cfg(not(feature = "metrics"))]
#[inline]
fn record_reload_rejected_counter() {}
#[cfg(feature = "telemetry")]
#[inline]
fn warn_reload_rejected(err: &OxCacheError) {
use crate::i18n::messages::{MSG_LOG_CONFERS_RELOAD_REJECTED, t};
tracing::warn!(
target = "oxcache::confers_config",
%err,
"{}",
t(MSG_LOG_CONFERS_RELOAD_REJECTED, &[])
);
}
#[cfg(not(feature = "telemetry"))]
#[inline]
fn warn_reload_rejected(_err: &OxCacheError) {}
#[inline]
fn oxcache_report_config_reload_rejected(err: &OxCacheError) {
record_reload_rejected_counter();
warn_reload_rejected(err);
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::CacheConfig;
use confers::{ConfigBus, ConfigChangeEvent, ConfigValue, InMemoryBus, SourceId};
use confers::{ConfigWriter, new_in_memory};
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::time::Duration;
#[cfg(feature = "metrics")]
fn dynamic_counter(metrics: &crate::infra::metrics::unified::UnifiedMetrics, key: &str) -> u64 {
use crate::infra::metrics::unified::MetricValue;
metrics
.get_dynamic_metrics()
.get(key)
.and_then(|v| match v {
MetricValue::Counter(c) => Some(*c),
_ => None,
})
.unwrap_or(0)
}
async fn memory_connector_with_capacity(capacity: u64) -> impl confers::ConfigConnector {
let conn = new_in_memory();
conn.set(
"cache.capacity",
confers::AnnotatedValue::new(
ConfigValue::U64(capacity),
SourceId::default(),
"cache.capacity",
),
)
.await
.unwrap();
conn.set(
"cache.default_ttl_ms",
confers::AnnotatedValue::new(
ConfigValue::U64(120_000),
SourceId::default(),
"cache.default_ttl_ms",
),
)
.await
.unwrap();
conn
}
#[tokio::test]
async fn load_from_extended_keys() {
let conn = new_in_memory();
let entries: &[(&str, ConfigValue)] = &[
("cache.capacity", ConfigValue::U64(2048)),
("cache.default_ttl_ms", ConfigValue::U64(90_000)),
("cache.tti_ms", ConfigValue::U64(15_000)),
("cache.null_cache_ttl_ms", ConfigValue::U64(500)),
("cache.ttl_jitter_factor", ConfigValue::F64(0.3)),
("cache.sync_mode", ConfigValue::Bool(true)),
("cache.backend", ConfigValue::String("redis".into())),
("cache.metrics_enabled", ConfigValue::Bool(false)),
(
"cache.redis_url",
ConfigValue::String("redis://cfg:6379".into()),
),
(
"cache.serialization_format",
ConfigValue::String("postcard".into()),
),
("cache.connection_pool_size", ConfigValue::U64(16)),
(
"cache.service_name",
ConfigValue::String("r9-confers".into()),
),
(
"cache.circuit_breaker.failure_threshold",
ConfigValue::U64(9),
),
(
"cache.circuit_breaker.recovery_timeout_ms",
ConfigValue::U64(45_000),
),
];
for (key, value) in entries {
conn.set(
key,
confers::AnnotatedValue::new(value.clone(), SourceId::default(), *key),
)
.await
.unwrap();
}
let cfg = OxcacheConfig::load_from(&conn).await.unwrap();
assert_eq!(cfg.capacity, 2048);
assert_eq!(cfg.default_ttl_ms, 90_000);
assert_eq!(cfg.tti_ms, Some(15_000));
assert_eq!(cfg.null_cache_ttl_ms, Some(500));
assert_eq!(cfg.ttl_jitter_factor, Some(0.3));
assert_eq!(cfg.sync_mode, Some(true));
assert_eq!(cfg.backend.as_deref(), Some("redis"));
assert_eq!(cfg.metrics_enabled, Some(false));
assert_eq!(cfg.redis_url.as_deref(), Some("redis://cfg:6379"));
assert_eq!(cfg.serialization_format.as_deref(), Some("postcard"));
assert_eq!(cfg.connection_pool_size, Some(16));
assert_eq!(cfg.service_name.as_deref(), Some("r9-confers"));
assert_eq!(cfg.circuit_breaker.failure_threshold, 9);
assert_eq!(cfg.circuit_breaker.recovery_timeout_ms, 45_000);
}
#[tokio::test]
async fn load_from_type_mismatch_fails_loudly() {
let conn = new_in_memory();
let key = "cache.sync_mode";
conn.set(
key,
confers::AnnotatedValue::new(ConfigValue::U64(1), SourceId::default(), key),
)
.await
.unwrap();
let err = OxcacheConfig::load_from(&conn).await.unwrap_err();
assert!(matches!(err, OxCacheError::InvalidInput(m) if m.contains(key)));
}
#[tokio::test]
async fn load_from_threshold_overflow_fails_loudly() {
let conn = new_in_memory();
let key = "cache.circuit_breaker.failure_threshold";
conn.set(
key,
confers::AnnotatedValue::new(
ConfigValue::U64(u32::MAX as u64 + 1),
SourceId::default(),
key,
),
)
.await
.unwrap();
let err = OxcacheConfig::load_from(&conn).await.unwrap_err();
assert!(matches!(err, OxCacheError::InvalidInput(m) if m.contains(key)));
}
#[cfg(target_pointer_width = "32")]
#[tokio::test]
async fn load_from_pool_size_overflow_fails_loudly() {
let conn = new_in_memory();
let key = "cache.connection_pool_size";
conn.set(
key,
confers::AnnotatedValue::new(
ConfigValue::U64(usize::MAX as u64 + 1),
SourceId::default(),
key,
),
)
.await
.unwrap();
let err = OxcacheConfig::load_from(&conn).await.unwrap_err();
assert!(matches!(err, OxCacheError::InvalidInput(m) if m.contains(key)));
}
#[test]
fn debug_redacts_redis_url() {
let cfg = OxcacheConfig {
redis_url: Some("redis://admin:s3cret@host:6379".into()),
..OxcacheConfig::default()
};
let debug = format!("{cfg:?}");
assert!(!debug.contains("s3cret"), "credentials leaked: {debug}");
assert!(
debug.contains("\"***\""),
"redaction marker missing: {debug}"
);
}
#[tokio::test]
async fn try_from_confers_maps_all_fields() {
let cfg = OxcacheConfig {
capacity: 4096,
default_ttl_ms: 60_000,
tti_ms: Some(12_000),
null_cache_ttl_ms: Some(400),
ttl_jitter_factor: Some(0.15),
sync_mode: Some(true),
backend: Some("moka".into()),
metrics_enabled: Some(false),
redis_url: None,
disk_path: Some("/tmp/cfg.redb".into()),
serialization_format: Some("bincode".into()),
connection_pool_size: Some(16),
service_name: Some("r9-confers-unified".into()),
circuit_breaker: CircuitBreakerSettings {
failure_threshold: 7,
recovery_timeout_ms: 25_000,
},
};
let unified = CacheConfig::try_from_confers(&cfg).unwrap();
assert_eq!(unified.capacity, Some(4096));
assert_eq!(unified.ttl, Some(Duration::from_millis(60_000)));
assert_eq!(unified.tti, Some(Duration::from_millis(12_000)));
assert_eq!(unified.null_cache_ttl, Some(Duration::from_millis(400)));
assert_eq!(unified.ttl_jitter_factor, Some(0.15));
assert_eq!(unified.sync_mode, Some(true));
assert_eq!(unified.backend, Some("moka".to_string()));
assert_eq!(unified.metrics_enabled, Some(false));
assert_eq!(unified.disk_path.as_deref(), Some("/tmp/cfg.redb"));
assert_eq!(unified.connection_pool_size, Some(16));
assert_eq!(unified.service_name.as_deref(), Some("r9-confers-unified"));
assert_eq!(unified.serialization_format.as_deref(), Some("bincode"));
assert_eq!(unified.circuit_breaker_failure_threshold, Some(7));
assert_eq!(
unified.circuit_breaker_reset_timeout,
Some(Duration::from_millis(25_000))
);
}
#[tokio::test]
async fn try_from_confers_unknown_backend_fails_loudly() {
let cfg = OxcacheConfig {
backend: Some("memcache".into()),
..OxcacheConfig::default()
};
let unified = CacheConfig::try_from_confers(&cfg).unwrap();
let err = unified.validate().unwrap_err();
assert!(matches!(err, OxCacheError::InvalidInput(m) if m.contains("memcache")));
}
#[tokio::test]
async fn try_from_confers_defaults_map_to_unset() {
let unified = CacheConfig::try_from_confers(&OxcacheConfig::default()).unwrap();
assert_eq!(unified.capacity, Some(10_000));
assert_eq!(unified.ttl, Some(Duration::from_millis(60_000)));
assert_eq!(unified.backend, None);
assert_eq!(unified.sync_mode, None);
assert_eq!(unified.serialization_format, None);
}
#[tokio::test]
async fn load_from_confers_connector() {
let conn = memory_connector_with_capacity(1234).await;
let cfg = OxcacheConfig::load_from(&conn).await.unwrap();
assert_eq!(cfg.capacity, 1234);
assert_eq!(cfg.default_ttl_ms, 120_000);
assert_eq!(cfg.circuit_breaker.failure_threshold, 5);
}
#[tokio::test]
async fn load_from_empty_connector_uses_defaults() {
let conn = new_in_memory();
let cfg = OxcacheConfig::load_from(&conn).await.unwrap();
assert_eq!(cfg, OxcacheConfig::default());
}
#[test]
fn config_serde_roundtrip() {
let cfg = OxcacheConfig {
capacity: 500,
default_ttl_ms: 30_000,
circuit_breaker: CircuitBreakerSettings {
failure_threshold: 3,
recovery_timeout_ms: 10_000,
},
..Default::default()
};
let json = serde_json::to_string(&cfg).unwrap();
let back: OxcacheConfig = serde_json::from_str(&json).unwrap();
assert_eq!(back, cfg);
let partial: OxcacheConfig = serde_json::from_str(r#"{"capacity": 42}"#).unwrap();
assert_eq!(partial.capacity, 42);
assert_eq!(partial.default_ttl_ms, 60_000);
}
#[cfg(feature = "memory")]
#[tokio::test]
async fn config_drives_l1_builder() {
let cfg = OxcacheConfig {
capacity: 777,
..Default::default()
};
let l1 = crate::cache::L1Builder::from(&cfg);
let backend = l1.build();
assert_eq!(backend.capacity().await.unwrap(), 777);
}
struct RecordingListener {
calls: AtomicU64,
last_capacity: AtomicU64,
}
impl RecordingListener {
fn new() -> Self {
Self {
calls: AtomicU64::new(0),
last_capacity: AtomicU64::new(0),
}
}
}
impl ConfigChangeListener for RecordingListener {
fn on_change(&self, _old: &OxcacheConfig, new: &OxcacheConfig) {
self.calls.fetch_add(1, Ordering::SeqCst);
self.last_capacity.store(new.capacity, Ordering::SeqCst);
}
}
#[tokio::test]
async fn watch_hot_reloads_on_confers_event() {
let connector = Arc::new(memory_connector_with_capacity(100).await);
connector
.set(
"cache.capacity",
confers::AnnotatedValue::new(
ConfigValue::U64(9999),
SourceId::default(),
"cache.capacity",
),
)
.await
.unwrap();
let bus = Arc::new(InMemoryBus::new());
let source = Arc::new(ConfersConfigSource::new(connector));
let initial = source.load().await.unwrap();
assert_eq!(initial.capacity, 9999, "重载前应读到新值");
let listener = Arc::new(RecordingListener::new());
let watcher = Arc::new(
ConfersConfigWatcher::new(OxcacheConfig {
capacity: 100,
..Default::default()
})
.on_change(listener.clone()),
);
assert_eq!(watcher.snapshot().get().capacity, 100);
let handle = watcher.watch(bus.clone(), source.clone());
tokio::time::sleep(Duration::from_millis(50)).await;
bus.publish(ConfigChangeEvent::new(
"test",
"unit-test",
vec!["cache.capacity".to_string()],
"",
))
.await
.unwrap();
for _ in 0..50 {
if watcher.snapshot().get().capacity == 9999 {
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
assert_eq!(
watcher.snapshot().get().capacity,
9999,
"热更新后快照应换装"
);
assert_eq!(
listener.calls.load(Ordering::SeqCst),
1,
"监听器应被回调 1 次"
);
assert_eq!(listener.last_capacity.load(Ordering::SeqCst), 9999);
handle.abort();
}
#[tokio::test]
async fn watch_ignores_non_cache_events() {
let connector = Arc::new(memory_connector_with_capacity(100).await);
let bus = Arc::new(InMemoryBus::new());
let source = Arc::new(ConfersConfigSource::new(connector));
let watcher = Arc::new(ConfersConfigWatcher::new(OxcacheConfig {
capacity: 100,
..Default::default()
}));
let handle = watcher.watch(bus.clone(), source);
tokio::time::sleep(Duration::from_millis(50)).await;
bus.publish(ConfigChangeEvent::new(
"test",
"unit-test",
vec!["other.key".to_string()],
"",
))
.await
.unwrap();
tokio::time::sleep(Duration::from_millis(100)).await;
assert_eq!(watcher.snapshot().get().capacity, 100, "无关事件不应换装");
handle.abort();
}
#[tokio::test]
#[serial_test::serial]
async fn watch_keeps_snapshot_when_reload_fails() {
struct FlakySource {
fail: AtomicBool,
}
#[async_trait::async_trait]
impl CacheConfigSource for FlakySource {
async fn load(&self) -> OxCacheResult<OxcacheConfig> {
if self.fail.load(Ordering::SeqCst) {
Err(OxCacheError::Operation("injected reload failure".into()))
} else {
Ok(OxcacheConfig {
capacity: 4321,
..OxcacheConfig::default()
})
}
}
}
let source = Arc::new(FlakySource {
fail: AtomicBool::new(false),
});
let bus = Arc::new(InMemoryBus::new());
let watcher = Arc::new(ConfersConfigWatcher::new(OxcacheConfig {
capacity: 100,
..OxcacheConfig::default()
}));
let handle = watcher.watch(bus.clone(), source.clone());
tokio::time::sleep(Duration::from_millis(50)).await;
bus.publish(ConfigChangeEvent::new(
"test",
"unit-test",
vec!["cache.capacity".to_string()],
"",
))
.await
.unwrap();
for _ in 0..50 {
if watcher.snapshot().get().capacity == 4321 {
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
assert_eq!(watcher.snapshot().get().capacity, 4321);
#[cfg(feature = "metrics")]
let rejected_before = dynamic_counter(
&crate::infra::metrics::unified::GLOBAL_UNIFIED_METRICS,
"oxcache_config_reload_rejected_total",
);
source.fail.store(true, Ordering::SeqCst);
bus.publish(ConfigChangeEvent::new(
"test",
"unit-test",
vec!["cache.capacity".to_string()],
"",
))
.await
.unwrap();
#[cfg(feature = "metrics")]
{
let deadline = std::time::Instant::now() + Duration::from_secs(2);
while dynamic_counter(
&crate::infra::metrics::unified::GLOBAL_UNIFIED_METRICS,
"oxcache_config_reload_rejected_total",
) <= rejected_before
&& std::time::Instant::now() < deadline
{
tokio::time::sleep(Duration::from_millis(10)).await;
}
assert!(
dynamic_counter(
&crate::infra::metrics::unified::GLOBAL_UNIFIED_METRICS,
"oxcache_config_reload_rejected_total",
) > rejected_before,
"reload rejection must increment oxcache_config_reload_rejected_total"
);
}
#[cfg(not(feature = "metrics"))]
tokio::time::sleep(Duration::from_millis(150)).await;
assert_eq!(
watcher.snapshot().get().capacity,
4321,
"重载失败应保留旧快照"
);
assert!(!handle.is_finished(), "watch 任务在重载失败后必须存活");
handle.abort();
}
}