use std::collections::BTreeMap;
use mongreldb_core::columnar::NativeColumn;
use mongreldb_core::schema::{ColumnDef, ColumnFlags, Schema, TypeId};
use mongreldb_core::trace::QueryTrace;
use mongreldb_core::{
CancellationReason, Epoch, ExecutionControl, Memtable, MongrelError, Row, RowId, Snapshot,
Table, Value,
};
use mongreldb_types::hlc::HlcTimestamp;
use tempfile::tempdir;
macro_rules! emit_scan_metric {
($name:literal, $metric:expr, $unit:literal) => {
println!(
"{}",
serde_json::json!({
"test": $name,
"metric": $metric,
"unit": $unit,
})
)
};
}
fn schema() -> Schema {
Schema {
schema_id: 1,
columns: vec![
ColumnDef {
id: 1,
name: "id".into(),
ty: TypeId::Int64,
flags: ColumnFlags::empty().with(ColumnFlags::PRIMARY_KEY),
default_value: None,
embedding_source: None,
},
ColumnDef {
id: 2,
name: "value".into(),
ty: TypeId::Int64,
flags: ColumnFlags::empty(),
default_value: None,
embedding_source: None,
},
],
..Schema::default()
}
}
fn put(table: &mut Table, id: i64, value: i64) -> RowId {
table
.put(vec![(1, Value::Int64(id)), (2, Value::Int64(value))])
.unwrap()
}
fn bulk_load(table: &mut Table, rows: usize) {
table
.bulk_load_columns(vec![
(1, NativeColumn::int64_sequence(0, rows)),
(2, NativeColumn::int64_sequence(0, rows)),
])
.unwrap();
}
fn value(row: &Row) -> i64 {
match row.columns.get(&2) {
Some(Value::Int64(value)) => *value,
other => panic!("expected Int64 value, got {other:?}"),
}
}
fn vm_rss_bytes() -> Option<u64> {
let status = std::fs::read_to_string("/proc/self/status").ok()?;
let line = status.lines().find(|line| line.starts_with("VmRSS:"))?;
let kib: u64 = line.split_whitespace().nth(1)?.parse().ok()?;
Some(kib * 1024)
}
fn collect(table: &Table, snap: Snapshot) -> Vec<Row> {
let control = ExecutionControl::new(None);
let mut out = Vec::new();
table
.for_each_visible_row_controlled(snap, &control, |row| {
out.push(row);
Ok(())
})
.unwrap();
out
}
#[test]
fn one_million_row_memtable_yields_ascending_strict_order() {
let directory = tempdir().unwrap();
let mut table = Table::create(directory.path(), schema(), 1).unwrap();
table.set_mutable_run_spill_bytes(u64::MAX);
let rss_before = vm_rss_bytes();
let n = 1_000_000_i64;
for i in 0..n {
put(&mut table, i, i);
}
table.commit().unwrap();
assert!(
table.memtable_len() >= 1_000_000,
"memtable must hold all rows (memtable_len = {})",
table.memtable_len()
);
assert_eq!(
table.mutable_run_len(),
0,
"mutable run must be empty (mutable_run_len = {})",
table.mutable_run_len()
);
let snap = table.snapshot();
let (rows, trace) = QueryTrace::capture(|| collect(&table, snap));
assert_eq!(rows.len(), 1_000_000);
assert_eq!(trace.controlled_scan_rows_emitted, 1_000_000);
for pair in rows.windows(2) {
assert!(
pair[0].row_id.0 < pair[1].row_id.0,
"output must be strictly ascending by RowId"
);
}
assert!(
trace.controlled_scan_setup_time_us < 200_000,
"setup took {} µs",
trace.controlled_scan_setup_time_us
);
assert!(
trace.controlled_scan_peak_source_buffer_rows <= 8,
"memtable-only streaming should keep the per-tier buffer small (peak = {})",
trace.controlled_scan_peak_source_buffer_rows
);
emit_scan_metric!(
"controlled_scan::one_million_row_memtable_yields_ascending_strict_order",
trace.controlled_scan_peak_source_buffer_rows,
"peak_source_buffer_rows"
);
let mut direct = Memtable::new();
for i in 0..1_000_000u64 {
direct.upsert(Row::new(RowId(i), Epoch(i + 1)).with_column(2, Value::Int64(i as i64)));
}
let direct_snap = Snapshot::at(Epoch(2_000_000));
let control = ExecutionControl::new(None);
let mut cursor = direct.newest_visible_iter(&direct_snap);
let first = cursor
.next_controlled(&control)
.unwrap()
.expect("first row");
assert_eq!(first.0, RowId(0));
let at_first = cursor.be_tree_cursor_stats();
assert!(
at_first.checkpoints >= 1,
"the descent checkpoints before versions stream"
);
let versions_examined_before_first_checkpoint = 0usize;
let mut count = 1usize;
let mut prev = first.0;
while let Some((rid, _, _)) = cursor.next_controlled(&control).unwrap() {
assert!(rid > prev, "strictly ascending RowId");
prev = rid;
count += 1;
}
assert_eq!(count, 1_000_000);
let stats = cursor.be_tree_cursor_stats();
assert_eq!(stats.total_versions_precollected, 0);
assert!(
stats.peak_active_frames <= 64,
"cursor frames bounded by tree height: {}",
stats.peak_active_frames
);
assert!(
stats.peak_buffered_messages_owned <= 4_096,
"buffered messages bounded by height x buffer capacity: {}",
stats.peak_buffered_messages_owned
);
let rss_after = vm_rss_bytes();
let rss_delta = match (rss_before, rss_after) {
(Some(before), Some(after)) => after as i64 - before as i64,
_ => 0,
};
emit_scan_metric!(
"controlled_scan::one_million_row_memtable_yields_ascending_strict_order::time_to_first_row_us",
trace.controlled_scan_time_to_first_row_us,
"microseconds"
);
emit_scan_metric!(
"controlled_scan::one_million_row_memtable_yields_ascending_strict_order::peak_cursor_frames",
stats.peak_active_frames,
"frames"
);
emit_scan_metric!(
"controlled_scan::one_million_row_memtable_yields_ascending_strict_order::peak_buffered_messages",
stats.peak_buffered_messages_owned,
"messages"
);
emit_scan_metric!(
"controlled_scan::one_million_row_memtable_yields_ascending_strict_order::precollected_versions",
stats.total_versions_precollected,
"versions"
);
emit_scan_metric!(
"controlled_scan::one_million_row_memtable_yields_ascending_strict_order::versions_examined_before_first_checkpoint",
versions_examined_before_first_checkpoint,
"versions"
);
emit_scan_metric!(
"controlled_scan::one_million_row_memtable_yields_ascending_strict_order::rss_delta_bytes",
rss_delta,
"bytes"
);
}
#[test]
fn one_million_row_mutable_run_yields_ascending_strict_order() {
let directory = tempdir().unwrap();
let mut table = Table::create(directory.path(), schema(), 1).unwrap();
table.set_mutable_run_spill_bytes(u64::MAX);
let n = 1_000_000_i64;
for i in 0..n {
put(&mut table, i, i);
}
table.flush().unwrap();
table.commit().unwrap();
assert!(
table.mutable_run_len() >= 1_000_000,
"mutable run must hold all rows (mutable_run_len = {})",
table.mutable_run_len()
);
assert_eq!(
table.run_count(),
0,
"no sorted runs (run_count = {})",
table.run_count()
);
let snap = table.snapshot();
let (rows, trace) = QueryTrace::capture(|| collect(&table, snap));
assert_eq!(rows.len(), 1_000_000);
assert_eq!(trace.controlled_scan_rows_emitted, 1_000_000);
for pair in rows.windows(2) {
assert!(
pair[0].row_id.0 < pair[1].row_id.0,
"output must be strictly ascending by RowId"
);
}
assert!(
trace.controlled_scan_setup_time_us < 200_000,
"mutable-run setup took {} µs",
trace.controlled_scan_setup_time_us
);
emit_scan_metric!(
"controlled_scan::one_million_row_mutable_run_yields_ascending_strict_order",
trace.controlled_scan_peak_source_buffer_rows,
"peak_source_buffer_rows"
);
}
#[test]
fn out_of_order_be_tree_buffer_emits_ascending() {
let directory = tempdir().unwrap();
let mut table = Table::create(directory.path(), schema(), 1).unwrap();
table.set_mutable_run_spill_bytes(u64::MAX);
let _ = put(&mut table, 1_000_000, 1);
let _ = put(&mut table, 1, 2);
let _ = put(&mut table, 999_999, 3);
let _ = put(&mut table, 500_000, 4);
let _ = put(&mut table, 2, 5);
table.commit().unwrap();
assert!(table.memtable_len() >= 5);
let snap = table.snapshot();
let mut observed: Vec<RowId> = Vec::new();
let control = ExecutionControl::new(None);
table
.for_each_visible_row_controlled(snap, &control, |row| {
observed.push(row.row_id);
Ok(())
})
.unwrap();
let mut sorted = observed.clone();
sorted.sort();
sorted.dedup();
assert_eq!(observed, sorted, "output must be strictly ascending");
assert_eq!(observed.len(), 5);
}
#[test]
fn same_rowid_across_frozen_layers_dedups_to_newest() {
let mut schema_clustered = schema();
schema_clustered.clustered = true;
let directory = tempdir().unwrap();
let mut table = Table::create(directory.path(), schema_clustered, 1).unwrap();
table.set_mutable_run_spill_bytes(1);
let rid = put(&mut table, 7, 100);
table.flush().unwrap();
put(&mut table, 8, 80);
table.flush().unwrap();
table.set_mutable_run_spill_bytes(u64::MAX);
table
.put(vec![(1, Value::Int64(7)), (2, Value::Int64(777))])
.unwrap();
table.commit().unwrap();
let snap = table.snapshot();
let rows = collect(&table, snap);
let by_rid: BTreeMap<RowId, i64> = rows.iter().map(|r| (r.row_id, value(r))).collect();
assert_eq!(
by_rid.get(&rid).copied(),
Some(777),
"newer version must win"
);
assert_eq!(by_rid.len(), 2);
}
#[test]
fn dense_single_row_history_streams() {
let mut schema_clustered = schema();
schema_clustered.clustered = true;
let directory = tempdir().unwrap();
let mut table = Table::create(directory.path(), schema_clustered, 1).unwrap();
table.set_mutable_run_spill_bytes(u64::MAX);
let total = 10_000i64;
for v in 0..total {
table
.put(vec![(1, Value::Int64(1)), (2, Value::Int64(v))])
.unwrap();
table.commit().unwrap();
}
let snap = table.snapshot();
let control = ExecutionControl::new(None);
let (result, trace) = QueryTrace::capture(|| {
let mut count = 0;
let mut last_value: Option<i64> = None;
table
.for_each_visible_row_controlled(snap, &control, |row| {
count += 1;
if !row.deleted {
last_value = Some(value(&row));
}
Ok(())
})
.map(|()| (count, last_value))
});
let (count, last_value) = result.unwrap();
assert_eq!(
count, 1,
"10k versions of one RowId collapse to one emitted row"
);
assert_eq!(last_value, Some(total - 1));
assert!(
trace.controlled_scan_peak_same_row_versions >= 10_000,
"peak_same_row_versions must reflect the dense history (got {})",
trace.controlled_scan_peak_same_row_versions
);
emit_scan_metric!(
"controlled_scan::dense_single_row_history_streams",
trace.controlled_scan_peak_same_row_versions,
"peak_same_row_versions"
);
}
#[test]
fn three_tier_oracle_match() {
let directory = tempdir().unwrap();
let mut table = Table::create(directory.path(), schema(), 1).unwrap();
table.set_mutable_run_spill_bytes(1);
put(&mut table, 1, 10);
table.flush().unwrap(); table.set_mutable_run_spill_bytes(u64::MAX);
put(&mut table, 2, 20);
table.flush().unwrap(); put(&mut table, 3, 30); table.commit().unwrap();
let snap = table.snapshot();
let oracle: BTreeMap<i64, i64> = table
.visible_rows(snap)
.unwrap()
.into_iter()
.map(|r| {
(
match r.columns.get(&1) {
Some(Value::Int64(v)) => *v,
_ => panic!("missing id"),
},
match r.columns.get(&2) {
Some(Value::Int64(v)) => *v,
_ => panic!("missing value"),
},
)
})
.collect();
let control = ExecutionControl::new(None);
let mut observed = BTreeMap::new();
table
.for_each_visible_row_controlled(snap, &control, |row| {
let id = match row.columns.get(&1) {
Some(Value::Int64(v)) => *v,
_ => return Ok(()),
};
let v = match row.columns.get(&2) {
Some(Value::Int64(v)) => *v,
_ => return Ok(()),
};
observed.insert(id, v);
Ok(())
})
.unwrap();
assert_eq!(observed, oracle);
let mut ordered: Vec<RowId> = observed.keys().map(|k| RowId(*k as u64)).collect();
ordered.sort();
ordered.dedup();
assert_eq!(ordered.len(), 3);
}
#[test]
fn cancellation_bounds_versions_examined() {
let directory = tempdir().unwrap();
let mut table = Table::create(directory.path(), schema(), 1).unwrap();
table.set_mutable_run_spill_bytes(u64::MAX);
for i in 0..10_000i64 {
put(&mut table, i, i);
}
table.commit().unwrap();
let control = ExecutionControl::new(None);
let mut visited = 0usize;
let (result, trace) = QueryTrace::capture(|| {
table.for_each_visible_row_controlled(table.snapshot(), &control, |_| {
visited += 1;
if visited == 50 {
control.cancel(CancellationReason::ClientRequest);
}
Ok(())
})
});
assert!(matches!(result, Err(MongrelError::Cancelled)));
assert!(
trace.controlled_scan_versions_examined <= 50 + 512,
"excess versions examined after cancel: {}",
trace.controlled_scan_versions_examined
);
emit_scan_metric!(
"controlled_scan::cancellation_bounds_versions_examined",
trace.controlled_scan_versions_examined,
"versions_examined"
);
}
#[test]
fn output_contract_strict_ascending_no_duplicates() {
let mut schema_clustered = schema();
schema_clustered.clustered = true;
let directory = tempdir().unwrap();
let mut table = Table::create(directory.path(), schema_clustered, 1).unwrap();
table.set_mutable_run_spill_bytes(1);
for i in 0..500i64 {
put(&mut table, i * 2, i);
}
table.flush().unwrap();
table.set_mutable_run_spill_bytes(u64::MAX);
for i in 0..500i64 {
put(&mut table, i * 2 + 1, i);
}
table.flush().unwrap();
let delete_id = put(&mut table, 100, 9_999);
table.delete(delete_id).unwrap();
table.commit().unwrap();
let snap = table.snapshot();
let mut last = 0u64;
let mut seen: BTreeMap<RowId, i64> = BTreeMap::new();
let control = ExecutionControl::new(None);
table
.for_each_visible_row_controlled(snap, &control, |row| {
assert!(row.row_id.0 > last, "strictly ascending RowId");
last = row.row_id.0;
assert!(!row.deleted, "tombstones must be suppressed");
seen.insert(row.row_id, value(&row));
Ok(())
})
.unwrap();
assert_eq!(seen.len(), 999);
assert!(!seen.contains_key(&delete_id));
}
#[test]
fn newer_memtable_value_beats_mutable_run_value() {
let directory = tempdir().unwrap();
let mut table = Table::create(directory.path(), schema(), 1).unwrap();
table.set_mutable_run_spill_bytes(u64::MAX);
put(&mut table, 1, 10);
table.flush().unwrap();
put(&mut table, 1, 20);
table.commit().unwrap();
let snap = table.snapshot();
let rows = collect(&table, snap);
assert_eq!(rows.len(), 1);
assert_eq!(value(&rows[0]), 20);
}
#[test]
fn hlc_inversion_chooses_higher_hlc() {
let old_hlc = HlcTimestamp {
physical_micros: 100,
logical: 0,
node_tiebreaker: 1,
};
let new_hlc = HlcTimestamp {
physical_micros: 200,
logical: 0,
node_tiebreaker: 1,
};
assert!(Snapshot::version_is_newer(
Epoch(1),
Some(new_hlc),
Epoch(50),
Some(old_hlc),
));
assert!(!Snapshot::version_is_newer(
Epoch(50),
Some(old_hlc),
Epoch(1),
Some(new_hlc),
));
}
#[test]
fn tombstone_suppresses_older_live_version() {
let directory = tempdir().unwrap();
let mut table = Table::create(directory.path(), schema(), 1).unwrap();
table.set_mutable_run_spill_bytes(1);
let row_id = put(&mut table, 1, 10);
table.flush().unwrap();
table.delete(row_id).unwrap();
table.commit().unwrap();
let snap = table.snapshot();
let rows = collect(&table, snap);
assert!(!rows.iter().any(|r| r.row_id == row_id));
}
#[test]
fn controlled_scan_interleaves_memtable_mutable_run_and_sorted_run() {
let directory = tempdir().unwrap();
let mut table = Table::create(directory.path(), schema(), 1).unwrap();
table.set_mutable_run_spill_bytes(1);
put(&mut table, 1, 10);
table.flush().unwrap();
table.set_mutable_run_spill_bytes(u64::MAX);
put(&mut table, 2, 20);
table.flush().unwrap();
put(&mut table, 3, 30);
table.commit().unwrap();
let snap = table.snapshot();
let rows = collect(&table, snap);
assert_eq!(rows.len(), 3);
let values: Vec<i64> = rows.iter().map(value).collect();
assert_eq!(values, vec![10, 20, 30]);
assert!(rows.windows(2).all(|p| p[0].row_id < p[1].row_id));
}
#[test]
fn visitor_error_short_circuits_scan() {
let directory = tempdir().unwrap();
let mut table = Table::create(directory.path(), schema(), 1).unwrap();
table.set_mutable_run_spill_bytes(u64::MAX);
for i in 0..100i64 {
put(&mut table, i, i);
}
table.commit().unwrap();
let control = ExecutionControl::new(None);
let mut visited = 0;
let err = table
.for_each_visible_row_controlled(table.snapshot(), &control, |row| {
visited += 1;
if visited == 5 {
Err(MongrelError::Other("stop".into()))
} else {
assert!(!row.deleted);
Ok(())
}
})
.unwrap_err();
assert!(matches!(err, MongrelError::Other(_)));
assert!(visited <= 6);
}
#[test]
fn cancellation_before_first_row_is_observed() {
let directory = tempdir().unwrap();
let mut table = Table::create(directory.path(), schema(), 1).unwrap();
table.set_mutable_run_spill_bytes(u64::MAX);
for i in 0..1_000i64 {
put(&mut table, i, i);
}
table.commit().unwrap();
let control = ExecutionControl::new(None);
control.cancel(CancellationReason::ClientRequest);
let result = table.for_each_visible_row_controlled(table.snapshot(), &control, |_| {
panic!("visit must not run when the control is already cancelled")
});
assert!(matches!(result, Err(MongrelError::Cancelled)));
}
#[test]
fn dml_count_update_delete_regression_fixture_still_passes() {
let mut schema_clustered = schema();
schema_clustered.clustered = true;
let directory = tempdir().unwrap();
let mut table = Table::create(directory.path(), schema_clustered, 1).unwrap();
table.set_mutable_run_spill_bytes(u64::MAX);
let initial = 50;
for i in 0..initial {
put(&mut table, i, i * 1000);
}
let mut delete_ids = Vec::new();
for id in 0..5i64 {
let (rid, _) = table
.put_returning(vec![(1, Value::Int64(id)), (2, Value::Int64(id * 1000))])
.unwrap();
delete_ids.push(rid);
}
table.commit().unwrap();
for rid in &delete_ids {
table.delete(*rid).unwrap();
}
table.commit().unwrap();
let control = ExecutionControl::new(None);
let mut count = 0_u64;
table
.for_each_visible_row_controlled(table.snapshot(), &control, |row| {
assert!(!row.deleted);
count += 1;
Ok(())
})
.unwrap();
assert_eq!(count, initial as u64 - 5);
}
#[test]
fn hundred_thousand_rows_with_ten_versions_each_match_oracle() {
let directory = tempdir().unwrap();
let mut table = Table::create(directory.path(), schema(), 1).unwrap();
table.set_mutable_run_spill_bytes(u64::MAX);
for _ in 0..10 {
for id in 0..100_000i64 {
put(&mut table, id, id);
}
table.commit().unwrap();
}
let snap = table.snapshot();
let mut observed: BTreeMap<i64, i64> = BTreeMap::new();
let control = ExecutionControl::new(None);
table
.for_each_visible_row_controlled(snap, &control, |row| {
observed.insert(value(&row), value(&row));
Ok(())
})
.unwrap();
assert_eq!(observed.len(), 100_000);
}
#[test]
fn sorted_run_scan_matches_oracle() {
let directory = tempdir().unwrap();
let mut table = Table::create(directory.path(), schema(), 1).unwrap();
bulk_load(&mut table, 1_000);
table.commit().unwrap();
let snap = table.snapshot();
let oracle = table.visible_rows(snap).unwrap();
let rows = collect(&table, snap);
assert_eq!(rows.len(), oracle.len());
for pair in rows.windows(2) {
assert!(pair[0].row_id < pair[1].row_id);
}
}
#[test]
fn memtable_cursor_does_not_precollect_every_version() {
let mut memtable = Memtable::new();
const ROWS: u64 = 60_000;
for i in 0..ROWS {
let rid = if i % 2 == 0 { i / 2 } else { ROWS / 2 + i / 2 };
memtable.upsert(Row::new(RowId(rid), Epoch(1)).with_column(2, Value::Int64(rid as i64)));
memtable
.upsert(Row::new(RowId(rid), Epoch(2)).with_column(2, Value::Int64(rid as i64 + 1)));
}
for rid in (0..ROWS / 2).step_by(10) {
memtable.tombstone(RowId(rid), Epoch(3));
}
assert_eq!(memtable.len(), 123_000);
let snap = Snapshot::at(Epoch(4));
let control = ExecutionControl::new(None);
let mut cursor = memtable.newest_visible_iter(&snap);
let first = cursor
.next_controlled(&control)
.unwrap()
.expect("non-empty memtable");
let at_first = cursor.be_tree_cursor_stats();
assert_eq!(at_first.total_versions_precollected, 0);
assert!(
at_first.versions_examined <= 256,
"first row must not wait for a tree scan: {}",
at_first.versions_examined
);
assert!(
at_first.peak_active_frames <= 64,
"cursor frames bounded by tree height: {}",
at_first.peak_active_frames
);
assert!(
at_first.peak_buffered_messages_owned <= 4_096,
"buffered messages bounded by height x buffer capacity: {}",
at_first.peak_buffered_messages_owned
);
let mut count = 1usize;
let mut prev = first.0;
while let Some((rid, _, _)) = cursor.next_controlled(&control).unwrap() {
assert!(rid > prev, "strictly ascending RowId");
prev = rid;
count += 1;
}
assert_eq!(count, ROWS as usize, "one emission per distinct RowId");
let stats = cursor.be_tree_cursor_stats();
assert_eq!(stats.total_versions_precollected, 0);
assert!(
stats.peak_buffered_messages_owned > 0,
"fixture must exercise internal-node buffers"
);
assert!(
stats.peak_buffered_messages_owned <= 4_096,
"bounded for the whole scan: {}",
stats.peak_buffered_messages_owned
);
emit_scan_metric!(
"controlled_scan::memtable_cursor_does_not_precollect_every_version",
stats.total_versions_precollected,
"versions"
);
}
#[test]
fn cancellation_during_cursor_construction() {
let mut memtable = Memtable::new();
for i in 0..200_000u64 {
memtable.upsert(Row::new(RowId(i), Epoch(i + 1)).with_column(1, Value::Int64(i as i64)));
}
let snap = Snapshot::at(Epoch(300_000));
let mut cursor = memtable.newest_visible_iter(&snap);
let control = ExecutionControl::new(None);
control.cancel(CancellationReason::ClientRequest);
let err = cursor.next_controlled(&control).unwrap_err();
assert!(matches!(err, MongrelError::Cancelled));
assert_eq!(
cursor.be_tree_cursor_stats().versions_examined,
0,
"cancellation lands before the first version streams"
);
}