Skip to main content

reifydb_runtime/
version_epoch.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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		// Several commits at the same wall-clock instant (e.g. a write and the flow processing it
97		// triggers) must collapse to the HIGHEST version, or a row written by the later same-instant
98		// commit would read as too young to ever expire.
99		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}