reifydb_store_multi/store/
drop.rs1use reifydb_codec::key::encoded::EncodedKey;
5use reifydb_core::{common::CommitVersion, interface::store::EntryKind};
6
7use crate::{Result, tier::TierStorage};
8
9#[derive(Debug, Clone)]
10pub struct DropEntry {
11 pub key: EncodedKey,
12
13 pub version: CommitVersion,
14
15 pub value_bytes: u64,
16}
17
18pub(crate) fn find_keys_to_drop<S: TierStorage>(
19 storage: &S,
20 table: EntryKind,
21 key: &[u8],
22 pending_version: Option<CommitVersion>,
23) -> Result<Vec<DropEntry>> {
24 let all_versions = storage.get_all_versions(table, key)?;
25
26 let mut versioned_entries: Vec<(CommitVersion, u64)> = all_versions
27 .into_iter()
28 .map(|(version, value)| {
29 let value_bytes = value.as_ref().map(|v| v.len() as u64).unwrap_or(0);
30 (version, value_bytes)
31 })
32 .collect();
33
34 if let Some(pending_ver) = pending_version
35 && !versioned_entries.iter().any(|(v, _)| *v == pending_ver)
36 {
37 versioned_entries.push((pending_ver, 0));
38 }
39
40 versioned_entries.sort_by(|a, b| b.0.cmp(&a.0));
41
42 let mut entries_to_drop = Vec::with_capacity(versioned_entries.len().saturating_sub(1));
43 let drop_key = EncodedKey::new(key.to_vec());
44
45 for (idx, (entry_version, value_bytes)) in versioned_entries.into_iter().enumerate() {
46 let should_drop = idx > 0;
47
48 if should_drop {
49 if Some(entry_version) == pending_version {
50 continue;
51 }
52
53 entries_to_drop.push(DropEntry {
54 key: drop_key.clone(),
55 version: entry_version,
56 value_bytes,
57 });
58 }
59 }
60
61 Ok(entries_to_drop)
62}
63
64#[cfg(test)]
65pub mod tests {
66 use std::collections::HashMap;
67
68 use reifydb_value::util::cowvec::CowVec;
69
70 use super::*;
71 use crate::tier::commit::buffer::MultiCommitBufferTier;
72
73 fn setup_versioned_entries(storage: &MultiCommitBufferTier, table: EntryKind, key: &[u8], versions: &[u64]) {
75 for v in versions {
76 let entries = vec![(EncodedKey::new(key.to_vec()), Some(CowVec::new(vec![*v as u8])))];
77 storage.set(CommitVersion(*v), HashMap::from([(table, entries)])).unwrap();
78 }
79 }
80
81 fn extract_dropped_versions(entries: &[DropEntry]) -> Vec<u64> {
83 entries.iter().map(|e| e.version.0).collect()
84 }
85
86 #[test]
87 fn test_drop_historical_versions() {
88 let storage = MultiCommitBufferTier::memory();
89 let table = EntryKind::Multi;
90 let key = b"test_key";
91
92 setup_versioned_entries(&storage, table, key, &[1, 5, 10, 20, 100]);
94
95 let to_drop = find_keys_to_drop(&storage, table, key, None).unwrap();
97
98 assert_eq!(to_drop.len(), 4);
99 let versions = extract_dropped_versions(&to_drop);
100 assert!(versions.contains(&1));
101 assert!(versions.contains(&5));
102 assert!(versions.contains(&10));
103 assert!(versions.contains(&20));
104 assert!(!versions.contains(&100));
105 }
106
107 #[test]
108 fn test_keep_latest_with_pending() {
109 let storage = MultiCommitBufferTier::memory();
110 let table = EntryKind::Multi;
111 let key = b"test_key";
112
113 setup_versioned_entries(&storage, table, key, &[1, 5, 10]);
115
116 let to_drop = find_keys_to_drop(&storage, table, key, Some(CommitVersion(20))).unwrap();
118
119 assert_eq!(to_drop.len(), 3);
120 let versions = extract_dropped_versions(&to_drop);
121 assert!(versions.contains(&1));
122 assert!(versions.contains(&5));
123 assert!(versions.contains(&10));
124 assert!(!versions.contains(&20));
125 }
126
127 #[test]
128 fn test_single_version_no_drop() {
129 let storage = MultiCommitBufferTier::memory();
130 let table = EntryKind::Multi;
131 let key = b"test_key";
132
133 setup_versioned_entries(&storage, table, key, &[42]);
134
135 let to_drop = find_keys_to_drop(&storage, table, key, None).unwrap();
137 assert!(to_drop.is_empty());
138 }
139
140 #[test]
141 fn test_empty_storage() {
142 let storage = MultiCommitBufferTier::memory();
143 let table = EntryKind::Multi;
144 let key = b"nonexistent";
145
146 let to_drop = find_keys_to_drop(&storage, table, key, None).unwrap();
147 assert!(to_drop.is_empty());
148 }
149}