use serde_json::json;
use std::time::{Duration, Instant};
use crate::engine::Engine;
use crate::session::SessionContext;
use crate::storage::ProjectedValue;
use crate::storage_adapter::{
Memory, SharedStorageAdapterRead, StorageAdapter, StorageAdapterRead, StorageBeginScanOptions,
StorageKey, StoragePrefix, StorageReadOptions, StorageSpace, StorageWriteOptions,
};
fn sizes_from_env(var: &str, default: &[usize]) -> Vec<usize> {
match std::env::var(var) {
Ok(raw) => raw
.split(',')
.filter(|part| !part.trim().is_empty())
.map(|part| part.trim().parse::<usize>().expect("size must parse"))
.collect(),
Err(_) => default.to_vec(),
}
}
fn reps_from_env(default: usize) -> usize {
std::env::var("LIX_TOMBSTONE_REPS")
.ok()
.and_then(|raw| raw.parse::<usize>().ok())
.unwrap_or(default)
}
async fn open_session() -> (Memory, SessionContext<Memory>) {
let storage = Memory::new();
Engine::initialize(storage.clone())
.await
.expect("engine should initialize");
let engine = Engine::new(storage.clone())
.await
.expect("engine should open");
let session = engine.open_session().await.expect("session should open");
(storage, session)
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
struct RowCensus {
entries: usize,
tombstones: usize,
packed_bases: usize,
root_bases: usize,
diff_records: usize,
}
impl RowCensus {
fn live(self) -> usize {
self.entries - self.tombstones
}
}
const HEAD_VALUE_DELETED_BIT: u8 = 0b0000_0001;
async fn space_values(
read: &(impl StorageAdapterRead + ?Sized),
space: StorageSpace,
) -> Vec<bytes::Bytes> {
let range = StoragePrefix {
bytes: bytes::Bytes::new(),
}
.to_range()
.expect("valid empty prefix");
let mut cursor = read
.begin_scan(space, range, StorageBeginScanOptions::default())
.await
.expect("scan the serving space");
cursor
.collect_all()
.await
.expect("collect serving entries")
.into_iter()
.map(|entry| match entry.value {
ProjectedValue::FullValue(bytes) => bytes,
ProjectedValue::KeyOnly => bytes::Bytes::new(),
})
.collect()
}
async fn row_census(storage: &Memory) -> RowCensus {
let adapter = StorageAdapter::new(storage.clone());
let read = adapter
.begin_read(StorageReadOptions::default())
.await
.expect("read the hot row plane");
let rows = space_values(&read, crate::hot_state::ROW_SPACE).await;
let tombstones = rows
.iter()
.filter(|value| value.len() > 1 && value[1] & HEAD_VALUE_DELETED_BIT != 0)
.count();
let packed_bases = space_values(&read, crate::hot_state::PACKED_CURRENT_BASE_SPACE)
.await
.len();
let root_bases = space_values(&read, crate::hot_state::ROOT_CURRENT_BASE_SPACE)
.await
.len();
let diff_records = space_values(&read, crate::hot_state::DIFF_SPACE)
.await
.len();
RowCensus {
entries: rows.len(),
tombstones,
packed_bases,
root_bases,
diff_records,
}
}
async fn reopen_session(storage: &Memory) -> SessionContext<Memory> {
let engine = Engine::new(storage.clone())
.await
.expect("engine should reopen");
engine.open_session().await.expect("session should reopen")
}
async fn space_entries(
read: &(impl StorageAdapterRead + ?Sized),
space: StorageSpace,
) -> Vec<(StorageKey, bytes::Bytes)> {
let range = StoragePrefix {
bytes: bytes::Bytes::new(),
}
.to_range()
.expect("valid empty prefix");
let mut cursor = read
.begin_scan(space, range, StorageBeginScanOptions::default())
.await
.expect("scan the serving space");
cursor
.collect_all()
.await
.expect("collect serving entries")
.into_iter()
.map(|entry| {
let value = match entry.value {
ProjectedValue::FullValue(bytes) => bytes,
ProjectedValue::KeyOnly => bytes::Bytes::new(),
};
(entry.key, value)
})
.collect()
}
async fn drop_all_tombstones(storage: &Memory) -> usize {
let adapter = StorageAdapter::new(storage.clone());
let read = adapter
.begin_read(StorageReadOptions::default())
.await
.expect("read the hot row plane");
let entries = space_entries(&read, crate::hot_state::ROW_SPACE).await;
drop(read);
let mut writes = adapter.new_write_set();
let mut removed = 0_usize;
for (key, value) in entries {
if value.len() > 1 && value[1] & HEAD_VALUE_DELETED_BIT != 0 {
writes.delete(crate::hot_state::ROW_SPACE, key);
removed += 1;
}
}
adapter
.commit_write_set(writes, StorageWriteOptions::default())
.await
.expect("tombstone removal should commit");
removed
}
async fn working_diff_rows(session: &SessionContext<Memory>, schema_key: &str) -> usize {
session
.execute(
"SELECT row_pk FROM lix_working_diff WHERE schema_key = $1",
&[crate::Value::Text(schema_key.to_string())],
)
.await
.expect("working diff should read")
.len()
}
async fn run_repository_gc(storage: &Memory) {
let adapter = StorageAdapter::new(storage.clone());
let read = SharedStorageAdapterRead::new(
adapter
.begin_read(StorageReadOptions::default())
.await
.expect("gc read"),
);
let mut gc_writes = adapter.new_write_set();
crate::gc::stage_repository_gc(read, &mut gc_writes)
.await
.expect("repository gc should plan");
adapter
.commit_write_set(gc_writes, StorageWriteOptions::default())
.await
.expect("gc write set should commit");
}
fn probe_schema(key: &str) -> serde_json::Value {
json!({
"$schema": "https://lix.dev/schema-v1.json",
"key": key,
"columns": [
{ "name": "id", "type": "text", "nullable": false },
{ "name": "locale", "type": "text", "nullable": false },
],
"primary_key": ["id"],
})
}
async fn register(session: &SessionContext<Memory>, schema: serde_json::Value) {
session
.execute(
"INSERT INTO lix_registered_schema (value) VALUES (CAST($1 AS JSONB))",
&[crate::Value::Text(schema.to_string())],
)
.await
.expect("schema should register");
}
const CHUNK: usize = 250;
async fn insert_rows(session: &SessionContext<Memory>, table: &str, count: usize) {
let mut index = 0;
while index < count {
let end = (index + CHUNK).min(count);
let values = (index..end)
.map(|i| {
let locale = if i == 0 { "keep" } else { "drop" };
format!("('row-{i}', '{locale}')")
})
.collect::<Vec<_>>()
.join(",");
session
.execute(
&format!("INSERT INTO {table} (id, locale) VALUES {values}"),
&[],
)
.await
.expect("rows should insert");
index = end;
}
}
async fn delete_all_but_first(session: &SessionContext<Memory>, table: &str, count: usize) {
let mut index = 1;
while index < count {
let end = (index + CHUNK).min(count);
let ids = (index..end)
.map(|i| format!("'row-{i}'"))
.collect::<Vec<_>>()
.join(",");
session
.execute(&format!("DELETE FROM {table} WHERE id IN ({ids})"), &[])
.await
.expect("bulk delete should run");
index = end;
}
}
async fn timed_scan(
session: &SessionContext<Memory>,
sql: &str,
expect_rows: usize,
reps: usize,
) -> Duration {
for _ in 0..2 {
let rows = session.execute(sql, &[]).await.expect("warmup should run");
assert_eq!(rows.len(), expect_rows, "warmup returned the wrong count");
}
let mut samples = Vec::new();
for _ in 0..reps {
let start = Instant::now();
let rows = session.execute(sql, &[]).await.expect("scan should run");
let elapsed = start.elapsed();
assert_eq!(rows.len(), expect_rows, "scan returned the wrong row count");
samples.push(elapsed);
}
samples.sort();
samples[samples.len() / 2]
}
fn scan_sql(table: &str) -> String {
format!("SELECT id FROM {table} WHERE locale = 'keep'")
}
#[tokio::test]
#[ignore = "measurement probe, not a gate"]
async fn hot_row_tombstone_accumulation_and_scan_cost() {
let sizes = sizes_from_env("LIX_TOMBSTONE_SIZES", &[100, 1_000, 10_000]);
let reps = reps_from_env(5);
println!(
"phase1 | arm,n,deletes,row_entries,tombstones,live_entries,packed_bases,root_bases,answer_rows,scan_us"
);
for n in sizes {
{
let (storage, session) = open_session().await;
register(&session, probe_schema("churnrow")).await;
insert_rows(&session, "churnrow", n).await;
delete_all_but_first(&session, "churnrow", n).await;
let census = row_census(&storage).await;
let scan = timed_scan(&session, &scan_sql("churnrow"), 1, reps).await;
println!(
"churn,{n},{},{},{},{},{},{},1,{}",
n - 1,
census.entries,
census.tombstones,
census.live(),
census.packed_bases,
census.root_bases,
scan.as_micros()
);
}
{
let (storage, session) = open_session().await;
register(&session, probe_schema("freshrow")).await;
insert_rows(&session, "freshrow", 1).await;
let census = row_census(&storage).await;
let scan = timed_scan(&session, &scan_sql("freshrow"), 1, reps).await;
println!(
"fresh,{n},0,{},{},{},{},{},1,{}",
census.entries,
census.tombstones,
census.live(),
census.packed_bases,
census.root_bases,
scan.as_micros()
);
}
}
}
#[tokio::test]
#[ignore = "measurement probe, not a gate"]
async fn hot_row_tombstone_reclamation_events() {
let n = sizes_from_env("LIX_TOMBSTONE_SIZES", &[1_000])[0];
let reps = reps_from_env(5);
let (storage, session) = open_session().await;
register(&session, probe_schema("evtrow")).await;
register(&session, probe_schema("otherrow")).await;
insert_rows(&session, "evtrow", n).await;
delete_all_but_first(&session, "evtrow", n).await;
println!("phase2 | event,row_entries,tombstones,live_entries,packed_bases,root_bases,scan_us");
let report = |label: &str, census: RowCensus, scan: Duration| {
println!(
"{label},{},{},{},{},{},{}",
census.entries,
census.tombstones,
census.live(),
census.packed_bases,
census.root_bases,
scan.as_micros()
);
};
let census = row_census(&storage).await;
let scan = timed_scan(&session, &scan_sql("evtrow"), 1, reps).await;
report("after_churn", census, scan);
for index in 0..20 {
session
.execute(
"INSERT INTO otherrow (id, locale) VALUES ($1, 'x')",
&[crate::Value::Text(format!("o-{index}"))],
)
.await
.expect("unrelated insert should commit");
}
let census = row_census(&storage).await;
let scan = timed_scan(&session, &scan_sql("evtrow"), 1, reps).await;
report("after_20_commits", census, scan);
session
.create_checkpoint()
.await
.expect("checkpoint should publish");
let census = row_census(&storage).await;
let scan = timed_scan(&session, &scan_sql("evtrow"), 1, reps).await;
report("after_checkpoint", census, scan);
run_repository_gc(&storage).await;
let census = row_census(&storage).await;
let scan = timed_scan(&session, &scan_sql("evtrow"), 1, reps).await;
report("after_gc", census, scan);
session
.execute(
"INSERT INTO otherrow (id, locale) VALUES ('post-ckpt', 'x')",
&[],
)
.await
.expect("post-checkpoint insert should commit");
session
.create_checkpoint()
.await
.expect("second checkpoint should publish");
run_repository_gc(&storage).await;
let census = row_census(&storage).await;
let scan = timed_scan(&session, &scan_sql("evtrow"), 1, reps).await;
report("after_ckpt2_gc", census, scan);
}
#[tokio::test]
#[ignore = "measurement probe, not a gate"]
async fn hot_row_tombstone_generation_rotation() {
let n = sizes_from_env("LIX_TOMBSTONE_SIZES", &[1_000])[0];
let reps = reps_from_env(5);
let (storage, session) = open_session().await;
register(&session, probe_schema("rotrow")).await;
insert_rows(&session, "rotrow", n).await;
delete_all_but_first(&session, "rotrow", n).await;
println!("phase3 | event,row_entries,tombstones,live_entries,packed_bases,root_bases,scan_us");
let report = |label: &str, census: RowCensus, scan: Duration| {
println!(
"{label},{},{},{},{},{},{}",
census.entries,
census.tombstones,
census.live(),
census.packed_bases,
census.root_bases,
scan.as_micros()
);
};
let main_branch_id = session
.active_branch_id()
.await
.expect("active branch should resolve");
let census = row_census(&storage).await;
let scan = timed_scan(&session, &scan_sql("rotrow"), 1, reps).await;
report("after_churn", census, scan);
let branch = session
.create_branch(crate::CreateBranchOptions {
id: None,
name: "e45-rotation".to_string(),
from_commit_id: None,
})
.await
.expect("branch should create");
let census = row_census(&storage).await;
let scan = timed_scan(&session, &scan_sql("rotrow"), 1, reps).await;
report("after_create_branch", census, scan);
session
.switch_branch(crate::SwitchBranchOptions {
branch_id: branch.id.clone(),
})
.await
.expect("branch should switch");
let census = row_census(&storage).await;
let scan = timed_scan(&session, &scan_sql("rotrow"), 1, reps).await;
report("after_switch_to_branch", census, scan);
session
.execute(
"INSERT INTO rotrow (id, locale) VALUES ('branch-row', 'x')",
&[],
)
.await
.expect("branch insert should commit");
session
.switch_branch(crate::SwitchBranchOptions {
branch_id: main_branch_id,
})
.await
.expect("branch should switch back");
let census = row_census(&storage).await;
let scan = timed_scan(&session, &scan_sql("rotrow"), 1, reps).await;
report("after_switch_back_to_main", census, scan);
run_repository_gc(&storage).await;
let census = row_census(&storage).await;
let scan = timed_scan(&session, &scan_sql("rotrow"), 1, reps).await;
report("after_gc", census, scan);
}
#[tokio::test]
#[ignore = "measurement probe, not a gate"]
async fn hot_row_tombstone_compaction_premise() {
let n = sizes_from_env("LIX_TOMBSTONE_SIZES", &[1_000])[0];
let reps = reps_from_env(5);
println!(
"phase4 | arm,event,row_entries,tombstones,packed_bases,root_bases,diff_records,working_diff_rows,answer_rows,collection_rows,scan_us"
);
{
let (storage, session) = open_session().await;
register(&session, probe_schema("safrow")).await;
insert_rows(&session, "safrow", n).await;
delete_all_but_first(&session, "safrow", n).await;
session
.create_checkpoint()
.await
.expect("checkpoint should publish");
let census = row_census(&storage).await;
assert_eq!(census.packed_bases, 0, "safe arm expects no packed base");
assert_eq!(census.root_bases, 0, "safe arm expects no root base");
let diff_before = working_diff_rows(&session, "safrow").await;
assert_eq!(
diff_before, 0,
"a checkpoint must discharge every delete's working-diff obligation"
);
let scan_before = timed_scan(&session, &scan_sql("safrow"), 1, reps).await;
let collection_before = session
.execute("SELECT id FROM safrow", &[])
.await
.expect("collection scan should run")
.len();
println!(
"safe,before_drop,{},{},{},{},{},{},1,{},{}",
census.entries,
census.tombstones,
census.packed_bases,
census.root_bases,
census.diff_records,
diff_before,
collection_before,
scan_before.as_micros()
);
let removed = drop_all_tombstones(&storage).await;
assert_eq!(removed, n - 1, "every tombstone should have been removable");
let session = reopen_session(&storage).await;
let census = row_census(&storage).await;
let diff_after = working_diff_rows(&session, "safrow").await;
let answer = session
.execute(&scan_sql("safrow"), &[])
.await
.expect("query should still run");
let collection_after = session
.execute("SELECT id FROM safrow", &[])
.await
.expect("collection scan should still run")
.len();
let scan_after = timed_scan(&session, &scan_sql("safrow"), 1, reps).await;
println!(
"safe,after_drop,{},{},{},{},{},{},{},{},{}",
census.entries,
census.tombstones,
census.packed_bases,
census.root_bases,
census.diff_records,
diff_after,
answer.len(),
collection_after,
scan_after.as_micros()
);
assert_eq!(answer.len(), 1, "the surviving row must still answer");
assert_eq!(
collection_after, 1,
"dropping tombstones must not resurrect a deleted row"
);
assert_eq!(
diff_after, diff_before,
"dropping a discharged tombstone must not move the working diff"
);
}
{
let (storage, session) = open_session().await;
register(&session, probe_schema("nmrow")).await;
insert_rows(&session, "nmrow", n).await;
session
.create_checkpoint()
.await
.expect("first checkpoint should publish");
delete_all_but_first(&session, "nmrow", n).await;
let census = row_census(&storage).await;
let diff_before = working_diff_rows(&session, "nmrow").await;
let scan_before = timed_scan(&session, &scan_sql("nmrow"), 1, reps).await;
println!(
"near_miss,before_drop,{},{},{},{},{},{},1,1,{}",
census.entries,
census.tombstones,
census.packed_bases,
census.root_bases,
census.diff_records,
diff_before,
scan_before.as_micros()
);
assert!(
diff_before > 0,
"an uncheckpointed delete must be visible in the working diff, \
otherwise this arm proves nothing"
);
let removed = drop_all_tombstones(&storage).await;
let session = reopen_session(&storage).await;
let census = row_census(&storage).await;
let diff_after = working_diff_rows(&session, "nmrow").await;
let answer = session
.execute(&scan_sql("nmrow"), &[])
.await
.expect("query should still run");
let collection_after = session
.execute("SELECT id FROM nmrow", &[])
.await
.expect("collection scan should still run")
.len();
println!(
"near_miss,after_drop,{},{},{},{},{},{},{},{},-",
census.entries,
census.tombstones,
census.packed_bases,
census.root_bases,
census.diff_records,
diff_after,
answer.len(),
collection_after
);
println!(
"near_miss | removed={removed} working_diff_rows_lost={}",
diff_before.saturating_sub(diff_after)
);
assert_eq!(
diff_after, diff_before,
"the working diff must survive tombstone removal; if it stops \
surviving, ROW_SPACE has become the working diff's authority and \
the compaction rule needs condition (b) back"
);
assert_eq!(
answer.len(),
1,
"the surviving row must still answer before any checkpoint"
);
assert_eq!(
collection_after, 1,
"removing an uncheckpointed tombstone must still not resurrect a row"
);
}
}
#[tokio::test]
#[ignore = "measurement probe, not a gate"]
async fn rotated_generation_read_scaling() {
let sizes = sizes_from_env("LIX_TOMBSTONE_SIZES", &[1, 100, 1_000, 10_000]);
let reps = reps_from_env(5);
println!(
"phase5 | arm,n,row_entries,tombstones,live_entries,packed_bases,root_bases,scan_us,point_us"
);
let point_sql = "SELECT id FROM liverow WHERE id = 'row-0'";
for n in sizes {
for rotate in [false, true] {
let (storage, session) = open_session().await;
register(&session, probe_schema("liverow")).await;
insert_rows(&session, "liverow", n).await;
if rotate {
let branch = session
.create_branch(crate::CreateBranchOptions {
id: None,
name: "e51-rotation".to_string(),
from_commit_id: None,
})
.await
.expect("branch should create");
session
.switch_branch(crate::SwitchBranchOptions {
branch_id: branch.id.clone(),
})
.await
.expect("branch should switch");
}
let census = row_census(&storage).await;
let scan = timed_scan(&session, &scan_sql("liverow"), 1, reps).await;
let point = timed_scan(&session, point_sql, 1, reps).await;
println!(
"{},{n},{},{},{},{},{},{},{}",
if rotate { "live_rot" } else { "live_main" },
census.entries,
census.tombstones,
census.live(),
census.packed_bases,
census.root_bases,
scan.as_micros(),
point.as_micros()
);
}
}
}
#[cfg(feature = "storage-benches")]
#[tokio::test]
async fn rotated_generation_serving_view_is_cached_after_the_first_read() {
let (_storage, session) = open_session().await;
register(&session, probe_schema("cachedrow")).await;
insert_rows(&session, "cachedrow", 50).await;
let branch = session
.create_branch(crate::CreateBranchOptions {
id: None,
name: "e51-cache-guard".to_string(),
from_commit_id: None,
})
.await
.expect("branch should create");
session
.switch_branch(crate::SwitchBranchOptions {
branch_id: branch.id.clone(),
})
.await
.expect("branch should switch");
let _ = crate::storage_bench::take_root_base_batch_cache_accounting();
for _ in 0..5 {
let rows = session
.execute(&scan_sql("cachedrow"), &[])
.await
.expect("rotated scan should run");
assert_eq!(rows.len(), 1, "the rotated generation must serve one row");
}
let (hits, misses) = crate::storage_bench::take_root_base_batch_cache_accounting();
assert!(
misses > 0,
"the first rotated read must materialize the serving view (hits={hits} misses={misses})"
);
assert!(
hits > 0,
"later rotated reads must be served from the materialized view, \
not re-derived from canonical records (hits={hits} misses={misses})"
);
}
#[tokio::test]
#[ignore = "measurement probe, not a gate"]
async fn undo_restores_a_row_whose_tombstone_was_removed() {
let (storage, session) = open_session().await;
register(&session, probe_schema("undorow")).await;
insert_rows(&session, "undorow", 3).await;
let before = session
.execute("SELECT id FROM undorow", &[])
.await
.expect("collection reads")
.len();
session
.execute("DELETE FROM undorow WHERE id = 'row-1'", &[])
.await
.expect("delete should commit");
let census_after_delete = row_census(&storage).await;
let after_delete = session
.execute("SELECT id FROM undorow", &[])
.await
.expect("collection reads")
.len();
let removed = drop_all_tombstones(&storage).await;
let session = reopen_session(&storage).await;
let census_after_drop = row_census(&storage).await;
let after_drop = session
.execute("SELECT id FROM undorow", &[])
.await
.expect("collection reads")
.len();
session.undo().await.expect("undo should publish");
let after_undo = session
.execute("SELECT id FROM undorow", &[])
.await
.expect("collection reads")
.len();
let restored = session
.execute("SELECT id, locale FROM undorow WHERE id = 'row-1'", &[])
.await
.expect("restored row reads");
println!(
"phase6 | before={before} after_delete={after_delete} tombstones_after_delete={} \
removed={removed} after_drop={after_drop} tombstones_after_drop={} after_undo={after_undo} \
restored_rows={}",
census_after_delete.tombstones,
census_after_drop.tombstones,
restored.len()
);
assert_eq!(before, 3, "fixture should start with three rows");
assert_eq!(after_delete, 2, "the delete should be visible");
assert_eq!(removed, 1, "the delete should have left exactly one tombstone");
assert_eq!(
after_drop, 2,
"removing the tombstone must not resurrect the row"
);
assert_eq!(
after_undo, 3,
"undo must restore the deleted row with its tombstone already gone; \
if this fails, undo depends on the ROW_SPACE tombstone and the \
never-write design is unsafe"
);
assert_eq!(restored.len(), 1, "the restored row must be readable");
}
#[tokio::test]
#[ignore = "measurement probe, not a gate"]
async fn retention_fence_durability_across_supported_operations() {
async fn untracked_insert(
session: &SessionContext<Memory>,
table: &str,
) -> Result<(), String> {
session
.execute(
&format!(
"INSERT INTO {table} (id, locale, lixcol_untracked) VALUES ('row-0', 'u', TRUE)"
),
&[],
)
.await
.map(|_| ())
.map_err(|error| error.to_string())
}
println!("phase6 | route,tombstones,untracked_insert_result");
for route in [
"base",
"checkpoint",
"branch_roundtrip",
"on_new_branch",
"tombstone_dropped",
] {
let (storage, session) = open_session().await;
register(&session, probe_schema("fencerow")).await;
session
.execute(
"INSERT INTO fencerow (id, locale) VALUES ('row-0', 'keep')",
&[],
)
.await
.expect("tracked insert should commit");
session
.execute("DELETE FROM fencerow WHERE id = 'row-0'", &[])
.await
.expect("tracked delete should commit");
let mut session = session;
match route {
"base" => {}
"checkpoint" => {
session
.create_checkpoint()
.await
.expect("checkpoint should publish");
}
"branch_roundtrip" => {
let main_branch_id = session
.active_branch_id()
.await
.expect("active branch resolves");
let branch = session
.create_branch(crate::CreateBranchOptions {
id: None,
name: "e45-fence-roundtrip".to_string(),
from_commit_id: None,
})
.await
.expect("branch should create");
session
.switch_branch(crate::SwitchBranchOptions {
branch_id: branch.id,
})
.await
.expect("switch to branch");
session
.switch_branch(crate::SwitchBranchOptions {
branch_id: main_branch_id,
})
.await
.expect("switch back to main");
}
"on_new_branch" => {
let branch = session
.create_branch(crate::CreateBranchOptions {
id: None,
name: "e45-fence-newbranch".to_string(),
from_commit_id: None,
})
.await
.expect("branch should create");
session
.switch_branch(crate::SwitchBranchOptions {
branch_id: branch.id,
})
.await
.expect("switch to branch");
}
"tombstone_dropped" => {
let removed = drop_all_tombstones(&storage).await;
assert_eq!(removed, 1, "the tracked delete should leave one tombstone");
session = reopen_session(&storage).await;
}
_ => unreachable!(),
}
let census = row_census(&storage).await;
let result = untracked_insert(&session, "fencerow").await;
assert!(
result.is_err(),
"route '{route}' let an untracked row take a tracked-deleted identity; \
the retention fence does not survive this state"
);
let verdict = match &result {
Ok(()) => "SUCCEEDED".to_string(),
Err(message) => format!("refused: {}", message.replace(',', ";")),
};
println!("{route},{},{verdict}", census.tombstones);
if result.is_ok() {
let undo = session.undo().await;
match undo {
Ok(_) => {
let rows = session
.execute(
"SELECT id, lixcol_untracked FROM fencerow WHERE id = 'row-0'",
&[],
)
.await
.expect("identity reads after undo");
println!(
"{route},POST_UNDO,rows_for_identity={} <- 2 means tracked+untracked \
coexist",
rows.len()
);
}
Err(error) => println!(
"{route},POST_UNDO,undo_refused: {}",
error.to_string().replace(',', ";")
),
}
}
drop(session);
}
}
#[cfg(feature = "storage-benches")]
#[tokio::test]
#[ignore = "measurement probe, not a gate"]
async fn rotated_generation_key_allocation_census() {
let sizes = sizes_from_env("LIX_TOMBSTONE_SIZES", &[100, 1_000, 10_000]);
let reps = reps_from_env(5);
println!(
"phase8 | arm,n,phase,scans,rows,decodes,dec_in_b,own_b,esc,acct_b,enc,enc_b,\
clones,clone_b,matkey,matkey_b,reverify,hit,miss,dur,exa,rep,us"
);
for n in sizes {
for rotate in [false, true] {
let (_storage, session) = open_session().await;
register(&session, probe_schema("censusrow")).await;
insert_rows(&session, "censusrow", n).await;
if rotate {
let branch = session
.create_branch(crate::CreateBranchOptions {
id: None,
name: "e52-census".to_string(),
from_commit_id: None,
})
.await
.expect("branch should create");
session
.switch_branch(crate::SwitchBranchOptions {
branch_id: branch.id.clone(),
})
.await
.expect("branch should switch");
}
let scan = scan_sql("censusrow");
let arm = if rotate { "rot" } else { "main" };
let _ = crate::storage_bench::take_tracked_key_allocation_census();
let _ = crate::storage_bench::take_root_base_batch_cache_accounting();
let _ = crate::storage_bench::take_tracked_scan_branch_accounting();
let started = Instant::now();
let rows = session.execute(&scan, &[]).await.expect("cold scan");
let cold_us = started.elapsed().as_micros();
assert_eq!(rows.len(), 1, "the cold scan must answer exactly one row");
print_census_line(arm, n, "cold", 1, cold_us);
let started = Instant::now();
for _ in 0..reps {
let rows = session.execute(&scan, &[]).await.expect("warm scan");
assert_eq!(rows.len(), 1, "the warm scan must answer exactly one row");
}
let warm_us = started.elapsed().as_micros();
print_census_line(arm, n, "warm", reps, warm_us);
}
}
}
#[cfg(feature = "storage-benches")]
fn print_census_line(arm: &str, n: usize, phase: &str, scans: usize, us: u128) {
let c = crate::storage_bench::take_tracked_key_allocation_census();
let (hit, miss) = crate::storage_bench::take_root_base_batch_cache_accounting();
let (dur, exa, rep) = crate::storage_bench::take_tracked_scan_branch_accounting();
println!(
"{arm},{n},{phase},{scans},{},{},{},{},{},{},{},{},{},{},{},{},{},{hit},{miss},{dur},{exa},{rep},{us}",
c.commit_delta_rows_loaded,
c.key_decode_calls,
c.key_decode_input_bytes,
c.key_decode_owned_string_bytes,
c.key_decode_escaped_strings,
c.commit_delta_account_id_bytes,
c.commit_delta_point_key_encodes,
c.commit_delta_point_key_encode_bytes,
c.commit_delta_request_key_clones,
c.commit_delta_request_key_clone_bytes,
c.materialize_owned_key_builds,
c.materialize_owned_key_bytes,
c.materialize_reverify_rows,
);
}
async fn created_at_of(session: &SessionContext<Memory>, table: &str, id: &str) -> String {
let result = session
.execute(
&format!("SELECT lixcol_created_at AS created_at FROM {table} WHERE id = '{id}'"),
&[],
)
.await
.expect("created_at reads");
result.rows()[0]
.get::<String>("created_at")
.expect("created_at is text")
}
#[tokio::test]
async fn recreated_identity_created_at_without_compaction() {
let (_storage, session) = open_session().await;
register(&session, probe_schema("c8row")).await;
session
.execute("INSERT INTO c8row (id, locale) VALUES ('row-0', 'first')", &[])
.await
.expect("first insert should commit");
session
.create_checkpoint()
.await
.expect("baseline checkpoint should publish");
let first = created_at_of(&session, "c8row", "row-0").await;
session
.execute("DELETE FROM c8row WHERE id = 'row-0'", &[])
.await
.expect("delete should commit");
session
.create_checkpoint()
.await
.expect("checkpoint should publish");
session
.execute(
"INSERT INTO c8row (id, locale) VALUES ('row-0', 'second')",
&[],
)
.await
.expect("re-insert should commit");
let second = created_at_of(&session, "c8row", "row-0").await;
println!(
"phase8 | arm=control first_created_at={first} second_created_at={second} \
inherited={}",
first == second
);
assert_eq!(
second, first,
"a re-inserted identity inherits its deleted predecessor's created_at"
);
}
#[tokio::test]
async fn recreated_identity_created_at_after_compaction() {
let (storage, session) = open_session().await;
register(&session, probe_schema("c8row")).await;
let branch_id = session
.active_branch_id()
.await
.expect("active branch id reads");
session
.execute("INSERT INTO c8row (id, locale) VALUES ('row-0', 'first')", &[])
.await
.expect("first insert should commit");
session
.create_checkpoint()
.await
.expect("baseline checkpoint should publish");
let first = created_at_of(&session, "c8row", "row-0").await;
session
.execute("DELETE FROM c8row WHERE id = 'row-0'", &[])
.await
.expect("delete should commit");
session
.create_checkpoint()
.await
.expect("checkpoint should publish");
drop(session);
let before = row_census(&storage).await;
let removed = drop_all_tombstones(&storage).await;
let after = row_census(&storage).await;
println!(
"phase8 | arm=compacted tombstones_removed={removed} \
entries {}->{} tombstones {}->{}",
before.entries, after.entries, before.tombstones, after.tombstones
);
assert_eq!(
after.tombstones, 0,
"the identity must carry no tombstone into the re-insert (removed={removed})"
);
let session = reopen_session(&storage).await;
let hits_before = crate::hot_state::BROAD_CANONICAL_CREATED_AT_HITS
.load(std::sync::atomic::Ordering::Relaxed);
let keys_before = crate::hot_state::BROAD_CANONICAL_CREATED_AT_KEYS
.load(std::sync::atomic::Ordering::Relaxed);
let lookups_before = crate::hot_state::BROAD_CANONICAL_CREATED_AT_LOOKUPS
.load(std::sync::atomic::Ordering::Relaxed);
session
.execute(
"INSERT INTO c8row (id, locale) VALUES ('row-0', 'second')",
&[],
)
.await
.expect("re-insert should commit");
let hits = crate::hot_state::BROAD_CANONICAL_CREATED_AT_HITS
.load(std::sync::atomic::Ordering::Relaxed)
- hits_before;
let keys = crate::hot_state::BROAD_CANONICAL_CREATED_AT_KEYS
.load(std::sync::atomic::Ordering::Relaxed)
- keys_before;
let lookups = crate::hot_state::BROAD_CANONICAL_CREATED_AT_LOOKUPS
.load(std::sync::atomic::Ordering::Relaxed)
- lookups_before;
let second = created_at_of(&session, "c8row", "row-0").await;
println!(
"phase8 | arm=compacted first_created_at={first} second_created_at={second} \
inherited={} lookups={lookups} keys={keys} hits={hits}",
first == second
);
drop(session);
assert!(lookups > 0, "the canonical lookup must have run");
assert!(
hits > 0,
"the re-insert over a compacted identity must have inherited from canonical"
);
assert_eq!(
second, first,
"canonical must supply the created_at the compacted tombstone used to carry"
);
let engine = Engine::new(storage.clone())
.await
.expect("engine should reopen for rebuild");
let validations_before = crate::tracked_state::DIFF_ROW_CREATED_AT_VALIDATIONS
.load(std::sync::atomic::Ordering::Relaxed);
let rebuild = engine.rebuild_tracked_state_for_branch(&branch_id).await;
let validations = crate::tracked_state::DIFF_ROW_CREATED_AT_VALIDATIONS
.load(std::sync::atomic::Ordering::Relaxed)
- validations_before;
match &rebuild {
Ok(()) => println!("phase8 | arm=compacted rebuild=accepted"),
Err(error) => println!(
"phase8 | arm=compacted rebuild=REJECTED {}",
error.to_string().replace('\n', " ")
),
}
println!(
"phase8 | verdict inherited={} rebuild_ok={} created_at_validations={validations}",
first == second,
rebuild.is_ok()
);
}
#[tokio::test]
async fn broad_canonical_created_at_recovery_misses_for_new_identities() {
let (_storage, session) = open_session().await;
register(&session, probe_schema("c8new")).await;
let keys_before = crate::hot_state::BROAD_CANONICAL_CREATED_AT_KEYS
.load(std::sync::atomic::Ordering::Relaxed);
let hits_before = crate::hot_state::BROAD_CANONICAL_CREATED_AT_HITS
.load(std::sync::atomic::Ordering::Relaxed);
session
.execute(
"INSERT INTO c8new (id, locale) VALUES ('new-0', 'a'), ('new-1', 'b')",
&[],
)
.await
.expect("new rows should commit");
let keys = crate::hot_state::BROAD_CANONICAL_CREATED_AT_KEYS
.load(std::sync::atomic::Ordering::Relaxed)
- keys_before;
let hits = crate::hot_state::BROAD_CANONICAL_CREATED_AT_HITS
.load(std::sync::atomic::Ordering::Relaxed)
- hits_before;
println!("phase8 | arm=new_identities keys={keys} hits={hits}");
assert!(
keys >= 2,
"both new identities must reach the canonical lookup"
);
assert_eq!(
hits, 0,
"a genuinely new identity has no canonical ancestor to inherit from"
);
}
#[derive(Clone, Copy, Debug, Default)]
struct CompactionCounters {
routes: u64,
offered: u64,
candidates: u64,
compacted: u64,
}
fn compaction_counters() -> CompactionCounters {
let load = |counter: &std::sync::atomic::AtomicU64| {
counter.load(std::sync::atomic::Ordering::Relaxed)
};
CompactionCounters {
routes: load(&crate::hot_state::COMPACTED_TOMBSTONE_ROUTES),
offered: load(&crate::hot_state::COMPACTED_TOMBSTONE_OFFERED),
candidates: load(&crate::hot_state::COMPACTED_TOMBSTONE_CANDIDATES),
compacted: load(&crate::hot_state::COMPACTED_TOMBSTONE_COMPACTED),
}
}
impl CompactionCounters {
fn since(self, before: Self) -> Self {
Self {
routes: self.routes - before.routes,
offered: self.offered - before.offered,
candidates: self.candidates - before.candidates,
compacted: self.compacted - before.compacted,
}
}
}
impl std::fmt::Display for CompactionCounters {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"routes={} offered={} candidates={} compacted={}",
self.routes, self.offered, self.candidates, self.compacted
)
}
}
#[tokio::test]
async fn checkpoint_compacts_discharged_branch_tombstones() {
const N: usize = 300;
let (storage, session) = open_session().await;
register(&session, probe_schema("cp10row")).await;
insert_rows(&session, "cp10row", N).await;
session
.create_checkpoint()
.await
.expect("baseline checkpoint should publish");
delete_all_but_first(&session, "cp10row", N).await;
let before = row_census(&storage).await;
assert!(
before.tombstones >= N - 1,
"the churn must leave one tombstone per delete, saw {}",
before.tombstones
);
let counters_before = compaction_counters();
session
.create_checkpoint()
.await
.expect("checkpoint should publish");
let counters = compaction_counters().since(counters_before);
let after = row_census(&storage).await;
println!(
"phase10 | compaction entries {}->{} tombstones {}->{} {counters}",
before.entries, after.entries, before.tombstones, after.tombstones
);
assert!(
counters.candidates >= (N - 1) as u64,
"every discharged tombstone must reach the compaction route, saw {}",
counters.candidates
);
assert!(
counters.compacted >= (N - 1) as u64,
"every gate must clear for a branch-local schema with no base, saw {}",
counters.compacted
);
assert!(
after.tombstones * 10 < before.tombstones,
"the checkpoint must reclaim the tombstones it discharged, {} -> {}",
before.tombstones,
after.tombstones
);
assert!(
after.entries < before.entries,
"reclaimed tombstones must leave the serving view smaller"
);
let session = reopen_session(&storage).await;
let survivors = session
.execute("SELECT id FROM cp10row", &[])
.await
.expect("collection scan should run");
assert_eq!(
survivors.len(),
1,
"compaction must not resurrect a deleted row"
);
let answer = session
.execute(&scan_sql("cp10row"), &[])
.await
.expect("scan should run");
assert_eq!(answer.len(), 1, "the surviving row must still answer");
assert_eq!(
working_diff_rows(&session, "cp10row").await,
0,
"a checkpoint discharges every delete, compacted or not"
);
}
#[tokio::test]
async fn recreated_identity_inherits_created_at_after_engine_compaction() {
let (storage, session) = open_session().await;
register(&session, probe_schema("ci10row")).await;
session
.execute(
"INSERT INTO ci10row (id, locale) VALUES ('row-0', 'first')",
&[],
)
.await
.expect("first insert should commit");
session
.create_checkpoint()
.await
.expect("baseline checkpoint should publish");
let first = created_at_of(&session, "ci10row", "row-0").await;
session
.execute("DELETE FROM ci10row WHERE id = 'row-0'", &[])
.await
.expect("delete should commit");
let counters_before = compaction_counters();
session
.create_checkpoint()
.await
.expect("checkpoint should publish");
let counters = compaction_counters().since(counters_before);
assert!(
counters.compacted > 0,
"this test is about a compacted identity; nothing was compacted"
);
let hits_before = crate::hot_state::BROAD_CANONICAL_CREATED_AT_HITS
.load(std::sync::atomic::Ordering::Relaxed);
session
.execute(
"INSERT INTO ci10row (id, locale) VALUES ('row-0', 'second')",
&[],
)
.await
.expect("re-insert should commit");
let hits = crate::hot_state::BROAD_CANONICAL_CREATED_AT_HITS
.load(std::sync::atomic::Ordering::Relaxed)
- hits_before;
let second = created_at_of(&session, "ci10row", "row-0").await;
println!(
"phase10 | created_at first={first} second={second} inherited={} \
canonical_hits={hits} compacted={}",
first == second,
counters.compacted
);
assert!(
hits > 0,
"the re-insert must have sourced its created_at from canonical state"
);
assert_eq!(
second, first,
"a re-insert over a compacted identity inherits its original created_at"
);
drop(session);
let engine = Engine::new(storage.clone())
.await
.expect("engine should reopen for rebuild");
let branch_id = engine
.open_session()
.await
.expect("session should open")
.active_branch_id()
.await
.expect("active branch id reads");
engine
.rebuild_tracked_state_for_branch(&branch_id)
.await
.expect("a compacted branch must still rebuild");
}
#[tokio::test]
#[ignore = "measurement probe, not a gate"]
async fn checkpoint_compaction_scan_cost() {
let sizes = sizes_from_env("LIX_TOMBSTONE_SIZES", &[100, 1_000, 10_000]);
let reps = reps_from_env(5);
println!("phase10 | n,event,row_entries,tombstones,scan_us,counters");
for n in sizes {
let (storage, session) = open_session().await;
register(&session, probe_schema("m10row")).await;
insert_rows(&session, "m10row", n).await;
session
.create_checkpoint()
.await
.expect("baseline checkpoint should publish");
delete_all_but_first(&session, "m10row", n).await;
let census = row_census(&storage).await;
let scan = timed_scan(&session, &scan_sql("m10row"), 1, reps).await;
println!(
"{n},after_churn,{},{},{},-",
census.entries,
census.tombstones,
scan.as_micros()
);
let counters_before = compaction_counters();
session
.create_checkpoint()
.await
.expect("checkpoint should publish");
let counters = compaction_counters().since(counters_before);
let census = row_census(&storage).await;
let scan = timed_scan(&session, &scan_sql("m10row"), 1, reps).await;
println!(
"{n},after_checkpoint,{},{},{},{}",
census.entries,
census.tombstones,
scan.as_micros(),
counters
);
}
}
#[tokio::test]
#[ignore = "measurement probe, not a gate"]
async fn hot_row_tombstone_decode_census() {
use std::sync::atomic::Ordering;
fn reset() {
crate::hot_state::HOT_SCAN_DECODED_ENTRIES.store(0, Ordering::Relaxed);
crate::hot_state::HOT_SCAN_MATCHED_ENTRIES.store(0, Ordering::Relaxed);
crate::hot_state::HOT_SCAN_TOMBSTONE_ENTRIES.store(0, Ordering::Relaxed);
}
fn read() -> (u64, u64, u64) {
(
crate::hot_state::HOT_SCAN_DECODED_ENTRIES.load(Ordering::Relaxed),
crate::hot_state::HOT_SCAN_MATCHED_ENTRIES.load(Ordering::Relaxed),
crate::hot_state::HOT_SCAN_TOMBSTONE_ENTRIES.load(Ordering::Relaxed),
)
}
let sizes = sizes_from_env("LIX_TOMBSTONE_SIZES", &[1_000, 10_000]);
let reps = reps_from_env(5);
println!(
"phase7 | arm,n,row_entries,space_tombstones,decoded,matched,decoded_tombstones,answer_rows,scan_us"
);
for n in sizes {
{
let (storage, session) = open_session().await;
register(&session, probe_schema("c7fresh")).await;
insert_rows(&session, "c7fresh", 1).await;
let census = row_census(&storage).await;
let _ = timed_scan(&session, &scan_sql("c7fresh"), 1, 1).await;
reset();
let rows = session
.execute(&scan_sql("c7fresh"), &[])
.await
.expect("scan");
let (decoded, matched, tombs) = read();
let scan = timed_scan(&session, &scan_sql("c7fresh"), 1, reps).await;
println!(
"fresh,{n},{},{},{decoded},{matched},{tombs},{},{}",
census.entries,
census.tombstones,
rows.len(),
scan.as_micros()
);
}
{
let (storage, session) = open_session().await;
register(&session, probe_schema("c7upd")).await;
insert_rows(&session, "c7upd", 2).await;
for i in 0..(n - 1) {
session
.execute(
"UPDATE c7upd SET locale = $1 WHERE id = 'row-1'",
&[crate::Value::Text(format!("v{i}"))],
)
.await
.expect("update should commit");
}
let census = row_census(&storage).await;
let _ = timed_scan(&session, &scan_sql("c7upd"), 1, 1).await;
reset();
let rows = session.execute(&scan_sql("c7upd"), &[]).await.expect("scan");
let (decoded, matched, tombs) = read();
let scan = timed_scan(&session, &scan_sql("c7upd"), 1, reps).await;
println!(
"update_churn,{n},{},{},{decoded},{matched},{tombs},{},{}",
census.entries,
census.tombstones,
rows.len(),
scan.as_micros()
);
}
{
let (storage, session) = open_session().await;
register(&session, probe_schema("c7clean")).await;
insert_rows(&session, "c7clean", n).await;
session
.create_checkpoint()
.await
.expect("checkpoint should publish");
delete_all_but_first(&session, "c7clean", n).await;
let census = row_census(&storage).await;
let _ = timed_scan(&session, &scan_sql("c7clean"), 1, 1).await;
reset();
let rows = session
.execute(&scan_sql("c7clean"), &[])
.await
.expect("scan");
let (decoded, matched, tombs) = read();
let scan = timed_scan(&session, &scan_sql("c7clean"), 1, reps).await;
println!(
"churn_clean,{n},{},{},{decoded},{matched},{tombs},{},{}",
census.entries,
census.tombstones,
rows.len(),
scan.as_micros()
);
assert!(
tombs > 0,
"a Clean pre-image keeps its tombstones, so the tombstone census must observe them here"
);
}
{
let (storage, session) = open_session().await;
register(&session, probe_schema("c7churn")).await;
insert_rows(&session, "c7churn", n).await;
delete_all_but_first(&session, "c7churn", n).await;
let census = row_census(&storage).await;
let _ = timed_scan(&session, &scan_sql("c7churn"), 1, 1).await;
reset();
let rows = session
.execute(&scan_sql("c7churn"), &[])
.await
.expect("scan");
let (decoded, matched, tombs) = read();
let scan = timed_scan(&session, &scan_sql("c7churn"), 1, reps).await;
println!(
"churn,{n},{},{},{decoded},{matched},{tombs},{},{}",
census.entries,
census.tombstones,
rows.len(),
scan.as_micros()
);
assert!(
decoded > 0,
"the churn arm must decode something, otherwise the census does not observe this site"
);
let _ = drop_all_tombstones(&storage).await;
let session = reopen_session(&storage).await;
let census = row_census(&storage).await;
let _ = timed_scan(&session, &scan_sql("c7churn"), 1, 1).await;
reset();
let rows = session
.execute(&scan_sql("c7churn"), &[])
.await
.expect("scan");
let (decoded, matched, tombs) = read();
let scan = timed_scan(&session, &scan_sql("c7churn"), 1, reps).await;
println!(
"churn_dropped,{n},{},{},{decoded},{matched},{tombs},{},{}",
census.entries,
census.tombstones,
rows.len(),
scan.as_micros()
);
}
}
}
#[tokio::test]
#[ignore = "measurement probe, not a gate"]
async fn hot_row_tombstone_churn_cycles() {
let n = sizes_from_env("LIX_TOMBSTONE_SIZES", &[500])[0];
let rounds = sizes_from_env("LIX_TOMBSTONE_ROUNDS", &[8])[0];
let reps = reps_from_env(5);
println!(
"phase11 | cadence,round,row_entries,tombstones,live_entries,packed_bases,root_bases,answer_rows,scan_us,compaction"
);
for (label, cadence) in [
("no_checkpoint", 0_usize),
("checkpoint_each_round", 1),
("checkpoint_every_4", 4),
("deferred_delete_ckpt", 9999),
] {
let (storage, session) = open_session().await;
let table = match cadence {
0 => "p11never",
1 => "p11each",
9999 => "p11defer",
_ => "p11four",
};
register(&session, probe_schema(table)).await;
session
.execute(
&format!("INSERT INTO {table} (id, locale) VALUES ('row-0', 'keep')"),
&[],
)
.await
.expect("survivor should insert");
if cadence == 9999 {
for round in 1..=rounds {
let ids = (0..n)
.map(|i| format!("('d{round}-{i}', 'drop')"))
.collect::<Vec<_>>()
.join(",");
session
.execute(
&format!("INSERT INTO {table} (id, locale) VALUES {ids}"),
&[],
)
.await
.expect("round insert should commit");
session
.create_checkpoint()
.await
.expect("insert checkpoint should publish");
let del = (0..n)
.map(|i| format!("'d{round}-{i}'"))
.collect::<Vec<_>>()
.join(",");
session
.execute(&format!("DELETE FROM {table} WHERE id IN ({del})"), &[])
.await
.expect("round delete should commit");
let before = compaction_counters();
session
.create_checkpoint()
.await
.expect("delete checkpoint should publish");
let counters = compaction_counters().since(before);
let census = row_census(&storage).await;
let scan = timed_scan(&session, &scan_sql(table), 1, reps).await;
let live = session
.execute(&format!("SELECT id FROM {table}"), &[])
.await
.expect("collection scan should run")
.len();
assert_eq!(live, 1, "every round must end with exactly one live row");
println!(
"{label},{round},{},{},{},{},{},1,{},{counters}",
census.entries,
census.tombstones,
census.live(),
census.packed_bases,
census.root_bases,
scan.as_micros()
);
}
continue;
}
for round in 1..=rounds {
let ids = (0..n)
.map(|i| format!("('c{round}-{i}', 'drop')"))
.collect::<Vec<_>>()
.join(",");
session
.execute(
&format!("INSERT INTO {table} (id, locale) VALUES {ids}"),
&[],
)
.await
.expect("round insert should commit");
let del = (0..n)
.map(|i| format!("'c{round}-{i}'"))
.collect::<Vec<_>>()
.join(",");
session
.execute(
&format!("DELETE FROM {table} WHERE id IN ({del})"),
&[],
)
.await
.expect("round delete should commit");
let mut counters = CompactionCounters::default();
if cadence > 0 && round % cadence == 0 {
let before = compaction_counters();
session
.create_checkpoint()
.await
.expect("checkpoint should publish");
counters = compaction_counters().since(before);
}
let census = row_census(&storage).await;
let scan = timed_scan(&session, &scan_sql(table), 1, reps).await;
let live = session
.execute(&format!("SELECT id FROM {table}"), &[])
.await
.expect("collection scan should run")
.len();
assert_eq!(live, 1, "every round must end with exactly one live row");
println!(
"{label},{round},{},{},{},{},{},1,{},{counters}",
census.entries,
census.tombstones,
census.live(),
census.packed_bases,
census.root_bases,
scan.as_micros()
);
}
}
}
async fn drop_tombstones_named(storage: &Memory, needle: &str) -> usize {
let adapter = StorageAdapter::new(storage.clone());
let read = adapter
.begin_read(StorageReadOptions::default())
.await
.expect("read the hot row plane");
let entries = space_entries(&read, crate::hot_state::ROW_SPACE).await;
drop(read);
let mut writes = adapter.new_write_set();
let mut removed = 0_usize;
for (key, value) in entries {
let is_tombstone = value.len() > 1 && value[1] & HEAD_VALUE_DELETED_BIT != 0;
let names = key
.0
.windows(needle.len())
.any(|window| window == needle.as_bytes());
if is_tombstone && names {
writes.delete(crate::hot_state::ROW_SPACE, key);
removed += 1;
}
}
adapter
.commit_write_set(writes, StorageWriteOptions::default())
.await
.expect("tombstone removal should commit");
removed
}
#[tokio::test]
#[ignore = "measurement probe, not a gate"]
async fn interval_local_tombstone_has_no_dependent_reader() {
const N: usize = 40;
{
let (storage, session) = open_session().await;
register(&session, probe_schema("p12row")).await;
session
.execute(
"INSERT INTO p12row (id, locale) VALUES ('row-0', 'keep')",
&[],
)
.await
.expect("survivor should insert");
session
.create_checkpoint()
.await
.expect("baseline checkpoint should publish");
for i in 0..N {
session
.execute(
"INSERT INTO p12row (id, locale) VALUES ($1, 'drop')",
&[crate::Value::Text(format!("ephem-{i}"))],
)
.await
.expect("ephemeral insert should commit");
}
for i in 0..N {
session
.execute(
"DELETE FROM p12row WHERE id = $1",
&[crate::Value::Text(format!("ephem-{i}"))],
)
.await
.expect("ephemeral delete should commit");
}
let census = row_census(&storage).await;
let diff_before = working_diff_rows(&session, "p12row").await;
let live_before = session
.execute("SELECT id FROM p12row", &[])
.await
.expect("collection read")
.len();
println!(
"phase12 | arm=local before: entries={} tombstones={} packed_bases={} root_bases={} working_diff_rows={diff_before} live={live_before}",
census.entries, census.tombstones, census.packed_bases, census.root_bases
);
assert_eq!(live_before, 1, "only the survivor is live before the removal");
let removed = drop_tombstones_named(&storage, "ephem-").await;
let session = reopen_session(&storage).await;
let census = row_census(&storage).await;
let diff_after = working_diff_rows(&session, "p12row").await;
let live_after = session
.execute("SELECT id FROM p12row", &[])
.await
.expect("collection read")
.len();
let point = session
.execute(
"SELECT id FROM p12row WHERE id = $1",
&[crate::Value::Text("ephem-0".to_string())],
)
.await
.expect("point read")
.len();
println!(
"phase12 | arm=local after: removed={removed} entries={} tombstones={} working_diff_rows={diff_after} live={live_after} point_hits_for_deleted={point}"
, census.entries, census.tombstones);
assert_eq!(removed, N, "every interval-local tombstone should be removable");
assert_eq!(live_after, 1, "removal must not resurrect an ephemeral row");
assert_eq!(point, 0, "a point read must not resurrect an ephemeral row");
assert_eq!(
diff_after, diff_before,
"the working diff must not move when an interval-local tombstone goes"
);
session
.create_checkpoint()
.await
.expect("closing checkpoint should publish after the removal");
let live_ckpt = session
.execute("SELECT id FROM p12row", &[])
.await
.expect("collection read")
.len();
let census = row_census(&storage).await;
println!(
"phase12 | arm=local after_closing_checkpoint: entries={} tombstones={} live={live_ckpt} working_diff_rows={}",
census.entries,
census.tombstones,
working_diff_rows(&session, "p12row").await
);
assert_eq!(live_ckpt, 1, "the closing checkpoint must not resurrect");
}
for remove in [false, true] {
let (storage, session) = open_session().await;
register(&session, probe_schema("p12fork")).await;
session
.execute(
"INSERT INTO p12fork (id, locale) VALUES ('row-0', 'keep')",
&[],
)
.await
.expect("survivor should insert");
session
.create_checkpoint()
.await
.expect("baseline checkpoint should publish");
let main_branch_id = session
.active_branch_id()
.await
.expect("active branch should resolve");
for i in 0..N {
session
.execute(
"INSERT INTO p12fork (id, locale) VALUES ($1, 'drop')",
&[crate::Value::Text(format!("ephem-{i}"))],
)
.await
.expect("ephemeral insert should commit");
}
let fork = session
.create_branch(crate::CreateBranchOptions {
id: None,
name: "p12-midinterval".to_string(),
from_commit_id: None,
})
.await
.expect("branch should create");
for i in 0..N {
session
.execute(
"DELETE FROM p12fork WHERE id = $1",
&[crate::Value::Text(format!("ephem-{i}"))],
)
.await
.expect("ephemeral delete should commit");
}
let commits_before = session
.execute("SELECT id FROM lix_commit", &[])
.await
.expect("commit graph read")
.len();
let census = row_census(&storage).await;
println!(
"phase12 | arm=fork(remove={remove}) before: entries={} tombstones={} root_bases={} commits={commits_before}",
census.entries, census.tombstones, census.root_bases
);
let removed = if remove {
drop_tombstones_named(&storage, "ephem-").await
} else {
0
};
let session = reopen_session(&storage).await;
let live_main = session
.execute("SELECT id FROM p12fork", &[])
.await
.expect("collection read")
.len();
session
.switch_branch(crate::SwitchBranchOptions {
branch_id: fork.id.clone(),
})
.await
.expect("switch to the fork");
let live_fork = session
.execute("SELECT id FROM p12fork", &[])
.await
.expect("fork collection read")
.len();
let commits_after = session
.execute("SELECT id FROM lix_commit", &[])
.await
.expect("commit graph read")
.len();
println!(
"phase12 | arm=fork(remove={remove}) after: removed={removed} live_main={live_main} live_fork={live_fork} commits={commits_after}"
);
assert_eq!(
removed,
if remove { N } else { 0 },
"the control must remove nothing and the treatment every tombstone"
);
assert_eq!(live_main, 1, "main must still show only the survivor");
assert_eq!(
live_fork,
N + 1,
"the fork forked while the ephemeral rows were alive and must still see them"
);
assert_eq!(
commits_after, commits_before,
"removing a serving-view tombstone must not change the commit graph"
);
session
.switch_branch(crate::SwitchBranchOptions {
branch_id: main_branch_id,
})
.await
.expect("switch back to main");
let merged = session
.merge_branch(crate::MergeBranchOptions {
source_branch_id: fork.id.clone(),
})
.await;
let live_merged = session
.execute("SELECT id FROM p12fork", &[])
.await
.expect("post-merge collection read")
.len();
println!(
"phase12 | arm=fork(remove={remove}) merge: ok={} live_after_merge={live_merged}",
merged.is_ok()
);
merged.expect("merge must not error after the tombstones are gone");
}
for remove in [false, true] {
let (storage, session) = open_session().await;
register(&session, probe_schema("p12undo")).await;
session
.execute(
"INSERT INTO p12undo (id, locale) VALUES ('row-0', 'keep')",
&[],
)
.await
.expect("survivor should insert");
session
.create_checkpoint()
.await
.expect("baseline checkpoint should publish");
session
.execute(
"INSERT INTO p12undo (id, locale) VALUES ('ephem-0', 'drop')",
&[],
)
.await
.expect("ephemeral insert should commit");
session
.execute("DELETE FROM p12undo WHERE id = 'ephem-0'", &[])
.await
.expect("ephemeral delete should commit");
let removed = if remove {
drop_tombstones_named(&storage, "ephem-").await
} else {
0
};
let session = reopen_session(&storage).await;
session.undo().await.expect("undo should run");
let live = session
.execute("SELECT id FROM p12undo", &[])
.await
.expect("post-undo collection read")
.len();
println!("phase12 | arm=undo(remove={remove}): removed={removed} live_after_undo={live}");
assert_eq!(
live, 2,
"undo of the delete must restore the ephemeral row even with its tombstone gone"
);
}
}
#[derive(Clone, Copy, Debug, Default)]
struct ElisionCounters {
routes: u64,
offered: u64,
candidates: u64,
elided: u64,
}
fn elision_counters() -> ElisionCounters {
let load =
|counter: &std::sync::atomic::AtomicU64| counter.load(std::sync::atomic::Ordering::Relaxed);
ElisionCounters {
routes: load(&crate::hot_state::INTERVAL_LOCAL_TOMBSTONE_ROUTES),
offered: load(&crate::hot_state::INTERVAL_LOCAL_TOMBSTONE_OFFERED),
candidates: load(&crate::hot_state::INTERVAL_LOCAL_TOMBSTONE_CANDIDATES),
elided: load(&crate::hot_state::INTERVAL_LOCAL_TOMBSTONE_ELIDED),
}
}
impl ElisionCounters {
fn since(self, before: Self) -> Self {
Self {
routes: self.routes - before.routes,
offered: self.offered - before.offered,
candidates: self.candidates - before.candidates,
elided: self.elided - before.elided,
}
}
}
impl std::fmt::Display for ElisionCounters {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"routes={} offered={} candidates={} elided={}",
self.routes, self.offered, self.candidates, self.elided
)
}
}
#[tokio::test]
async fn interval_local_delete_publishes_no_tombstone() {
const N: usize = 200;
let (storage, session) = open_session().await;
register(&session, probe_schema("p13row")).await;
session
.execute(
"INSERT INTO p13row (id, locale) VALUES ('row-0', 'keep')",
&[],
)
.await
.expect("survivor should insert");
session
.create_checkpoint()
.await
.expect("baseline checkpoint should publish");
let before = elision_counters();
insert_rows_named(&session, "p13row", "ephem", N).await;
let after_insert = row_census(&storage).await;
delete_rows_named(&session, "p13row", "ephem", N).await;
let counters = elision_counters().since(before);
let census = row_census(&storage).await;
println!(
"phase13 | after_insert entries={} tombstones={} | after_delete entries={} tombstones={} packed={} root={} {counters}",
after_insert.entries,
after_insert.tombstones,
census.entries,
census.tombstones,
census.packed_bases,
census.root_bases
);
assert!(
counters.candidates >= N as u64,
"every interval-local delete must reach the elision route, saw {}",
counters.candidates
);
assert!(
counters.elided >= N as u64,
"every gate must clear for a branch-local schema with no base, saw {}",
counters.elided
);
assert_eq!(
census.tombstones, 0,
"an interval-local delete must publish no tombstone"
);
assert!(
census.entries < after_insert.entries,
"eliding must leave the serving view smaller than it was with the rows live"
);
let session = reopen_session(&storage).await;
let live = session
.execute("SELECT id FROM p13row", &[])
.await
.expect("collection scan should run");
assert_eq!(live.len(), 1, "eliding must not resurrect a deleted row");
let point = session
.execute(
"SELECT id FROM p13row WHERE id = $1",
&[crate::Value::Text("ephem-0".to_string())],
)
.await
.expect("point read should run");
assert_eq!(point.len(), 0, "a point read must not resurrect either");
assert_eq!(
working_diff_rows(&session, "p13row").await,
0,
"an identity created and deleted in one interval owes the working diff nothing"
);
session
.create_checkpoint()
.await
.expect("closing checkpoint should publish");
let live = session
.execute("SELECT id FROM p13row", &[])
.await
.expect("collection scan should run");
assert_eq!(live.len(), 1, "the closing checkpoint must not resurrect");
}
#[tokio::test]
async fn clean_pre_image_delete_still_publishes_a_tombstone() {
const N: usize = 200;
let (storage, session) = open_session().await;
register(&session, probe_schema("p13cln")).await;
insert_rows_named(&session, "p13cln", "keeprow", N).await;
session
.create_checkpoint()
.await
.expect("baseline checkpoint should publish");
let before = elision_counters();
delete_rows_named(&session, "p13cln", "keeprow", N).await;
let counters = elision_counters().since(before);
let census = row_census(&storage).await;
println!(
"phase13 | inversion=clean_pre_image entries={} tombstones={} working_diff_rows={} {counters}",
census.entries,
census.tombstones,
working_diff_rows(&session, "p13cln").await
);
assert_eq!(
census.tombstones, N,
"the delete of a checkpointed row must still publish its tombstone"
);
assert_eq!(
working_diff_rows(&session, "p13cln").await,
N,
"and the working diff must report every one of those deletes"
);
}
#[tokio::test]
async fn interval_local_delete_over_a_base_still_publishes_a_tombstone() {
const N: usize = 50;
let (storage, session) = open_session().await;
register(&session, probe_schema("p13base")).await;
insert_rows_named(&session, "p13base", "based", N).await;
session
.create_checkpoint()
.await
.expect("baseline checkpoint should publish");
let branch = session
.create_branch(crate::CreateBranchOptions {
id: None,
name: "p13-based".to_string(),
from_commit_id: None,
})
.await
.expect("branch should create");
session
.switch_branch(crate::SwitchBranchOptions {
branch_id: branch.id.clone(),
})
.await
.expect("branch should switch");
let census = row_census(&storage).await;
assert!(
census.packed_bases > 0 || census.root_bases > 0,
"this inversion needs a base in the generation, saw packed={} root={}",
census.packed_bases,
census.root_bases
);
let before = elision_counters();
insert_rows_named(&session, "p13base", "ephem", N).await;
delete_rows_named(&session, "p13base", "ephem", N).await;
let counters = elision_counters().since(before);
let census = row_census(&storage).await;
println!(
"phase13 | inversion=has_base entries={} tombstones={} packed={} root={} {counters}",
census.entries, census.tombstones, census.packed_bases, census.root_bases
);
assert!(
counters.candidates >= N as u64,
"the deltas must reach the route, or the refusal below is vacuous, saw {}",
counters.candidates
);
assert_eq!(
census.tombstones, N,
"gate (a) must keep every tombstone while a base is visible"
);
let session = reopen_session(&storage).await;
let live = session
.execute("SELECT id FROM p13base", &[])
.await
.expect("collection scan should run");
assert_eq!(live.len(), N, "the based rows must survive, the ephemerals must not");
}
async fn insert_rows_named(
session: &SessionContext<Memory>,
table: &str,
prefix: &str,
count: usize,
) {
let mut index = 0;
while index < count {
let end = (index + CHUNK).min(count);
let values = (index..end)
.map(|i| format!("('{prefix}-{i}', 'drop')"))
.collect::<Vec<_>>()
.join(",");
session
.execute(
&format!("INSERT INTO {table} (id, locale) VALUES {values}"),
&[],
)
.await
.expect("named rows should insert");
index = end;
}
}
async fn delete_rows_named(
session: &SessionContext<Memory>,
table: &str,
prefix: &str,
count: usize,
) {
let mut index = 0;
while index < count {
let end = (index + CHUNK).min(count);
let ids = (index..end)
.map(|i| format!("'{prefix}-{i}'"))
.collect::<Vec<_>>()
.join(",");
session
.execute(&format!("DELETE FROM {table} WHERE id IN ({ids})"), &[])
.await
.expect("named bulk delete should run");
index = end;
}
}