1use crate::engine::Table;
18use crate::epoch::{Epoch, MaintenanceReceipt, Snapshot};
19use crate::manifest::RunRef;
20use crate::memtable::Row;
21use crate::sorted_run::RunWriter;
22use crate::{ExecutionControl, MongrelError, Result};
23use std::collections::HashMap;
24use std::path::Path;
25
26impl Table {
27 pub const AUTO_COMPACT_RUN_THRESHOLD: usize = 8;
32
33 pub const MAX_L0_OVERLAPPING_RUNS: usize = 64;
41
42 pub fn should_compact(&self) -> bool {
48 if self.run_refs().len() >= Self::AUTO_COMPACT_RUN_THRESHOLD {
49 return true;
50 }
51 if self.l0_run_count() >= Self::MAX_L0_OVERLAPPING_RUNS {
55 return true;
56 }
57 self.ttl().is_some()
58 && !self.run_refs().is_empty()
59 && self.has_expired_run_rows().unwrap_or(false)
60 }
61
62 pub(crate) fn l0_run_count(&self) -> usize {
64 self.run_refs().iter().filter(|rr| rr.level == 0).count()
65 }
66
67 fn l1_range_overlaps(
74 &self,
75 min_row_id: u64,
76 max_row_id: u64,
77 retire: &std::collections::HashSet<u128>,
78 ) -> bool {
79 if min_row_id > max_row_id {
80 return false;
81 }
82 for rr in self.run_refs() {
83 if rr.level < 1 || retire.contains(&rr.run_id) {
84 continue;
85 }
86 if let Some((lo, hi)) = self.run_row_id_range(rr.run_id) {
87 if rr.row_count == 0 {
91 continue;
92 }
93 if min_row_id <= hi && lo <= max_row_id {
94 return true;
95 }
96 }
97 }
98 false
99 }
100
101 fn has_expired_run_rows(&self) -> Result<bool> {
102 self.has_expired_run_rows_inner(None)
103 }
104
105 fn has_expired_run_rows_inner(&self, control: Option<&ExecutionControl>) -> Result<bool> {
106 let now_nanos = crate::engine::unix_nanos_now();
107 for (run_index, run) in self.run_refs().iter().enumerate() {
108 if run_index % 256 == 0 {
109 if let Some(control) = control {
110 control.checkpoint()?;
111 }
112 }
113 let mut reader = self.open_reader(run.run_id)?;
114 for (row_index, row) in reader.all_rows()?.iter().enumerate() {
115 if row_index % 256 == 0 {
116 if let Some(control) = control {
117 control.checkpoint()?;
118 }
119 }
120 if self.row_expired_at(row, now_nanos) {
121 return Ok(true);
122 }
123 }
124 }
125 Ok(false)
126 }
127
128 pub fn maybe_compact(&mut self) -> Result<bool> {
134 if !self.should_compact() {
135 return Ok(false);
136 }
137 self.compact()?;
138 Ok(true)
139 }
140
141 pub fn compact(&mut self) -> Result<()> {
146 let control = ExecutionControl::new(None);
147 self.compact_controlled(&control, || true).map(|_| ())
148 }
149
150 #[doc(hidden)]
154 pub fn compact_controlled<F>(
155 &mut self,
156 control: &ExecutionControl,
157 before_publish: F,
158 ) -> Result<bool>
159 where
160 F: FnOnce() -> bool,
161 {
162 self.compact_controlled_with_receipt(control, before_publish)
163 .map(|(changed, _)| changed)
164 }
165
166 #[doc(hidden)]
169 pub fn compact_controlled_with_receipt<F>(
170 &mut self,
171 control: &ExecutionControl,
172 before_publish: F,
173 ) -> Result<(bool, Option<MaintenanceReceipt>)>
174 where
175 F: FnOnce() -> bool,
176 {
177 control.checkpoint()?;
178 let maintenance_epoch = self.current_epoch();
179 let reclaim_ttl = self.ttl().is_some() && self.has_expired_run_rows_inner(Some(control))?;
180 if self.run_refs().len() < 2 && !reclaim_ttl {
181 return Ok((false, None));
182 }
183 let min_active = self.min_active_snapshot();
184 let old_refs: Vec<RunRef> = self.run_refs().to_vec();
185 let now_nanos = crate::engine::unix_nanos_now();
186 let mutable_rows = if self.mutable_run_len() > 0 {
187 self.snapshot_mutable_run()
188 } else {
189 Vec::new()
190 };
191
192 let mut all: HashMap<u64, Vec<Row>> = HashMap::new();
193 let mut scanned = 0_usize;
194 for rr in &old_refs {
195 control.checkpoint()?;
196 let mut reader = self.open_reader(rr.run_id)?;
197 for row in reader.all_rows()? {
198 if scanned.is_multiple_of(256) {
199 control.checkpoint()?;
200 }
201 scanned += 1;
202 all.entry(row.row_id.0).or_default().push(row);
203 }
204 }
205 for row in mutable_rows {
206 if scanned.is_multiple_of(256) {
207 control.checkpoint()?;
208 }
209 scanned += 1;
210 all.entry(row.row_id.0).or_default().push(row);
211 }
212
213 let mut rows = Vec::new();
214 let mut current_live_count = 0u64;
215 for (row_index, (_, mut versions)) in all.into_iter().enumerate() {
216 if row_index % 256 == 0 {
217 control.checkpoint()?;
218 }
219 versions.sort_by_key(|row| row.committed_epoch);
223 let Some(newest) = versions
224 .iter()
225 .max_by(|a, b| cmp_version_order(a, b))
226 .cloned()
227 else {
228 continue;
229 };
230 let newest_epoch = newest.committed_epoch;
231 if !newest.deleted && !self.row_expired_at(&newest, now_nanos) {
232 current_live_count += 1;
233 }
234 for row in select_keep(&versions, min_active) {
235 if self.row_expired_at(&row, now_nanos) {
236 if is_same_version(&row, &newest)
237 && min_active.is_some_and(|epoch| newest_epoch > epoch)
238 {
239 let mut tombstone = row;
240 tombstone.deleted = true;
241 tombstone.columns.clear();
242 rows.push(tombstone);
243 }
244 } else {
245 rows.push(row);
246 }
247 }
248 }
249 rows.sort_by_key(|row| (row.row_id, row.committed_epoch));
250
251 let mut staged_run = None;
252 if !rows.is_empty() {
253 let retire: std::collections::HashSet<u128> =
263 old_refs.iter().map(|rr| rr.run_id).collect();
264 let new_min = rows.first().map(|r| r.row_id.0).unwrap_or(0);
265 let new_max = rows.last().map(|r| r.row_id.0).unwrap_or(0);
266 if self.l1_range_overlaps(new_min, new_max, &retire) {
267 return Err(MongrelError::InvalidArgument(format!(
268 "compaction would publish an L1 run overlapping an existing L1+ \
269 range [{new_min}, {new_max}]; wait for prior compactions to retire"
270 )));
271 }
272
273 let run_id = self.alloc_run_id()?;
274 let final_name = format!("r-{run_id}.sr");
275 let stage_name = format!(
276 "{final_name}.compact-stage-{}-{}",
277 std::process::id(),
278 std::time::SystemTime::now()
279 .duration_since(std::time::UNIX_EPOCH)
280 .unwrap_or_default()
281 .as_nanos()
282 );
283 let kek = self.kek();
284 let mut writer = RunWriter::new(self.schema(), run_id as u128, maintenance_epoch, 1)
285 .clean(min_active.is_none())
286 .with_zstd_level(self.compaction_zstd_level());
287 if let Some(k) = &kek {
288 writer = writer.with_encryption(k.as_ref(), self.indexable_column_specs());
289 }
290 let header = match self.create_run_entry(Path::new(&stage_name))? {
291 Some(file) => writer.write_file(file, &rows),
292 None => writer.write(self.runs_dir().join(&stage_name), &rows),
293 };
294 let header = match header {
295 Ok(header) => header,
296 Err(error) => {
297 let _ = self.remove_run_entry(Path::new(&stage_name));
298 return Err(error);
299 }
300 };
301 staged_run = Some((
302 stage_name,
303 final_name,
304 RunRef {
305 run_id: run_id as u128,
306 level: 1,
307 epoch_created: header.epoch_created,
308 row_count: header.row_count,
309 },
310 ));
311 }
312
313 if let Err(error) = control.checkpoint() {
314 if let Some((stage_name, _, _)) = staged_run {
315 let _ = self.remove_run_entry(Path::new(&stage_name));
316 }
317 return Err(error);
318 }
319 if !before_publish() {
320 if let Some((stage_name, _, _)) = staged_run {
321 let _ = self.remove_run_entry(Path::new(&stage_name));
322 }
323 return Err(MongrelError::Cancelled);
324 }
325
326 let replacement_ref = if let Some((stage_name, final_name, staged_ref)) = staged_run {
330 self.publish_run_entry(Path::new(&stage_name), Path::new(&final_name))?;
331 Some(staged_ref)
332 } else {
333 None
334 };
335
336 if self.mutable_run_len() > 0 {
337 self.drain_mutable_run();
338 }
339 self.live_count = current_live_count;
340 let retire_epoch = maintenance_epoch.0;
341 if let Some(replacement_ref) = replacement_ref {
342 self.set_run_refs(vec![replacement_ref]);
343 for run in &old_refs {
344 self.retire_run(run.run_id, retire_epoch);
345 }
346 } else {
347 self.set_run_refs(Vec::new());
348 for run in &old_refs {
349 self.retire_run(run.run_id, retire_epoch);
350 }
351 }
352
353 self.prepare_indexes_for_run_replacement();
358 if let Err(error) = self.persist_manifest(maintenance_epoch) {
359 self.poison_after_maintenance_publish_failure();
360 return Err(MongrelError::CommitOutcomeUnknown {
361 epoch: maintenance_epoch.0,
362 message: format!("compaction manifest publication failed: {error}"),
363 });
364 }
365 self.clear_result_cache();
366 self.bump_data_generation();
367
368 if let Err(error) = self
372 .rebuild_indexes_from_runs()
373 .and_then(|_| self.build_learned_ranges())
374 {
375 return Err(MongrelError::DurableCommit {
376 epoch: maintenance_epoch.0,
377 message: format!("compaction committed but index rebuild failed: {error}"),
378 });
379 }
380 self.finish_indexes_for_run_replacement();
381 self.checkpoint_indexes(maintenance_epoch);
382 Ok((
383 true,
384 Some(MaintenanceReceipt {
385 epoch: maintenance_epoch,
386 }),
387 ))
388 }
389}
390
391fn cmp_version_order(a: &Row, b: &Row) -> std::cmp::Ordering {
393 if Snapshot::version_is_newer(
394 a.committed_epoch,
395 a.commit_ts,
396 b.committed_epoch,
397 b.commit_ts,
398 ) {
399 std::cmp::Ordering::Greater
400 } else if Snapshot::version_is_newer(
401 b.committed_epoch,
402 b.commit_ts,
403 a.committed_epoch,
404 a.commit_ts,
405 ) {
406 std::cmp::Ordering::Less
407 } else {
408 a.committed_epoch
409 .cmp(&b.committed_epoch)
410 .then_with(|| a.commit_ts.cmp(&b.commit_ts))
411 }
412}
413
414fn is_same_version(a: &Row, b: &Row) -> bool {
415 a.committed_epoch == b.committed_epoch && a.commit_ts == b.commit_ts
416}
417
418fn select_keep(vers: &[Row], min_active: Option<Epoch>) -> Vec<Row> {
428 let Some(newest) = vers.iter().max_by(|a, b| cmp_version_order(a, b)).cloned() else {
429 return Vec::new();
430 };
431 match min_active {
432 None => {
433 if newest.deleted {
434 Vec::new()
435 } else {
436 vec![newest]
437 }
438 }
439 Some(min_e) => {
440 let recent_start = vers.partition_point(|row| row.committed_epoch < min_e);
441 let mut keep = vers[recent_start..].to_vec();
442 if recent_start > 0 && keep.first().is_none_or(|row| row.committed_epoch > min_e) {
443 let boundary = vers[recent_start - 1].clone();
444 if keep.is_empty() && boundary.deleted {
445 return Vec::new();
446 }
447 keep.insert(0, boundary);
448 }
449 keep
450 }
451 }
452}
453
454#[cfg(test)]
455mod tests {
456 use super::*;
457 use crate::schema::{ColumnDef, ColumnFlags, Schema, TypeId};
458 use crate::{Database, Snapshot, Value};
459 use tempfile::tempdir;
460
461 fn schema() -> Schema {
462 Schema {
463 schema_id: 1,
464 columns: vec![ColumnDef {
465 id: 1,
466 name: "v".into(),
467 ty: TypeId::Int64,
468 flags: ColumnFlags::empty().with(ColumnFlags::PRIMARY_KEY),
469 default_value: None,
470 embedding_source: None,
471 }],
472 indexes: Vec::new(),
473 colocation: vec![],
474 constraints: Default::default(),
475 clustered: false,
476 }
477 }
478
479 #[test]
480 fn compaction_merges_runs_and_gcs_tombstoned_row() {
481 let dir = tempdir().unwrap();
482 let mut db = Table::create(dir.path(), schema(), 1).unwrap();
483 db.set_mutable_run_spill_bytes(1);
485 let mut ids = Vec::new();
486 for i in 1..=5i64 {
487 ids.push(db.put(vec![(1, Value::Int64(i))]).unwrap());
488 }
489 db.flush().unwrap();
490 db.delete(ids[2]).unwrap();
491 db.flush().unwrap();
492 db.put(vec![(1, Value::Int64(60))]).unwrap();
493 db.flush().unwrap();
494 assert_eq!(db.run_count(), 3);
495
496 db.compact().unwrap();
497 assert_eq!(db.run_count(), 1);
498 let rows = db.visible_rows(db.snapshot()).unwrap();
499 let row_ids: Vec<u64> = rows.iter().map(|r| r.row_id.0).collect();
500 assert!(!row_ids.contains(&ids[2].0), "tombstoned row must be GC'd");
501 assert_eq!(rows.len(), 5);
502 }
503
504 #[test]
505 fn pinned_snapshot_survives_compaction() {
506 let dir = tempdir().unwrap();
507 let mut db = Table::create(dir.path(), schema(), 1).unwrap();
508 let r = db.put(vec![(1, Value::Int64(1))]).unwrap();
509 db.flush().unwrap(); let pinned = db.pin_snapshot();
513 assert_eq!(
514 db.get(r, pinned)
515 .and_then(|row| row.columns.get(&1).cloned()),
516 Some(Value::Int64(1))
517 );
518
519 db.delete(r).unwrap();
521 db.commit().unwrap();
522 db.flush().unwrap(); db.compact().unwrap(); assert_eq!(
527 db.get(r, pinned)
528 .and_then(|row| row.columns.get(&1).cloned()),
529 Some(Value::Int64(1))
530 );
531 assert_eq!(
533 db.get(r, db.snapshot())
534 .and_then(|row| row.columns.get(&1).cloned()),
535 None
536 );
537
538 db.unpin_snapshot(pinned);
540 db.compact().unwrap();
541 assert_eq!(
542 db.get(r, db.snapshot())
543 .and_then(|row| row.columns.get(&1).cloned()),
544 None
545 );
546 }
547
548 #[test]
549 fn controlled_compaction_cancel_before_publish_preserves_live_state() {
550 let dir = tempdir().unwrap();
551 let mut table = Table::create(dir.path(), schema(), 1).unwrap();
552 table.set_mutable_run_spill_bytes(1);
553 table.put(vec![(1, Value::Int64(1))]).unwrap();
554 table.flush().unwrap();
555 table.put(vec![(1, Value::Int64(2))]).unwrap();
556 table.flush().unwrap();
557 let before_refs: Vec<_> = table
558 .run_refs()
559 .iter()
560 .map(|run| (run.run_id, run.level, run.epoch_created, run.row_count))
561 .collect();
562
563 let error = table
564 .compact_controlled(&ExecutionControl::new(None), || false)
565 .unwrap_err();
566 assert!(matches!(error, MongrelError::Cancelled));
567 let after_refs: Vec<_> = table
568 .run_refs()
569 .iter()
570 .map(|run| (run.run_id, run.level, run.epoch_created, run.row_count))
571 .collect();
572 assert_eq!(after_refs, before_refs);
573 assert_eq!(table.visible_rows(table.snapshot()).unwrap().len(), 2);
574 assert!(std::fs::read_dir(table.runs_dir())
575 .unwrap()
576 .all(|entry| !entry
577 .unwrap()
578 .file_name()
579 .to_string_lossy()
580 .contains("compact-stage")));
581 }
582
583 #[test]
584 fn compaction_manifest_failure_poisons_standalone_table_until_reopen() {
585 let dir = tempdir().unwrap();
586 let mut table = Table::create(dir.path(), schema(), 1).unwrap();
587 table.set_mutable_run_spill_bytes(1);
588 table.put(vec![(1, Value::Int64(1))]).unwrap();
589 table.flush().unwrap();
590 table.put(vec![(1, Value::Int64(2))]).unwrap();
591 table.flush().unwrap();
592 let manifest = dir.path().join(crate::manifest::MANIFEST_FILENAME);
593 let saved_manifest = dir.path().join("_mf.saved");
594 std::fs::rename(&manifest, &saved_manifest).unwrap();
595 std::fs::create_dir(&manifest).unwrap();
596
597 let error = table.compact().unwrap_err();
598 assert!(matches!(error, MongrelError::CommitOutcomeUnknown { .. }));
599 assert_eq!(table.visible_rows(table.snapshot()).unwrap().len(), 2);
600 assert!(table
601 .put(vec![(1, Value::Int64(3))])
602 .unwrap_err()
603 .to_string()
604 .contains("reopen required"));
605
606 drop(table);
607 std::fs::remove_dir(&manifest).unwrap();
608 std::fs::rename(saved_manifest, manifest).unwrap();
609 let reopened = Table::open(dir.path()).unwrap();
610 assert_eq!(reopened.visible_rows(reopened.snapshot()).unwrap().len(), 2);
611 }
612
613 #[test]
614 fn compaction_manifest_failure_poisons_mounted_database() {
615 let dir = tempdir().unwrap();
616 let db = Database::create(dir.path()).unwrap();
617 let table_id = db.create_table("items", schema()).unwrap();
618 let handle = db.table("items").unwrap();
619 {
620 let mut table = handle.lock();
621 table.set_mutable_run_spill_bytes(1);
622 table.put(vec![(1, Value::Int64(1))]).unwrap();
623 table.flush().unwrap();
624 table.put(vec![(1, Value::Int64(2))]).unwrap();
625 table.flush().unwrap();
626 }
627 let manifest = dir
628 .path()
629 .join("tables")
630 .join(table_id.to_string())
631 .join(crate::manifest::MANIFEST_FILENAME);
632 let saved_manifest = manifest.with_extension("saved");
633 std::fs::rename(&manifest, &saved_manifest).unwrap();
634 std::fs::create_dir(&manifest).unwrap();
635
636 let error = handle.lock().compact().unwrap_err();
637 assert!(matches!(error, MongrelError::CommitOutcomeUnknown { .. }));
638 assert!(db
639 .create_table("blocked", schema())
640 .unwrap_err()
641 .to_string()
642 .contains("database poisoned"));
643 }
644
645 #[test]
646 fn _snapshot_import_used() {
647 let _ = Snapshot::at(Epoch(0));
648 }
649
650 fn hlc(physical_micros: u64) -> mongreldb_types::hlc::HlcTimestamp {
651 mongreldb_types::hlc::HlcTimestamp {
652 physical_micros,
653 logical: 0,
654 node_tiebreaker: 1,
655 }
656 }
657
658 fn stamped(id: u64, epoch: u64, ts: mongreldb_types::hlc::HlcTimestamp, v: i64) -> Row {
659 crate::memtable::Row::new_with_hlc(crate::rowid::RowId(id), Epoch(epoch), ts)
660 .with_column(1, Value::Int64(v))
661 }
662
663 #[test]
667 fn select_keep_preserves_hlc_visibility_under_inverted_epoch_order() {
668 let early = hlc(100);
669 let mid = hlc(150);
670 let late = hlc(200);
671 let mut versions = vec![
673 stamped(1, 1, late, 99), stamped(1, 2, mid, 50), stamped(1, 50, early, 1), ];
677 versions.sort_by_key(|row| row.committed_epoch);
678 let kept = select_keep(&versions, None);
680 assert_eq!(kept.len(), 1);
681 assert_eq!(kept[0].commit_ts, Some(late));
682 assert_eq!(kept[0].committed_epoch, Epoch(1));
683 assert_eq!(kept[0].columns.get(&1), Some(&Value::Int64(99)));
684
685 let kept = select_keep(&versions, Some(Epoch(2)));
689 assert!(kept.iter().any(|r| r.commit_ts == Some(mid)));
690 assert!(kept.iter().any(|r| r.commit_ts == Some(early)));
691 assert!(!kept.iter().any(|r| r.commit_ts == Some(late)));
692
693 let kept = select_keep(&versions, Some(Epoch(10)));
696 assert!(kept.iter().any(|r| r.commit_ts == Some(early)));
697 assert!(
698 kept.iter().any(|r| r.commit_ts == Some(mid)),
699 "boundary version below floor must be retained"
700 );
701 assert!(!kept.iter().any(|r| r.commit_ts == Some(late)));
702 }
703
704 #[test]
705 fn cmp_version_order_prefers_hlc_when_both_stamped() {
706 let early = stamped(1, 99, hlc(100), 1);
707 let late = stamped(1, 1, hlc(200), 2);
708 assert_eq!(cmp_version_order(&early, &late), std::cmp::Ordering::Less);
709 assert_eq!(
710 cmp_version_order(&late, &early),
711 std::cmp::Ordering::Greater
712 );
713 let a = crate::memtable::Row::new(crate::rowid::RowId(1), Epoch(1));
715 let b = crate::memtable::Row::new(crate::rowid::RowId(1), Epoch(3));
716 assert_eq!(cmp_version_order(&a, &b), std::cmp::Ordering::Less);
717 }
718}