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
//! Cache 字节操作方法(用于宏兼容)

use super::Cache;
use crate::error::OxCacheResult;
use crate::traits::CacheKey;
use std::sync::Arc;
use std::time::Duration;

#[cfg(any(feature = "serialization", feature = "full"))]
use crate::infra::Serializer;

impl<K, V> Cache<K, V>
where
    K: CacheKey,
    V: serde::Serialize + for<'de> serde::Deserialize<'de>,
{
    pub async fn get_bytes(&self, key: &str) -> OxCacheResult<Option<Vec<u8>>> {
        #[cfg(feature = "metrics")]
        let __start = std::time::Instant::now();
        let result = self.backend.get(key).await;
        #[cfg(feature = "metrics")]
        {
            let latency = __start.elapsed();
            let layer = self.metrics_layer();
            match &result {
                Ok(Some(_)) => self.metrics.record_hit(layer, latency),
                Ok(None) => self.metrics.record_miss(layer, latency),
                Err(_) => {}
            }
            self.record_backend_op();
        }
        result
    }

    pub async fn set_bytes(
        &self,
        key: &str,
        value: Vec<u8>,
        ttl: Option<u64>,
    ) -> OxCacheResult<()> {
        let ttl_duration = ttl.map(Duration::from_secs);
        #[cfg(feature = "metrics")]
        let __start = std::time::Instant::now();
        let result = self
            .backend
            .set(Arc::from(key), Arc::new(value), ttl_duration)
            .await;
        #[cfg(feature = "metrics")]
        {
            self.metrics
                .record_set(self.metrics_layer(), __start.elapsed());
            self.record_backend_op();
        }
        result
    }

    /// Synchronously get raw bytes from the cache (macro-compatible sync path).
    ///
    /// Returns `Err(NotSupported)` if `sync_mode(true)` was not set on the
    /// builder (i.e., `backend_sync` is `None`).
    pub fn get_bytes_sync(&self, key: &str) -> OxCacheResult<Option<Vec<u8>>> {
        let backend = self
            .backend_sync
            .as_ref()
            .ok_or_else(Self::sync_mode_error)?;
        #[cfg(feature = "metrics")]
        let __start = std::time::Instant::now();
        let result = backend.get(key);
        #[cfg(feature = "metrics")]
        {
            let latency = __start.elapsed();
            let layer = self.metrics_layer();
            match &result {
                Ok(Some(_)) => self.metrics.record_hit(layer, latency),
                Ok(None) => self.metrics.record_miss(layer, latency),
                Err(_) => {}
            }
            self.record_backend_op();
        }
        result
    }

    /// Synchronously set raw bytes in the cache (macro-compatible sync path).
    ///
    /// `ttl` is in seconds (matching the async `set_bytes` signature for
    /// macro symmetry). Returns `Err(NotSupported)` if `sync_mode(true)` was
    /// not set on the builder.
    pub fn set_bytes_sync(&self, key: &str, value: Vec<u8>, ttl: Option<u64>) -> OxCacheResult<()> {
        let backend = self
            .backend_sync
            .as_ref()
            .ok_or_else(Self::sync_mode_error)?;
        let ttl_duration = ttl.map(Duration::from_secs);
        #[cfg(feature = "metrics")]
        let __start = std::time::Instant::now();
        let result = backend.set(Arc::from(key), Arc::new(value), ttl_duration);
        #[cfg(feature = "metrics")]
        {
            self.metrics
                .record_set(self.metrics_layer(), __start.elapsed());
            self.record_backend_op();
        }
        result
    }

    #[cfg(any(feature = "serialization", feature = "full"))]
    pub fn serializer(&self) -> Arc<dyn Serializer> {
        self.serializer.clone()
    }

    // unified_serializer() 仅在 serialization/full feature 下可用
    #[cfg(any(feature = "serialization", feature = "full"))]
    pub fn unified_serializer(&self) -> crate::infra::UnifiedSerializer {
        self.unified_serializer.clone()
    }
}

#[cfg(test)]
mod tests {
    use crate::cache::Cache;

    // ========================================================================
    // get_bytes tests
    // ========================================================================

    #[tokio::test]
    async fn test_get_bytes_returns_none_for_missing_key() {
        let cache: Cache<String, Vec<u8>> = Cache::memory().await.unwrap();
        let result = cache.get_bytes("nonexistent_key").await.unwrap();
        assert!(result.is_none());
    }

    #[tokio::test]
    async fn test_get_bytes_returns_stored_value() {
        let cache: Cache<String, Vec<u8>> = Cache::memory().await.unwrap();
        let data = vec![1, 2, 3, 4, 5];
        cache
            .set_bytes("test_key", data.clone(), None)
            .await
            .unwrap();
        let result = cache.get_bytes("test_key").await.unwrap();
        assert_eq!(result, Some(data));
    }

    // ========================================================================
    // set_bytes tests
    // ========================================================================

    #[tokio::test]
    async fn test_set_bytes_without_ttl() {
        let cache: Cache<String, Vec<u8>> = Cache::memory().await.unwrap();
        let data = b"hello world".to_vec();
        let result = cache.set_bytes("hello", data.clone(), None).await;
        assert!(result.is_ok());
        let stored = cache.get_bytes("hello").await.unwrap();
        assert_eq!(stored, Some(data));
    }

    #[tokio::test]
    async fn test_set_bytes_with_ttl() {
        let cache: Cache<String, Vec<u8>> = Cache::memory().await.unwrap();
        let data = b"expiring".to_vec();
        let result = cache.set_bytes("temp", data.clone(), Some(3600)).await;
        assert!(result.is_ok());
        let stored = cache.get_bytes("temp").await.unwrap();
        assert_eq!(stored, Some(data));
    }

    #[tokio::test]
    async fn test_set_bytes_overwrites_existing() {
        let cache: Cache<String, Vec<u8>> = Cache::memory().await.unwrap();
        cache.set_bytes("counter", vec![1], None).await.unwrap();
        cache.set_bytes("counter", vec![2], None).await.unwrap();
        let result = cache.get_bytes("counter").await.unwrap();
        assert_eq!(result, Some(vec![2]));
    }

    #[tokio::test]
    async fn test_set_bytes_empty_value() {
        let cache: Cache<String, Vec<u8>> = Cache::memory().await.unwrap();
        let empty_data: Vec<u8> = vec![];
        cache
            .set_bytes("empty", empty_data.clone(), None)
            .await
            .unwrap();
        let result = cache.get_bytes("empty").await.unwrap();
        assert_eq!(result, Some(empty_data));
    }

    #[tokio::test]
    #[cfg(any(feature = "serialization", feature = "full"))]
    async fn test_unified_serializer_accessible() {
        let cache: Cache<String, Vec<u8>> = Cache::memory().await.unwrap();
        let serializer = cache.unified_serializer();
        // Verify it's a valid UnifiedSerializer by serializing and deserializing
        let original = "hello";
        let serialized = serializer.serialize(&original).unwrap();
        let deserialized: String = serializer.deserialize(&serialized).unwrap();
        assert_eq!(deserialized, original);
    }

    // ========================================================================
    // serializer() feature-gated test
    // ========================================================================

    #[tokio::test]
    #[cfg(any(feature = "serialization", feature = "full"))]
    async fn test_serializer_returns_json() {
        let cache: Cache<String, Vec<u8>> = Cache::memory().await.unwrap();
        let serializer = cache.serializer();
        // JsonSerializer should be returned - serialize then deserialize to verify
        let original = b"hello world";
        let serialized = serializer.serialize("bytes", original).unwrap();
        let deserialized = serializer.deserialize("bytes", &serialized).unwrap();
        assert_eq!(deserialized, original);
    }

    // ========================================================================
    // Large data tests
    // ========================================================================

    #[tokio::test]
    async fn test_set_bytes_large_data() {
        let cache: Cache<String, Vec<u8>> = Cache::memory().await.unwrap();
        let large_data = vec![0xAB; 1024 * 100]; // 100KB
        cache
            .set_bytes("large", large_data.clone(), None)
            .await
            .unwrap();
        let result = cache.get_bytes("large").await.unwrap();
        assert_eq!(result, Some(large_data));
    }

    // ========================================================================
    // get_bytes_sync / set_bytes_sync tests
    // ========================================================================

    // NOTE: multi_thread flavor required — MokaMemoryBackend's sync_block_on
    // uses block_in_place, which panics on current_thread runtimes.
    #[tokio::test(flavor = "multi_thread")]
    async fn test_get_bytes_sync_returns_none_for_missing_key() {
        let cache: Cache<String, Vec<u8>> = Cache::builder().sync_mode(true).build().await.unwrap();
        let result = cache.get_bytes_sync("nonexistent_key").unwrap();
        assert!(result.is_none(), "missing key should return None");
    }

    #[tokio::test(flavor = "multi_thread")]
    async fn test_set_bytes_sync_then_get_bytes_sync_roundtrip() {
        let cache: Cache<String, Vec<u8>> = Cache::builder().sync_mode(true).build().await.unwrap();
        let data = vec![1, 2, 3, 4, 5];
        cache
            .set_bytes_sync("test_key", data.clone(), None)
            .unwrap();
        let result = cache.get_bytes_sync("test_key").unwrap();
        assert_eq!(result, Some(data));
    }

    #[tokio::test(flavor = "multi_thread")]
    async fn test_set_bytes_sync_with_ttl_roundtrip() {
        let cache: Cache<String, Vec<u8>> = Cache::builder().sync_mode(true).build().await.unwrap();
        let data = b"expiring".to_vec();
        cache
            .set_bytes_sync("temp", data.clone(), Some(3600))
            .unwrap();
        let result = cache.get_bytes_sync("temp").unwrap();
        assert_eq!(result, Some(data));
    }

    #[tokio::test(flavor = "multi_thread")]
    async fn test_get_bytes_sync_without_sync_mode_returns_not_supported() {
        // Default builder (sync_mode=false) should return Err(NotSupported)
        let cache: Cache<String, Vec<u8>> = Cache::builder().build().await.unwrap();
        let result = cache.get_bytes_sync("any_key");
        assert!(
            matches!(result, Err(crate::error::OxCacheError::NotSupported(_))),
            "expected Err(NotSupported) when sync_mode is false, got {:?}",
            result
        );
    }

    #[tokio::test(flavor = "multi_thread")]
    async fn test_set_bytes_sync_without_sync_mode_returns_not_supported() {
        let cache: Cache<String, Vec<u8>> = Cache::builder().build().await.unwrap();
        let result = cache.set_bytes_sync("any_key", vec![1], None);
        assert!(
            matches!(result, Err(crate::error::OxCacheError::NotSupported(_))),
            "expected Err(NotSupported) when sync_mode is false, got {:?}",
            result
        );
    }

    // ========================================================================
    // 字节 API 指标埋点
    // ========================================================================

    /// 字节 API 与泛型路径同口径:get_bytes hit/miss、set_bytes 计入 unified。
    /// 全局计数器与 metrics 重置类测试互斥执行(serial 组),避免 reset 竞态。
    #[cfg(feature = "metrics")]
    #[tokio::test]
    #[serial_test::serial]
    async fn bytes_api_records_unified_metrics() {
        let before = crate::infra::GLOBAL_UNIFIED_METRICS.get_counters();

        let cache: Cache<String, Vec<u8>> = Cache::memory().await.unwrap();
        cache
            .set_bytes("metrics-bytes", vec![1, 2], None)
            .await
            .unwrap();
        let _ = cache.get_bytes("metrics-bytes").await.unwrap(); // hit
        let _ = cache.get_bytes("metrics-bytes-miss").await.unwrap(); // miss

        let after = crate::infra::GLOBAL_UNIFIED_METRICS.get_counters();
        assert!(after.l1_sets > before.l1_sets, "set_bytes must record sets");
        assert!(after.l1_hits > before.l1_hits, "get_bytes hit must record");
        assert!(
            after.l1_misses > before.l1_misses,
            "get_bytes miss must record"
        );
    }

    /// sync 字节 API 同样埋点。
    #[cfg(feature = "metrics")]
    #[tokio::test(flavor = "multi_thread")]
    #[serial_test::serial]
    async fn sync_bytes_api_records_unified_metrics() {
        let before = crate::infra::GLOBAL_UNIFIED_METRICS.get_counters();

        let cache: Cache<String, Vec<u8>> = Cache::builder().sync_mode(true).build().await.unwrap();
        cache
            .set_bytes_sync("metrics-bytes-sync", vec![9], None)
            .unwrap();
        let _ = cache.get_bytes_sync("metrics-bytes-sync").unwrap();

        let after = crate::infra::GLOBAL_UNIFIED_METRICS.get_counters();
        assert!(after.l1_sets > before.l1_sets, "set_bytes_sync must record");
        assert!(
            after.l1_hits > before.l1_hits,
            "get_bytes_sync hit must record"
        );
    }
}