oxcache 0.5.0-rc.7

A production-grade multi-level cache library for Rust: L1 memory (Moka/DashMap) + L2 distributed (Redis/Valkey/Dragonfly/Aerospike) + optional L3 disk (redb).
// Copyright (c) 2025-2026 Kirky.X🌠
// SPDX-License-Identifier: MIT
//! 写路径失效广播装饰器
//!
//! [`InvalidatingBackend`] 包装任意 [`CacheBackend`]:写操作(set/delete/
//! set_many/delete_many/clear)成功后经 [`InvalidationBus`] 广播失效事件,
//! 其他实例的监听任务失效各自本地 L1。发布失败不影响写结果(fire-and-forget:
//! 底层写入已成功,广播为尽力而为)。
//!
//! `expire` 默认仅透传(TTL 变更不改变值的一致性);需要 TTL 变更也触发
//! 失效(如心跳/续期键形)时,用
//! [`with_expire_broadcast`](InvalidatingBackend::with_expire_broadcast)
//! 开启(默认 `false`,关闭时与不开启该装饰器的既有行为逐位一致)。

use super::InvalidationBus;
use crate::backend::CacheBackend;
use crate::backend::interface::{BackendKind, CacheSetItem};
use crate::error::OxCacheResult;
use async_trait::async_trait;
use std::collections::HashMap;
use std::sync::Arc;
use std::time::Duration;

/// 写路径失效广播装饰器
pub struct InvalidatingBackend {
    inner: Arc<dyn CacheBackend>,
    bus: Arc<InvalidationBus>,
    /// expire 广播开关(默认 `false`:TTL 变更不改变值的一致性,不广播)
    expire_broadcast: bool,
}

impl InvalidatingBackend {
    /// 包装内部后端与失效总线
    ///
    /// expire 默认不广播;需要 TTL 变更也触发失效(如心跳/续期键形)时,
    /// 用 [`Self::with_expire_broadcast`] 开启。
    pub fn new(inner: Arc<dyn CacheBackend>, bus: Arc<InvalidationBus>) -> Self {
        Self {
            inner,
            bus,
            expire_broadcast: false,
        }
    }

    /// 开启 expire 广播(默认关闭)。
    ///
    /// 开启后,`expire` 调用成功即广播该键的失效事件(与 `delete` 的
    /// 语义一致,不区分键是否存在);`false`(默认)时 `expire` 仅透传。
    ///
    /// **流量提示**:开启后每次 expire 都多一次 PUBLISH,高频键形
    /// (心跳/续期类)会放大 Pub/Sub 流量,请按命名空间评估限流后再开启。
    pub fn with_expire_broadcast(mut self) -> Self {
        self.expire_broadcast = true;
        self
    }

    /// 内部后端
    pub fn inner(&self) -> &Arc<dyn CacheBackend> {
        &self.inner
    }
}

#[async_trait]
impl crate::backend::CacheReader for InvalidatingBackend {
    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_trait]
impl crate::backend::CacheWriter for InvalidatingBackend {
    async fn set(
        &self,
        key: Arc<str>,
        value: Arc<Vec<u8>>,
        ttl: Option<Duration>,
    ) -> OxCacheResult<()> {
        self.inner.set(key.clone(), value, ttl).await?;
        // 写已成功,广播尽力而为:发布失败不回滚写
        let _ = self.bus.invalidate_key(&key).await;
        Ok(())
    }

    async fn delete(&self, key: &str) -> OxCacheResult<()> {
        self.inner.delete(key).await?;
        let _ = self.bus.invalidate_key(key).await;
        Ok(())
    }

    async fn clear(&self) -> OxCacheResult<()> {
        self.inner.clear().await?;
        let _ = self.bus.invalidate_namespace("*").await;
        Ok(())
    }

    async fn expire(&self, key: &str, ttl: Duration) -> OxCacheResult<bool> {
        let result = self.inner.expire(key, ttl).await?;
        // 默认关闭:TTL 变更不改变值的一致性,仅透传;开启后广播尽力而为
        if self.expire_broadcast {
            let _ = self.bus.invalidate_key(key).await;
        }
        Ok(result)
    }

    async fn set_many(&self, items: &[CacheSetItem]) -> OxCacheResult<()> {
        self.inner.set_many(items).await?;
        for (key, _, _) in items {
            let _ = self.bus.invalidate_key(key).await;
        }
        Ok(())
    }

    async fn delete_many(&self, keys: &[String]) -> OxCacheResult<()> {
        self.inner.delete_many(keys).await?;
        for key in keys {
            let _ = self.bus.invalidate_key(key).await;
        }
        Ok(())
    }
}

#[async_trait]
impl crate::backend::CacheConnector for InvalidatingBackend {
    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()
    }
}

// `CacheBackend` 由 blanket impl 自动提供(Reader+Writer+Connector 均已实现),
// 不可再显式 impl(会与 blanket impl 冲突)。

#[cfg(test)]
mod tests {
    use super::*;
    use crate::backend::interface::{CacheConnector, CacheReader, CacheWriter};
    use crate::backend::{CacheBackend, MockBackend};
    use crate::features::invalidation::{
        DEFAULT_CHANNEL, InMemoryPubSubTransport, InvalidationConfig,
    };

    async fn setup() -> (
        Arc<InvalidatingBackend>,
        Arc<dyn CacheBackend>,
        Arc<dyn CacheBackend>,
    ) {
        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 inner_a: Arc<dyn CacheBackend> = Arc::new(MockBackend::new("mock-a", 100, false));
        let inner_b: Arc<dyn CacheBackend> = Arc::new(MockBackend::new("mock-b", 100, false));

        let decorated = Arc::new(InvalidatingBackend::new(inner_a.clone(), bus_a));
        // B 实例监听总线,失效自己的 L1
        let _handle = bus_b.spawn_listener(inner_b.clone()).await.unwrap();
        tokio::time::sleep(Duration::from_millis(50)).await;
        (decorated, inner_a, inner_b)
    }

    #[tokio::test]
    async fn set_publishes_invalidation_to_other_instances() {
        let (decorated, _inner_a, inner_b) = setup().await;

        // 预先在 B 放一个旧值,模拟 B 上一次读回填
        inner_b
            .set(Arc::from("user:1"), Arc::new(b"stale".to_vec()), None)
            .await
            .unwrap();
        assert!(inner_b.exists("user:1").await.unwrap());

        // A 经装饰器写入 → 广播 → B 失效
        decorated
            .set(Arc::from("user:1"), Arc::new(b"fresh".to_vec()), None)
            .await
            .unwrap();

        for _ in 0..50 {
            if !inner_b.exists("user:1").await.unwrap() {
                break;
            }
            tokio::time::sleep(Duration::from_millis(10)).await;
        }
        assert!(
            !inner_b.exists("user:1").await.unwrap(),
            "装饰器 set 后 B 实例的旧条目应被失效"
        );
    }

    #[tokio::test]
    async fn delete_publishes_invalidation_to_other_instances() {
        let (decorated, _inner_a, inner_b) = setup().await;

        inner_b
            .set(Arc::from("user:2"), Arc::new(b"stale".to_vec()), None)
            .await
            .unwrap();

        decorated.delete("user:2").await.unwrap();

        for _ in 0..50 {
            if !inner_b.exists("user:2").await.unwrap() {
                break;
            }
            tokio::time::sleep(Duration::from_millis(10)).await;
        }
        assert!(!inner_b.exists("user:2").await.unwrap());
    }

    #[tokio::test]
    async fn read_path_is_passthrough() {
        let (decorated, inner_a, _inner_b) = setup().await;
        inner_a
            .set(Arc::from("k"), Arc::new(b"v".to_vec()), None)
            .await
            .unwrap();
        assert_eq!(decorated.get("k").await.unwrap(), Some(b"v".to_vec()));
        assert!(decorated.exists("k").await.unwrap());
        assert_eq!(
            decorated.backend_kind(),
            inner_a.backend_kind(),
            "装饰器透传 backend_kind"
        );
    }

    #[tokio::test]
    async fn publish_failure_does_not_fail_write() {
        // 独立 transport 无订阅者:publish 为 no-op,写必须照常成功
        let transport = Arc::new(InMemoryPubSubTransport::new());
        let bus = Arc::new(InvalidationBus::new(
            transport,
            InvalidationConfig::new("solo"),
        ));
        let inner: Arc<dyn CacheBackend> = Arc::new(MockBackend::new("mock", 100, false));
        let decorated = InvalidatingBackend::new(inner.clone(), bus);

        decorated
            .set(Arc::from("k"), Arc::new(b"v".to_vec()), None)
            .await
            .unwrap();
        assert_eq!(inner.get("k").await.unwrap(), Some(b"v".to_vec()));
    }

    // ========================================================================
    // expire 广播开关(默认关)
    // ========================================================================

    /// 双实例总线:bus_b 监听并失效 remote 上的条目
    async fn setup_two_instances() -> (
        Arc<InvalidatingBackend>,
        Arc<InvalidatingBackend>,
        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 inner_a: Arc<dyn CacheBackend> = Arc::new(MockBackend::new("mock-a", 100, false));
        let remote: Arc<dyn CacheBackend> = Arc::new(MockBackend::new("remote", 100, false));
        let handle = bus_b.spawn_listener(remote.clone()).await.unwrap();
        tokio::time::sleep(Duration::from_millis(50)).await;
        let plain = Arc::new(InvalidatingBackend::new(inner_a.clone(), bus_a.clone()));
        let broadcasting =
            Arc::new(InvalidatingBackend::new(inner_a, bus_a).with_expire_broadcast());
        (plain, broadcasting, remote, handle)
    }

    /// 默认关闭:expire 透传但不得广播——与既有行为逐位一致
    #[tokio::test]
    async fn expire_broadcast_disabled_by_default_keeps_passthrough() {
        let (plain, _broadcasting, remote, _handle) = setup_two_instances().await;

        // 键经 inner() 直写内层(不经装饰器,不产生广播)
        plain
            .inner()
            .set(Arc::from("user:1"), Arc::new(b"v".to_vec()), None)
            .await
            .unwrap();
        remote
            .set(Arc::from("user:1"), Arc::new(b"stale".to_vec()), None)
            .await
            .unwrap();

        // expire 在内层成功(键存在),但默认开关下不得广播
        assert!(
            plain
                .expire("user:1", Duration::from_secs(30))
                .await
                .unwrap()
        );

        // 否定断言:窗口内轮询确认远端条目始终存在,任何时点消失即失败
        for _ in 0..15 {
            assert!(
                remote.exists("user:1").await.unwrap(),
                "默认关闭时 expire 不得广播失效(与既有行为一致)"
            );
            tokio::time::sleep(Duration::from_millis(10)).await;
        }
    }

    /// 开关开启:expire 成功后广播失效,其他实例的条目被删除
    #[tokio::test]
    async fn expire_broadcast_enabled_notifies_other_instances() {
        let (_plain, broadcasting, remote, _handle) = setup_two_instances().await;

        // 键经 inner() 直写内层(不经装饰器,不产生广播)
        broadcasting
            .inner()
            .set(Arc::from("user:2"), Arc::new(b"v".to_vec()), None)
            .await
            .unwrap();
        remote
            .set(Arc::from("user:2"), Arc::new(b"stale".to_vec()), None)
            .await
            .unwrap();

        assert!(
            broadcasting
                .expire("user:2", Duration::from_secs(30))
                .await
                .unwrap()
        );

        for _ in 0..50 {
            if !remote.exists("user:2").await.unwrap() {
                return;
            }
            tokio::time::sleep(Duration::from_millis(10)).await;
        }
        panic!("开关开启后 expire 应广播失效远端条目");
    }
}