Skip to main content

hightower_kv/
compactor.rs

1use std::sync::Arc;
2use std::time::Duration;
3
4use crate::error::Result;
5use crate::storage::{CompactionOptions, Storage};
6
7/// Coordinates background segment merges to keep disk usage and read
8/// amplification under control.
9///
10/// The intended lifecycle for a single compaction run is:
11/// 1. Inspect storage metadata and pick cold segments whose combined size
12///    exceeds `min_bytes` or whose live-data ratio falls below a future
13///    heuristic. We cap the batch using `max_segments` to avoid long stalls.
14/// 2. Stream the chosen segments in log order, retaining the newest value per
15///    key while skipping tombstones older than `tombstone_grace`.
16/// 3. Write the surviving entries into a fresh segment, fsync it, rebuild its
17///    sparse index/Bloom filter, and atomically swap it into Storage while
18///    removing the old segment files.
19/// 4. Optionally trigger snapshotting so restart time stays stable even as log
20///    history churns.
21/// Configuration for compaction behavior
22#[derive(Debug, Clone)]
23pub struct CompactionConfig {
24    /// Minimum bytes needed before compaction runs
25    pub min_bytes: u64,
26    /// Maximum number of segments to compact at once
27    pub max_segments: usize,
28    /// How long to keep tombstones before eviction
29    pub tombstone_grace: Duration,
30    /// Whether to emit a snapshot after compaction
31    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/// Compactor that merges log segments to reduce disk usage and read amplification
46#[derive(Debug, Clone)]
47pub struct Compactor {
48    storage: Arc<Storage>,
49    config: CompactionConfig,
50}
51
52impl Compactor {
53    /// Creates a new compactor with the given storage and configuration
54    pub fn new(storage: Arc<Storage>, config: CompactionConfig) -> Self {
55        Self { storage, config }
56    }
57
58    /// Runs a single compaction cycle if conditions are met
59    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    /// Returns a clone of the storage Arc
85    pub fn storage(&self) -> Arc<Storage> {
86        Arc::clone(&self.storage)
87    }
88
89    /// Returns a reference to the compaction configuration
90    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}