reifydb_cdc/storage/
recent_cache.rs1use std::{collections::BTreeMap, sync::Arc};
5
6use reifydb_core::{common::CommitVersion, interface::cdc::Cdc};
7use reifydb_runtime::sync::mutex::Mutex;
8
9#[derive(Clone)]
10pub struct RecentCdcCache {
11 inner: Arc<Mutex<BTreeMap<CommitVersion, Arc<Cdc>>>>,
12 capacity: usize,
13}
14
15impl RecentCdcCache {
16 pub const DEFAULT_CAPACITY: usize = 1024;
17
18 pub fn new(capacity: usize) -> Self {
19 Self {
20 inner: Arc::new(Mutex::new(BTreeMap::new())),
21 capacity: capacity.max(1),
22 }
23 }
24
25 pub fn insert(&self, cdc: &Cdc) {
26 let mut entries = self.inner.lock();
27 entries.insert(cdc.version, Arc::new(cdc.clone()));
28 while entries.len() > self.capacity {
29 let Some(lowest) = entries.keys().next().copied() else {
30 break;
31 };
32 entries.remove(&lowest);
33 }
34 }
35
36 pub fn get(&self, version: CommitVersion) -> Option<Arc<Cdc>> {
37 self.inner.lock().get(&version).cloned()
38 }
39
40 pub fn try_serve_range(
41 &self,
42 lo_inc: CommitVersion,
43 hi_inc: CommitVersion,
44 limit: usize,
45 ) -> Option<(Vec<Cdc>, bool)> {
46 let entries = self.inner.lock();
47 let min = entries.keys().next().copied()?;
48 let max = entries.keys().next_back().copied()?;
49 if lo_inc < min || hi_inc > max {
50 return None;
51 }
52 let mut range = entries.range(lo_inc..=hi_inc);
53 let mut items = Vec::new();
54 for (_, cdc) in range.by_ref().take(limit) {
55 items.push((**cdc).clone());
56 }
57 let has_more = range.next().is_some();
58 Some((items, has_more))
59 }
60
61 pub fn clear(&self) {
62 self.inner.lock().clear();
63 }
64}
65
66#[cfg(test)]
67mod tests {
68 use reifydb_value::value::datetime::DateTime;
69
70 use super::*;
71
72 fn cv(n: u64) -> CommitVersion {
73 CommitVersion(n)
74 }
75
76 fn cdc(version: u64) -> Cdc {
77 Cdc::new(cv(version), DateTime::default(), Vec::new(), Vec::new())
78 }
79
80 #[test]
81 fn insert_then_get_returns_entry() {
82 let cache = RecentCdcCache::new(4);
83 cache.insert(&cdc(1));
84 assert_eq!(cache.get(cv(1)).expect("present").version, cv(1));
85 assert!(cache.get(cv(2)).is_none());
86 }
87
88 #[test]
89 fn eviction_drops_lowest_version_when_over_capacity() {
90 let cache = RecentCdcCache::new(2);
91 cache.insert(&cdc(1));
92 cache.insert(&cdc(2));
93 cache.insert(&cdc(3));
94 assert!(cache.get(cv(1)).is_none(), "lowest version must be evicted");
95 assert!(cache.get(cv(2)).is_some());
96 assert!(cache.get(cv(3)).is_some());
97 }
98
99 #[test]
100 fn serve_range_returns_none_when_not_fully_covered() {
101 let cache = RecentCdcCache::new(2);
103 cache.insert(&cdc(5));
104 cache.insert(&cdc(6));
105 assert!(cache.try_serve_range(cv(3), cv(6), 100).is_none());
106 }
107
108 #[test]
109 fn serve_range_returns_none_when_above_cache_max() {
110 let cache = RecentCdcCache::new(8);
114 cache.insert(&cdc(5));
115 cache.insert(&cdc(6));
116 assert!(cache.try_serve_range(cv(5), cv(8), 100).is_none(), "must not serve past cache max");
117 assert!(cache.try_serve_range(cv(7), cv(7), 100).is_none(), "must not serve a gap above max as empty");
118 let (items, _) = cache.try_serve_range(cv(5), cv(6), 100).expect("fully covered up to max");
119 assert_eq!(items.len(), 2);
120 }
121
122 #[test]
123 fn serve_range_returns_none_when_empty() {
124 let cache = RecentCdcCache::new(4);
125 assert!(cache.try_serve_range(cv(1), cv(10), 100).is_none());
126 }
127
128 #[test]
129 fn serve_range_serves_covered_range_in_order() {
130 let cache = RecentCdcCache::new(8);
131 for v in 4..=8 {
132 cache.insert(&cdc(v));
133 }
134 let (items, has_more) = cache.try_serve_range(cv(5), cv(7), 100).expect("covered");
135 assert_eq!(items.iter().map(|c| c.version).collect::<Vec<_>>(), vec![cv(5), cv(6), cv(7)]);
136 assert!(!has_more);
137 }
138
139 #[test]
140 fn serve_range_reports_has_more_when_limited() {
141 let cache = RecentCdcCache::new(8);
142 for v in 1..=5 {
143 cache.insert(&cdc(v));
144 }
145 let (items, has_more) = cache.try_serve_range(cv(1), cv(5), 2).expect("covered");
146 assert_eq!(items.len(), 2);
147 assert!(has_more, "more entries remain in range beyond the limit");
148 }
149
150 #[test]
151 fn serve_range_at_exactly_min_is_covered() {
152 let cache = RecentCdcCache::new(4);
153 cache.insert(&cdc(10));
154 cache.insert(&cdc(11));
155 let (items, _) = cache.try_serve_range(cv(10), cv(11), 100).expect("covered at min");
156 assert_eq!(items.len(), 2);
157 }
158
159 #[test]
160 fn clone_shares_backing_storage() {
161 let a = RecentCdcCache::new(4);
162 let b = a.clone();
163 a.insert(&cdc(1));
164 assert!(b.get(cv(1)).is_some(), "clone observes writes from original");
165 }
166}