Skip to main content

reifydb_cdc/storage/
recent_cache.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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		// lo below the cache's min version => caller must fall back to storage.
102		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		// hi above the cache's max version: versions in (max, hi] may exist in
111		// durable storage but not in the cache, so claiming the range is empty
112		// would make the caller miss them. The caller must fall back to storage.
113		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}