Skip to main content

alopex_core/compaction/
leveled.rs

1//! Leveled compaction planning and a reference in-memory compaction routine.
2//!
3//! This is an algorithmic implementation that is usable without the full single-file
4//! in-place compaction plumbing. It selects compaction candidates based on the design rules
5//! in spec §3.4.2/§3.4.3 and can merge SSTable entry streams into new SSTable runs.
6
7use crate::compaction::merge_iter::MergingIterator;
8use crate::error::{Error, Result};
9use crate::lsm::sstable::SSTableEntry;
10use crate::types::Key;
11
12/// Leveled compaction configuration.
13#[derive(Debug, Clone, Copy, PartialEq, Eq)]
14pub struct LeveledCompactionConfig {
15    /// L0 → L1 trigger count (default: 4).
16    pub l0_compaction_trigger: usize,
17    /// Size multiplier for each level (default: 10).
18    pub level_size_multiplier: usize,
19    /// Maximum number of levels (default: 7).
20    pub max_levels: usize,
21    /// Target SSTable file size (bytes, default: 64MB).
22    pub target_file_size: usize,
23}
24
25impl Default for LeveledCompactionConfig {
26    fn default() -> Self {
27        Self {
28            l0_compaction_trigger: 4,
29            level_size_multiplier: 10,
30            max_levels: 7,
31            target_file_size: 64 * 1024 * 1024,
32        }
33    }
34}
35
36/// Inclusive key range for an SSTable.
37#[derive(Debug, Clone, PartialEq, Eq)]
38pub struct KeyRange {
39    /// Smallest key in the table.
40    pub first_key: Key,
41    /// Largest key in the table.
42    pub last_key: Key,
43}
44
45impl KeyRange {
46    /// Returns true if the ranges overlap.
47    pub fn overlaps(&self, other: &KeyRange) -> bool {
48        !(self.last_key < other.first_key || other.last_key < self.first_key)
49    }
50}
51
52/// Metadata for an SSTable used by the compaction planner.
53#[derive(Debug, Clone, PartialEq, Eq)]
54pub struct SSTableMeta {
55    /// Stable identifier (logical).
56    pub id: u64,
57    /// Level number (0..max_levels-1).
58    pub level: usize,
59    /// Approximate size in bytes.
60    pub size_bytes: u64,
61    /// Key range covered by this table.
62    pub key_range: KeyRange,
63}
64
65/// A compaction plan selecting input tables and output level.
66#[derive(Debug, Clone, PartialEq, Eq)]
67pub struct CompactionPlan {
68    /// Input level.
69    pub input_level: usize,
70    /// Output level (= input_level + 1).
71    pub output_level: usize,
72    /// SSTable IDs selected from input_level.
73    pub input_ids: Vec<u64>,
74    /// SSTable IDs selected from output_level that overlap input key-range.
75    pub overlapping_output_ids: Vec<u64>,
76}
77
78/// Leveled compaction planner + reference compaction routine.
79#[derive(Debug, Clone)]
80pub struct LeveledCompaction {
81    config: LeveledCompactionConfig,
82}
83
84impl LeveledCompaction {
85    /// Create a new compaction planner.
86    pub fn new(config: LeveledCompactionConfig) -> Result<Self> {
87        if config.max_levels < 2 {
88            return Err(Error::InvalidFormat("max_levels must be >= 2".into()));
89        }
90        if config.l0_compaction_trigger == 0 {
91            return Err(Error::InvalidFormat(
92                "l0_compaction_trigger must be >= 1".into(),
93            ));
94        }
95        if config.level_size_multiplier < 2 {
96            return Err(Error::InvalidFormat(
97                "level_size_multiplier must be >= 2".into(),
98            ));
99        }
100        if config.target_file_size == 0 {
101            return Err(Error::InvalidFormat("target_file_size must be > 0".into()));
102        }
103        Ok(Self { config })
104    }
105
106    /// Pick a compaction plan from per-level SSTable metadata.
107    ///
108    /// The planner chooses:
109    /// - L0 → L1 when `L0.len() >= l0_compaction_trigger`.
110    /// - Otherwise, the first level `n` where `size(Ln) > target_file_size * multiplier^n`.
111    pub fn pick_plan(&self, levels: &[Vec<SSTableMeta>]) -> Option<CompactionPlan> {
112        if levels.len() < 2 {
113            return None;
114        }
115
116        // L0 trigger.
117        let l0 = &levels[0];
118        if l0.len() >= self.config.l0_compaction_trigger {
119            let input_ids: Vec<u64> = l0.iter().map(|m| m.id).collect();
120            let input_range = merge_key_ranges(l0.iter().map(|m| &m.key_range));
121            let overlaps = levels[1]
122                .iter()
123                .filter(|m| m.key_range.overlaps(&input_range))
124                .map(|m| m.id)
125                .collect();
126            return Some(CompactionPlan {
127                input_level: 0,
128                output_level: 1,
129                input_ids,
130                overlapping_output_ids: overlaps,
131            });
132        }
133
134        // Size trigger for levels 1..max_levels-2.
135        let max = self.config.max_levels.min(levels.len());
136        for level in 1..max.saturating_sub(1) {
137            let total: u64 = levels[level].iter().map(|m| m.size_bytes).sum();
138            let limit = self.level_size_limit(level);
139            if total > limit {
140                // TODO: choose a better candidate selection strategy (oldest run, smallest overlap,
141                // largest size, etc.). For now we select the first table as a deterministic
142                // placeholder.
143                let input_meta = levels[level].first()?;
144                let input_ids = vec![input_meta.id];
145                let overlaps = levels[level + 1]
146                    .iter()
147                    .filter(|m| m.key_range.overlaps(&input_meta.key_range))
148                    .map(|m| m.id)
149                    .collect();
150                return Some(CompactionPlan {
151                    input_level: level,
152                    output_level: level + 1,
153                    input_ids,
154                    overlapping_output_ids: overlaps,
155                });
156            }
157        }
158
159        None
160    }
161
162    fn level_size_limit(&self, level: usize) -> u64 {
163        // Interpret "base_size" as target_file_size for level 1, then grow by multiplier^level.
164        let mult = self.config.level_size_multiplier as u64;
165        let base = self.config.target_file_size as u64;
166        let pow = mult.saturating_pow(level as u32);
167        base.saturating_mul(pow)
168    }
169
170    /// Compact multiple sorted entry streams and return output runs split by `target_file_size`.
171    ///
172    /// If `output_level` is the last level (`max_levels-1`), tombstones and older versions for that
173    /// key are removed completely (spec §3.4.3).
174    pub fn compact_entries(
175        &self,
176        sources: Vec<Box<dyn Iterator<Item = SSTableEntry>>>,
177        output_level: usize,
178    ) -> Result<Vec<Vec<SSTableEntry>>> {
179        if output_level >= self.config.max_levels {
180            return Err(Error::InvalidFormat("output_level out of range".into()));
181        }
182
183        let mut out_files: Vec<Vec<SSTableEntry>> = Vec::new();
184        let mut current: Vec<SSTableEntry> = Vec::new();
185        let mut current_bytes: usize = 0;
186
187        let mut iter = MergingIterator::new(sources).peekable();
188        let drop_tombstones = output_level + 1 >= self.config.max_levels;
189
190        while let Some(entry) = iter.next() {
191            if !drop_tombstones {
192                push_splitting(
193                    &mut out_files,
194                    &mut current,
195                    &mut current_bytes,
196                    entry,
197                    self.config.target_file_size,
198                );
199                continue;
200            }
201
202            // Bottom-level: if the newest visible version is a tombstone, drop the entire key.
203            let key = entry.key.clone();
204            if entry.value.is_none() {
205                // Skip all remaining versions for this key.
206                while let Some(next) = iter.peek() {
207                    if next.key == key {
208                        let _ = iter.next();
209                    } else {
210                        break;
211                    }
212                }
213                continue;
214            }
215
216            push_splitting(
217                &mut out_files,
218                &mut current,
219                &mut current_bytes,
220                entry,
221                self.config.target_file_size,
222            );
223        }
224
225        if !current.is_empty() {
226            out_files.push(current);
227        }
228        Ok(out_files)
229    }
230}
231
232fn merge_key_ranges<'a, I>(ranges: I) -> KeyRange
233where
234    I: IntoIterator<Item = &'a KeyRange>,
235{
236    let mut it = ranges.into_iter();
237    let first = it.next().expect("non-empty ranges");
238    let mut min = first.first_key.clone();
239    let mut max = first.last_key.clone();
240    for r in it {
241        if r.first_key < min {
242            min = r.first_key.clone();
243        }
244        if r.last_key > max {
245            max = r.last_key.clone();
246        }
247    }
248    KeyRange {
249        first_key: min,
250        last_key: max,
251    }
252}
253
254fn estimate_entry_size(entry: &SSTableEntry) -> usize {
255    // Rough estimate for splitting; not the actual on-disk size.
256    let val_len = entry.value.as_ref().map(|v| v.len()).unwrap_or(0);
257    1 + 8 + 8 + 4 + 4 + entry.key.len() + val_len
258}
259
260fn push_splitting(
261    out_files: &mut Vec<Vec<SSTableEntry>>,
262    current: &mut Vec<SSTableEntry>,
263    current_bytes: &mut usize,
264    entry: SSTableEntry,
265    target_bytes: usize,
266) {
267    let size = estimate_entry_size(&entry);
268    if !current.is_empty() && current_bytes.saturating_add(size) > target_bytes {
269        out_files.push(std::mem::take(current));
270        *current_bytes = 0;
271    }
272    *current_bytes = current_bytes.saturating_add(size);
273    current.push(entry);
274}
275
276#[cfg(all(test, not(target_arch = "wasm32")))]
277mod tests {
278    use super::*;
279
280    fn meta(id: u64, level: usize, size: u64, first: &[u8], last: &[u8]) -> SSTableMeta {
281        SSTableMeta {
282            id,
283            level,
284            size_bytes: size,
285            key_range: KeyRange {
286                first_key: first.to_vec(),
287                last_key: last.to_vec(),
288            },
289        }
290    }
291
292    fn put(key: &[u8], value: &[u8], ts: u64, seq: u64) -> SSTableEntry {
293        SSTableEntry {
294            key: key.to_vec(),
295            value: Some(value.to_vec()),
296            timestamp: ts,
297            sequence: seq,
298        }
299    }
300
301    fn del(key: &[u8], ts: u64, seq: u64) -> SSTableEntry {
302        SSTableEntry {
303            key: key.to_vec(),
304            value: None,
305            timestamp: ts,
306            sequence: seq,
307        }
308    }
309
310    #[test]
311    fn picks_l0_trigger() {
312        let comp = LeveledCompaction::new(LeveledCompactionConfig {
313            l0_compaction_trigger: 2,
314            ..Default::default()
315        })
316        .unwrap();
317        let levels = vec![
318            vec![meta(1, 0, 10, b"a", b"b"), meta(2, 0, 10, b"c", b"d")],
319            vec![meta(10, 1, 10, b"b", b"c"), meta(11, 1, 10, b"x", b"z")],
320        ];
321        let plan = comp.pick_plan(&levels).unwrap();
322        assert_eq!(plan.input_level, 0);
323        assert_eq!(plan.output_level, 1);
324        assert_eq!(plan.input_ids, vec![1, 2]);
325        assert_eq!(plan.overlapping_output_ids, vec![10]);
326    }
327
328    #[test]
329    fn picks_size_trigger_for_l1() {
330        let comp = LeveledCompaction::new(LeveledCompactionConfig {
331            l0_compaction_trigger: 100,
332            target_file_size: 10,
333            level_size_multiplier: 2,
334            max_levels: 4,
335        })
336        .unwrap();
337        let levels = vec![
338            vec![],
339            vec![meta(5, 1, 100, b"a", b"z")],
340            vec![meta(6, 2, 10, b"m", b"n")],
341            vec![],
342        ];
343        let plan = comp.pick_plan(&levels).unwrap();
344        assert_eq!(plan.input_level, 1);
345        assert_eq!(plan.output_level, 2);
346        assert_eq!(plan.input_ids, vec![5]);
347        assert_eq!(plan.overlapping_output_ids, vec![6]);
348    }
349
350    #[test]
351    fn bottom_level_drops_tombstone_and_older_versions() {
352        let comp = LeveledCompaction::new(LeveledCompactionConfig {
353            max_levels: 3,
354            target_file_size: 1024,
355            ..Default::default()
356        })
357        .unwrap();
358        let sources: Vec<Box<dyn Iterator<Item = SSTableEntry>>> = vec![
359            Box::new(
360                vec![
361                    del(b"k", 20, 1),
362                    put(b"k", b"old", 10, 1),
363                    put(b"x", b"v", 1, 1),
364                ]
365                .into_iter(),
366            ),
367            Box::new(vec![put(b"a", b"1", 5, 1)].into_iter()),
368        ];
369        // output_level=2 is bottom for max_levels=3.
370        let out = comp.compact_entries(sources, 2).unwrap();
371        let flat: Vec<_> = out.into_iter().flatten().collect();
372        assert!(flat.iter().all(|e| e.key != b"k".to_vec()));
373        assert!(flat.iter().any(|e| e.key == b"x".to_vec()));
374    }
375}