Skip to main content

sz_orm_core/
l2_cache.rs

1//! L2 二级缓存(Level-2 Cache)
2//!
3//! 对应文档 6.8 节改进项 21(L2 二级缓存)。
4//!
5//! # 核心概念
6//!
7//! - **L2Cache**:跨 Session 共享的二级缓存(与 Hibernate L2 Cache / MyBatis 二级缓存对应)
8//! - **CacheKey**:统一缓存键构造(table + pk 或 table + query_hash)
9//! - **L2CacheStats**:命中率统计(hits/misses/evictions/sets)
10//! - **表级失效**:`invalidate_table(table)` 一次失效某表的所有缓存项
11//!
12//! 与 L1 缓存(Session 级别)的区别:
13//! - L1:单次 Session/请求 内有效,事务结束自动清空
14//! - L2:跨 Session 共享,进程级缓存,需显式失效
15//!
16//! # 设计灵感
17//!
18//! - Hibernate L2 Cache(`@Cache` / `@Cacheable`)
19//! - MyBatis 二级缓存(`<cache>` 标签)
20//! - Rails `Rails.cache`
21//! - Django cache framework
22//!
23//! # 使用示例
24//!
25//! ```no_run
26//! use sz_orm_core::l2_cache::{L2Cache, CacheKey};
27//! use sz_orm_core::Value;
28//!
29//! // 1. 创建 L2 缓存
30//! let cache = L2Cache::new();
31//!
32//! // 2. 缓存单行(pk 维度)
33//! let key = CacheKey::by_pk("users", 1);
34//! cache.put(&key, Value::String("Alice".to_string()), None);
35//!
36//! // 3. 读取
37//! let val = cache.get(&key);
38//! assert!(val.is_some());
39//!
40//! // 4. 表级失效(用户表更新后)
41//! cache.invalidate_table("users");
42//! assert!(cache.get(&key).is_none());
43//!
44//! // 5. 命中率统计
45//! let stats = cache.stats();
46//! println!("hit rate: {:.2}%", stats.hit_rate() * 100.0);
47//! ```
48
49use crate::cache::Cache;
50use crate::error::CacheError;
51use crate::value::Value;
52use std::collections::HashMap;
53use std::future::Future;
54use std::pin::Pin;
55use std::sync::{Arc, RwLock};
56use std::time::Duration;
57// #72 修复:使用 tokio::time::Instant 替代 std::time::Instant
58// tokio::time::Instant 支持 tokio::time::pause() 测试辅助,
59// 允许测试在不真实睡眠的情况下控制时间流逝。
60// 在非测试环境(未调用 pause)下,行为与 std::time::Instant 完全一致。
61use tokio::time::Instant;
62
63// ============================================================================
64// InvalidationBus — 缓存失效消息总线(跨实例失效)
65// ============================================================================
66
67/// 缓存失效消息
68#[derive(Debug, Clone)]
69pub enum InvalidationMessage {
70    /// 失效单个 key
71    InvalidateKey(String),
72    /// 失效整张表
73    InvalidateTable(String),
74    /// 失效所有缓存
75    InvalidateAll,
76}
77
78/// 缓存失效总线 trait
79///
80/// 用于跨实例缓存失效:当一个实例失效了某张表的缓存时,
81/// 通过总线通知其他订阅者同步失效。
82pub trait InvalidationBus: Send + Sync {
83    /// 发布失效消息
84    fn publish(&self, message: InvalidationMessage);
85    /// 订阅失效消息(返回一个迭代器,drain 当前缓冲的消息)
86    fn subscribe(&self) -> Box<dyn Iterator<Item = InvalidationMessage> + Send>;
87}
88
89/// 进程内失效总线(单实例用)
90///
91/// 基于 `tokio::sync::broadcast` 实现多订阅者广播。
92/// `subscribe()` 返回的迭代器会 drain 当前已缓冲但未消费的消息。
93pub struct LocalInvalidationBus {
94    tx: tokio::sync::broadcast::Sender<InvalidationMessage>,
95}
96
97impl LocalInvalidationBus {
98    /// 创建进程内失效总线,`capacity` 为广播缓冲区容量
99    pub fn new(capacity: usize) -> Self {
100        let (tx, _rx) = tokio::sync::broadcast::channel(capacity.max(1));
101        Self { tx }
102    }
103}
104
105impl Default for LocalInvalidationBus {
106    fn default() -> Self {
107        Self::new(256)
108    }
109}
110
111impl InvalidationBus for LocalInvalidationBus {
112    fn publish(&self, message: InvalidationMessage) {
113        // 忽略无订阅者的错误
114        let _ = self.tx.send(message);
115    }
116
117    fn subscribe(&self) -> Box<dyn Iterator<Item = InvalidationMessage> + Send> {
118        let mut rx = self.tx.subscribe();
119        Box::new(std::iter::from_fn(move || loop {
120            match rx.try_recv() {
121                Ok(msg) => return Some(msg),
122                // 缓冲区为空或通道已关闭 → 终止迭代
123                Err(tokio::sync::broadcast::error::TryRecvError::Empty)
124                | Err(tokio::sync::broadcast::error::TryRecvError::Closed) => return None,
125                // 滞后(订阅者落后太多)→ 跳过丢失的消息,继续读下一条
126                Err(tokio::sync::broadcast::error::TryRecvError::Lagged(_)) => continue,
127            }
128        }))
129    }
130}
131
132// ============================================================================
133// CacheKey — 统一缓存键
134// ============================================================================
135
136/// 统一缓存键
137///
138/// 通过 `table` + `kind` + `identifier` 三元组唯一标识一个缓存项:
139/// - `table`:表名(用于表级失效)
140/// - `kind`:缓存类型(ByPk / ByQuery / ByRelation)
141/// - `identifier`:具体标识(pk 值 / 查询哈希 / 关联键)
142#[derive(Debug, Clone, PartialEq, Eq, Hash)]
143pub struct CacheKey {
144    /// 表名
145    pub table: String,
146    /// 缓存类型
147    pub kind: CacheKeyKind,
148    /// 具体标识
149    pub identifier: String,
150}
151
152/// 缓存键类型
153#[derive(Debug, Clone, PartialEq, Eq, Hash)]
154pub enum CacheKeyKind {
155    /// 按主键缓存
156    ByPk,
157    /// 按查询条件缓存
158    ByQuery,
159    /// 按关联关系缓存
160    ByRelation,
161}
162
163impl CacheKey {
164    /// 构造主键维度的缓存键
165    pub fn by_pk(table: impl Into<String>, pk: impl std::fmt::Display) -> Self {
166        Self {
167            table: table.into(),
168            kind: CacheKeyKind::ByPk,
169            identifier: pk.to_string(),
170        }
171    }
172
173    /// 构造查询维度的缓存键(identifier 通常是 SQL + params 的哈希)
174    pub fn by_query(table: impl Into<String>, query_hash: impl std::fmt::Display) -> Self {
175        Self {
176            table: table.into(),
177            kind: CacheKeyKind::ByQuery,
178            identifier: query_hash.to_string(),
179        }
180    }
181
182    /// 构造关联维度的缓存键
183    pub fn by_relation(table: impl Into<String>, relation: impl std::fmt::Display) -> Self {
184        Self {
185            table: table.into(),
186            kind: CacheKeyKind::ByRelation,
187            identifier: relation.to_string(),
188        }
189    }
190
191    /// 序列化为字符串(用于底层存储键)
192    pub fn to_string_key(&self) -> String {
193        let kind_str = match self.kind {
194            CacheKeyKind::ByPk => "pk",
195            CacheKeyKind::ByQuery => "q",
196            CacheKeyKind::ByRelation => "rel",
197        };
198        format!("l2:{}:{}:{}", self.table, kind_str, self.identifier)
199    }
200}
201
202impl std::fmt::Display for CacheKey {
203    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
204        write!(f, "{}", self.to_string_key())
205    }
206}
207
208// ============================================================================
209// L2CacheStats — 命中率统计
210// ============================================================================
211
212/// L2 缓存命中率统计
213#[derive(Debug, Clone, Default)]
214pub struct L2CacheStats {
215    /// 命中次数
216    pub hits: u64,
217    /// 未命中次数
218    pub misses: u64,
219    /// 设置次数
220    pub sets: u64,
221    /// 失效次数(含单键和表级失效)
222    pub evictions: u64,
223    /// 当前缓存项数量
224    pub size: usize,
225}
226
227/// 按表分桶的命中率统计
228///
229/// 用于细粒度观察每张表的缓存命中情况,
230/// 识别"热点表"与"低命中表",指导缓存策略调整。
231#[derive(Debug, Clone, Default)]
232pub struct PerTableStats {
233    /// 命中次数
234    pub hits: u64,
235    /// 未命中次数
236    pub misses: u64,
237    /// 设置次数
238    pub sets: u64,
239    /// 失效次数
240    pub evictions: u64,
241}
242
243impl PerTableStats {
244    /// 总查询次数(hits + misses)
245    pub fn total_lookups(&self) -> u64 {
246        self.hits + self.misses
247    }
248
249    /// 命中率(0.0 ~ 1.0)
250    pub fn hit_rate(&self) -> f64 {
251        let total = self.total_lookups();
252        if total == 0 {
253            0.0
254        } else {
255            self.hits as f64 / total as f64
256        }
257    }
258}
259
260impl L2CacheStats {
261    /// 总查询次数(hits + misses)
262    pub fn total_lookups(&self) -> u64 {
263        self.hits + self.misses
264    }
265
266    /// 命中率(0.0 ~ 1.0)
267    pub fn hit_rate(&self) -> f64 {
268        let total = self.total_lookups();
269        if total == 0 {
270            0.0
271        } else {
272            self.hits as f64 / total as f64
273        }
274    }
275
276    /// 未命中率(0.0 ~ 1.0)
277    pub fn miss_rate(&self) -> f64 {
278        1.0 - self.hit_rate()
279    }
280
281    /// 合并两个统计(用于多分片汇总)
282    pub fn merge(&mut self, other: &L2CacheStats) {
283        self.hits += other.hits;
284        self.misses += other.misses;
285        self.sets += other.sets;
286        self.evictions += other.evictions;
287        self.size += other.size;
288    }
289}
290
291// ============================================================================
292// CacheEntry — 缓存项
293// ============================================================================
294
295/// 缓存项(值 + 过期时间)
296#[derive(Debug, Clone)]
297struct CacheEntry {
298    /// 缓存值
299    value: Value,
300    /// 过期时间(None 表示永不过期)
301    expires_at: Option<Instant>,
302}
303
304impl CacheEntry {
305    fn new(value: Value, ttl: Option<Duration>) -> Self {
306        // Duration::MAX 会导致 Instant::now() + Duration::MAX 溢出
307        // 将其视为永不过期(expires_at = None),与 None 语义一致
308        let expires_at = ttl.and_then(|d| {
309            if d == Duration::MAX {
310                None
311            } else {
312                Some(Instant::now() + d)
313            }
314        });
315        Self { value, expires_at }
316    }
317
318    fn is_expired(&self) -> bool {
319        self.expires_at
320            .map(|t| t <= Instant::now())
321            .unwrap_or(false)
322    }
323}
324
325// ============================================================================
326// LruOrder — O(1) LRU 顺序跟踪器(arena 双向链表 + HashMap)
327// ============================================================================
328
329/// LRU 顺序跟踪器 — 所有操作 O(1)
330///
331/// 基于 arena(Vec<LruNode>)的双向链表 + HashMap 索引实现:
332/// - `touch(key)`:将 key 移到 MRU 端(已存在则摘链+追加,新 key 直接追加)— O(1)
333/// - `remove(key)`:从链表中摘除并回收节点 — O(1)
334/// - `lru_key()`:返回 LRU 端的 key — O(1)
335/// - `iter_keys()`:从 LRU 到 MRU 遍历 — O(n)
336///
337/// 相比 `Vec<String>` + `retain` 方案(touch/remove 为 O(n)),本实现将高频操作
338/// 降为 O(1),仅遍历(用于查找过期 key)保持 O(n)。
339struct LruOrder {
340    /// 节点池(arena):节点索引即数组下标
341    nodes: Vec<LruNode>,
342    /// 空闲节点列表(复用已删除节点的槽位,避免 Vec 无限增长)
343    free_list: Vec<usize>,
344    /// key → 节点索引
345    index: HashMap<String, usize>,
346    /// 链表头(LRU 端,淘汰时从此处取)
347    head: Option<usize>,
348    /// 链表尾(MRU 端,新访问的加入此处)
349    tail: Option<usize>,
350}
351
352/// 双向链表节点
353struct LruNode {
354    key: String,
355    prev: Option<usize>,
356    next: Option<usize>,
357}
358
359impl LruOrder {
360    fn new() -> Self {
361        Self {
362            nodes: Vec::new(),
363            free_list: Vec::new(),
364            index: HashMap::new(),
365            head: None,
366            tail: None,
367        }
368    }
369
370    /// 触碰 key:已存在则移到尾部,不存在则创建并追加到尾部 — O(1)
371    fn touch(&mut self, key: &str) {
372        if let Some(&idx) = self.index.get(key) {
373            self.unlink(idx);
374            self.link_tail(idx);
375        } else {
376            let idx = self.alloc_node(key.to_string());
377            self.link_tail(idx);
378            self.index.insert(key.to_string(), idx);
379        }
380    }
381
382    /// 移除 key — O(1)
383    fn remove(&mut self, key: &str) {
384        if let Some(idx) = self.index.remove(key) {
385            self.unlink(idx);
386            self.free_node(idx);
387        }
388    }
389
390    /// 返回 LRU 端的 key(最久未访问) — O(1)
391    fn lru_key(&self) -> Option<&str> {
392        self.head.map(|idx| self.nodes[idx].key.as_str())
393    }
394
395    /// 从 LRU 到 MRU 遍历所有 key — O(n)
396    fn iter_keys(&self) -> impl Iterator<Item = &str> {
397        LruIter {
398            nodes: &self.nodes,
399            current: self.head,
400        }
401    }
402
403    /// 清空所有 — O(n)(需释放 Vec/HashMap 内存)
404    fn clear(&mut self) {
405        self.nodes.clear();
406        self.free_list.clear();
407        self.index.clear();
408        self.head = None;
409        self.tail = None;
410    }
411
412    /// 当前元素数量 — O(1)
413    #[allow(dead_code)]
414    fn len(&self) -> usize {
415        self.index.len()
416    }
417
418    /// 分配节点(优先复用空闲槽位)
419    fn alloc_node(&mut self, key: String) -> usize {
420        if let Some(idx) = self.free_list.pop() {
421            self.nodes[idx] = LruNode {
422                key,
423                prev: None,
424                next: None,
425            };
426            idx
427        } else {
428            self.nodes.push(LruNode {
429                key,
430                prev: None,
431                next: None,
432            });
433            self.nodes.len() - 1
434        }
435    }
436
437    /// 回收节点到空闲列表
438    fn free_node(&mut self, idx: usize) {
439        self.free_list.push(idx);
440    }
441
442    /// 从链表中摘除节点(仅修改前后指针,不释放节点)
443    fn unlink(&mut self, idx: usize) {
444        let prev = self.nodes[idx].prev;
445        let next = self.nodes[idx].next;
446        match prev {
447            Some(p) => self.nodes[p].next = next,
448            None => self.head = next,
449        }
450        match next {
451            Some(n) => self.nodes[n].prev = prev,
452            None => self.tail = prev,
453        }
454        self.nodes[idx].prev = None;
455        self.nodes[idx].next = None;
456    }
457
458    /// 将节点链接到链表尾部(MRU 端)
459    fn link_tail(&mut self, idx: usize) {
460        match self.tail {
461            Some(t) => {
462                self.nodes[t].next = Some(idx);
463                self.nodes[idx].prev = Some(t);
464            }
465            None => self.head = Some(idx),
466        }
467        self.nodes[idx].next = None;
468        self.tail = Some(idx);
469    }
470}
471
472/// LRU 链表迭代器(从 LRU 端到 MRU 端)
473struct LruIter<'a> {
474    nodes: &'a [LruNode],
475    current: Option<usize>,
476}
477
478impl<'a> Iterator for LruIter<'a> {
479    type Item = &'a str;
480
481    fn next(&mut self) -> Option<Self::Item> {
482        let idx = self.current?;
483        let node = &self.nodes[idx];
484        self.current = node.next;
485        Some(node.key.as_str())
486    }
487}
488
489// ============================================================================
490// L2Cache — 跨 Session 共享的二级缓存
491// ============================================================================
492
493/// L2 二级缓存 — 跨 Session 共享
494///
495/// 线程安全:内部使用 RwLock,可在多线程环境下共享。
496///
497/// # 示例
498///
499/// ```
500/// use sz_orm_core::l2_cache::{L2Cache, CacheKey};
501/// use sz_orm_core::Value;
502/// use std::time::Duration;
503///
504/// let cache = L2Cache::new();
505///
506/// // 缓存单行
507/// let key = CacheKey::by_pk("users", 1);
508/// cache.put(&key, Value::String("Alice".to_string()), None);
509///
510/// // 读取
511/// assert!(cache.get(&key).is_some());
512///
513/// // 表级失效
514/// cache.invalidate_table("users");
515/// assert!(cache.get(&key).is_none());
516/// ```
517pub struct L2Cache {
518    /// 缓存数据
519    data: RwLock<HashMap<String, CacheEntry>>,
520    /// 表名索引(用于表级失效)— table -> Vec<key_string>(去重)
521    table_index: RwLock<HashMap<String, Vec<String>>>,
522    /// LRU 访问顺序跟踪器(O(1) touch/remove/lru_key,arena 双向链表 + HashMap)
523    ///
524    /// # 锁顺序约定
525    ///
526    /// 跨字段持锁时遵循:`data` → `access_order` → `table_index` → `stats`,
527    /// 避免死锁。本字段不允许在持 `data` 写锁时获取其他写锁。
528    access_order: RwLock<LruOrder>,
529    /// 全局统计信息
530    stats: RwLock<L2CacheStats>,
531    /// 按表分桶的统计信息(table -> PerTableStats)
532    table_stats: RwLock<HashMap<String, PerTableStats>>,
533    /// 默认 TTL(`put` 传 `None` 时使用,要"永不失效"请传 `Some(Duration::MAX)`)
534    default_ttl: Option<Duration>,
535    /// 最大容量(LRU 淘汰)
536    max_size: usize,
537    /// 缓存失效总线(可选,用于跨实例失效通知)
538    invalidation_bus: Option<Arc<dyn InvalidationBus>>,
539}
540
541impl Default for L2Cache {
542    fn default() -> Self {
543        Self::new()
544    }
545}
546
547impl L2Cache {
548    /// 创建 L2 缓存(默认容量 10000,无 TTL)
549    pub fn new() -> Self {
550        Self {
551            data: RwLock::new(HashMap::new()),
552            table_index: RwLock::new(HashMap::new()),
553            access_order: RwLock::new(LruOrder::new()),
554            stats: RwLock::new(L2CacheStats::default()),
555            table_stats: RwLock::new(HashMap::new()),
556            default_ttl: None,
557            max_size: 10_000,
558            invalidation_bus: None,
559        }
560    }
561
562    /// 设置默认 TTL
563    pub fn with_default_ttl(mut self, ttl: Duration) -> Self {
564        self.default_ttl = Some(ttl);
565        self
566    }
567
568    /// 设置最大容量
569    pub fn with_max_size(mut self, max_size: usize) -> Self {
570        self.max_size = max_size;
571        self
572    }
573
574    /// 设置缓存失效总线(用于跨实例失效通知)
575    pub fn with_invalidation_bus(mut self, bus: Arc<dyn InvalidationBus>) -> Self {
576        self.invalidation_bus = Some(bus);
577        self
578    }
579
580    /// 存入缓存项
581    ///
582    /// # TTL 语义
583    ///
584    /// - `ttl = Some(d)`:使用 `d` 作为过期时间
585    /// - `ttl = None`:使用 `default_ttl`(若未设置则永不过期)
586    /// - 要显式表示"永不失效",请传 `Some(Duration::MAX)`
587    pub fn put(&self, key: &CacheKey, value: Value, ttl: Option<Duration>) {
588        let actual_ttl = ttl.or(self.default_ttl);
589        let entry = CacheEntry::new(value, actual_ttl);
590        let key_str = key.to_string_key();
591
592        // 1. 写入数据 + LRU 淘汰
593        {
594            let mut data = self.data.write().expect("L2Cache data lock poisoned (put)");
595            let exists = data.contains_key(&key_str);
596            if !exists && data.len() >= self.max_size {
597                // LRU 淘汰:优先淘汰已过期的 key,否则淘汰 LRU 端(access_order 头部)
598                let victim = {
599                    // 不在持 data 写锁时获取 access_order 写锁,先读 access_order
600                    let order = self
601                        .access_order
602                        .read()
603                        .expect("L2Cache access_order lock poisoned (put-victim-read)");
604                    // 优先找已过期的 key(O(n) 遍历,仅缓存满时触发)
605                    // 分两步计算,避免闭包捕获 order 导致生命周期问题
606                    let expired = order
607                        .iter_keys()
608                        .find(|k| data.get(*k).map(|e| e.is_expired()).unwrap_or(false))
609                        .map(|s| s.to_string());
610                    let lru = order.lru_key().map(|s| s.to_string());
611                    expired.or(lru)
612                };
613                if let Some(victim) = victim {
614                    data.remove(&victim);
615                    // 同步清理 access_order(O(1) remove)
616                    let mut order = self
617                        .access_order
618                        .write()
619                        .expect("L2Cache access_order lock poisoned (put-victim-remove)");
620                    order.remove(&victim);
621                }
622            }
623            data.insert(key_str.clone(), entry);
624        };
625
626        // 2. 更新 LRU 访问顺序(O(1) touch:新 key 追加尾部,已存在 key 移到尾部)
627        {
628            let mut order = self
629                .access_order
630                .write()
631                .expect("L2Cache access_order lock poisoned (put-touch)");
632            order.touch(&key_str);
633        }
634
635        // 3. 更新表索引(去重,避免重复 push 导致 invalidate_table 统计错误)
636        {
637            let mut idx = self
638                .table_index
639                .write()
640                .expect("L2Cache table_index lock poisoned (put)");
641            let keys = idx.entry(key.table.clone()).or_default();
642            if !keys.contains(&key_str) {
643                keys.push(key_str);
644            }
645        }
646
647        // 4. 更新统计(不在此处读取 data.len(),避免锁顺序敏感)
648        {
649            let mut stats = self
650                .stats
651                .write()
652                .expect("L2Cache stats lock poisoned (put)");
653            stats.sets += 1;
654        }
655        // 4.1 更新按表分桶统计
656        {
657            if let Ok(mut tbl_stats) = self.table_stats.write() {
658                tbl_stats.entry(key.table.clone()).or_default().sets += 1;
659            }
660        }
661    }
662
663    /// 读取缓存项(不存在或已过期返回 None)
664    ///
665    /// 命中时会更新 LRU 访问顺序(移到尾部)。
666    pub fn get(&self, key: &CacheKey) -> Option<Value> {
667        let key_str = key.to_string_key();
668        let table_name = key.table.clone();
669        let result = {
670            let data = self.data.read().ok()?;
671            if let Some(entry) = data.get(&key_str) {
672                if entry.is_expired() {
673                    None
674                } else {
675                    Some(entry.value.clone())
676                }
677            } else {
678                None
679            }
680        };
681
682        // 命中时更新 LRU 顺序(O(1) touch:移到尾部)
683        if result.is_some() {
684            let mut order = self
685                .access_order
686                .write()
687                .expect("L2Cache access_order lock poisoned (get)");
688            order.touch(&key_str);
689        }
690
691        // 更新全局统计
692        if let Ok(mut stats) = self.stats.write() {
693            if result.is_some() {
694                stats.hits += 1;
695            } else {
696                stats.misses += 1;
697            }
698        }
699        // 更新按表分桶统计
700        if let Ok(mut tbl_stats) = self.table_stats.write() {
701            let entry = tbl_stats.entry(table_name).or_default();
702            if result.is_some() {
703                entry.hits += 1;
704            } else {
705                entry.misses += 1;
706            }
707        }
708
709        result
710    }
711
712    /// 失效单个缓存项
713    pub fn invalidate(&self, key: &CacheKey) {
714        let key_str = key.to_string_key();
715        let table_name = key.table.clone();
716        let removed = {
717            let mut data = self
718                .data
719                .write()
720                .expect("L2Cache data lock poisoned (invalidate)");
721            data.remove(&key_str).is_some()
722        };
723        if removed {
724            let mut order = self
725                .access_order
726                .write()
727                .expect("L2Cache access_order lock poisoned (invalidate)");
728            order.remove(&key_str);
729        }
730        if removed {
731            let mut stats = self
732                .stats
733                .write()
734                .expect("L2Cache stats lock poisoned (invalidate)");
735            stats.evictions += 1;
736            if let Ok(mut tbl_stats) = self.table_stats.write() {
737                tbl_stats.entry(table_name).or_default().evictions += 1;
738            }
739        }
740    }
741
742    /// 失效整张表的所有缓存项
743    ///
744    /// 仅统计实际从缓存中删除的 key 数量,避免 evictions 偏大。
745    /// 若设置了失效总线,会同时发布 `InvalidateTable` 消息通知其他实例。
746    pub fn invalidate_table(&self, table: &str) {
747        let keys_to_remove: Vec<String> = {
748            let idx = match self.table_index.read() {
749                Ok(i) => i,
750                Err(_) => return,
751            };
752            idx.get(table).cloned().unwrap_or_default()
753        };
754
755        let mut actually_removed: usize = 0;
756        {
757            let mut data = self
758                .data
759                .write()
760                .expect("L2Cache data lock poisoned (invalidate_table)");
761            for k in &keys_to_remove {
762                if data.remove(k).is_some() {
763                    actually_removed += 1;
764                }
765            }
766        }
767
768        // O(m) 批量移除(m = keys_to_remove),而非旧实现的 O(n*m) retain
769        if actually_removed > 0 {
770            let mut order = self
771                .access_order
772                .write()
773                .expect("L2Cache access_order lock poisoned (invalidate_table)");
774            for k in &keys_to_remove {
775                order.remove(k);
776            }
777        }
778
779        if let Ok(mut idx) = self.table_index.write() {
780            idx.remove(table);
781        }
782        if actually_removed > 0 {
783            let mut stats = self
784                .stats
785                .write()
786                .expect("L2Cache stats lock poisoned (invalidate_table)");
787            stats.evictions += actually_removed as u64;
788            if let Ok(mut tbl_stats) = self.table_stats.write() {
789                tbl_stats.entry(table.to_string()).or_default().evictions +=
790                    actually_removed as u64;
791            }
792        }
793
794        // 发布失效消息到总线(通知其他订阅实例)
795        if let Some(bus) = &self.invalidation_bus {
796            bus.publish(InvalidationMessage::InvalidateTable(table.to_string()));
797        }
798    }
799
800    /// 清空所有缓存
801    pub fn clear(&self) {
802        let removed = {
803            let mut data = self
804                .data
805                .write()
806                .expect("L2Cache data lock poisoned (clear)");
807            let n = data.len();
808            data.clear();
809            n
810        };
811        if let Ok(mut order) = self.access_order.write() {
812            order.clear();
813        }
814        if let Ok(mut idx) = self.table_index.write() {
815            idx.clear();
816        }
817        if let Ok(mut tbl_stats) = self.table_stats.write() {
818            tbl_stats.clear();
819        }
820        if removed > 0 {
821            let mut stats = self
822                .stats
823                .write()
824                .expect("L2Cache stats lock poisoned (clear)");
825            stats.evictions += removed as u64;
826            stats.size = 0;
827        }
828    }
829
830    /// 获取当前缓存项数量
831    pub fn size(&self) -> usize {
832        self.data.read().map(|d| d.len()).unwrap_or(0)
833    }
834
835    /// 获取统计信息
836    pub fn stats(&self) -> L2CacheStats {
837        let mut s = self.stats.read().map(|s| s.clone()).unwrap_or_default();
838        // 实时同步 size 字段(不写入 stats,避免持锁读 data)
839        s.size = self.size();
840        s
841    }
842
843    /// 重置统计信息(含全局和按表分桶)
844    pub fn reset_stats(&self) {
845        if let Ok(mut stats) = self.stats.write() {
846            *stats = L2CacheStats::default();
847        }
848        if let Ok(mut tbl_stats) = self.table_stats.write() {
849            tbl_stats.clear();
850        }
851    }
852
853    /// 获取指定表的命中率统计
854    pub fn table_stats(&self, table: &str) -> Option<PerTableStats> {
855        self.table_stats
856            .read()
857            .ok()
858            .and_then(|s| s.get(table).cloned())
859    }
860
861    /// 获取所有表的命中率统计快照
862    pub fn all_table_stats(&self) -> HashMap<String, PerTableStats> {
863        self.table_stats
864            .read()
865            .map(|s| s.clone())
866            .unwrap_or_default()
867    }
868
869    /// 检查缓存项是否存在(不更新统计与 LRU 顺序)
870    pub fn contains(&self, key: &CacheKey) -> bool {
871        let key_str = key.to_string_key();
872        self.data
873            .read()
874            .map(|d| d.get(&key_str).map(|e| !e.is_expired()).unwrap_or(false))
875            .unwrap_or(false)
876    }
877
878    /// 手动清理所有过期项
879    pub fn evict_expired(&self) -> usize {
880        let expired_keys: Vec<String> = {
881            let data = self
882                .data
883                .read()
884                .expect("L2Cache data lock poisoned (evict_expired-read)");
885            data.iter()
886                .filter(|(_, e)| e.is_expired())
887                .map(|(k, _)| k.clone())
888                .collect()
889        };
890
891        // 反向查找 key_str -> table_name,用于按表分桶统计
892        let key_to_table: HashMap<String, String> = {
893            let idx = self
894                .table_index
895                .read()
896                .expect("L2Cache table_index lock poisoned (evict_expired-idx)");
897            let mut map = HashMap::new();
898            for (table, keys) in idx.iter() {
899                for k in keys {
900                    map.insert(k.clone(), table.clone());
901                }
902            }
903            map
904        };
905
906        let mut removed = 0;
907        if !expired_keys.is_empty() {
908            let mut data = self
909                .data
910                .write()
911                .expect("L2Cache data lock poisoned (evict_expired-write)");
912            for k in &expired_keys {
913                if data.remove(k).is_some() {
914                    removed += 1;
915                }
916            }
917        }
918
919        if removed > 0 {
920            // O(m) 批量移除(m = expired_keys),而非旧实现的 O(n*m) retain
921            let mut order = self
922                .access_order
923                .write()
924                .expect("L2Cache access_order lock poisoned (evict_expired)");
925            for k in &expired_keys {
926                order.remove(k);
927            }
928            {
929                let mut stats = self
930                    .stats
931                    .write()
932                    .expect("L2Cache stats lock poisoned (evict_expired)");
933                stats.evictions += removed as u64;
934            }
935            // 更新按表分桶统计(单独持锁,避免与 stats 锁同时持有)
936            if let Ok(mut tbl_stats) = self.table_stats.write() {
937                for k in &expired_keys {
938                    if let Some(table) = key_to_table.get(k) {
939                        tbl_stats.entry(table.clone()).or_default().evictions += 1;
940                    }
941                }
942            }
943        }
944        removed
945    }
946
947    /// 更新缓存项的 TTL(若 key 不存在返回 false)
948    ///
949    /// 用于 `Cache` trait 的 `expire` 方法实现。
950    pub fn update_ttl(&self, key: &CacheKey, ttl: Duration) -> bool {
951        let key_str = key.to_string_key();
952        let mut data = match self.data.write() {
953            Ok(d) => d,
954            Err(_) => return false,
955        };
956        if let Some(entry) = data.get_mut(&key_str) {
957            entry.expires_at = Some(Instant::now() + ttl);
958            true
959        } else {
960            false
961        }
962    }
963
964    /// 获取缓存项的剩余 TTL
965    ///
966    /// 返回值:
967    /// - `None`:key 不存在或已过期
968    /// - `Some(None)`:key 存在但无 TTL(永不过期)
969    /// - `Some(Some(d))`:key 存在且剩余 TTL 为 d
970    ///
971    /// 用于 `Cache` trait 的 `ttl` 方法实现。
972    pub fn remaining_ttl(&self, key: &CacheKey) -> Option<Option<Duration>> {
973        let key_str = key.to_string_key();
974        let data = self.data.read().ok()?;
975        let entry = data.get(&key_str)?;
976        match entry.expires_at {
977            Some(expires_at) => {
978                let now = Instant::now();
979                if expires_at <= now {
980                    None
981                } else {
982                    Some(Some(expires_at.duration_since(now)))
983                }
984            }
985            None => Some(None),
986        }
987    }
988}
989
990// ============================================================================
991// Cache trait 实现 — 让 L2Cache 可作为通用 Cache 使用
992// ============================================================================
993
994/// 为 L2Cache 实现 `Cache` trait
995///
996/// 通过 `CacheKey::by_pk("__cache__", key)` 将字符串 key 映射到 L2Cache 的 CacheKey 体系,
997/// 所有通过 `Cache` trait 写入的缓存项归入 `__cache__` 表,与业务缓存项隔离。
998///
999/// 值以 `Value::Bytes(Vec<u8>)` 存储;若通过 `Cache::get` 读取到的 Value 非 Bytes 类型
1000/// (如直接通过 `L2Cache::put` 写入的其他类型),则回退为 JSON 序列化。
1001impl Cache for L2Cache {
1002    fn get(&self, key: &str) -> Result<Option<Vec<u8>>, CacheError> {
1003        let cache_key = CacheKey::by_pk("__cache__", key);
1004        match L2Cache::get(self, &cache_key) {
1005            Some(Value::Bytes(bytes)) => Ok(Some(bytes)),
1006            Some(other) => {
1007                let json = serde_json::to_vec(&other)
1008                    .map_err(|e| CacheError::SerializationError(e.to_string()))?;
1009                Ok(Some(json))
1010            }
1011            None => Ok(None),
1012        }
1013    }
1014
1015    fn set(&self, key: &str, value: Vec<u8>, ttl: Option<Duration>) -> Result<(), CacheError> {
1016        let cache_key = CacheKey::by_pk("__cache__", key);
1017        self.put(&cache_key, Value::Bytes(value), ttl);
1018        Ok(())
1019    }
1020
1021    fn delete(&self, key: &str) -> Result<(), CacheError> {
1022        let cache_key = CacheKey::by_pk("__cache__", key);
1023        self.invalidate(&cache_key);
1024        Ok(())
1025    }
1026
1027    fn clear(&self) -> Result<(), CacheError> {
1028        // 仅清除通过 Cache trait 写入的原始缓存项(__cache__ 表),
1029        // 不影响通过 CacheKey 直接写入的业务缓存项。
1030        self.invalidate_table("__cache__");
1031        Ok(())
1032    }
1033
1034    fn exists(&self, key: &str) -> Result<bool, CacheError> {
1035        let cache_key = CacheKey::by_pk("__cache__", key);
1036        Ok(self.contains(&cache_key))
1037    }
1038
1039    fn expire(&self, key: &str, ttl: Duration) -> Result<(), CacheError> {
1040        let cache_key = CacheKey::by_pk("__cache__", key);
1041        if self.update_ttl(&cache_key, ttl) {
1042            Ok(())
1043        } else {
1044            Err(CacheError::NotFound(key.to_string()))
1045        }
1046    }
1047
1048    fn ttl(&self, key: &str) -> Result<Option<Duration>, CacheError> {
1049        let cache_key = CacheKey::by_pk("__cache__", key);
1050        match self.remaining_ttl(&cache_key) {
1051            None => Err(CacheError::NotFound(key.to_string())),
1052            Some(None) => Ok(None),
1053            Some(Some(d)) => Ok(Some(d)),
1054        }
1055    }
1056}
1057
1058// ============================================================================
1059// L2CacheBackend — 分布式缓存后端抽象(trait + InMemoryBackend + RedisBackend stub)
1060// ============================================================================
1061
1062/// L2 缓存异步 Future 类型别名
1063///
1064/// 用于简化 `L2CacheBackend` trait 中方法的返回类型签名,
1065/// 避免重复书写复杂的 `Pin<Box<dyn Future<...> + Send + 'a>>`。
1066pub type L2CacheFuture<'a, T> = Pin<Box<dyn Future<Output = Result<T, CacheError>> + Send + 'a>>;
1067
1068/// L2 缓存后端 trait(分布式抽象)
1069///
1070/// 定义跨进程共享的二级缓存后端接口,支持进程内内存、Redis 等实现。
1071/// 手动解糖 async 方法(不使用 `#[async_trait]`),与 `Connection` trait 风格一致。
1072///
1073/// # 设计要点
1074///
1075/// - **键值以 `&[u8]` 传输**:后端无关的序列化格式(由调用方决定 bincode/json 等)
1076/// - **TTL 可选**:`Some(Duration)` 设置过期时间,`None` 表示永不过期
1077/// - **前缀失效**:`invalidate_prefix` 批量失效某前缀的所有键(用于表级失效)
1078///
1079/// # 实现方
1080///
1081/// - [`InMemoryBackend`]:进程内内存后端(默认,单机场景)
1082/// - [`RedisBackend`]:Redis 分布式后端(stub,需启用 `redis` feature 并补充依赖)
1083pub trait L2CacheBackend: Send + Sync {
1084    /// 获取缓存值,不存在或已过期返回 `None`
1085    fn get<'a>(&'a self, key: &'a str) -> L2CacheFuture<'a, Option<Vec<u8>>>;
1086
1087    /// 设置缓存值,`ttl` 为 `None` 表示永不过期
1088    fn set<'a>(
1089        &'a self,
1090        key: &'a str,
1091        value: &'a [u8],
1092        ttl: Option<Duration>,
1093    ) -> L2CacheFuture<'a, ()>;
1094
1095    /// 删除单个缓存键
1096    fn delete<'a>(&'a self, key: &'a str) -> L2CacheFuture<'a, ()>;
1097
1098    /// 按前缀批量失效缓存项(用于表级失效)
1099    fn invalidate_prefix<'a>(&'a self, prefix: &'a str) -> L2CacheFuture<'a, ()>;
1100}
1101
1102/// 进程内内存后端(默认实现)
1103///
1104/// 适用于单机场景,不跨进程共享。内部使用 `RwLock<HashMap>` 存储,
1105/// `invalidate_prefix` 通过遍历键前缀匹配实现(O(n),单机场景可接受)。
1106///
1107/// # 线程安全
1108///
1109/// 所有操作通过 `RwLock` 保护,可在多线程环境下共享。
1110///
1111/// # 注意
1112///
1113/// 由于 `std::sync::RwLock` 的 guard 是 `!Send`,所有同步操作在创建
1114/// `Future` 之前完成,guard 在 block 退出时释放,避免跨 `.await` 持锁。
1115pub struct InMemoryBackend {
1116    /// 缓存数据:key -> (value, expiry_time)
1117    /// 使用类型别名降低类型复杂度(clippy::type_complexity)
1118    data: RwLock<InMemoryCacheData>,
1119}
1120
1121/// 内存缓存条目类型别名
1122type InMemoryCacheData = HashMap<String, (Vec<u8>, Option<Instant>)>;
1123
1124impl Default for InMemoryBackend {
1125    fn default() -> Self {
1126        Self::new()
1127    }
1128}
1129
1130impl InMemoryBackend {
1131    /// 创建空的内存后端
1132    pub fn new() -> Self {
1133        Self {
1134            data: RwLock::new(HashMap::new()),
1135        }
1136    }
1137}
1138
1139impl L2CacheBackend for InMemoryBackend {
1140    fn get<'a>(&'a self, key: &'a str) -> L2CacheFuture<'a, Option<Vec<u8>>> {
1141        // 同步完成读操作,guard 在 block 退出时释放,避免跨 await 持锁
1142        let result = {
1143            let data = match self.data.read() {
1144                Ok(d) => d,
1145                Err(e) => {
1146                    let err = CacheError::from(e);
1147                    return Box::pin(async move { Err(err) });
1148                }
1149            };
1150            match data.get(key) {
1151                Some((value, expiry)) => {
1152                    // 过期检查:expiry 为 None 表示永不过期
1153                    if expiry.map(|t| t <= Instant::now()).unwrap_or(false) {
1154                        Ok(None)
1155                    } else {
1156                        Ok(Some(value.clone()))
1157                    }
1158                }
1159                None => Ok(None),
1160            }
1161        };
1162        Box::pin(async move { result })
1163    }
1164
1165    fn set<'a>(
1166        &'a self,
1167        key: &'a str,
1168        value: &'a [u8],
1169        ttl: Option<Duration>,
1170    ) -> L2CacheFuture<'a, ()> {
1171        let result = {
1172            let mut data = match self.data.write() {
1173                Ok(d) => d,
1174                Err(e) => {
1175                    let err = CacheError::from(e);
1176                    return Box::pin(async move { Err(err) });
1177                }
1178            };
1179            let expiry = ttl.map(|d| Instant::now() + d);
1180            data.insert(key.to_string(), (value.to_vec(), expiry));
1181            Ok(())
1182        };
1183        Box::pin(async move { result })
1184    }
1185
1186    fn delete<'a>(&'a self, key: &'a str) -> L2CacheFuture<'a, ()> {
1187        let result = {
1188            let mut data = match self.data.write() {
1189                Ok(d) => d,
1190                Err(e) => {
1191                    let err = CacheError::from(e);
1192                    return Box::pin(async move { Err(err) });
1193                }
1194            };
1195            data.remove(key);
1196            Ok(())
1197        };
1198        Box::pin(async move { result })
1199    }
1200
1201    fn invalidate_prefix<'a>(&'a self, prefix: &'a str) -> L2CacheFuture<'a, ()> {
1202        let result = {
1203            let mut data = match self.data.write() {
1204                Ok(d) => d,
1205                Err(e) => {
1206                    let err = CacheError::from(e);
1207                    return Box::pin(async move { Err(err) });
1208                }
1209            };
1210            // O(n) 遍历,删除所有以 prefix 开头的键
1211            let keys_to_remove: Vec<String> = data
1212                .keys()
1213                .filter(|k| k.starts_with(prefix))
1214                .cloned()
1215                .collect();
1216            for k in keys_to_remove {
1217                data.remove(&k);
1218            }
1219            Ok(())
1220        };
1221        Box::pin(async move { result })
1222    }
1223}
1224
1225/// Redis 分布式缓存后端
1226///
1227/// 基于 `redis` crate 0.27 + `tokio-comp` 异步 IO + `connection-manager` 自动重连。
1228///
1229/// # 实现要点
1230///
1231/// - **连接管理**:使用 `redis::aio::ConnectionManager`(内部自动重连的连接池)
1232/// - **`get`** → `redis::cmd("GET")` 异步执行
1233/// - **`set`** → `SET key value` + 可选 `EX seconds`(合并为单次 `SET` 命令,原子性保证)
1234/// - **`delete`** → `redis::cmd("DEL")`
1235/// - **`invalidate_prefix`** → `SCAN` + 批量 `DEL`(避免 `KEYS` 阻塞 Redis 主线程)
1236///   - 使用 `COUNT 100` 分批扫描,避免单次 SCAN 拉取过多 key 导致阻塞
1237///   - 多次 DEL 调用合并为单次 pipeline 批量执行,减少 RTT 开销
1238///
1239/// # 错误处理
1240///
1241/// - 连接失败 → `CacheError::Internal`,由调用方决定是否重试
1242/// - Redis 命令错误 → 原始错误字符串包装为 `CacheError::Internal`
1243///
1244/// # 启用方式
1245///
1246/// 在 `Cargo.toml` 中启用 `redis` feature:
1247/// ```toml
1248/// [dependencies]
1249/// sz-orm-core = { version = "1.0", features = ["redis"] }
1250/// ```
1251///
1252/// # 使用示例
1253///
1254/// ```no_run
1255/// # use sz_orm_core::l2_cache::{RedisBackend, L2CacheBackend};
1256/// # use std::time::Duration;
1257/// # #[tokio::main]
1258/// # async fn main() -> Result<(), Box<dyn std::error::Error>> {
1259/// let backend = RedisBackend::new("redis://127.0.0.1:6379/0").await?;
1260/// backend.set("user:1", b"alice", Some(Duration::from_secs(60))).await?;
1261/// let val = backend.get("user:1").await?;
1262/// assert_eq!(val, Some(b"alice".to_vec()));
1263/// backend.delete("user:1").await?;
1264/// # Ok(())
1265/// # }
1266/// ```
1267#[cfg(feature = "redis")]
1268pub struct RedisBackend {
1269    /// Redis 异步连接管理器(自动重连)
1270    manager: redis::aio::ConnectionManager,
1271}
1272
1273#[cfg(feature = "redis")]
1274impl RedisBackend {
1275    /// 创建 Redis 后端
1276    ///
1277    /// `url` 格式:`redis://[:password@]host:port[/db]`
1278    /// - `redis://127.0.0.1:6379/0` — 默认 DB 0
1279    /// - `redis://:secret@127.0.0.1:6379/1` — 带密码,DB 1
1280    ///
1281    /// # 错误
1282    ///
1283    /// - 连接失败 → `CacheError::Internal`
1284    pub async fn new(url: impl Into<String>) -> Result<Self, CacheError> {
1285        let url = url.into();
1286        let client = redis::Client::open(url.as_str())
1287            .map_err(|e| CacheError::Internal(format!("Redis client create failed: {}", e)))?;
1288        let manager = redis::aio::ConnectionManager::new(client)
1289            .await
1290            .map_err(|e| CacheError::Internal(format!("Redis connect failed: {}", e)))?;
1291        Ok(Self { manager })
1292    }
1293
1294    /// 使用已有 ConnectionManager 创建后端(用于复用连接池)
1295    pub fn from_manager(manager: redis::aio::ConnectionManager) -> Self {
1296        Self { manager }
1297    }
1298
1299    /// SCAN + 批量 DEL 实现前缀失效
1300    ///
1301    /// 使用 `SCAN cursor MATCH prefix* COUNT 100` 迭代扫描所有匹配的 key,
1302    /// 累计到本地 Vec 后通过 pipeline 批量 DEL,避免:
1303    /// 1. `KEYS pattern` 阻塞 Redis 主线程(O(N) 全表扫描)
1304    /// 2. 单次 DEL 调用过多导致 RTT 累积
1305    ///
1306    /// # 参数
1307    /// - `prefix`:key 前缀(不含通配符,函数内部追加 `*`)
1308    ///
1309    /// # 返回
1310    /// - `Ok(())`:扫描完成,无论是否删除了 key
1311    /// - `Err(_)`:连接错误或命令执行失败
1312    async fn invalidate_prefix_inner(&self, prefix: &str) -> Result<(), CacheError> {
1313        let pattern = format!("{}*", prefix);
1314        let mut cursor: u64 = 0;
1315        loop {
1316            // SCAN 返回 (next_cursor, Vec<key>)
1317            // 注意:必须先 clone 出独立的 conn,避免 &mut temporary 借用问题
1318            let mut conn = self.manager.clone();
1319            let scan_result: redis::RedisResult<(u64, Vec<String>)> = redis::cmd("SCAN")
1320                .arg(cursor)
1321                .arg("MATCH")
1322                .arg(&pattern)
1323                .arg("COUNT")
1324                .arg(100usize)
1325                .query_async(&mut conn)
1326                .await;
1327            let (next_cursor, keys): (u64, Vec<String>) = scan_result
1328                .map_err(|e| CacheError::Internal(format!("Redis SCAN failed: {}", e)))?;
1329
1330            if !keys.is_empty() {
1331                // 批量 DEL:使用 pipeline 减少往返次数
1332                let mut pipe = redis::pipe();
1333                for k in &keys {
1334                    pipe.del(k);
1335                }
1336                // 显式指定 RedisResult<()> 类型,避免 never type fallback 警告
1337                let del_result: redis::RedisResult<()> = pipe.query_async(&mut conn).await;
1338                del_result.map_err(|e| {
1339                    CacheError::Internal(format!("Redis DEL pipeline failed: {}", e))
1340                })?;
1341            }
1342
1343            // cursor == 0 表示扫描完成
1344            if next_cursor == 0 {
1345                break;
1346            }
1347            cursor = next_cursor;
1348        }
1349        Ok(())
1350    }
1351}
1352
1353#[cfg(feature = "redis")]
1354impl L2CacheBackend for RedisBackend {
1355    fn get<'a>(&'a self, key: &'a str) -> L2CacheFuture<'a, Option<Vec<u8>>> {
1356        Box::pin(async move {
1357            use redis::AsyncCommands;
1358            let mut conn = self.manager.clone();
1359            let value: Option<Vec<u8>> = conn
1360                .get(key)
1361                .await
1362                .map_err(|e| CacheError::Internal(format!("Redis GET failed: {}", e)))?;
1363            Ok(value)
1364        })
1365    }
1366
1367    fn set<'a>(
1368        &'a self,
1369        key: &'a str,
1370        value: &'a [u8],
1371        ttl: Option<Duration>,
1372    ) -> L2CacheFuture<'a, ()> {
1373        Box::pin(async move {
1374            use redis::AsyncCommands;
1375            let mut conn = self.manager.clone();
1376            // 合并 SET + EX 为单次原子操作(SET key value EX seconds)
1377            // 避免 SET 后 EXPIRE 之间的窗口期 key 无 TTL
1378            match ttl {
1379                Some(d) => {
1380                    let secs = d.as_secs();
1381                    if secs > 0 {
1382                        let _: () = conn.set_ex(key, value, secs).await.map_err(|e| {
1383                            CacheError::Internal(format!("Redis SET EX failed: {}", e))
1384                        })?;
1385                    } else {
1386                        // TTL < 1s:退化为 SET + PEXPIRE(毫秒精度)
1387                        let _: () = conn.set(key, value).await.map_err(|e| {
1388                            CacheError::Internal(format!("Redis SET failed: {}", e))
1389                        })?;
1390                        // 毫秒数转换为 i64(u128 → i64,实际值不会超过 i64 范围)
1391                        let ms: i64 = d.as_millis().min(i64::MAX as u128) as i64;
1392                        let _: () = conn.pexpire(key, ms).await.map_err(|e| {
1393                            CacheError::Internal(format!("Redis PEXPIRE failed: {}", e))
1394                        })?;
1395                    }
1396                }
1397                None => {
1398                    let _: () = conn
1399                        .set(key, value)
1400                        .await
1401                        .map_err(|e| CacheError::Internal(format!("Redis SET failed: {}", e)))?;
1402                }
1403            }
1404            Ok(())
1405        })
1406    }
1407
1408    fn delete<'a>(&'a self, key: &'a str) -> L2CacheFuture<'a, ()> {
1409        Box::pin(async move {
1410            use redis::AsyncCommands;
1411            let mut conn = self.manager.clone();
1412            let _: () = conn
1413                .del(key)
1414                .await
1415                .map_err(|e| CacheError::Internal(format!("Redis DEL failed: {}", e)))?;
1416            Ok(())
1417        })
1418    }
1419
1420    fn invalidate_prefix<'a>(&'a self, prefix: &'a str) -> L2CacheFuture<'a, ()> {
1421        Box::pin(async move { self.invalidate_prefix_inner(prefix).await })
1422    }
1423}
1424
1425/// Redis 分布式缓存后端(stub,未启用 `redis` feature 时使用)
1426///
1427/// 当未启用 `redis` feature 时,所有操作返回 `CacheError::Internal`,
1428/// 提示用户在 `Cargo.toml` 中启用 `redis` feature。
1429#[cfg(not(feature = "redis"))]
1430pub struct RedisBackend {
1431    /// Redis 连接字符串(保留字段用于错误提示)
1432    url: String,
1433}
1434
1435#[cfg(not(feature = "redis"))]
1436impl RedisBackend {
1437    /// 创建 Redis 后端 stub
1438    ///
1439    /// 返回 stub 实例,所有操作将返回 `CacheError::Internal`。
1440    /// 启用 `redis` feature 后自动切换为真实实现。
1441    pub fn new(_url: impl Into<String>) -> Self {
1442        Self { url: _url.into() }
1443    }
1444}
1445
1446#[cfg(not(feature = "redis"))]
1447impl L2CacheBackend for RedisBackend {
1448    fn get<'a>(&'a self, _key: &'a str) -> L2CacheFuture<'a, Option<Vec<u8>>> {
1449        let url = self.url.clone();
1450        Box::pin(async move {
1451            Err(CacheError::Internal(format!(
1452                "RedisBackend not compiled: enable 'redis' feature in sz-orm-core. URL: {}",
1453                url
1454            )))
1455        })
1456    }
1457
1458    fn set<'a>(
1459        &'a self,
1460        _key: &'a str,
1461        _value: &'a [u8],
1462        _ttl: Option<Duration>,
1463    ) -> L2CacheFuture<'a, ()> {
1464        let url = self.url.clone();
1465        Box::pin(async move {
1466            Err(CacheError::Internal(format!(
1467                "RedisBackend not compiled: enable 'redis' feature in sz-orm-core. URL: {}",
1468                url
1469            )))
1470        })
1471    }
1472
1473    fn delete<'a>(&'a self, _key: &'a str) -> L2CacheFuture<'a, ()> {
1474        let url = self.url.clone();
1475        Box::pin(async move {
1476            Err(CacheError::Internal(format!(
1477                "RedisBackend not compiled: enable 'redis' feature in sz-orm-core. URL: {}",
1478                url
1479            )))
1480        })
1481    }
1482
1483    fn invalidate_prefix<'a>(&'a self, _prefix: &'a str) -> L2CacheFuture<'a, ()> {
1484        let url = self.url.clone();
1485        Box::pin(async move {
1486            Err(CacheError::Internal(format!(
1487                "RedisBackend not compiled: enable 'redis' feature in sz-orm-core. URL: {}",
1488                url
1489            )))
1490        })
1491    }
1492}
1493
1494// ============================================================================
1495// WriteBehind — 异步写回缓存模式(Fix #40)
1496// ============================================================================
1497//
1498// Write-Behind 模式:写操作立即更新缓存,并异步批量刷新到后端存储(如数据库)。
1499// 适用于写吞吐高、可容忍短暂数据不一致的场景。
1500//
1501// # 工作流程
1502//
1503// 1. `write()` / `delete()` → 立即更新 L2CacheBackend,同时将操作入队
1504// 2. 后台任务每 `flush_interval` 触发一次 `flush()`,或显式调用 `flush()`
1505// 3. `flush()` 将队列中的操作批量应用回调 `on_flush`
1506//
1507// # 失败处理
1508//
1509// - 缓存写入失败:立即返回错误给调用方
1510// - 队列写入失败(锁中毒):返回 `CacheError::Internal`
1511// - 后端刷新失败:调用 `on_error` 回调,操作**保留在队列中**等待下次重试
1512//
1513// # 注意
1514//
1515// - 不保证写入顺序与刷新顺序一致(多生产者并发入队)
1516// - 同一 key 的多次写入会按入队顺序刷新(FIFO)
1517// - 调用方需自行处理幂等性(如使用 upsert)
1518
1519/// 写回操作类型
1520#[derive(Debug, Clone)]
1521pub enum WriteOp {
1522    /// SET 操作(key, value, ttl)
1523    Set {
1524        /// 缓存键
1525        key: String,
1526        /// 缓存值(已序列化的字节)
1527        value: Vec<u8>,
1528        /// TTL(与 set 调用一致)
1529        ttl: Option<Duration>,
1530    },
1531    /// DELETE 操作
1532    Delete {
1533        /// 缓存键
1534        key: String,
1535    },
1536}
1537
1538/// 写回刷新回调类型
1539///
1540/// 接收一批待刷新的操作,调用方需将其应用到后端存储(如执行 SQL)。
1541/// 返回 `Err` 表示刷新失败,操作将保留在队列中等待重试。
1542pub type FlushCallback = Arc<
1543    dyn Fn(Vec<WriteOp>) -> Pin<Box<dyn Future<Output = Result<(), CacheError>> + Send>>
1544        + Send
1545        + Sync,
1546>;
1547
1548/// 写回错误回调类型
1549pub type ErrorCallback = Arc<dyn Fn(Vec<WriteOp>, CacheError) + Send + Sync>;
1550
1551/// Write-Behind 写入器
1552///
1553/// 包装一个 `L2CacheBackend`,将写操作同时写入缓存与内存队列,
1554/// 后台任务或显式 `flush()` 触发批量刷新到后端存储。
1555///
1556/// # 线程安全
1557///
1558/// 内部使用 `tokio::sync::Mutex` 保护队列,可被多线程并发调用。
1559///
1560/// # 示例
1561///
1562/// ```no_run
1563/// use sz_orm_core::l2_cache::{WriteBehindWriter, WriteOp, InMemoryBackend};
1564/// use std::sync::Arc;
1565/// use std::time::Duration;
1566///
1567/// # async fn example() -> Result<(), Box<dyn std::error::Error>> {
1568/// let backend = Arc::new(InMemoryBackend::new());
1569/// let on_flush = Arc::new(|ops: Vec<WriteOp>| {
1570///     Box::pin(async move {
1571///         // 这里将 ops 应用到数据库(如批量 INSERT/UPDATE)
1572///         for op in &ops {
1573///             println!("flushing: {:?}", op);
1574///         }
1575///         Ok(())
1576///     }) as std::pin::Pin<Box<dyn std::future::Future<Output = Result<(), _>> + Send>>
1577///     as _
1578/// });
1579/// let writer = WriteBehindWriter::new(backend.clone(), on_flush);
1580///
1581/// // 立即更新缓存,并异步刷新到数据库
1582/// writer.write(b"key1", b"value1", None).await?;
1583///
1584/// // 显式刷新所有待处理操作
1585/// writer.flush().await?;
1586/// # Ok(())
1587/// # }
1588/// ```
1589pub struct WriteBehindWriter {
1590    /// 被包装的 L2 缓存后端
1591    backend: Arc<dyn L2CacheBackend>,
1592    /// 待刷新操作队列
1593    queue: tokio::sync::Mutex<Vec<WriteOp>>,
1594    /// 刷新回调
1595    on_flush: FlushCallback,
1596    /// 错误回调(可选)
1597    on_error: Option<ErrorCallback>,
1598}
1599
1600impl WriteBehindWriter {
1601    /// 创建 Write-Behind 写入器
1602    ///
1603    /// # 参数
1604    /// - `backend`:被包装的 L2 缓存后端(如 `InMemoryBackend`、`RedisBackend`)
1605    /// - `on_flush`:刷新回调,接收一批操作并应用到后端存储
1606    pub fn new(backend: Arc<dyn L2CacheBackend>, on_flush: FlushCallback) -> Self {
1607        Self {
1608            backend,
1609            queue: tokio::sync::Mutex::new(Vec::new()),
1610            on_flush,
1611            on_error: None,
1612        }
1613    }
1614
1615    /// 设置错误回调
1616    ///
1617    /// 当 `flush()` 失败时调用,传入失败的操作和错误信息。
1618    /// 注意:失败的操作会保留在队列中等待下次重试。
1619    pub fn with_error_callback(mut self, on_error: ErrorCallback) -> Self {
1620        self.on_error = Some(on_error);
1621        self
1622    }
1623
1624    /// 写入缓存(立即更新后端缓存 + 入队待刷新)
1625    ///
1626    /// # 参数
1627    /// - `key`:缓存键
1628    /// - `value`:缓存值(字节切片)
1629    /// - `ttl`:TTL,`None` 表示永不过期
1630    pub async fn write(
1631        &self,
1632        key: &[u8],
1633        value: &[u8],
1634        ttl: Option<Duration>,
1635    ) -> Result<(), CacheError> {
1636        let key_str = String::from_utf8_lossy(key).into_owned();
1637        // 1. 立即更新缓存(同步可见性优先)
1638        self.backend.set(&key_str, value, ttl).await?;
1639        // 2. 入队待刷新
1640        let op = WriteOp::Set {
1641            key: key_str,
1642            value: value.to_vec(),
1643            ttl,
1644        };
1645        self.queue.lock().await.push(op);
1646        Ok(())
1647    }
1648
1649    /// 删除缓存项(立即从后端缓存删除 + 入队待刷新)
1650    pub async fn delete(&self, key: &[u8]) -> Result<(), CacheError> {
1651        let key_str = String::from_utf8_lossy(key).into_owned();
1652        // 1. 立即从缓存删除
1653        self.backend.delete(&key_str).await?;
1654        // 2. 入队待刷新
1655        let op = WriteOp::Delete { key: key_str };
1656        self.queue.lock().await.push(op);
1657        Ok(())
1658    }
1659
1660    /// 刷新所有待处理操作到后端存储
1661    ///
1662    /// 将队列中的操作一次性传给 `on_flush` 回调。
1663    /// 若回调返回错误,操作保留在队列中等待下次重试。
1664    pub async fn flush(&self) -> Result<(), CacheError> {
1665        // 1. 取出所有待刷新操作(drain)
1666        let ops: Vec<WriteOp> = {
1667            let mut guard = self.queue.lock().await;
1668            std::mem::take(&mut *guard)
1669        };
1670        if ops.is_empty() {
1671            return Ok(());
1672        }
1673        // 2. 调用刷新回调
1674        match (self.on_flush)(ops.clone()).await {
1675            Ok(()) => Ok(()),
1676            Err(e) => {
1677                // 刷新失败:将操作放回队列,等待下次重试
1678                let mut guard = self.queue.lock().await;
1679                guard.extend(ops.clone());
1680                // 触发错误回调(如有)
1681                if let Some(ref on_error) = self.on_error {
1682                    on_error(ops, e.clone());
1683                }
1684                Err(e)
1685            }
1686        }
1687    }
1688
1689    /// 当前队列中待刷新的操作数(用于监控)
1690    pub async fn pending_count(&self) -> usize {
1691        self.queue.lock().await.len()
1692    }
1693
1694    /// 启动后台自动刷新任务
1695    ///
1696    /// 每 `interval` 触发一次 `flush()`,直到 `WriteBehindWriter` 被丢弃。
1697    /// 返回 `JoinHandle`,调用方可用于等待任务结束。
1698    ///
1699    /// # 注意
1700    ///
1701    /// 调用方需保证 `WriteBehindWriter` 的生命周期长于后台任务,
1702    /// 否则在 writer 被丢弃后,后台任务会因 Arc 引用计数归零而停止。
1703    pub fn spawn_auto_flush(self: Arc<Self>, interval: Duration) -> tokio::task::JoinHandle<()> {
1704        tokio::spawn(async move {
1705            let mut ticker = tokio::time::interval(interval);
1706            // 跳过首次立即触发(首次 tick 会立即返回)
1707            ticker.tick().await;
1708            loop {
1709                ticker.tick().await;
1710                // 刷新失败时记录日志(不中断循环)
1711                if let Err(e) = self.flush().await {
1712                    eprintln!("[WriteBehind] auto flush failed: {}", e);
1713                }
1714            }
1715        })
1716    }
1717}
1718
1719// ============================================================================
1720// 单元测试
1721// ============================================================================
1722
1723#[cfg(test)]
1724mod tests {
1725    use super::*;
1726    use crate::Value;
1727    use std::thread;
1728    use std::time::Duration;
1729
1730    // ===== CacheKey 测试 =====
1731
1732    #[test]
1733    fn test_cache_key_by_pk() {
1734        let key = CacheKey::by_pk("users", 1);
1735        assert_eq!(key.table, "users");
1736        assert_eq!(key.kind, CacheKeyKind::ByPk);
1737        assert_eq!(key.identifier, "1");
1738        assert_eq!(key.to_string_key(), "l2:users:pk:1");
1739    }
1740
1741    #[test]
1742    fn test_cache_key_by_query() {
1743        let key = CacheKey::by_query("orders", "abc123");
1744        assert_eq!(key.kind, CacheKeyKind::ByQuery);
1745        assert_eq!(key.to_string_key(), "l2:orders:q:abc123");
1746    }
1747
1748    #[test]
1749    fn test_cache_key_by_relation() {
1750        let key = CacheKey::by_relation("users", "posts:1");
1751        assert_eq!(key.kind, CacheKeyKind::ByRelation);
1752        assert_eq!(key.to_string_key(), "l2:users:rel:posts:1");
1753    }
1754
1755    #[test]
1756    fn test_cache_key_equality() {
1757        let k1 = CacheKey::by_pk("users", 1);
1758        let k2 = CacheKey::by_pk("users", 1);
1759        let k3 = CacheKey::by_pk("users", 2);
1760        assert_eq!(k1, k2);
1761        assert_ne!(k1, k3);
1762    }
1763
1764    #[test]
1765    fn test_cache_key_display() {
1766        let key = CacheKey::by_pk("users", 42);
1767        assert_eq!(format!("{}", key), "l2:users:pk:42");
1768    }
1769
1770    // ===== L2CacheStats 测试 =====
1771
1772    #[test]
1773    fn test_stats_hit_rate_empty() {
1774        let stats = L2CacheStats::default();
1775        assert_eq!(stats.hit_rate(), 0.0);
1776        assert_eq!(stats.total_lookups(), 0);
1777    }
1778
1779    #[test]
1780    fn test_stats_hit_rate_calculation() {
1781        let stats = L2CacheStats {
1782            hits: 80,
1783            misses: 20,
1784            ..Default::default()
1785        };
1786        assert_eq!(stats.total_lookups(), 100);
1787        assert!((stats.hit_rate() - 0.8).abs() < 0.001);
1788        assert!((stats.miss_rate() - 0.2).abs() < 0.001);
1789    }
1790
1791    #[test]
1792    fn test_stats_merge() {
1793        let mut s1 = L2CacheStats {
1794            hits: 10,
1795            misses: 5,
1796            sets: 15,
1797            evictions: 2,
1798            size: 100,
1799        };
1800        let s2 = L2CacheStats {
1801            hits: 20,
1802            misses: 10,
1803            sets: 30,
1804            evictions: 5,
1805            size: 200,
1806        };
1807        s1.merge(&s2);
1808        assert_eq!(s1.hits, 30);
1809        assert_eq!(s1.misses, 15);
1810        assert_eq!(s1.sets, 45);
1811        assert_eq!(s1.evictions, 7);
1812        assert_eq!(s1.size, 300);
1813    }
1814
1815    // ===== L2Cache 基本操作 =====
1816
1817    #[test]
1818    fn test_put_and_get() {
1819        let cache = L2Cache::new();
1820        let key = CacheKey::by_pk("users", 1);
1821
1822        cache.put(&key, Value::String("Alice".to_string()), None);
1823        let val = cache.get(&key);
1824        assert_eq!(val, Some(Value::String("Alice".to_string())));
1825    }
1826
1827    #[test]
1828    fn test_get_missing_returns_none() {
1829        let cache = L2Cache::new();
1830        let key = CacheKey::by_pk("users", 999);
1831        assert_eq!(cache.get(&key), None);
1832    }
1833
1834    #[test]
1835    fn test_overwrite_existing_key() {
1836        let cache = L2Cache::new();
1837        let key = CacheKey::by_pk("users", 1);
1838
1839        cache.put(&key, Value::String("Alice".to_string()), None);
1840        cache.put(&key, Value::String("Bob".to_string()), None);
1841        assert_eq!(cache.get(&key), Some(Value::String("Bob".to_string())));
1842    }
1843
1844    #[test]
1845    fn test_invalidate_single_key() {
1846        let cache = L2Cache::new();
1847        let key = CacheKey::by_pk("users", 1);
1848
1849        cache.put(&key, Value::I64(42), None);
1850        assert!(cache.get(&key).is_some());
1851
1852        cache.invalidate(&key);
1853        assert!(cache.get(&key).is_none());
1854    }
1855
1856    // ===== 表级失效 =====
1857
1858    #[test]
1859    fn test_invalidate_table_removes_all_entries_for_table() {
1860        let cache = L2Cache::new();
1861
1862        let k1 = CacheKey::by_pk("users", 1);
1863        let k2 = CacheKey::by_pk("users", 2);
1864        let k3 = CacheKey::by_query("users", "hash1");
1865        let k4 = CacheKey::by_pk("orders", 1); // 不同表
1866
1867        cache.put(&k1, Value::I64(1), None);
1868        cache.put(&k2, Value::I64(2), None);
1869        cache.put(&k3, Value::I64(3), None);
1870        cache.put(&k4, Value::I64(4), None);
1871
1872        cache.invalidate_table("users");
1873
1874        // users 表的所有缓存项应被失效
1875        assert!(cache.get(&k1).is_none());
1876        assert!(cache.get(&k2).is_none());
1877        assert!(cache.get(&k3).is_none());
1878        // orders 表的缓存项应保留
1879        assert!(cache.get(&k4).is_some());
1880    }
1881
1882    #[test]
1883    fn test_invalidate_table_no_op_for_unknown_table() {
1884        let cache = L2Cache::new();
1885        let k1 = CacheKey::by_pk("users", 1);
1886        cache.put(&k1, Value::I64(1), None);
1887
1888        cache.invalidate_table("nonexistent");
1889        assert!(cache.get(&k1).is_some());
1890    }
1891
1892    // ===== TTL 测试 =====
1893
1894    #[test]
1895    fn test_ttl_expiration() {
1896        let cache = L2Cache::new();
1897        let key = CacheKey::by_pk("users", 1);
1898
1899        cache.put(&key, Value::I64(42), Some(Duration::from_millis(50)));
1900        assert!(cache.get(&key).is_some());
1901
1902        // 等待 TTL 过期
1903        thread::sleep(Duration::from_millis(100));
1904        assert!(cache.get(&key).is_none());
1905    }
1906
1907    #[test]
1908    fn test_default_ttl_applied_when_no_explicit_ttl() {
1909        let cache = L2Cache::new().with_default_ttl(Duration::from_millis(50));
1910        let key = CacheKey::by_pk("users", 1);
1911
1912        cache.put(&key, Value::I64(42), None); // 不显式传 TTL
1913        assert!(cache.get(&key).is_some());
1914
1915        thread::sleep(Duration::from_millis(100));
1916        assert!(cache.get(&key).is_none());
1917    }
1918
1919    #[test]
1920    fn test_explicit_ttl_overrides_default() {
1921        // 语义验证:ttl=Some(Duration::MAX) 表示永不失效,覆盖默认 TTL
1922        let cache = L2Cache::new().with_default_ttl(Duration::from_millis(50));
1923        let key = CacheKey::by_pk("users", 1);
1924
1925        // 显式传 Some(Duration::MAX) 覆盖默认 TTL(永不失效)
1926        cache.put(&key, Value::I64(42), Some(Duration::MAX));
1927
1928        // 等待默认 TTL 已过期的时间
1929        thread::sleep(Duration::from_millis(100));
1930        // 由于显式传 Some(Duration::MAX),应仍然有效
1931        assert!(cache.get(&key).is_some());
1932    }
1933
1934    #[test]
1935    fn test_none_ttl_uses_default_ttl() {
1936        // 语义验证:ttl=None 时使用 default_ttl
1937        let cache = L2Cache::new().with_default_ttl(Duration::from_millis(50));
1938        let key = CacheKey::by_pk("users", 1);
1939
1940        cache.put(&key, Value::I64(42), None);
1941        assert!(cache.get(&key).is_some());
1942
1943        thread::sleep(Duration::from_millis(100));
1944        // None 使用了 default_ttl,应已过期
1945        assert!(cache.get(&key).is_none());
1946    }
1947
1948    // ===== 命中率统计 =====
1949
1950    #[test]
1951    fn test_stats_tracks_hits_and_misses() {
1952        let cache = L2Cache::new();
1953
1954        let k1 = CacheKey::by_pk("users", 1);
1955        let k2 = CacheKey::by_pk("users", 2);
1956
1957        cache.put(&k1, Value::I64(1), None);
1958
1959        // 1 次命中
1960        cache.get(&k1);
1961        // 2 次未命中
1962        cache.get(&k2);
1963        cache.get(&k2);
1964
1965        let stats = cache.stats();
1966        assert_eq!(stats.hits, 1);
1967        assert_eq!(stats.misses, 2);
1968        assert_eq!(stats.sets, 1);
1969    }
1970
1971    #[test]
1972    fn test_stats_tracks_evictions() {
1973        let cache = L2Cache::new();
1974        let k1 = CacheKey::by_pk("users", 1);
1975        let k2 = CacheKey::by_pk("users", 2);
1976
1977        cache.put(&k1, Value::I64(1), None);
1978        cache.put(&k2, Value::I64(2), None);
1979
1980        cache.invalidate(&k1); // evictions = 1
1981        cache.invalidate_table("users"); // 仅 k2 实际被删除,evictions = 2
1982
1983        let stats = cache.stats();
1984        // invalidate(k1) 删除 1 项;invalidate_table("users") 仅删除 k2(k1 已不存在)
1985        assert_eq!(stats.evictions, 2);
1986    }
1987
1988    #[test]
1989    fn test_stats_reset() {
1990        let cache = L2Cache::new();
1991        let k1 = CacheKey::by_pk("users", 1);
1992
1993        cache.put(&k1, Value::I64(1), None);
1994        cache.get(&k1);
1995        cache.get(&k1);
1996
1997        let stats_before = cache.stats();
1998        assert!(stats_before.hits > 0);
1999
2000        cache.reset_stats();
2001        let stats_after = cache.stats();
2002        assert_eq!(stats_after.hits, 0);
2003        assert_eq!(stats_after.misses, 0);
2004        assert_eq!(stats_after.sets, 0);
2005    }
2006
2007    // ===== 容量管理 =====
2008
2009    #[test]
2010    fn test_max_size_eviction() {
2011        let cache = L2Cache::new().with_max_size(3);
2012
2013        for i in 0..5 {
2014            let k = CacheKey::by_pk("users", i);
2015            cache.put(&k, Value::I64(i), None);
2016        }
2017
2018        // 真正的 LRU:容量严格不超过 max_size
2019        let size = cache.size();
2020        assert_eq!(
2021            size, 3,
2022            "size should be exactly max_size after LRU eviction, got {}",
2023            size
2024        );
2025    }
2026
2027    #[test]
2028    fn test_lru_eviction_order() {
2029        // 验证 LRU 顺序:访问 k0 后,下次淘汰应跳过 k0 而淘汰 k1
2030        let cache = L2Cache::new().with_max_size(3);
2031
2032        let k0 = CacheKey::by_pk("users", 0);
2033        let k1 = CacheKey::by_pk("users", 1);
2034        let k2 = CacheKey::by_pk("users", 2);
2035        let k3 = CacheKey::by_pk("users", 3);
2036
2037        cache.put(&k0, Value::I64(0), None);
2038        cache.put(&k1, Value::I64(1), None);
2039        cache.put(&k2, Value::I64(2), None);
2040
2041        // 访问 k0,使其成为最近使用
2042        let _ = cache.get(&k0);
2043
2044        // 插入 k3,应淘汰 k1(最久未访问)
2045        cache.put(&k3, Value::I64(3), None);
2046
2047        assert!(
2048            cache.get(&k0).is_some(),
2049            "k0 should survive (recently accessed)"
2050        );
2051        assert!(
2052            cache.get(&k1).is_none(),
2053            "k1 should be evicted (LRU victim)"
2054        );
2055        assert!(cache.get(&k2).is_some(), "k2 should survive");
2056        assert!(
2057            cache.get(&k3).is_some(),
2058            "k3 should survive (just inserted)"
2059        );
2060    }
2061
2062    #[test]
2063    fn test_clear_all() {
2064        let cache = L2Cache::new();
2065        cache.put(&CacheKey::by_pk("users", 1), Value::I64(1), None);
2066        cache.put(&CacheKey::by_pk("users", 2), Value::I64(2), None);
2067        cache.put(&CacheKey::by_pk("orders", 1), Value::I64(3), None);
2068
2069        assert_eq!(cache.size(), 3);
2070        cache.clear();
2071        assert_eq!(cache.size(), 0);
2072    }
2073
2074    // ===== contains(不更新统计)=====
2075
2076    #[test]
2077    fn test_contains_does_not_update_stats() {
2078        let cache = L2Cache::new();
2079        let k1 = CacheKey::by_pk("users", 1);
2080        cache.put(&k1, Value::I64(1), None);
2081
2082        let exists = cache.contains(&k1);
2083        assert!(exists);
2084
2085        let stats = cache.stats();
2086        assert_eq!(stats.hits, 0);
2087        assert_eq!(stats.misses, 0);
2088    }
2089
2090    #[test]
2091    fn test_contains_returns_false_for_missing() {
2092        let cache = L2Cache::new();
2093        let k = CacheKey::by_pk("users", 999);
2094        assert!(!cache.contains(&k));
2095    }
2096
2097    #[test]
2098    fn test_contains_returns_false_for_expired() {
2099        let cache = L2Cache::new();
2100        let k = CacheKey::by_pk("users", 1);
2101        cache.put(&k, Value::I64(1), Some(Duration::from_millis(10)));
2102
2103        thread::sleep(Duration::from_millis(50));
2104        assert!(!cache.contains(&k));
2105    }
2106
2107    // ===== evict_expired 手动清理 =====
2108
2109    #[test]
2110    fn test_evict_expired_removes_only_expired_entries() {
2111        let cache = L2Cache::new();
2112
2113        let k1 = CacheKey::by_pk("users", 1);
2114        let k2 = CacheKey::by_pk("users", 2);
2115
2116        cache.put(&k1, Value::I64(1), Some(Duration::from_millis(10)));
2117        cache.put(&k2, Value::I64(2), None); // 永不过期
2118
2119        thread::sleep(Duration::from_millis(50));
2120        let removed = cache.evict_expired();
2121
2122        assert_eq!(removed, 1);
2123        assert!(cache.get(&k1).is_none());
2124        assert!(cache.get(&k2).is_some());
2125    }
2126
2127    #[test]
2128    fn test_evict_expired_returns_zero_if_no_expired() {
2129        let cache = L2Cache::new();
2130        let k1 = CacheKey::by_pk("users", 1);
2131        cache.put(&k1, Value::I64(1), None);
2132
2133        let removed = cache.evict_expired();
2134        assert_eq!(removed, 0);
2135    }
2136
2137    // ===== 多线程测试 =====
2138
2139    #[test]
2140    fn test_concurrent_access() {
2141        let cache = std::sync::Arc::new(L2Cache::new());
2142        let mut handles = Vec::new();
2143
2144        // 多线程写入
2145        for i in 0..4 {
2146            let c = cache.clone();
2147            handles.push(thread::spawn(move || {
2148                for j in 0..10 {
2149                    let k = CacheKey::by_pk("users", i * 10 + j);
2150                    c.put(&k, Value::I64(i * 10 + j), None);
2151                }
2152            }));
2153        }
2154        for h in handles {
2155            h.join().unwrap();
2156        }
2157
2158        assert_eq!(cache.size(), 40);
2159
2160        // 多线程读取
2161        let mut handles = Vec::new();
2162        for i in 0..4 {
2163            let c = cache.clone();
2164            handles.push(thread::spawn(move || {
2165                for j in 0..10 {
2166                    let k = CacheKey::by_pk("users", i * 10 + j);
2167                    let v = c.get(&k);
2168                    assert!(v.is_some());
2169                }
2170            }));
2171        }
2172        for h in handles {
2173            h.join().unwrap();
2174        }
2175
2176        let stats = cache.stats();
2177        assert_eq!(stats.hits, 40);
2178    }
2179
2180    // ===== Default 测试 =====
2181
2182    #[test]
2183    fn test_default() {
2184        let cache = L2Cache::default();
2185        assert_eq!(cache.size(), 0);
2186    }
2187
2188    // ===== 综合场景 =====
2189
2190    #[test]
2191    fn test_realistic_scenario() {
2192        let cache = L2Cache::new();
2193
2194        // 1. 缓存用户表数据
2195        for i in 1..=5 {
2196            cache.put(
2197                &CacheKey::by_pk("users", i),
2198                Value::String(format!("user_{}", i)),
2199                None,
2200            );
2201        }
2202
2203        // 2. 缓存查询结果
2204        cache.put(
2205            &CacheKey::by_query("users", "active_users_hash"),
2206            Value::I64(5),
2207            None,
2208        );
2209
2210        // 3. 读取(部分命中、部分未命中)
2211        for i in 1..=10 {
2212            let _ = cache.get(&CacheKey::by_pk("users", i));
2213        }
2214
2215        let stats = cache.stats();
2216        assert_eq!(stats.hits, 5); // 1-5 命中
2217        assert_eq!(stats.misses, 5); // 6-10 未命中
2218        assert_eq!(stats.sets, 6); // 5 pk + 1 query
2219
2220        // 4. 用户表更新,失效所有缓存
2221        cache.invalidate_table("users");
2222
2223        // 5. 再次读取应全部未命中
2224        cache.reset_stats();
2225        for i in 1..=5 {
2226            let _ = cache.get(&CacheKey::by_pk("users", i));
2227        }
2228        let stats2 = cache.stats();
2229        assert_eq!(stats2.hits, 0);
2230        assert_eq!(stats2.misses, 5);
2231    }
2232
2233    // ===== WriteBehindWriter 测试(Fix #40) =====
2234
2235    #[tokio::test]
2236    async fn test_write_behind_basic_write_and_flush() {
2237        use std::sync::atomic::{AtomicUsize, Ordering};
2238        // 计数刷新调用次数
2239        let counter = Arc::new(AtomicUsize::new(0));
2240        let counter_clone = counter.clone();
2241        let on_flush: FlushCallback = Arc::new(move |ops: Vec<WriteOp>| {
2242            let c = counter_clone.clone();
2243            Box::pin(async move {
2244                c.fetch_add(ops.len(), Ordering::SeqCst);
2245                Ok(())
2246            })
2247        });
2248        let backend = Arc::new(InMemoryBackend::new());
2249        let writer = WriteBehindWriter::new(backend.clone(), on_flush);
2250
2251        // 写入 3 个键
2252        writer.write(b"k1", b"v1", None).await.unwrap();
2253        writer.write(b"k2", b"v2", None).await.unwrap();
2254        writer.write(b"k3", b"v3", None).await.unwrap();
2255
2256        // 缓存应立即可见
2257        let v1 = backend.get("k1").await.unwrap();
2258        assert_eq!(v1, Some(b"v1".to_vec()));
2259
2260        // 队列应有 3 个待刷新
2261        assert_eq!(writer.pending_count().await, 3);
2262
2263        // 刷新
2264        writer.flush().await.unwrap();
2265        assert_eq!(counter.load(Ordering::SeqCst), 3);
2266        assert_eq!(writer.pending_count().await, 0);
2267    }
2268
2269    #[tokio::test]
2270    async fn test_write_behind_delete() {
2271        let on_flush: FlushCallback =
2272            Arc::new(|_ops: Vec<WriteOp>| Box::pin(async move { Ok(()) }));
2273        let backend = Arc::new(InMemoryBackend::new());
2274        let writer = WriteBehindWriter::new(backend.clone(), on_flush);
2275
2276        // 写入后删除
2277        writer.write(b"k1", b"v1", None).await.unwrap();
2278        assert!(backend.get("k1").await.unwrap().is_some());
2279        writer.delete(b"k1").await.unwrap();
2280        // 删除后缓存中应不存在
2281        assert!(backend.get("k1").await.unwrap().is_none());
2282
2283        // flush 应处理 2 个操作(Set + Delete)
2284        writer.flush().await.unwrap();
2285        assert_eq!(writer.pending_count().await, 0);
2286    }
2287
2288    #[tokio::test]
2289    async fn test_write_behind_flush_failure_retries() {
2290        // 模拟刷新总是失败
2291        let on_flush: FlushCallback = Arc::new(|_ops: Vec<WriteOp>| {
2292            Box::pin(async move { Err(CacheError::Internal("backend down".to_string())) })
2293        });
2294        let backend = Arc::new(InMemoryBackend::new());
2295        let writer = WriteBehindWriter::new(backend.clone(), on_flush);
2296
2297        writer.write(b"k1", b"v1", None).await.unwrap();
2298        // flush 失败,操作应保留在队列中
2299        let result = writer.flush().await;
2300        assert!(result.is_err());
2301        assert_eq!(writer.pending_count().await, 1);
2302    }
2303
2304    #[tokio::test]
2305    async fn test_write_behind_empty_flush_noop() {
2306        let on_flush: FlushCallback =
2307            Arc::new(|_ops: Vec<WriteOp>| Box::pin(async move { Ok(()) }));
2308        let backend = Arc::new(InMemoryBackend::new());
2309        let writer = WriteBehindWriter::new(backend, on_flush);
2310        // 空队列 flush 应立即成功
2311        writer.flush().await.unwrap();
2312        assert_eq!(writer.pending_count().await, 0);
2313    }
2314
2315    #[tokio::test]
2316    async fn test_write_behind_error_callback_invoked() {
2317        use std::sync::atomic::{AtomicUsize, Ordering};
2318        let error_counter = Arc::new(AtomicUsize::new(0));
2319        let ec = error_counter.clone();
2320        let on_error: ErrorCallback = Arc::new(move |_ops, _err| {
2321            ec.fetch_add(1, Ordering::SeqCst);
2322        });
2323        let on_flush: FlushCallback = Arc::new(|_ops: Vec<WriteOp>| {
2324            Box::pin(async move { Err(CacheError::Internal("fail".to_string())) })
2325        });
2326        let backend = Arc::new(InMemoryBackend::new());
2327        let writer = WriteBehindWriter::new(backend, on_flush).with_error_callback(on_error);
2328
2329        writer.write(b"k1", b"v1", None).await.unwrap();
2330        let _ = writer.flush().await;
2331        assert_eq!(error_counter.load(Ordering::SeqCst), 1);
2332    }
2333}