Skip to main content

reifydb_cdc/storage/
cached.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::collections::Bound;
5
6use reifydb_core::{
7	common::CommitVersion,
8	interface::cdc::{Cdc, CdcBatch},
9};
10use reifydb_value::value::datetime::DateTime;
11
12use super::{CdcStorage, CdcStorageResult, DropBeforeResult, normalize_range_inclusive, recent_cache::RecentCdcCache};
13
14#[derive(Clone)]
15pub struct CachedCdcStorage<S: CdcStorage> {
16	inner: S,
17	cache: RecentCdcCache,
18}
19
20impl<S: CdcStorage> CachedCdcStorage<S> {
21	pub fn new(inner: S, capacity: usize) -> Self {
22		Self {
23			inner,
24			cache: RecentCdcCache::new(capacity),
25		}
26	}
27
28	pub fn inner(&self) -> &S {
29		&self.inner
30	}
31}
32
33impl<S: CdcStorage> CdcStorage for CachedCdcStorage<S> {
34	fn write(&self, cdc: &Cdc) -> CdcStorageResult<()> {
35		self.inner.write(cdc)?;
36		self.cache.insert(cdc);
37		Ok(())
38	}
39
40	fn read(&self, version: CommitVersion) -> CdcStorageResult<Option<Cdc>> {
41		if let Some(cdc) = self.cache.get(version) {
42			return Ok(Some((*cdc).clone()));
43		}
44		self.inner.read(version)
45	}
46
47	fn read_range(
48		&self,
49		start: Bound<CommitVersion>,
50		end: Bound<CommitVersion>,
51		batch_size: u64,
52	) -> CdcStorageResult<CdcBatch> {
53		if let Some((lo_inc, hi_inc)) = normalize_range_inclusive(start, end)
54			&& let Some((items, has_more)) = self.cache.try_serve_range(lo_inc, hi_inc, batch_size as usize)
55		{
56			return Ok(CdcBatch {
57				items,
58				has_more,
59			});
60		}
61		self.inner.read_range(start, end, batch_size)
62	}
63
64	fn count(&self, version: CommitVersion) -> CdcStorageResult<usize> {
65		self.inner.count(version)
66	}
67
68	fn min_version(&self) -> CdcStorageResult<Option<CommitVersion>> {
69		self.inner.min_version()
70	}
71
72	fn max_version(&self) -> CdcStorageResult<Option<CommitVersion>> {
73		self.inner.max_version()
74	}
75
76	fn drop_before(&self, version: CommitVersion, limit: usize) -> CdcStorageResult<DropBeforeResult> {
77		self.inner.drop_before(version, limit)
78	}
79
80	fn find_ttl_cutoff(&self, cutoff: DateTime) -> CdcStorageResult<Option<CommitVersion>> {
81		self.inner.find_ttl_cutoff(cutoff)
82	}
83}
84
85#[cfg(test)]
86mod tests {
87	use std::collections::Bound;
88
89	use reifydb_core::{common::CommitVersion, interface::cdc::Cdc};
90	use reifydb_value::value::datetime::DateTime;
91
92	use super::*;
93	use crate::storage::memory::MemoryCdcStorage;
94
95	fn cv(n: u64) -> CommitVersion {
96		CommitVersion(n)
97	}
98
99	fn cdc(version: u64) -> Cdc {
100		Cdc::new(cv(version), DateTime::default(), Vec::new(), Vec::new())
101	}
102
103	#[test]
104	fn write_is_persisted_to_inner_and_served_from_cache() {
105		let cached = CachedCdcStorage::new(MemoryCdcStorage::new(), 16);
106		cached.write(&cdc(1)).unwrap();
107		// inner has it durably
108		assert!(cached.inner().read(cv(1)).unwrap().is_some());
109		// and the cache serves the read
110		assert_eq!(cached.read(cv(1)).unwrap().unwrap().version, cv(1));
111	}
112
113	#[test]
114	fn read_range_served_from_cache_when_covered() {
115		let cached = CachedCdcStorage::new(MemoryCdcStorage::new(), 16);
116		for v in 1..=5 {
117			cached.write(&cdc(v)).unwrap();
118		}
119		let batch = cached.read_range(Bound::Excluded(cv(1)), Bound::Included(cv(4)), 100).unwrap();
120		assert_eq!(batch.items.iter().map(|c| c.version).collect::<Vec<_>>(), vec![cv(2), cv(3), cv(4)]);
121		assert!(!batch.has_more);
122	}
123
124	#[test]
125	fn read_range_falls_back_to_inner_when_below_cache_window() {
126		// Capacity 2 keeps only versions {4,5}; a request starting at 1 is not covered, so the
127		// decorator must fall through to the backend, which still has the full history.
128		let inner = MemoryCdcStorage::new();
129		let cached = CachedCdcStorage::new(inner, 2);
130		for v in 1..=5 {
131			cached.write(&cdc(v)).unwrap();
132		}
133		let batch = cached.read_range(Bound::Included(cv(1)), Bound::Included(cv(5)), 100).unwrap();
134		assert_eq!(batch.items.len(), 5, "fallback must serve the full range from the backend");
135	}
136}