#[cfg(feature = "write-support")]
use super::model::{MergeEntry, RowData};
#[cfg(feature = "write-support")]
use crate::storage::write_engine::mutation::{DecoratedKey, RangeTombstone};
#[cfg(feature = "write-support")]
#[derive(Debug, Default)]
pub(super) struct PartitionCarriers {
pub range_tombstones: Vec<(DecoratedKey, RangeTombstone)>,
pub max_partition_deletion: Option<(i64, i32)>,
pub partition_delete_key: Option<DecoratedKey>,
}
#[cfg(feature = "write-support")]
pub(super) fn is_range_marker_carrier(entry: &MergeEntry) -> bool {
entry.range_deletion.is_some()
&& entry.complex_deletions.is_empty()
&& entry.row_deletion.is_none()
&& matches!(&entry.row_data, RowData::Live { cells } if cells.is_empty())
}
#[cfg(feature = "write-support")]
pub(super) fn scan_partition_carriers(rows: &[MergeEntry]) -> PartitionCarriers {
let mut carriers = PartitionCarriers::default();
for row in rows {
if row.is_partition_delete_carrier() {
if let Some((mfda, ldt)) = row.partition_deletion {
carriers
.partition_delete_key
.get_or_insert_with(|| row.key.clone());
match carriers.max_partition_deletion {
Some((cur_mfda, _)) if cur_mfda >= mfda => {}
_ => carriers.max_partition_deletion = Some((mfda, ldt)),
}
continue;
}
}
if is_range_marker_carrier(row) {
if let Some(rt) = row.range_deletion.clone() {
carriers.range_tombstones.push((row.key.clone(), rt));
}
}
}
carriers
}
#[cfg(all(test, feature = "write-support"))]
mod tests {
use super::*;
use crate::storage::write_engine::merge::CellData;
use crate::storage::write_engine::mutation::ClusteringBound;
use crate::types::Value;
fn key(token: i64) -> DecoratedKey {
DecoratedKey::new(token, vec![token as u8])
}
fn range_tombstone(deletion_time: i64, local_deletion_time: i32) -> RangeTombstone {
RangeTombstone {
start: ClusteringBound::Bottom,
end: ClusteringBound::Top,
deletion_time,
local_deletion_time,
}
}
fn live_row(token: i64) -> MergeEntry {
MergeEntry::new(
0,
key(token),
None,
100,
RowData::Live {
cells: vec![CellData::new("c".to_string(), Value::Integer(1), 100)],
},
)
}
fn range_carrier(token: i64, deletion_time: i64, ldt: i32) -> MergeEntry {
MergeEntry::new(
0,
key(token),
None,
deletion_time,
RowData::Live { cells: Vec::new() },
)
.with_range_deletion(range_tombstone(deletion_time, ldt))
}
fn partition_carrier(token: i64, mfda: i64, ldt: i32) -> MergeEntry {
MergeEntry::new(
0,
key(token),
None,
mfda,
RowData::Live { cells: Vec::new() },
)
.with_partition_deletion((mfda, ldt))
}
#[test]
fn empty_input_yields_no_carriers() {
let carriers = scan_partition_carriers(&[]);
assert!(carriers.range_tombstones.is_empty());
assert_eq!(carriers.max_partition_deletion, None);
assert!(carriers.partition_delete_key.is_none());
}
#[test]
fn plain_rows_yield_no_carriers() {
let rows = vec![live_row(1), live_row(1), live_row(1)];
let carriers = scan_partition_carriers(&rows);
assert!(carriers.range_tombstones.is_empty());
assert_eq!(carriers.max_partition_deletion, None);
assert!(carriers.partition_delete_key.is_none());
}
#[test]
fn extracts_range_carriers_in_first_seen_order() {
let rows = vec![
range_carrier(1, 200, 20),
live_row(1),
range_carrier(1, 100, 10),
];
let carriers = scan_partition_carriers(&rows);
assert_eq!(carriers.range_tombstones.len(), 2);
assert_eq!(carriers.range_tombstones[0].1.deletion_time, 200);
assert_eq!(carriers.range_tombstones[1].1.deletion_time, 100);
assert_eq!(carriers.max_partition_deletion, None);
assert!(carriers.partition_delete_key.is_none());
}
#[test]
fn keeps_max_partition_deletion_across_carriers() {
let rows = vec![
partition_carrier(7, 50, 5),
live_row(7),
partition_carrier(7, 300, 30),
partition_carrier(7, 200, 20),
];
let carriers = scan_partition_carriers(&rows);
assert_eq!(carriers.max_partition_deletion, Some((300, 30)));
assert_eq!(carriers.partition_delete_key, Some(key(7)));
assert!(carriers.range_tombstones.is_empty());
}
#[test]
fn equal_mfda_keeps_first_seen_local_deletion_time() {
let rows = vec![partition_carrier(3, 100, 11), partition_carrier(3, 100, 22)];
let carriers = scan_partition_carriers(&rows);
assert_eq!(carriers.max_partition_deletion, Some((100, 11)));
}
#[test]
fn mixed_carriers_and_rows_split_correctly() {
let rows = vec![
live_row(5),
range_carrier(5, 400, 40),
partition_carrier(5, 250, 25),
live_row(5),
range_carrier(5, 150, 15),
];
let carriers = scan_partition_carriers(&rows);
assert_eq!(carriers.range_tombstones.len(), 2);
assert_eq!(carriers.range_tombstones[0].1.deletion_time, 400);
assert_eq!(carriers.range_tombstones[1].1.deletion_time, 150);
assert_eq!(carriers.max_partition_deletion, Some((250, 25)));
assert_eq!(carriers.partition_delete_key, Some(key(5)));
}
#[test]
fn is_range_marker_carrier_rejects_non_empty_live_row() {
let entry = live_row(1).with_range_deletion(range_tombstone(100, 10));
assert!(!is_range_marker_carrier(&entry));
assert!(is_range_marker_carrier(&range_carrier(1, 100, 10)));
}
}