hightower_kv/
compactor.rs1use std::sync::Arc;
2use std::time::Duration;
3
4use crate::error::Result;
5use crate::storage::{CompactionOptions, Storage};
6
7#[derive(Debug, Clone)]
23pub struct CompactionConfig {
24 pub min_bytes: u64,
26 pub max_segments: usize,
28 pub tombstone_grace: Duration,
30 pub emit_snapshot: bool,
32}
33
34impl Default for CompactionConfig {
35 fn default() -> Self {
36 Self {
37 min_bytes: 32 * 1024 * 1024,
38 max_segments: 4,
39 tombstone_grace: Duration::from_secs(300),
40 emit_snapshot: false,
41 }
42 }
43}
44
45#[derive(Debug, Clone)]
47pub struct Compactor {
48 storage: Arc<Storage>,
49 config: CompactionConfig,
50}
51
52impl Compactor {
53 pub fn new(storage: Arc<Storage>, config: CompactionConfig) -> Self {
55 Self { storage, config }
56 }
57
58 pub fn run_once(&self) -> Result<()> {
60 let sealed = self.storage.sealed_segments_snapshot();
61 if sealed.is_empty() {
62 return Ok(());
63 }
64
65 let total_bytes: u64 = sealed.iter().map(|segment| segment.bytes_written()).sum();
66 if total_bytes < self.config.min_bytes {
67 return Ok(());
68 }
69
70 let options = CompactionOptions {
71 tombstone_grace: self.config.tombstone_grace,
72 min_bytes: self.config.min_bytes,
73 max_segments: self.config.max_segments,
74 emit_snapshot: self.config.emit_snapshot,
75 };
76 if !self.storage.compact_all(options)? {
77 return Ok(());
78 }
79
80 self.storage.sync()?;
81 Ok(())
82 }
83
84 pub fn storage(&self) -> Arc<Storage> {
86 Arc::clone(&self.storage)
87 }
88
89 pub fn config(&self) -> &CompactionConfig {
91 &self.config
92 }
93}
94
95#[cfg(test)]
96mod tests {
97 use super::*;
98 use crate::command::Command;
99 use crate::config::StoreConfig;
100 use std::path::PathBuf;
101 use tempfile::tempdir;
102
103 #[test]
104 fn run_once_is_noop_when_under_threshold() {
105 let temp = tempdir().unwrap();
106 let mut cfg = StoreConfig::default();
107 cfg.data_dir = temp.path().join("compactor").to_string_lossy().into_owned();
108 cfg.max_segment_size = 1_000_000;
109 let storage = Arc::new(Storage::new(&cfg).unwrap());
110 let compactor = Compactor::new(storage.clone(), CompactionConfig::default());
111 compactor.run_once().unwrap();
112 assert_eq!(storage.segment_snapshot().len(), 1);
113 }
114
115 #[test]
116 fn compacts_multiple_segments_into_one() {
117 let temp = tempdir().unwrap();
118 let mut cfg = StoreConfig::default();
119 cfg.data_dir = temp.path().join("compactor").to_string_lossy().into_owned();
120 cfg.max_segment_size = 64;
121 let storage = Arc::new(Storage::new(&cfg).unwrap());
122 for i in 0..5 {
123 let command = Command::Set {
124 key: format!("key{i}").into_bytes(),
125 value: vec![b'x'; 32],
126 version: i + 1,
127 timestamp: i as i64,
128 };
129 storage.apply(&command).unwrap();
130 }
131 storage
132 .apply(&Command::Delete {
133 key: b"key1".to_vec(),
134 version: 99,
135 timestamp: 0,
136 })
137 .unwrap();
138 assert!(storage.segment_snapshot().len() > 1);
139
140 let sealed_before = storage.sealed_segments_snapshot();
141 assert!(sealed_before.len() >= 2);
142 let bytes_target: u64 = sealed_before
143 .iter()
144 .take(2)
145 .map(|segment| segment.bytes_written())
146 .sum();
147 let segments_to_merge = sealed_before.len().min(2);
148
149 let mut config = CompactionConfig::default();
150 config.min_bytes = bytes_target;
151 config.max_segments = 2;
152 let compactor = Compactor::new(storage.clone(), config);
153 compactor.run_once().unwrap();
154
155 let sealed_after = storage.sealed_segments_snapshot();
156 assert_eq!(
157 sealed_after.len(),
158 sealed_before.len() - segments_to_merge + 1
159 );
160 let entry = storage.lookup(b"key1").unwrap();
161 assert!(entry.is_tombstone);
162 let command = storage.fetch_command(&entry).unwrap().unwrap();
163 assert!(matches!(command, Command::Delete { .. }));
164 }
165
166 #[test]
167 fn tombstone_grace_evicts_old_deletes() {
168 let temp = tempdir().unwrap();
169 let mut cfg = StoreConfig::default();
170 cfg.data_dir = temp.path().join("compactor").to_string_lossy().into_owned();
171 cfg.max_segment_size = 64;
172 let storage = Arc::new(Storage::new(&cfg).unwrap());
173
174 storage
175 .apply(&Command::Set {
176 key: b"victim".to_vec(),
177 value: b"value".to_vec(),
178 version: 1,
179 timestamp: 1,
180 })
181 .unwrap();
182 storage
183 .apply(&Command::Delete {
184 key: b"victim".to_vec(),
185 version: 2,
186 timestamp: 0,
187 })
188 .unwrap();
189
190 for i in 0..5 {
191 storage
192 .apply(&Command::Set {
193 key: format!("filler{i}").into_bytes(),
194 value: vec![b'x'; 32],
195 version: 10 + i,
196 timestamp: 10 + i as i64,
197 })
198 .unwrap();
199 }
200
201 let victim_segment = storage.lookup(b"victim").unwrap().segment_id;
202 let sealed_before = storage.sealed_segments_snapshot();
203 let mut target_bytes = 0u64;
204 let mut segments_needed = 0usize;
205 for segment in &sealed_before {
206 segments_needed += 1;
207 target_bytes += segment.bytes_written();
208 if segment.id() == victim_segment {
209 break;
210 }
211 }
212
213 let mut config = CompactionConfig::default();
214 config.min_bytes = target_bytes;
215 config.max_segments = segments_needed;
216 config.tombstone_grace = Duration::from_secs(1);
217 let compactor = Compactor::new(storage.clone(), config);
218 compactor.run_once().unwrap();
219
220 assert!(storage.lookup(b"victim").is_none());
221 }
222
223 #[test]
224 fn respects_max_segment_limit() {
225 let temp = tempdir().unwrap();
226 let mut cfg = StoreConfig::default();
227 cfg.data_dir = temp.path().join("compactor").to_string_lossy().into_owned();
228 cfg.max_segment_size = 64;
229 let storage = Arc::new(Storage::new(&cfg).unwrap());
230
231 for i in 0..8 {
232 storage
233 .apply(&Command::Set {
234 key: format!("hot{i}").into_bytes(),
235 value: vec![b'x'; 32],
236 version: i + 1,
237 timestamp: i as i64,
238 })
239 .unwrap();
240 }
241
242 let sealed_before = storage.sealed_segments_snapshot();
243 assert!(sealed_before.len() >= 2);
244
245 let mut config = CompactionConfig::default();
246 config.min_bytes = 0;
247 config.max_segments = 1;
248 let compactor = Compactor::new(storage.clone(), config);
249 compactor.run_once().unwrap();
250
251 let sealed_after = storage.sealed_segments_snapshot();
252 let before_ids: Vec<u64> = sealed_before.iter().map(|segment| segment.id()).collect();
253 let after_ids: Vec<u64> = sealed_after.iter().map(|segment| segment.id()).collect();
254
255 let new_ids: Vec<u64> = after_ids
256 .iter()
257 .copied()
258 .filter(|id| !before_ids.contains(id))
259 .collect();
260 assert_eq!(new_ids.len(), 1);
261 assert!(!after_ids.contains(&before_ids[0]));
262 for id in before_ids.iter().skip(1) {
263 assert!(after_ids.contains(id));
264 }
265 }
266
267 #[test]
268 fn emit_snapshot_is_best_effort() {
269 let temp = tempdir().unwrap();
270 let mut cfg = StoreConfig::default();
271 cfg.data_dir = temp.path().join("compactor").to_string_lossy().into_owned();
272 cfg.max_segment_size = 64;
273 let storage = Arc::new(Storage::new(&cfg).unwrap());
274
275 for i in 0..8 {
276 storage
277 .apply(&Command::Set {
278 key: format!("snap{i}").into_bytes(),
279 value: vec![b'x'; 32],
280 version: i + 1,
281 timestamp: i as i64,
282 })
283 .unwrap();
284 }
285
286 let mut config = CompactionConfig::default();
287 config.min_bytes = 0;
288 config.emit_snapshot = true;
289 let snapshot_path = PathBuf::from(cfg.data_dir.clone()).join("snapshot.bin");
290 let compactor = Compactor::new(storage, config);
291 compactor.run_once().unwrap();
292
293 assert!(snapshot_path.exists());
294 }
295
296 #[test]
297 fn config_defaults_are_reasonable() {
298 let cfg = CompactionConfig::default();
299 assert!(cfg.min_bytes >= 4 * 1024 * 1024);
300 assert!(cfg.max_segments >= 1);
301 assert!(cfg.tombstone_grace >= Duration::from_secs(60));
302 }
303}