1use crate::compaction::merge_iter::MergingIterator;
8use crate::error::{Error, Result};
9use crate::lsm::sstable::SSTableEntry;
10use crate::types::Key;
11
12#[derive(Debug, Clone, Copy, PartialEq, Eq)]
14pub struct LeveledCompactionConfig {
15 pub l0_compaction_trigger: usize,
17 pub level_size_multiplier: usize,
19 pub max_levels: usize,
21 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#[derive(Debug, Clone, PartialEq, Eq)]
38pub struct KeyRange {
39 pub first_key: Key,
41 pub last_key: Key,
43}
44
45impl KeyRange {
46 pub fn overlaps(&self, other: &KeyRange) -> bool {
48 !(self.last_key < other.first_key || other.last_key < self.first_key)
49 }
50}
51
52#[derive(Debug, Clone, PartialEq, Eq)]
54pub struct SSTableMeta {
55 pub id: u64,
57 pub level: usize,
59 pub size_bytes: u64,
61 pub key_range: KeyRange,
63}
64
65#[derive(Debug, Clone, PartialEq, Eq)]
67pub struct CompactionPlan {
68 pub input_level: usize,
70 pub output_level: usize,
72 pub input_ids: Vec<u64>,
74 pub overlapping_output_ids: Vec<u64>,
76}
77
78#[derive(Debug, Clone)]
80pub struct LeveledCompaction {
81 config: LeveledCompactionConfig,
82}
83
84impl LeveledCompaction {
85 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 pub fn pick_plan(&self, levels: &[Vec<SSTableMeta>]) -> Option<CompactionPlan> {
112 if levels.len() < 2 {
113 return None;
114 }
115
116 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 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 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 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 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 let key = entry.key.clone();
204 if entry.value.is_none() {
205 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 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 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}