reifydb_runtime/
version_epoch.rs1use std::{collections::BTreeMap, sync::Arc};
5
6use crate::sync::rwlock::RwLock;
7
8const MAX_SAMPLES: usize = 100_000;
9
10#[derive(Clone)]
11pub struct VersionEpoch {
12 inner: Arc<RwLock<BTreeMap<u64, u64>>>,
13}
14
15impl Default for VersionEpoch {
16 fn default() -> Self {
17 Self::new()
18 }
19}
20
21impl VersionEpoch {
22 pub fn new() -> Self {
23 Self {
24 inner: Arc::new(RwLock::new(BTreeMap::new())),
25 }
26 }
27
28 pub fn record(&self, bucket_nanos: u64, version: u64) {
29 let mut samples = self.inner.write();
30 if let Some((&last_bucket, &last_version)) = samples.iter().next_back() {
31 if bucket_nanos < last_bucket || version < last_version {
32 return;
33 }
34 if bucket_nanos == last_bucket {
35 samples.insert(bucket_nanos, version);
36 return;
37 }
38 }
39 samples.insert(bucket_nanos, version);
40 while samples.len() > MAX_SAMPLES {
41 let oldest = *samples.keys().next().expect("samples is non-empty after insert");
42 samples.remove(&oldest);
43 }
44 }
45
46 pub fn floor_version_at(&self, target_nanos: u64) -> Option<u64> {
47 self.inner.read().range(..=target_nanos).next_back().map(|(_, version)| *version)
48 }
49
50 #[cfg(test)]
51 pub fn sample_count(&self) -> usize {
52 self.inner.read().len()
53 }
54}
55
56#[cfg(test)]
57mod tests {
58 use super::VersionEpoch;
59
60 #[test]
61 fn cold_epoch_returns_none_so_gc_deletes_nothing() {
62 let epoch = VersionEpoch::new();
63 assert_eq!(
64 epoch.floor_version_at(1_000),
65 None,
66 "an empty epoch must yield no cutoff; otherwise a cold start would evict the whole store"
67 );
68 }
69
70 #[test]
71 fn floor_returns_latest_sample_at_or_before_target() {
72 let epoch = VersionEpoch::new();
73 epoch.record(100, 10);
74 epoch.record(200, 20);
75 epoch.record(300, 30);
76
77 assert_eq!(epoch.floor_version_at(50), None, "target older than every sample -> no cutoff");
78 assert_eq!(epoch.floor_version_at(100), Some(10), "exact bucket is included");
79 assert_eq!(epoch.floor_version_at(250), Some(20), "floor is the latest bucket <= target");
80 assert_eq!(epoch.floor_version_at(9_999), Some(30), "target after all samples -> newest");
81 }
82
83 #[test]
84 fn record_drops_non_monotonic_samples() {
85 let epoch = VersionEpoch::new();
86 epoch.record(200, 20);
87 epoch.record(100, 10);
88 epoch.record(300, 15);
89
90 assert_eq!(epoch.sample_count(), 1, "a stale bucket and a regressed version must both be rejected");
91 assert_eq!(epoch.floor_version_at(9_999), Some(20));
92 }
93
94 #[test]
95 fn record_keeps_highest_version_within_a_bucket() {
96 let epoch = VersionEpoch::new();
100 epoch.record(100, 5);
101 epoch.record(100, 9);
102 epoch.record(100, 7);
103
104 assert_eq!(epoch.sample_count(), 1, "one bucket holds a single sample");
105 assert_eq!(epoch.floor_version_at(100), Some(9), "the highest version committed at this instant wins");
106 assert_eq!(epoch.floor_version_at(9_999), Some(9));
107 }
108}