use std::sync::{
atomic::{AtomicU64, Ordering},
Arc,
};
use crate::types::{ColumnId, TombstoneType, Value};
#[derive(Debug, Clone)]
pub struct ScanSummary {
pub element_tombstones_detected: u64,
}
#[derive(Debug, Clone)]
pub struct ScanSummaryHandle {
element_tombstones: Arc<AtomicU64>,
}
impl ScanSummaryHandle {
fn new() -> Self {
Self {
element_tombstones: Arc::new(AtomicU64::new(0)),
}
}
pub(super) fn add_element_tombstones(&self, n: u64) {
self.element_tombstones.fetch_add(n, Ordering::Relaxed);
}
pub fn read(&self) -> ScanSummary {
ScanSummary {
element_tombstones_detected: self.element_tombstones.load(Ordering::Relaxed),
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct RowKeys {
pub partition: Vec<Value>,
pub clustering: Vec<Value>,
}
impl RowKeys {
pub fn partition_only(partition: Vec<Value>) -> Self {
Self {
partition,
clustering: Vec::new(),
}
}
pub fn new(partition: Vec<Value>, clustering: Vec<Value>) -> Self {
Self {
partition,
clustering,
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct CellMeta {
pub writetime: i64,
pub expires_at: Option<i64>,
}
impl CellMeta {
pub fn new(writetime: i64) -> Self {
Self {
writetime,
expires_at: None,
}
}
pub fn with_ttl(writetime: i64, expires_at: i64) -> Self {
Self {
writetime,
expires_at: Some(expires_at),
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct CellDelta {
pub value: Option<Value>,
pub writetime: i64,
pub expires_at: Option<i64>,
pub replaced: bool,
}
impl CellDelta {
pub fn value(value: Value, writetime: i64) -> Self {
Self {
value: Some(value),
writetime,
expires_at: None,
replaced: false,
}
}
pub fn tombstone(writetime: i64) -> Self {
Self {
value: None,
writetime,
expires_at: None,
replaced: false,
}
}
pub fn value_with_ttl(value: Value, writetime: i64, expires_at: i64) -> Self {
Self {
value: Some(value),
writetime,
expires_at: Some(expires_at),
replaced: false,
}
}
pub fn collection_replace(value: Value, writetime: i64) -> Self {
Self {
value: Some(value),
writetime,
expires_at: None,
replaced: true,
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct RangeBound {
pub values: Vec<Value>,
pub inclusive: bool,
}
impl RangeBound {
pub fn new(values: Vec<Value>, inclusive: bool) -> Self {
Self { values, inclusive }
}
pub fn inclusive(values: Vec<Value>) -> Self {
Self {
values,
inclusive: true,
}
}
pub fn exclusive(values: Vec<Value>) -> Self {
Self {
values,
inclusive: false,
}
}
pub fn open() -> Self {
Self {
values: Vec::new(),
inclusive: false,
}
}
pub fn is_prefix(&self) -> bool {
!self.values.is_empty()
}
}
#[derive(Debug, Clone, PartialEq)]
pub enum DeltaRecord {
Upsert {
keys: RowKeys,
liveness: Option<CellMeta>,
cells: Vec<(ColumnId, CellDelta)>,
},
StaticUpsert {
partition_key: RowKeys,
cells: Vec<(ColumnId, CellDelta)>,
},
RowDelete {
keys: RowKeys,
deleted_at: i64,
},
RangeDelete {
partition_key: RowKeys,
start: RangeBound,
end: RangeBound,
deleted_at: i64,
},
PartitionDelete {
partition_key: RowKeys,
deleted_at: i64,
},
}
impl DeltaRecord {
pub fn partition_key(&self) -> &[Value] {
match self {
DeltaRecord::Upsert { keys, .. } => &keys.partition,
DeltaRecord::StaticUpsert { partition_key, .. } => &partition_key.partition,
DeltaRecord::RowDelete { keys, .. } => &keys.partition,
DeltaRecord::RangeDelete { partition_key, .. } => &partition_key.partition,
DeltaRecord::PartitionDelete { partition_key, .. } => &partition_key.partition,
}
}
pub fn op_name(&self) -> &'static str {
match self {
DeltaRecord::Upsert { .. } => "upsert",
DeltaRecord::StaticUpsert { .. } => "static_upsert",
DeltaRecord::RowDelete { .. } => "row_delete",
DeltaRecord::RangeDelete { .. } => "range_delete",
DeltaRecord::PartitionDelete { .. } => "partition_delete",
}
}
}
pub type ScanDeltaOutput = (
tokio::sync::mpsc::Receiver<crate::Result<DeltaRecord>>,
ScanSummaryHandle,
);
pub fn scan_delta(
sstable_dir: std::path::PathBuf,
schema: crate::schema::TableSchema,
buffer_size: usize,
) -> ScanDeltaOutput {
let (tx, rx) = tokio::sync::mpsc::channel(buffer_size.max(1));
let summary = ScanSummaryHandle::new();
let summary_for_task = summary.clone();
tokio::spawn(async move {
if let Err(e) = run_scan_delta(sstable_dir, schema, tx.clone(), summary_for_task).await {
let _ = tx.send(Err(e)).await;
}
});
(rx, summary)
}
async fn run_scan_delta(
sstable_dir: std::path::PathBuf,
schema: crate::schema::TableSchema,
tx: tokio::sync::mpsc::Sender<crate::Result<DeltaRecord>>,
summary: ScanSummaryHandle,
) -> crate::Result<()> {
use crate::storage::sstable::reader::SSTableReader;
let data_db = find_data_db(&sstable_dir)?;
let config = crate::Config::default();
let platform =
std::sync::Arc::new(crate::Platform::new(&config).await.map_err(|e| {
crate::Error::corruption(format!("scan_delta: platform init failed: {e}"))
})?);
let reader = std::sync::Arc::new(
SSTableReader::open(&data_db, &config, platform)
.await
.map_err(|e| {
crate::Error::corruption(format!("scan_delta: failed to open {:?}: {e}", data_db))
})?,
);
let schema_arc = std::sync::Arc::new(schema);
let (stitched, parser) = reader.prepare_delta_scan().await?;
let schema_for_parse = std::sync::Arc::clone(&schema_arc);
let reader_arc = std::sync::Arc::clone(&reader);
let parse_result = tokio::task::spawn_blocking(move || -> crate::Result<()> {
parser.parse_block_emit_delta(
&stitched,
Some(&schema_for_parse),
&reader_arc,
|(
partition_key_raw,
cells,
cell_meta,
row_liveness_ts,
is_static,
is_row_tombstone,
marked_for_delete_at,
range_info, // Issue #699: Some((start_vals,start_incl,end_vals,end_incl,del_at)) for range tombstone
is_partition_tombstone, col_complex_meta, liveness_expires_at_micros, )| {
let pk_columns = crate::storage::partition_key_codec::decode_partition_key_columns(
&partition_key_raw.0,
&schema_arc,
)
.map_err(|e| crate::Error::corruption(format!(
"scan_delta: failed to decode partition key (raw bytes {:?}): {e}",
partition_key_raw.0
)))?;
let partition_values: Vec<Value> = pk_columns.into_iter().map(|(_, v)| v).collect();
if is_partition_tombstone {
let deleted_at = marked_for_delete_at.ok_or_else(|| {
crate::Error::corruption(format!(
"scan_delta: partition tombstone for pk={:?} has no markedForDeleteAt \
— cannot represent faithfully (no-heuristics, issue #28)",
partition_values
))
})?;
let record = DeltaRecord::PartitionDelete {
partition_key: RowKeys::partition_only(partition_values),
deleted_at,
};
return match tx.blocking_send(Ok(record)) {
Ok(()) => Ok(std::ops::ControlFlow::Continue(())),
Err(_) => Ok(std::ops::ControlFlow::Break(())),
};
}
if let Some((start_vals, start_incl, end_vals, end_incl, del_at)) = range_info {
let record = DeltaRecord::RangeDelete {
partition_key: RowKeys::partition_only(partition_values),
start: RangeBound::new(start_vals, start_incl),
end: RangeBound::new(end_vals, end_incl),
deleted_at: del_at,
};
return match tx.blocking_send(Ok(record)) {
Ok(()) => Ok(std::ops::ControlFlow::Continue(())),
Err(_) => Ok(std::ops::ControlFlow::Break(())),
};
}
let pk_col_names: std::collections::HashSet<&str> = schema_arc
.partition_keys.iter().map(|k| k.name.as_str()).collect();
let clustering_col_names: std::collections::HashSet<&str> = schema_arc
.clustering_keys.iter().map(|ck| ck.name.as_str()).collect();
let static_col_names: std::collections::HashSet<&str> = schema_arc
.columns.iter().filter(|c| c.is_static).map(|c| c.name.as_str()).collect();
let clustering_values: Vec<Value> = schema_arc
.clustering_keys.iter()
.filter_map(|ck| cells.get(&ck.name).cloned())
.collect();
if is_row_tombstone {
let deleted_at = marked_for_delete_at.ok_or_else(|| {
crate::Error::corruption(format!(
"scan_delta: row tombstone for pk={:?} ck={:?} has no markedForDeleteAt \
— cannot represent faithfully (no-heuristics, issue #28)",
partition_values, clustering_values
))
})?;
let record = DeltaRecord::RowDelete {
keys: RowKeys::new(partition_values, clustering_values),
deleted_at,
};
return match tx.blocking_send(Ok(record)) {
Ok(()) => Ok(std::ops::ControlFlow::Continue(())),
Err(_) => Ok(std::ops::ControlFlow::Break(())),
};
}
{
let mut row_element_tombstones: u64 = 0;
for (col_name, ccm) in &col_complex_meta {
if ccm.element_tombstone_count > 0 {
row_element_tombstones += ccm.element_tombstone_count;
log::warn!(
"scan_delta DS4: collection column '{}' has {} element-level tombstone(s) \
that cannot be represented in v1 delta semantics (Issue #493 follow-up). \
These removals are counted in the scan summary but not in the emitted records.",
col_name, ccm.element_tombstone_count
);
}
}
if row_element_tombstones > 0 {
summary.add_element_tombstones(row_element_tombstones);
}
}
let mut cell_deltas: Vec<(ColumnId, CellDelta)> = Vec::new();
for (col_name, value) in &cells {
if pk_col_names.contains(col_name.as_str())
|| clustering_col_names.contains(col_name.as_str())
{
continue;
}
if is_static && !static_col_names.contains(col_name.as_str()) {
continue;
}
if !is_static && static_col_names.contains(col_name.as_str()) {
continue;
}
let meta = cell_meta.get(col_name.as_str());
let (writetime, expires_at) = match meta {
Some(m) => {
let exp = m.expiration.as_ref().map(|e| {
e.expires_at_seconds.saturating_mul(1_000_000)
});
(m.write_timestamp_micros, exp)
}
None => (row_liveness_ts.unwrap_or(0), None),
};
let (replaced, writetime) = match col_complex_meta.get(col_name.as_str()) {
Some(ccm) => {
let effective_wt = if ccm.max_element_writetime != 0 {
ccm.max_element_writetime
} else {
writetime
};
(ccm.has_collection_tombstone, effective_wt)
}
None => (false, writetime),
};
let cell = match value {
Value::Tombstone(info)
if info.tombstone_type == TombstoneType::CellTombstone =>
{
CellDelta {
value: None,
writetime: info.deletion_time,
expires_at: None,
replaced: false,
}
}
_ => CellDelta {
value: Some(value.clone()),
writetime,
expires_at,
replaced,
},
};
cell_deltas.push((ColumnId::new(col_name), cell));
}
let record = if is_static {
DeltaRecord::StaticUpsert {
partition_key: RowKeys::partition_only(partition_values),
cells: cell_deltas,
}
} else {
let liveness = row_liveness_ts.map(|ts| {
match liveness_expires_at_micros {
Some(exp) => CellMeta::with_ttl(ts, exp),
None => CellMeta::new(ts),
}
});
DeltaRecord::Upsert {
keys: RowKeys::new(partition_values, clustering_values),
liveness,
cells: cell_deltas,
}
};
match tx.blocking_send(Ok(record)) {
Ok(()) => Ok(std::ops::ControlFlow::Continue(())),
Err(_) => {
Ok(std::ops::ControlFlow::Break(()))
}
}
},
)
})
.await;
match parse_result {
Ok(result) => result,
Err(join_err) => Err(crate::Error::corruption(format!(
"scan_delta: parse task panicked: {join_err}"
))),
}
}
fn find_data_db(dir: &std::path::Path) -> crate::Result<std::path::PathBuf> {
if !dir.exists() {
return Err(crate::Error::corruption(format!(
"scan_delta: SSTable directory does not exist: {:?}",
dir
)));
}
let entries = std::fs::read_dir(dir).map_err(|e| {
crate::Error::corruption(format!("scan_delta: cannot read directory {:?}: {e}", dir))
})?;
let mut candidates: Vec<std::path::PathBuf> = entries
.filter_map(|e| e.ok())
.filter(|e| {
e.file_name()
.to_str()
.map(|n| n.ends_with("-Data.db"))
.unwrap_or(false)
})
.map(|e| e.path())
.collect();
if candidates.is_empty() {
return Err(crate::Error::corruption(format!(
"scan_delta: no Data.db file found in {:?}",
dir
)));
}
if candidates.len() > 1 {
candidates.sort();
log::warn!(
"scan_delta: {:?} contains {} Data.db files (expected 1 per generation); \
using lexicographically first: {:?}. Consider compacting before scanning.",
dir,
candidates.len(),
candidates[0]
);
}
Ok(candidates.remove(0))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::types::Value;
#[test]
fn row_keys_partition_only() {
let keys = RowKeys::partition_only(vec![Value::Integer(42)]);
assert_eq!(keys.partition, vec![Value::Integer(42)]);
assert!(keys.clustering.is_empty());
}
#[test]
fn row_keys_full() {
let keys = RowKeys::new(vec![Value::Integer(1)], vec![Value::Text("a".into())]);
assert_eq!(keys.partition.len(), 1);
assert_eq!(keys.clustering.len(), 1);
}
#[test]
fn cell_meta_no_ttl() {
let m = CellMeta::new(1_000_000);
assert_eq!(m.writetime, 1_000_000);
assert!(m.expires_at.is_none());
}
#[test]
fn cell_meta_with_ttl() {
let m = CellMeta::with_ttl(1_000_000, 2_000_000);
assert_eq!(m.expires_at, Some(2_000_000));
}
#[test]
fn cell_delta_value() {
let d = CellDelta::value(Value::Text("hello".into()), 100);
assert!(d.value.is_some());
assert_eq!(d.writetime, 100);
assert!(d.expires_at.is_none());
assert!(!d.replaced);
}
#[test]
fn cell_delta_tombstone() {
let d = CellDelta::tombstone(200);
assert!(d.value.is_none());
assert_eq!(d.writetime, 200);
assert!(!d.replaced);
}
#[test]
fn cell_delta_with_ttl() {
let d = CellDelta::value_with_ttl(Value::Integer(7), 100, 9999);
assert_eq!(d.expires_at, Some(9999));
assert!(!d.replaced);
}
#[test]
fn cell_delta_collection_replace() {
let d = CellDelta::collection_replace(Value::Text("x".into()), 300);
assert!(d.replaced);
assert!(d.value.is_some());
}
#[test]
fn range_bound_inclusive() {
let b = RangeBound::inclusive(vec![Value::Text("a".into())]);
assert!(b.inclusive);
assert_eq!(b.values.len(), 1);
}
#[test]
fn range_bound_exclusive() {
let b = RangeBound::exclusive(vec![Value::Text("m".into())]);
assert!(!b.inclusive);
}
#[test]
fn range_bound_open() {
let b = RangeBound::open();
assert!(b.values.is_empty());
assert!(!b.inclusive);
assert!(!b.is_prefix());
}
#[test]
fn range_bound_prefix() {
let b = RangeBound::inclusive(vec![Value::Integer(2024)]);
assert!(b.is_prefix());
}
fn sample_pk() -> RowKeys {
RowKeys::partition_only(vec![Value::Integer(1)])
}
fn sample_row_keys() -> RowKeys {
RowKeys::new(vec![Value::Integer(1)], vec![Value::Text("ck1".into())])
}
fn sample_cell() -> (ColumnId, CellDelta) {
(
ColumnId::new("val"),
CellDelta::value(Value::Text("hello".into()), 1_700_000_000_000_000),
)
}
#[test]
fn delta_record_upsert() {
let rec = DeltaRecord::Upsert {
keys: sample_row_keys(),
liveness: Some(CellMeta::new(1_700_000_000_000_000)),
cells: vec![sample_cell()],
};
assert_eq!(rec.op_name(), "upsert");
assert_eq!(rec.partition_key(), &[Value::Integer(1)]);
if let DeltaRecord::Upsert {
keys,
liveness,
cells,
} = &rec
{
assert_eq!(keys.clustering, vec![Value::Text("ck1".into())]);
assert!(liveness.is_some());
assert_eq!(cells.len(), 1);
} else {
panic!("expected Upsert");
}
}
#[test]
fn delta_record_upsert_no_liveness() {
let rec = DeltaRecord::Upsert {
keys: sample_row_keys(),
liveness: None,
cells: vec![sample_cell()],
};
if let DeltaRecord::Upsert { liveness, .. } = &rec {
assert!(liveness.is_none());
} else {
panic!("expected Upsert");
}
}
#[test]
fn delta_record_static_upsert() {
let rec = DeltaRecord::StaticUpsert {
partition_key: sample_pk(),
cells: vec![(
ColumnId::new("st"),
CellDelta::value(Value::Text("S".into()), 1_700_000_000_000_000),
)],
};
assert_eq!(rec.op_name(), "static_upsert");
assert_eq!(rec.partition_key(), &[Value::Integer(1)]);
if let DeltaRecord::StaticUpsert {
partition_key,
cells,
} = &rec
{
assert!(partition_key.clustering.is_empty());
assert_eq!(cells.len(), 1);
} else {
panic!("expected StaticUpsert");
}
}
#[test]
fn delta_record_row_delete() {
let rec = DeltaRecord::RowDelete {
keys: sample_row_keys(),
deleted_at: 1_700_000_000_000_000,
};
assert_eq!(rec.op_name(), "row_delete");
if let DeltaRecord::RowDelete { deleted_at, .. } = &rec {
assert_eq!(*deleted_at, 1_700_000_000_000_000);
} else {
panic!("expected RowDelete");
}
}
#[test]
fn delta_record_range_delete() {
let rec = DeltaRecord::RangeDelete {
partition_key: sample_pk(),
start: RangeBound::inclusive(vec![Value::Text("a".into())]),
end: RangeBound::exclusive(vec![Value::Text("m".into())]),
deleted_at: 1_700_000_000_000_001,
};
assert_eq!(rec.op_name(), "range_delete");
if let DeltaRecord::RangeDelete {
start,
end,
deleted_at,
..
} = &rec
{
assert!(start.inclusive);
assert!(!end.inclusive);
assert_eq!(*deleted_at, 1_700_000_000_000_001);
} else {
panic!("expected RangeDelete");
}
}
#[test]
fn delta_record_range_delete_prefix_bound() {
let rec = DeltaRecord::RangeDelete {
partition_key: sample_pk(),
start: RangeBound::inclusive(vec![Value::Integer(2024)]),
end: RangeBound::inclusive(vec![Value::Integer(2024)]),
deleted_at: 1_700_000_000_000_002,
};
if let DeltaRecord::RangeDelete { start, end, .. } = &rec {
assert!(start.is_prefix());
assert!(end.is_prefix());
assert_eq!(start.values.len(), 1);
assert_eq!(end.values.len(), 1);
} else {
panic!("expected RangeDelete");
}
}
#[test]
fn delta_record_partition_delete() {
let rec = DeltaRecord::PartitionDelete {
partition_key: sample_pk(),
deleted_at: 1_700_000_000_000_003,
};
assert_eq!(rec.op_name(), "partition_delete");
if let DeltaRecord::PartitionDelete {
partition_key,
deleted_at,
} = &rec
{
assert_eq!(partition_key.partition, vec![Value::Integer(1)]);
assert!(partition_key.clustering.is_empty());
assert_eq!(*deleted_at, 1_700_000_000_000_003);
} else {
panic!("expected PartitionDelete");
}
}
#[test]
fn op_names_are_distinct() {
let ops = [
DeltaRecord::Upsert {
keys: sample_row_keys(),
liveness: None,
cells: vec![],
},
DeltaRecord::StaticUpsert {
partition_key: sample_pk(),
cells: vec![],
},
DeltaRecord::RowDelete {
keys: sample_row_keys(),
deleted_at: 0,
},
DeltaRecord::RangeDelete {
partition_key: sample_pk(),
start: RangeBound::open(),
end: RangeBound::open(),
deleted_at: 0,
},
DeltaRecord::PartitionDelete {
partition_key: sample_pk(),
deleted_at: 0,
},
];
let names: Vec<&str> = ops.iter().map(|r| r.op_name()).collect();
let mut unique = names.clone();
unique.sort_unstable();
unique.dedup();
assert_eq!(unique.len(), names.len(), "duplicate op_name: {:?}", names);
}
#[test]
fn partition_tombstone_type_exists() {
use crate::types::TombstoneType;
let t = TombstoneType::PartitionTombstone;
let s = format!("{:?}", t);
assert!(s.contains("Partition"), "unexpected debug: {}", s);
}
#[test]
fn cell_tombstone_value_is_none_untouched_columns_absent() {
let tombstone_cell = CellDelta::tombstone(1_700_000_000_000_000);
let rec = DeltaRecord::Upsert {
keys: RowKeys::new(vec![Value::Integer(1)], vec![Value::Text("a".into())]),
liveness: None,
cells: vec![(ColumnId::new("val"), tombstone_cell.clone())],
};
if let DeltaRecord::Upsert { cells, .. } = &rec {
let val_entry = cells.iter().find(|(id, _)| id.0 == "val");
assert!(val_entry.is_some(), "val should be present");
let (_, cell) = val_entry.unwrap();
assert!(
cell.value.is_none(),
"cell tombstone must have value == None"
);
assert_eq!(cell.writetime, 1_700_000_000_000_000);
let other_entry = cells.iter().find(|(id, _)| id.0 == "other");
assert!(other_entry.is_none(), "untouched column must be absent");
} else {
panic!("expected Upsert");
}
}
#[test]
fn liveness_none_for_update_some_for_insert() {
let ts = 1_700_000_000_000_000_i64;
let update_rec = DeltaRecord::Upsert {
keys: RowKeys::new(vec![Value::Integer(1)], vec![Value::Text("a".into())]),
liveness: None,
cells: vec![(
ColumnId::new("val"),
CellDelta::value(Value::Text("x".into()), ts),
)],
};
if let DeltaRecord::Upsert { liveness, .. } = &update_rec {
assert!(liveness.is_none(), "UPDATE must have liveness == None");
}
let insert_rec = DeltaRecord::Upsert {
keys: RowKeys::new(vec![Value::Integer(2)], vec![Value::Text("b".into())]),
liveness: Some(CellMeta::new(ts)),
cells: vec![(
ColumnId::new("val"),
CellDelta::value(Value::Text("y".into()), ts),
)],
};
if let DeltaRecord::Upsert { liveness, .. } = &insert_rec {
let lv = liveness.as_ref().expect("INSERT must have liveness");
assert_eq!(lv.writetime, ts);
assert!(lv.expires_at.is_none());
}
}
#[test]
fn insert_with_ttl_liveness_carries_expires_at() {
let ts: i64 = 1_700_000_000_000_000;
let exp: i64 = 1_700_000_086_400_000_000;
let rec = DeltaRecord::Upsert {
keys: RowKeys::new(vec![Value::Integer(1)], vec![Value::Text("a".into())]),
liveness: Some(CellMeta::with_ttl(ts, exp)),
cells: vec![(
ColumnId::new("val"),
CellDelta::value_with_ttl(Value::Text("x".into()), ts, exp),
)],
};
if let DeltaRecord::Upsert {
liveness, cells, ..
} = &rec
{
let lv = liveness.as_ref().unwrap();
assert_eq!(lv.writetime, ts);
assert_eq!(lv.expires_at, Some(exp));
let (_, cell) = &cells[0];
assert_eq!(cell.expires_at, Some(exp));
}
}
#[test]
fn static_column_update_emits_static_upsert() {
let ts: i64 = 1_700_000_000_000_000;
let rec = DeltaRecord::StaticUpsert {
partition_key: RowKeys::partition_only(vec![Value::Integer(42)]),
cells: vec![(
ColumnId::new("static_col"),
CellDelta::value(Value::Text("S".into()), ts),
)],
};
assert_eq!(rec.op_name(), "static_upsert");
assert_eq!(rec.partition_key(), &[Value::Integer(42)]);
if let DeltaRecord::StaticUpsert {
partition_key,
cells,
} = &rec
{
assert!(
partition_key.clustering.is_empty(),
"StaticUpsert must have empty clustering"
);
assert_eq!(cells.len(), 1);
let (col_id, cell) = &cells[0];
assert_eq!(col_id.0, "static_col");
assert!(cell.value.is_some());
} else {
panic!("expected StaticUpsert");
}
}
#[test]
fn cell_tombstone_writetime_is_deletion_time_not_row_ts() {
let row_ts: i64 = 1_000_000;
let del_ts: i64 = 2_000_000;
let cell = CellDelta::tombstone(del_ts);
assert_eq!(cell.writetime, del_ts);
assert_ne!(
cell.writetime, row_ts,
"tombstone writetime must be the deletion timestamp, not the row timestamp"
);
}
#[test]
fn ds4_replaced_false_for_scalar_column() {
let cell = CellDelta::value(Value::Integer(42), 1_700_000_000_000_000);
assert!(
!cell.replaced,
"scalar column must always have replaced=false; got replaced=true"
);
}
#[test]
fn ds4_collection_append_replaced_false() {
let cell = CellDelta {
value: Some(Value::List(vec![Value::Text("new_element".into())])),
writetime: 1_700_000_000_000_000,
expires_at: None,
replaced: false,
};
assert!(
!cell.replaced,
"collection append must have replaced=false (no collection tombstone)"
);
}
#[test]
fn ds4_collection_overwrite_replaced_true() {
let cell = CellDelta::collection_replace(
Value::List(vec![Value::Text("a".into()), Value::Text("b".into())]),
1_700_000_000_000_000,
);
assert!(
cell.replaced,
"collection overwrite must have replaced=true (collection tombstone present)"
);
}
#[test]
fn ds4_collection_writetime_equals_max_element_writetime() {
let ts_early: i64 = 1_700_000_000_000_000;
let ts_late: i64 = 1_700_000_100_000_000;
let cell = CellDelta {
value: Some(Value::List(vec![
Value::Text("a".into()),
Value::Text("b".into()),
])),
writetime: ts_late, expires_at: None,
replaced: false,
};
assert_eq!(
cell.writetime, ts_late,
"writetime must equal max element writetime; expected {ts_late}, got {}",
cell.writetime
);
assert!(
cell.writetime > ts_early,
"max element writetime {ts_late} must be greater than an earlier element writetime {ts_early}"
);
}
#[test]
fn ds4_scan_summary_handle_accumulates_element_tombstones() {
let handle = ScanSummaryHandle::new();
assert_eq!(
handle.read().element_tombstones_detected,
0,
"initial element_tombstones_detected must be 0"
);
handle.add_element_tombstones(3);
handle.add_element_tombstones(5);
let summary = handle.read();
assert_eq!(
summary.element_tombstones_detected, 8,
"element_tombstones_detected must accumulate: expected 8, got {}",
summary.element_tombstones_detected
);
}
#[test]
fn ds4_scan_summary_handle_clone_shares_counter() {
let handle = ScanSummaryHandle::new();
let clone = handle.clone();
clone.add_element_tombstones(7);
assert_eq!(
handle.read().element_tombstones_detected,
7,
"original handle must reflect counter updated via clone"
);
}
#[test]
fn ds4_cell_tombstone_replaced_false() {
let cell = CellDelta::tombstone(1_700_000_000_000_000);
assert!(
!cell.replaced,
"cell tombstones must have replaced=false regardless of column type"
);
}
#[tokio::test]
async fn scan_delta_yields_upserts_from_simple_table() {
let root = match std::env::var("CQLITE_DATASETS_ROOT") {
Ok(r) => std::path::PathBuf::from(r),
Err(_) => {
eprintln!("CQLITE_DATASETS_ROOT not set — skipping scan_delta integration test");
return;
}
};
let base = root.join("sstables/test_basic");
if !base.exists() {
eprintln!("test_basic not found — skipping");
return;
}
let table_dir = std::fs::read_dir(&base).ok().and_then(|mut it| {
it.find_map(|e| {
e.ok()
.filter(|e| {
e.file_name()
.to_str()
.map(|n| n.starts_with("simple_table"))
.unwrap_or(false)
})
.map(|e| e.path())
})
});
let Some(table_dir) = table_dir else {
eprintln!("simple_table dir not found — skipping");
return;
};
let has_data_db = std::fs::read_dir(&table_dir)
.ok()
.map(|it| {
it.filter_map(|e| e.ok()).any(|e| {
e.file_name()
.to_str()
.map(|n| n.ends_with("-Data.db"))
.unwrap_or(false)
})
})
.unwrap_or(false);
if !has_data_db {
eprintln!("No Data.db in simple_table — skipping (run fetch-datasets.sh)");
return;
}
let schema = crate::schema::TableSchema {
keyspace: "test_basic".to_string(),
table: "simple_table".to_string(),
partition_keys: vec![crate::schema::KeyColumn {
name: "id".to_string(),
data_type: "uuid".to_string(),
position: 0,
}],
clustering_keys: vec![],
columns: vec![
crate::schema::Column {
name: "name".to_string(),
data_type: "text".to_string(),
nullable: true,
default: None,
is_static: false,
},
crate::schema::Column {
name: "value".to_string(),
data_type: "int".to_string(),
nullable: true,
default: None,
is_static: false,
},
],
comments: std::collections::HashMap::new(),
dropped_columns: std::collections::HashMap::new(),
};
let (mut rx, _scan_summary) = scan_delta(table_dir, schema, 64);
let mut upsert_count = 0_usize;
let mut total = 0_usize;
while let Some(result) = rx.recv().await {
total += 1;
match result {
Ok(DeltaRecord::Upsert { .. }) => upsert_count += 1,
Ok(DeltaRecord::StaticUpsert { .. }) => {}
Ok(other) => {
panic!(
"simple_table should have no tombstones; got {:?}",
other.op_name()
);
}
Err(e) => panic!("scan_delta error: {e}"),
}
}
eprintln!(
"scan_delta simple_table: {} total records, {} upserts",
total, upsert_count
);
assert!(
upsert_count > 0,
"expected at least one Upsert from simple_table"
);
}
#[tokio::test]
async fn scan_delta_cells_have_nonzero_writetime() {
let root = match std::env::var("CQLITE_DATASETS_ROOT") {
Ok(r) => std::path::PathBuf::from(r),
Err(_) => return,
};
let base = root.join("sstables/test_basic");
if !base.exists() {
return;
}
let table_dir = std::fs::read_dir(&base).ok().and_then(|mut it| {
it.find_map(|e| {
e.ok()
.filter(|e| {
e.file_name()
.to_str()
.map(|n| n.starts_with("simple_table"))
.unwrap_or(false)
})
.map(|e| e.path())
})
});
let Some(table_dir) = table_dir else {
return;
};
let has_data_db = std::fs::read_dir(&table_dir)
.ok()
.map(|it| {
it.filter_map(|e| e.ok()).any(|e| {
e.file_name()
.to_str()
.map(|n| n.ends_with("-Data.db"))
.unwrap_or(false)
})
})
.unwrap_or(false);
if !has_data_db {
return;
}
let schema = crate::schema::TableSchema {
keyspace: "test_basic".to_string(),
table: "simple_table".to_string(),
partition_keys: vec![crate::schema::KeyColumn {
name: "id".to_string(),
data_type: "uuid".to_string(),
position: 0,
}],
clustering_keys: vec![],
columns: vec![
crate::schema::Column {
name: "name".to_string(),
data_type: "text".to_string(),
nullable: true,
default: None,
is_static: false,
},
crate::schema::Column {
name: "value".to_string(),
data_type: "int".to_string(),
nullable: true,
default: None,
is_static: false,
},
],
comments: std::collections::HashMap::new(),
dropped_columns: std::collections::HashMap::new(),
};
let (mut rx, _scan_summary) = scan_delta(table_dir, schema, 64);
let mut checked = 0_usize;
while let Some(result) = rx.recv().await {
if let Ok(DeltaRecord::Upsert { cells, .. }) = result {
for (_, cell) in &cells {
assert!(
cell.writetime > 1_262_304_000_000_000,
"writetime {} is suspiciously small (cell {:?})",
cell.writetime,
cell.value
);
checked += 1;
}
}
}
if checked > 0 {
eprintln!("scan_delta writetime check: verified {} cells", checked);
}
}
#[tokio::test]
async fn scan_delta_emits_static_upsert_from_static_columns_table() {
let root = match std::env::var("CQLITE_DATASETS_ROOT") {
Ok(r) => std::path::PathBuf::from(r),
Err(_) => {
eprintln!("CQLITE_DATASETS_ROOT not set — skipping StaticUpsert e2e test");
return;
}
};
let base = root.join("sstables/test_basic");
if !base.exists() {
eprintln!("test_basic not found — skipping StaticUpsert e2e test");
return;
}
let table_dir = std::fs::read_dir(&base).ok().and_then(|mut it| {
it.find_map(|e| {
e.ok()
.filter(|e| {
e.file_name()
.to_str()
.map(|n| n.starts_with("static_columns_table"))
.unwrap_or(false)
})
.map(|e| e.path())
})
});
let Some(table_dir) = table_dir else {
eprintln!("static_columns_table dir not found — skipping StaticUpsert e2e test");
return;
};
let has_data_db = std::fs::read_dir(&table_dir)
.ok()
.map(|it| {
it.filter_map(|e| e.ok()).any(|e| {
e.file_name()
.to_str()
.map(|n| n.ends_with("-Data.db"))
.unwrap_or(false)
})
})
.unwrap_or(false);
if !has_data_db {
eprintln!("No Data.db in static_columns_table — skipping (run fetch-datasets.sh)");
return;
}
let schema = crate::schema::TableSchema {
keyspace: "test_basic".to_string(),
table: "static_columns_table".to_string(),
partition_keys: vec![crate::schema::KeyColumn {
name: "partition_key".to_string(),
data_type: "uuid".to_string(),
position: 0,
}],
clustering_keys: vec![crate::schema::ClusteringColumn {
name: "clustering_key".to_string(),
data_type: "timestamp".to_string(),
position: 0,
order: crate::schema::ClusteringOrder::Asc,
}],
columns: vec![
crate::schema::Column {
name: "static_data".to_string(),
data_type: "text".to_string(),
nullable: true,
default: None,
is_static: true,
},
crate::schema::Column {
name: "row_data".to_string(),
data_type: "text".to_string(),
nullable: true,
default: None,
is_static: false,
},
crate::schema::Column {
name: "row_value".to_string(),
data_type: "int".to_string(),
nullable: true,
default: None,
is_static: false,
},
],
comments: std::collections::HashMap::new(),
dropped_columns: std::collections::HashMap::new(),
};
let (mut rx, _scan_summary) = scan_delta(table_dir, schema, 64);
let mut static_upsert_count = 0_usize;
let mut upsert_count = 0_usize;
while let Some(result) = rx.recv().await {
match result {
Ok(DeltaRecord::StaticUpsert { ref cells, .. }) => {
static_upsert_count += 1;
assert!(
!cells.is_empty(),
"StaticUpsert must have at least one cell delta"
);
}
Ok(DeltaRecord::Upsert { .. }) => {
upsert_count += 1;
}
Ok(other) => {
eprintln!(
"scan_delta static_columns_table: unexpected record: {}",
other.op_name()
);
}
Err(e) => panic!("scan_delta error on static_columns_table: {e}"),
}
}
eprintln!(
"scan_delta static_columns_table: {} StaticUpserts, {} Upserts",
static_upsert_count, upsert_count
);
assert!(
static_upsert_count > 0,
"expected at least one StaticUpsert from static_columns_table; \
got {} StaticUpserts and {} Upserts",
static_upsert_count,
upsert_count
);
}
#[tokio::test]
async fn scan_delta_emits_cell_tombstone_from_cell_tombstones_table() {
let root = match std::env::var("CQLITE_DATASETS_ROOT") {
Ok(r) => std::path::PathBuf::from(r),
Err(_) => {
eprintln!("CQLITE_DATASETS_ROOT not set — skipping cell-tombstone e2e test");
return;
}
};
let deltas_dir = root.join("sstables/test_deltas");
if !deltas_dir.exists() {
eprintln!(
"test_deltas not found at {:?} — skipping cell-tombstone e2e test \
(run `bash test-data/scripts/generate-deltas.sh` to regenerate)",
deltas_dir
);
return;
}
let table_dir = std::fs::read_dir(&deltas_dir).ok().and_then(|mut it| {
it.find_map(|e| {
e.ok()
.filter(|e| {
e.file_name()
.to_str()
.map(|n| n.starts_with("cell_tombstones"))
.unwrap_or(false)
})
.map(|e| e.path())
})
});
let Some(table_dir) = table_dir else {
eprintln!("cell_tombstones dir not found — skipping cell-tombstone e2e test");
return;
};
let has_data_db = std::fs::read_dir(&table_dir)
.ok()
.map(|it| {
it.filter_map(|e| e.ok()).any(|e| {
let name = e.file_name();
let n = name.to_string_lossy();
n.ends_with("-Data.db") && !n.ends_with(".db.jsonl")
})
})
.unwrap_or(false);
if !has_data_db {
eprintln!(
"No binary Data.db in cell_tombstones — skipping cell-tombstone e2e test \
(run `bash test-data/scripts/generate-deltas.sh` to regenerate binaries; \
test_deltas binaries are not in the published dataset asset)"
);
return;
}
let schema = crate::schema::TableSchema {
keyspace: "test_deltas".to_string(),
table: "cell_tombstones".to_string(),
partition_keys: vec![crate::schema::KeyColumn {
name: "pk".to_string(),
data_type: "int".to_string(),
position: 0,
}],
clustering_keys: vec![crate::schema::ClusteringColumn {
name: "ck".to_string(),
data_type: "int".to_string(),
position: 0,
order: crate::schema::ClusteringOrder::Asc,
}],
columns: vec![
crate::schema::Column {
name: "col_a".to_string(),
data_type: "text".to_string(),
nullable: true,
default: None,
is_static: false,
},
crate::schema::Column {
name: "col_b".to_string(),
data_type: "text".to_string(),
nullable: true,
default: None,
is_static: false,
},
],
comments: std::collections::HashMap::new(),
dropped_columns: std::collections::HashMap::new(),
};
let (mut rx, _scan_summary) = scan_delta(table_dir, schema, 64);
let mut cell_tombstone_count = 0_usize;
let mut total_cells = 0_usize;
while let Some(result) = rx.recv().await {
match result {
Ok(DeltaRecord::Upsert { cells, .. }) => {
for (col_id, cell) in &cells {
total_cells += 1;
if cell.value.is_none() {
cell_tombstone_count += 1;
eprintln!(
"cell-tombstone e2e: column {:?} has CellDelta {{ value: None, writetime: {} }}",
col_id.0, cell.writetime
);
}
}
}
Ok(DeltaRecord::StaticUpsert { .. }) => {}
Ok(other) => {
eprintln!(
"scan_delta cell_tombstones: got {} (out of #698 scope)",
other.op_name()
);
}
Err(e) => panic!("scan_delta error on cell_tombstones: {e}"),
}
}
eprintln!(
"scan_delta cell_tombstones e2e: {} cell tombstones out of {} total cells",
cell_tombstone_count, total_cells
);
assert!(
cell_tombstone_count > 0,
"expected at least one CellDelta {{ value: None }} from cell_tombstones; \
got {} total cells with 0 tombstones",
total_cells
);
}
fn find_test_deltas_table_dir(
root: &std::path::Path,
table_prefix: &str,
) -> Option<std::path::PathBuf> {
let deltas_dir = root.join("sstables/test_deltas");
if !deltas_dir.exists() {
eprintln!(
"test_deltas not found at {:?} — skipping e2e test \
(run `bash test-data/scripts/generate-deltas.sh` to regenerate)",
deltas_dir
);
return None;
}
let table_dir = std::fs::read_dir(&deltas_dir).ok()?.find_map(|e| {
let entry = e.ok()?;
let name = entry.file_name();
let n = name.to_string_lossy();
if !n.starts_with(table_prefix) {
return None;
}
let path = entry.path();
let has_data_db = std::fs::read_dir(&path)
.ok()
.map(|it| {
it.filter_map(|e| e.ok()).any(|e| {
let fname = e.file_name();
let fn_str = fname.to_string_lossy();
fn_str.ends_with("-Data.db") && !fn_str.ends_with(".db.jsonl")
})
})
.unwrap_or(false);
if has_data_db {
Some(path)
} else {
None
}
});
match table_dir {
Some(dir) => Some(dir),
None => {
eprintln!(
"No binary Data.db found in any {}-* directory under test_deltas — \
skipping e2e test (run `bash test-data/scripts/generate-deltas.sh` \
to regenerate binaries)",
table_prefix
);
None
}
}
}
#[test]
fn partition_header_full_returns_deleted_at_for_dead_partition_nb_format() {
use crate::storage::sstable::reader::parsing::PublicV5CompressedLegacyParser;
let parser = PublicV5CompressedLegacyParser::new(
"test_deltas".to_string(),
"partition_tombstones".to_string(),
0, 0, None,
);
let dead_ts: i64 = 0x0000_617C_AF98_CC00_u64 as i64;
let mut buf = vec![
0x00_u8, 0x04_u8, 0x00, 0x00, 0x00, 0x01, 0x61, 0x23, 0x45, 0x67, ];
buf.extend_from_slice(&dead_ts.to_be_bytes());
let (_row_key, _next_offset, partition_deletion) = parser
.parse_partition_header_full(&buf, 0)
.expect("parse should succeed");
assert!(
partition_deletion.is_some(),
"DEAD partition must yield Some(deleted_at); got None"
);
assert_eq!(
partition_deletion.unwrap(),
dead_ts,
"deleted_at must equal the markedForDeleteAt bytes in the header"
);
}
#[test]
fn parse_range_tombstone_marker_full_surfaces_unknown_bound_kind() {
use crate::schema::{KeyColumn, TableSchema};
use crate::storage::sstable::reader::parsing::PublicV5CompressedLegacyParser;
let parser = PublicV5CompressedLegacyParser::new(
"test_deltas".to_string(),
"range_tombstones".to_string(),
0, 0, None,
);
let schema = TableSchema {
keyspace: "test_deltas".to_string(),
table: "range_tombstones".to_string(),
partition_keys: vec![KeyColumn {
name: "pk".to_string(),
data_type: "int".to_string(),
position: 0,
}],
clustering_keys: vec![],
columns: vec![],
comments: std::collections::HashMap::new(),
dropped_columns: std::collections::HashMap::new(),
};
let crafted_marker: Vec<u8> = vec![
0x02, 99u8, 0x00, 0x00, 0x03, 0x00, 0x00, 0x00, ];
let (bound_values, bound_kind, deleted_at_primary, deleted_at_secondary, next_offset) =
parser
.parse_range_tombstone_marker_full(&crafted_marker, 0, &schema)
.expect("parse_range_tombstone_marker_full should succeed on crafted buffer");
assert_eq!(
bound_kind, 99,
"parse_range_tombstone_marker_full must return bound_kind=99 unchanged; \
if this fails the production hard-error branch in parse_block_emit_delta \
is no longer reachable for unknown kinds"
);
assert!(bound_values.is_empty(), "cluster_count=0 → no bound values");
assert_eq!(deleted_at_primary, 0);
assert!(deleted_at_secondary.is_none());
assert_eq!(next_offset, crafted_marker.len());
}
#[test]
fn unknown_range_tombstone_kind_error_message_format() {
let unknown_kind: u8 = 99;
let offset: usize = 8; let pk_raw: &[u8] = &[0x00, 0x00, 0x00, 0x01]; let err = crate::Error::corruption(format!(
"delta-scan: unknown range tombstone bound kind {} at offset {} \
in test_deltas.range_tombstones (partition key {:?}) — cannot represent faithfully \
(no-heuristics mandate, issue #28)",
unknown_kind, offset, pk_raw
));
let msg = format!("{}", err);
assert!(
msg.contains("unknown range tombstone bound kind"),
"error message must name the unknown-kind problem: {}",
msg
);
assert!(
msg.contains("99"),
"error message must include the bad kind value: {}",
msg
);
assert!(
msg.contains("issue #28"),
"error message must cite the no-heuristics mandate: {}",
msg
);
}
#[test]
fn row_delete_and_upsert_can_coexist_for_same_key() {
let ts: i64 = 1_700_000_000_000_000;
let del_ts: i64 = 1_700_000_000_100_000;
let pk = RowKeys::new(vec![Value::Integer(1)], vec![Value::Text("ck".into())]);
let row_delete = DeltaRecord::RowDelete {
keys: pk.clone(),
deleted_at: del_ts,
};
let upsert = DeltaRecord::Upsert {
keys: pk.clone(),
liveness: None,
cells: vec![(
ColumnId::new("name"),
CellDelta::value(Value::Text("Alice".into()), ts),
)],
};
assert_ne!(row_delete.op_name(), upsert.op_name());
assert_eq!(row_delete.partition_key(), upsert.partition_key());
assert_eq!(row_delete.op_name(), "row_delete");
assert_eq!(upsert.op_name(), "upsert");
if let DeltaRecord::RowDelete { deleted_at, .. } = &row_delete {
assert_eq!(*deleted_at, del_ts);
}
if let DeltaRecord::Upsert { cells, .. } = &upsert {
assert_eq!(cells.len(), 1);
let (_, cell) = &cells[0];
assert_eq!(cell.writetime, ts);
assert!(cell.value.is_some());
}
}
#[test]
fn range_delete_prefix_bound_multi_column_clustering() {
let rec = DeltaRecord::RangeDelete {
partition_key: RowKeys::partition_only(vec![Value::Integer(1)]),
start: RangeBound::new(vec![Value::Integer(2)], true),
end: RangeBound::new(vec![Value::Integer(2)], true),
deleted_at: 1_700_000_000_000_000,
};
if let DeltaRecord::RangeDelete { start, end, .. } = &rec {
assert_eq!(start.values.len(), 1);
assert_eq!(end.values.len(), 1);
assert!(start.is_prefix()); assert!(end.is_prefix());
assert!(start.inclusive);
assert!(end.inclusive);
} else {
panic!("expected RangeDelete");
}
}
#[test]
fn range_delete_mixed_inclusivity() {
let rec = DeltaRecord::RangeDelete {
partition_key: RowKeys::partition_only(vec![Value::Integer(2)]),
start: RangeBound::inclusive(vec![Value::Integer(2)]),
end: RangeBound::exclusive(vec![Value::Integer(4)]),
deleted_at: 1_700_000_000_000_001,
};
if let DeltaRecord::RangeDelete {
start,
end,
deleted_at,
..
} = &rec
{
assert!(start.inclusive, "start should be inclusive (>=)");
assert!(!end.inclusive, "end should be exclusive (<)");
assert_eq!(*deleted_at, 1_700_000_000_000_001);
} else {
panic!("expected RangeDelete");
}
}
#[tokio::test]
async fn scan_delta_emits_row_delete_from_row_tombstones_table() {
let root = match std::env::var("CQLITE_DATASETS_ROOT") {
Ok(r) => std::path::PathBuf::from(r),
Err(_) => {
eprintln!("CQLITE_DATASETS_ROOT not set — skipping row-tombstone e2e test");
return;
}
};
let Some(table_dir) = find_test_deltas_table_dir(&root, "row_tombstones") else {
return;
};
let schema = crate::schema::TableSchema {
keyspace: "test_deltas".to_string(),
table: "row_tombstones".to_string(),
partition_keys: vec![crate::schema::KeyColumn {
name: "pk".to_string(),
data_type: "int".to_string(),
position: 0,
}],
clustering_keys: vec![crate::schema::ClusteringColumn {
name: "ck".to_string(),
data_type: "int".to_string(),
position: 0,
order: crate::schema::ClusteringOrder::Asc,
}],
columns: vec![crate::schema::Column {
name: "val".to_string(),
data_type: "text".to_string(),
nullable: true,
default: None,
is_static: false,
}],
comments: std::collections::HashMap::new(),
dropped_columns: std::collections::HashMap::new(),
};
let (mut rx, _scan_summary) = scan_delta(table_dir, schema, 64);
let mut row_delete_count = 0_usize;
let mut upsert_count = 0_usize;
while let Some(result) = rx.recv().await {
match result {
Ok(DeltaRecord::RowDelete { keys, deleted_at }) => {
row_delete_count += 1;
assert!(
!keys.clustering.is_empty(),
"RowDelete must have non-empty clustering key; got pk={:?}",
keys.partition
);
assert!(
deleted_at > 1_262_304_000_000_000,
"RowDelete deleted_at={} is suspiciously small",
deleted_at
);
eprintln!(
"row-tombstone e2e: RowDelete pk={:?} ck={:?} deleted_at={}",
keys.partition, keys.clustering, deleted_at
);
}
Ok(DeltaRecord::Upsert { .. }) => upsert_count += 1,
Ok(DeltaRecord::StaticUpsert { .. }) => {}
Ok(other) => {
panic!(
"row_tombstones should only have Upsert and RowDelete; got {}",
other.op_name()
);
}
Err(e) => panic!("scan_delta error on row_tombstones: {e}"),
}
}
eprintln!(
"scan_delta row_tombstones e2e: {} RowDelete + {} Upsert",
row_delete_count, upsert_count
);
assert!(
row_delete_count > 0,
"expected at least one RowDelete from row_tombstones; got 0 (with {} upserts)",
upsert_count
);
}
#[tokio::test]
async fn scan_delta_emits_range_delete_from_range_tombstones_table() {
let root = match std::env::var("CQLITE_DATASETS_ROOT") {
Ok(r) => std::path::PathBuf::from(r),
Err(_) => {
eprintln!("CQLITE_DATASETS_ROOT not set — skipping range-tombstone e2e test");
return;
}
};
let Some(table_dir) = find_test_deltas_table_dir(&root, "range_tombstones") else {
return;
};
let schema = crate::schema::TableSchema {
keyspace: "test_deltas".to_string(),
table: "range_tombstones".to_string(),
partition_keys: vec![crate::schema::KeyColumn {
name: "pk".to_string(),
data_type: "int".to_string(),
position: 0,
}],
clustering_keys: vec![
crate::schema::ClusteringColumn {
name: "ck1".to_string(),
data_type: "int".to_string(),
position: 0,
order: crate::schema::ClusteringOrder::Asc,
},
crate::schema::ClusteringColumn {
name: "ck2".to_string(),
data_type: "text".to_string(),
position: 1,
order: crate::schema::ClusteringOrder::Asc,
},
],
columns: vec![crate::schema::Column {
name: "val".to_string(),
data_type: "text".to_string(),
nullable: true,
default: None,
is_static: false,
}],
comments: std::collections::HashMap::new(),
dropped_columns: std::collections::HashMap::new(),
};
let (mut rx, _scan_summary) = scan_delta(table_dir, schema, 64);
let mut range_delete_count = 0_usize;
let mut upsert_count = 0_usize;
while let Some(result) = rx.recv().await {
match result {
Ok(DeltaRecord::RangeDelete {
partition_key,
start,
end,
deleted_at,
}) => {
range_delete_count += 1;
assert!(
!partition_key.partition.is_empty(),
"RangeDelete must have a partition key"
);
assert!(
partition_key.clustering.is_empty(),
"RangeDelete partition_key must have empty clustering (bounds carry it)"
);
assert!(
deleted_at > 1_262_304_000_000_000,
"RangeDelete deleted_at={} is suspiciously small",
deleted_at
);
eprintln!(
"range-tombstone e2e: RangeDelete pk={:?} start=({:?}, incl={}) \
end=({:?}, incl={}) deleted_at={}",
partition_key.partition,
start.values,
start.inclusive,
end.values,
end.inclusive,
deleted_at
);
}
Ok(DeltaRecord::Upsert { .. }) => upsert_count += 1,
Ok(DeltaRecord::StaticUpsert { .. }) => {}
Ok(other) => {
panic!(
"range_tombstones should only have Upsert and RangeDelete; got {}",
other.op_name()
);
}
Err(e) => panic!("scan_delta error on range_tombstones: {e}"),
}
}
eprintln!(
"scan_delta range_tombstones e2e: {} RangeDelete + {} Upsert",
range_delete_count, upsert_count
);
assert!(
range_delete_count > 0,
"expected at least one RangeDelete from range_tombstones; got 0 (with {} upserts)",
upsert_count
);
}
#[tokio::test]
async fn scan_delta_emits_partition_delete_from_partition_tombstones_table() {
let root = match std::env::var("CQLITE_DATASETS_ROOT") {
Ok(r) => std::path::PathBuf::from(r),
Err(_) => {
eprintln!("CQLITE_DATASETS_ROOT not set — skipping partition-tombstone e2e test");
return;
}
};
let Some(table_dir) = find_test_deltas_table_dir(&root, "partition_tombstones") else {
return;
};
let schema = crate::schema::TableSchema {
keyspace: "test_deltas".to_string(),
table: "partition_tombstones".to_string(),
partition_keys: vec![crate::schema::KeyColumn {
name: "pk".to_string(),
data_type: "int".to_string(),
position: 0,
}],
clustering_keys: vec![crate::schema::ClusteringColumn {
name: "ck".to_string(),
data_type: "int".to_string(),
position: 0,
order: crate::schema::ClusteringOrder::Asc,
}],
columns: vec![crate::schema::Column {
name: "val".to_string(),
data_type: "text".to_string(),
nullable: true,
default: None,
is_static: false,
}],
comments: std::collections::HashMap::new(),
dropped_columns: std::collections::HashMap::new(),
};
let (mut rx, _scan_summary) = scan_delta(table_dir, schema, 64);
let mut partition_delete_count = 0_usize;
let mut upsert_count = 0_usize;
while let Some(result) = rx.recv().await {
match result {
Ok(DeltaRecord::PartitionDelete {
partition_key,
deleted_at,
}) => {
partition_delete_count += 1;
assert!(
!partition_key.partition.is_empty(),
"PartitionDelete must have a partition key"
);
assert!(
partition_key.clustering.is_empty(),
"PartitionDelete must have empty clustering key"
);
assert!(
deleted_at > 1_262_304_000_000_000,
"PartitionDelete deleted_at={} is suspiciously small",
deleted_at
);
eprintln!(
"partition-tombstone e2e: PartitionDelete pk={:?} deleted_at={}",
partition_key.partition, deleted_at
);
}
Ok(DeltaRecord::Upsert { .. }) => upsert_count += 1,
Ok(DeltaRecord::StaticUpsert { .. }) => {}
Ok(other) => {
panic!(
"partition_tombstones should only have Upsert and PartitionDelete; got {}",
other.op_name()
);
}
Err(e) => panic!("scan_delta error on partition_tombstones: {e}"),
}
}
eprintln!(
"scan_delta partition_tombstones e2e: {} PartitionDelete + {} Upsert",
partition_delete_count, upsert_count
);
assert!(
partition_delete_count > 0,
"expected at least one PartitionDelete from partition_tombstones; got 0 (with {} upserts)",
upsert_count
);
assert!(
partition_delete_count >= 2,
"expected at least 2 PartitionDeletes (pk=2 and pk=4); got {} (with {} upserts)",
partition_delete_count,
upsert_count
);
}
#[tokio::test]
async fn scan_delta_emits_both_range_deletes_from_adjacent_ranges_table() {
let root = match std::env::var("CQLITE_DATASETS_ROOT") {
Ok(r) => std::path::PathBuf::from(r),
Err(_) => {
eprintln!("CQLITE_DATASETS_ROOT not set — skipping adjacent-ranges e2e test");
return;
}
};
let Some(table_dir) = find_test_deltas_table_dir(&root, "adjacent_ranges") else {
return;
};
let schema = crate::schema::TableSchema {
keyspace: "test_deltas".to_string(),
table: "adjacent_ranges".to_string(),
partition_keys: vec![crate::schema::KeyColumn {
name: "pk".to_string(),
data_type: "int".to_string(),
position: 0,
}],
clustering_keys: vec![crate::schema::ClusteringColumn {
name: "ck".to_string(),
data_type: "int".to_string(),
position: 0,
order: crate::schema::ClusteringOrder::Asc,
}],
columns: vec![crate::schema::Column {
name: "val".to_string(),
data_type: "text".to_string(),
nullable: true,
default: None,
is_static: false,
}],
comments: std::collections::HashMap::new(),
dropped_columns: std::collections::HashMap::new(),
};
let (mut rx, _scan_summary) = scan_delta(table_dir, schema, 128);
let mut range_deletes_by_pk: std::collections::HashMap<
i32,
Vec<(RangeBound, RangeBound, i64)>,
> = std::collections::HashMap::new();
let mut upsert_count = 0_usize;
while let Some(result) = rx.recv().await {
match result {
Ok(DeltaRecord::RangeDelete {
partition_key,
start,
end,
deleted_at,
}) => {
assert!(
!partition_key.partition.is_empty(),
"RangeDelete must have a non-empty partition key"
);
assert!(
deleted_at > 0,
"RangeDelete deleted_at must be positive; got {}",
deleted_at
);
let pk_int = match &partition_key.partition[0] {
Value::Integer(n) => *n,
other => panic!("expected Integer pk; got {:?}", other),
};
eprintln!(
"adjacent-ranges e2e: RangeDelete pk={} start=({:?}, incl={}) \
end=({:?}, incl={}) deleted_at={}",
pk_int,
start.values,
start.inclusive,
end.values,
end.inclusive,
deleted_at
);
range_deletes_by_pk
.entry(pk_int)
.or_default()
.push((start, end, deleted_at));
}
Ok(DeltaRecord::Upsert { .. }) => upsert_count += 1,
Ok(DeltaRecord::StaticUpsert { .. }) => {}
Ok(DeltaRecord::RowDelete { .. }) => {} Ok(DeltaRecord::PartitionDelete { .. }) => {}
Err(e) => panic!("scan_delta error on adjacent_ranges: {e}"),
}
}
let total_range_deletes: usize = range_deletes_by_pk.values().map(|v| v.len()).sum();
eprintln!(
"adjacent_ranges e2e: {} total RangeDeletes across {} partitions, {} Upserts",
total_range_deletes,
range_deletes_by_pk.len(),
upsert_count
);
assert!(
total_range_deletes >= 2,
"expected at least 2 RangeDeletes (one per adjacent range); got {} (with {} upserts)",
total_range_deletes,
upsert_count
);
for (pk, records) in &range_deletes_by_pk {
if records.len() >= 2 {
let timestamps: std::collections::HashSet<i64> =
records.iter().map(|(_, _, ts)| *ts).collect();
assert!(
timestamps.len() >= 2,
"pk={}: expected at least 2 distinct deleted_at values from adjacent ranges \
with different USING TIMESTAMP values; all {} records share the same timestamp. \
This indicates the boundary-marker secondary deletion time is not decoded.",
pk,
records.len()
);
eprintln!(
"pk={}: {} RangeDeletes with {} distinct timestamps — boundary marker correctly decoded",
pk, records.len(), timestamps.len()
);
}
}
}
#[tokio::test]
async fn ds4_scan_delta_collection_table_e2e() {
let root = match std::env::var("CQLITE_DATASETS_ROOT") {
Ok(r) => std::path::PathBuf::from(r),
Err(_) => {
eprintln!("CQLITE_DATASETS_ROOT not set — skipping DS4 collection e2e test");
return;
}
};
let base = root.join("sstables/test_collections");
if !base.exists() {
eprintln!("test_collections not found — skipping DS4 e2e");
return;
}
let table_dir = std::fs::read_dir(&base).ok().and_then(|mut it| {
it.find_map(|e| {
e.ok()
.filter(|e| {
e.file_name()
.to_str()
.map(|n| n.starts_with("collection_table"))
.unwrap_or(false)
})
.map(|e| e.path())
})
});
let Some(table_dir) = table_dir else {
eprintln!("collection_table dir not found — skipping DS4 e2e");
return;
};
let has_data_db = std::fs::read_dir(&table_dir)
.ok()
.map(|it| {
it.filter_map(|e| e.ok()).any(|e| {
e.file_name()
.to_str()
.map(|n| n.ends_with("-Data.db"))
.unwrap_or(false)
})
})
.unwrap_or(false);
if !has_data_db {
eprintln!("No Data.db in collection_table — skipping DS4 e2e (run fetch-datasets.sh)");
return;
}
let schema = crate::schema::TableSchema {
keyspace: "test_collections".to_string(),
table: "collection_table".to_string(),
partition_keys: vec![crate::schema::KeyColumn {
name: "id".to_string(),
data_type: "uuid".to_string(),
position: 0,
}],
clustering_keys: vec![],
columns: vec![
crate::schema::Column {
name: "tags".to_string(),
data_type: "set<text>".to_string(),
nullable: true,
default: None,
is_static: false,
},
crate::schema::Column {
name: "scores".to_string(),
data_type: "list<int>".to_string(),
nullable: true,
default: None,
is_static: false,
},
crate::schema::Column {
name: "properties".to_string(),
data_type: "map<text, text>".to_string(),
nullable: true,
default: None,
is_static: false,
},
crate::schema::Column {
name: "numbers_set".to_string(),
data_type: "set<int>".to_string(),
nullable: true,
default: None,
is_static: false,
},
crate::schema::Column {
name: "ordered_values".to_string(),
data_type: "list<timestamp>".to_string(),
nullable: true,
default: None,
is_static: false,
},
crate::schema::Column {
name: "metadata_map".to_string(),
data_type: "map<text, bigint>".to_string(),
nullable: true,
default: None,
is_static: false,
},
],
comments: std::collections::HashMap::new(),
dropped_columns: std::collections::HashMap::new(),
};
let (mut rx, summary_handle) = scan_delta(table_dir, schema, 64);
let mut upsert_count = 0_usize;
let mut total = 0_usize;
let mut collection_cells_seen = 0_usize;
while let Some(result) = rx.recv().await {
total += 1;
match result {
Ok(DeltaRecord::Upsert { ref cells, .. }) => {
upsert_count += 1;
for (col_id, cell) in cells {
let col_name = col_id.name();
if matches!(
col_name,
"tags"
| "scores"
| "properties"
| "numbers_set"
| "ordered_values"
| "metadata_map"
) && cell.value.is_some()
{
collection_cells_seen += 1;
assert!(
cell.writetime > 1_577_836_800_000_000,
"DS4: collection cell '{}' writetime {} is suspiciously small — \
expected max element writetime",
col_name,
cell.writetime
);
}
}
}
Ok(_) => {}
Err(e) => panic!("scan_delta DS4 collection e2e error: {e}"),
}
}
let summary = summary_handle.read();
eprintln!(
"DS4 collection_table e2e: {} total records, {} upserts, {} collection cells, \
{} element tombstones detected",
total, upsert_count, collection_cells_seen, summary.element_tombstones_detected
);
assert!(
upsert_count > 0,
"DS4 e2e: expected at least one Upsert from collection_table"
);
assert!(
collection_cells_seen > 0,
"DS4 e2e: expected at least one collection cell (tags/scores/properties/…) in Upsert records"
);
assert_eq!(
summary.element_tombstones_detected, 0,
"DS4 e2e: collection_table fixture uses appends only — expected 0 element tombstones, \
got {}",
summary.element_tombstones_detected
);
}
}