#[cfg(feature = "write-support")]
use super::model::{MergeEntry, RowData};
#[cfg(feature = "write-support")]
use super::{carriers, KWayMerger, PurgeCounts};
#[cfg(feature = "write-support")]
use crate::error::{Error, Result};
#[cfg(feature = "write-support")]
use crate::storage::write_engine::mutation::{ClusteringKey, DecoratedKey, RangeTombstone};
#[cfg(feature = "write-support")]
use std::cmp::Reverse;
#[cfg(feature = "write-support")]
use std::collections::VecDeque;
#[cfg(all(test, feature = "write-support", feature = "dhat-heap"))]
mod streaming_dhat_test;
#[cfg(feature = "write-support")]
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum StreamingStep {
ClusterGroup {
key: DecoratedKey,
row: Box<MergeEntry>,
},
PartitionEnd {
key: DecoratedKey,
},
Complete,
}
#[cfg(feature = "write-support")]
#[derive(Debug)]
pub(crate) struct PartitionReconcileCheckpoint {
pending_rows: VecDeque<MergeEntry>,
range_tombstones: Vec<(DecoratedKey, RangeTombstone)>,
max_partition_deletion: Option<(i64, i32)>,
partition_delete_key: Option<DecoratedKey>,
partition_deletion_resolved: bool,
current_cluster: Option<(Option<ClusteringKey>, Vec<MergeEntry>)>,
purges: PurgeCounts,
awaiting_flush: VecDeque<MergeEntry>,
range_tombstones_emitted: bool,
}
#[cfg(feature = "write-support")]
impl PartitionReconcileCheckpoint {
fn fresh() -> Self {
Self {
pending_rows: VecDeque::new(),
range_tombstones: Vec::new(),
max_partition_deletion: None,
partition_delete_key: None,
partition_deletion_resolved: false,
current_cluster: None,
purges: PurgeCounts::default(),
awaiting_flush: VecDeque::new(),
range_tombstones_emitted: false,
}
}
}
#[cfg(feature = "write-support")]
pub struct StreamingMerger<'a> {
merger: &'a mut KWayMerger,
partition_key: Option<DecoratedKey>,
state: PartitionReconcileCheckpoint,
}
#[cfg(feature = "write-support")]
impl<'a> StreamingMerger<'a> {
pub fn new(merger: &'a mut KWayMerger) -> Self {
Self {
merger,
partition_key: None,
state: PartitionReconcileCheckpoint::fresh(),
}
}
pub(crate) fn resume(
merger: &'a mut KWayMerger,
saved: Option<(DecoratedKey, PartitionReconcileCheckpoint)>,
) -> Self {
let (partition_key, state) = match saved {
Some((key, checkpoint)) => (Some(key), checkpoint),
None => (None, PartitionReconcileCheckpoint::fresh()),
};
Self {
merger,
partition_key,
state,
}
}
pub(crate) fn into_paused_state(self) -> Option<(DecoratedKey, PartitionReconcileCheckpoint)> {
self.partition_key.map(|key| (key, self.state))
}
pub(crate) fn current_partition_key(&self) -> Option<DecoratedKey> {
self.partition_key.clone()
}
fn pull_one(&mut self) -> Result<Option<MergeEntry>> {
let Some(current_key) = &self.partition_key else {
return Ok(None);
};
match self.merger.heap.peek() {
Some(Reverse(top)) if &top.entry.key == current_key => {}
_ => return Ok(None),
}
let Reverse(top) = self
.merger
.heap
.pop()
.ok_or_else(|| Error::InvalidInput("Merge heap unexpectedly empty".to_string()))?;
let entry = top.entry;
self.merger.refill_heap(entry.run_index)?;
Ok(Some(entry))
}
fn resolve_partition_prefix(&mut self) {
if self.state.partition_deletion_resolved {
return;
}
self.state.partition_deletion_resolved = true;
let (effective_gc_before, max_purgeable_timestamp) = self.merger.effective_gc_settings();
if let (Some((pmfda, pldt)), Some(key)) = (
self.state.max_partition_deletion,
self.state.partition_delete_key.take(),
) {
let purge = match effective_gc_before {
Some(gc_before) => {
i64::from(pldt as u32) < gc_before && pmfda < max_purgeable_timestamp
}
None => false,
};
if purge {
self.state.purges.partition_tombstones += 1;
} else {
self.state.purges.emitted += 1;
self.state.pending_rows.push_back(
MergeEntry::new(
usize::MAX,
key,
None,
pmfda,
RowData::Tombstone {
deletion_time: pmfda,
local_deletion_time: pldt,
},
)
.with_partition_deletion((pmfda, pldt)),
);
}
}
}
fn flush_range_tombstones(
&mut self,
effective_gc_before: Option<i64>,
max_purgeable_timestamp: i64,
) {
self.state.range_tombstones_emitted = true;
KWayMerger::coalesce_range_tombstones(
&mut self.state.range_tombstones,
&self.merger.schema,
);
if let Some((pmfda, _)) = self.state.max_partition_deletion {
self.state
.range_tombstones
.retain(|(_, rt)| rt.deletion_time > pmfda);
}
for (key, rt) in self.state.range_tombstones.clone() {
if let Some(gc_before) = effective_gc_before {
if i64::from(rt.local_deletion_time as u32) < gc_before
&& rt.deletion_time < max_purgeable_timestamp
{
self.state.purges.range_tombstones += 1;
continue;
}
}
self.state.purges.emitted += 1;
self.state.pending_rows.push_back(
MergeEntry::new(
usize::MAX,
key,
None,
rt.deletion_time,
RowData::Live { cells: Vec::new() },
)
.with_range_deletion(rt),
);
}
}
fn reshadow_and_flush_awaiting(&mut self) {
if self.state.awaiting_flush.is_empty() {
return;
}
KWayMerger::coalesce_range_tombstones(
&mut self.state.range_tombstones,
&self.merger.schema,
);
for entry in std::mem::take(&mut self.state.awaiting_flush) {
if let Some(shadowed) = KWayMerger::apply_range_shadowing(
entry,
&self.state.range_tombstones,
&self.merger.schema,
) {
if let Some(survivor) = KWayMerger::apply_partition_shadowing(
shadowed,
self.state.max_partition_deletion,
) {
self.state.pending_rows.push_back(survivor);
}
}
}
}
fn finalize_current_cluster(&mut self) -> Result<()> {
let Some((ck, rows)) = self.state.current_cluster.take() else {
self.resolve_partition_prefix();
return Ok(());
};
self.resolve_partition_prefix();
let (effective_gc_before, max_purgeable_timestamp) = self.merger.effective_gc_settings();
if let Some(entry) = KWayMerger::reconcile_cluster_with_overlap_counted(
ck,
rows,
&self.merger.schema.dropped_columns,
effective_gc_before,
max_purgeable_timestamp,
self.merger.now_secs,
&mut self.state.purges,
) {
if self.state.range_tombstones.is_empty() {
if let Some(survivor) =
KWayMerger::apply_partition_shadowing(entry, self.state.max_partition_deletion)
{
self.state.pending_rows.push_back(survivor);
}
} else {
self.state.awaiting_flush.push_back(entry);
}
}
Ok(())
}
pub fn step_streaming(&mut self) -> Result<StreamingStep> {
loop {
if let Some(row) = self.state.pending_rows.pop_front() {
let Some(key) = self.partition_key.clone() else {
return Err(Error::InvalidInput(
"pending row without an active partition key".to_string(),
));
};
return Ok(StreamingStep::ClusterGroup {
key,
row: Box::new(row),
});
}
if self.partition_key.is_none() {
if self.merger.heap.is_empty() && self.merger.current_partition.is_none() {
self.merger.initialize_heap()?;
}
let Some(Reverse(top)) = self.merger.heap.peek() else {
return Ok(StreamingStep::Complete);
};
self.partition_key = Some(top.entry.key.clone());
self.state = PartitionReconcileCheckpoint::fresh();
}
match self.pull_one()? {
Some(entry) => {
if entry.is_partition_delete_carrier() {
if let Some((mfda, ldt)) = entry.partition_deletion {
self.state
.partition_delete_key
.get_or_insert_with(|| entry.key.clone());
match self.state.max_partition_deletion {
Some((cur_mfda, _)) if cur_mfda >= mfda => {}
_ => self.state.max_partition_deletion = Some((mfda, ldt)),
}
continue;
}
}
if carriers::is_range_marker_carrier(&entry) {
if let Some(rt) = entry.range_deletion.clone() {
self.state.range_tombstones.push((entry.key.clone(), rt));
if self.state.partition_deletion_resolved {
self.reshadow_and_flush_awaiting();
}
}
continue;
}
let starts_new_cluster = match &self.state.current_cluster {
Some((ck, _)) => *ck != entry.clustering_key,
None => false,
};
if starts_new_cluster {
self.finalize_current_cluster()?;
}
match &mut self.state.current_cluster {
Some((_, rows)) => rows.push(entry),
None => {
self.state.current_cluster =
Some((entry.clustering_key.clone(), vec![entry]));
}
}
}
None => {
self.finalize_current_cluster()?;
self.reshadow_and_flush_awaiting();
if !self.state.range_tombstones.is_empty()
&& !self.state.range_tombstones_emitted
{
let (effective_gc_before, max_purgeable_timestamp) =
self.merger.effective_gc_settings();
self.flush_range_tombstones(effective_gc_before, max_purgeable_timestamp);
}
if !self.state.pending_rows.is_empty() {
continue;
}
let purged = self.state.purges.total();
if purged > 0 {
crate::observability::add_counter(
crate::observability::catalog::COMPACTION_TOMBSTONES_PURGED,
purged,
&[],
);
}
let Some(key) = self.partition_key.take() else {
return Err(Error::InvalidInput(
"partition ended without an active key".to_string(),
));
};
return Ok(StreamingStep::PartitionEnd { key });
}
}
}
}
}
#[cfg(all(test, feature = "write-support"))]
mod tests {
use super::*;
use super::super::MergeStep;
use crate::schema::{ClusteringColumn, ClusteringOrder, KeyColumn, TableSchema};
use crate::storage::write_engine::merge::model::{CellData, RowData};
use crate::storage::write_engine::merge::{RunReader, SSTableRowIterator};
use crate::storage::write_engine::mutation::{ClusteringKey, RangeTombstone};
use crate::types::Value;
use std::cmp::Reverse;
use std::collections::{BinaryHeap, HashMap};
fn key(token: i64) -> DecoratedKey {
DecoratedKey::new(token, vec![token as u8])
}
fn ck(v: i32) -> ClusteringKey {
ClusteringKey::single("ck", Value::Integer(v))
}
fn live_entry(run_index: usize, token: i64, cluster: i32, ts: i64) -> MergeEntry {
MergeEntry::new(
run_index,
key(token),
Some(ck(cluster)),
ts,
RowData::Live {
cells: vec![CellData::new("c".to_string(), Value::Integer(cluster), ts)],
},
)
}
fn range_carrier(run_index: usize, token: i64, deletion_time: i64, ldt: i32) -> MergeEntry {
let rt = RangeTombstone {
start: crate::storage::write_engine::mutation::ClusteringBound::Bottom,
end: crate::storage::write_engine::mutation::ClusteringBound::Inclusive(ck(1)),
deletion_time,
local_deletion_time: ldt,
};
MergeEntry::new(
run_index,
key(token),
None,
deletion_time,
RowData::Live { cells: Vec::new() },
)
.with_range_deletion(rt)
}
fn partition_carrier(run_index: usize, token: i64, mfda: i64, ldt: i32) -> MergeEntry {
MergeEntry::new(
run_index,
key(token),
None,
mfda,
RowData::Live { cells: Vec::new() },
)
.with_partition_deletion((mfda, ldt))
}
fn test_schema(order: ClusteringOrder) -> TableSchema {
TableSchema {
keyspace: "ks_1668".to_string(),
table: "t_1668".to_string(),
partition_keys: vec![KeyColumn {
name: "id".to_string(),
data_type: "int".to_string(),
position: 0,
}],
clustering_keys: vec![ClusteringColumn {
name: "ck".to_string(),
data_type: "int".to_string(),
position: 0,
order,
}],
columns: vec![],
comments: HashMap::new(),
dropped_columns: HashMap::new(),
}
}
fn two_col_schema() -> TableSchema {
TableSchema {
keyspace: "ks_1668".to_string(),
table: "t_1668_multi".to_string(),
partition_keys: vec![KeyColumn {
name: "id".to_string(),
data_type: "int".to_string(),
position: 0,
}],
clustering_keys: vec![
ClusteringColumn {
name: "ck1".to_string(),
data_type: "int".to_string(),
position: 0,
order: ClusteringOrder::Asc,
},
ClusteringColumn {
name: "ck2".to_string(),
data_type: "int".to_string(),
position: 1,
order: ClusteringOrder::Desc,
},
],
columns: vec![],
comments: HashMap::new(),
dropped_columns: HashMap::new(),
}
}
fn ck_multi(pairs: &[(&str, i32)]) -> ClusteringKey {
ClusteringKey::new(
pairs
.iter()
.map(|(n, v)| (n.to_string(), Value::Integer(*v)))
.collect(),
)
}
fn live_entry_ck(run_index: usize, token: i64, ck: ClusteringKey, ts: i64) -> MergeEntry {
MergeEntry::new(
run_index,
key(token),
Some(ck),
ts,
RowData::Live {
cells: vec![CellData::new("c".to_string(), Value::Integer(1), ts)],
},
)
}
struct VecIterator(std::vec::IntoIter<MergeEntry>);
impl SSTableRowIterator for VecIterator {
fn next(&mut self) -> Option<Result<MergeEntry>> {
self.0.next().map(Ok)
}
}
fn merger_over(entries: Vec<MergeEntry>, schema: TableSchema) -> KWayMerger {
merger_over_runs(vec![entries], schema)
}
fn merger_over_runs(runs: Vec<Vec<MergeEntry>>, schema: TableSchema) -> KWayMerger {
KWayMerger {
runs: runs
.into_iter()
.map(|entries| RunReader::new(Box::new(VecIterator(entries.into_iter()))))
.collect(),
heap: BinaryHeap::new(),
current_partition: None,
gc_before_secs: None,
now_secs: None,
purge_safe: false,
max_purgeable_timestamp: None,
schema_arc: std::sync::Arc::new(schema.clone()),
schema,
}
}
#[test]
fn heap_groups_contiguously_by_ord() {
let mut heap: BinaryHeap<Reverse<MergeEntry>> = BinaryHeap::new();
for &(run, cluster, ts) in &[
(2, 1, 100),
(0, 0, 300),
(1, 2, 50),
(1, 0, 200),
(2, 2, 60),
(0, 1, 150),
(0, 2, 70),
(2, 0, 250),
(1, 1, 120),
] {
heap.push(Reverse(live_entry(run, 1, cluster, ts)));
}
let mut popped_clusters = Vec::new();
while let Some(Reverse(entry)) = heap.pop() {
popped_clusters.push(entry.clustering_key.clone());
}
let mut closed: Vec<Option<ClusteringKey>> = Vec::new();
let mut current: Option<Option<ClusteringKey>> = None;
for c in &popped_clusters {
if current.as_ref() != Some(c) {
if let Some(prev) = current.take() {
assert!(
!closed.contains(&prev),
"clustering key reappeared non-contiguously after the heap moved past it"
);
closed.push(prev);
}
current = Some(c.clone());
}
}
let mut distinct: Vec<Option<ClusteringKey>> = Vec::new();
for c in &popped_clusters {
if !distinct.contains(c) {
distinct.push(c.clone());
}
}
assert_eq!(distinct.len(), 3, "expected 3 distinct clustering keys");
}
#[test]
fn empty_merger_yields_complete() {
let mut merger = merger_over(vec![], test_schema(ClusteringOrder::Asc));
let mut stream = StreamingMerger::new(&mut merger);
assert!(matches!(
stream.step_streaming().unwrap(),
StreamingStep::Complete
));
}
fn drain_streaming(merger: &mut KWayMerger) -> Vec<MergeEntry> {
let mut stream = StreamingMerger::new(merger);
let mut rows = Vec::new();
let mut in_partition = false;
loop {
match stream.step_streaming().unwrap() {
StreamingStep::ClusterGroup { row, .. } => {
in_partition = true;
rows.push(*row);
}
StreamingStep::PartitionEnd { .. } => {
assert!(in_partition, "PartitionEnd with no preceding ClusterGroup");
in_partition = false;
}
StreamingStep::Complete => {
assert!(!in_partition, "Complete while a partition was still open");
return rows;
}
}
}
}
fn drain_whole_partition(merger: &mut KWayMerger) -> Vec<MergeEntry> {
let mut rows = Vec::new();
loop {
match merger.step().unwrap() {
MergeStep::Partition {
rows: partition_rows,
..
} => rows.extend(partition_rows),
MergeStep::Complete => return rows,
}
}
}
fn split_carriers_and_rows(rows: Vec<MergeEntry>) -> (Vec<MergeEntry>, Vec<MergeEntry>) {
rows.into_iter()
.partition(|entry| entry.clustering_key.is_none())
}
fn carrier_sort_key(entry: &MergeEntry) -> (bool, i64) {
(entry.partition_deletion.is_some(), entry.timestamp)
}
#[test]
fn step_streaming_matches_step_for_mixed_tombstone_fixture() {
let schema = test_schema(ClusteringOrder::Asc);
let entries = vec![
partition_carrier(0, 1, 10, 1),
live_entry(0, 1, 0, 100),
live_entry(0, 1, 1, 100),
range_carrier(0, 1, 50, 5),
live_entry(0, 1, 2, 100),
];
let mut old_merger = merger_over(entries.clone(), schema.clone());
let old_rows = drain_whole_partition(&mut old_merger);
let mut new_merger = merger_over(entries, schema);
let new_rows = drain_streaming(&mut new_merger);
assert!(
!old_rows.is_empty(),
"fixture must produce at least one row"
);
let (mut old_carriers, old_clustering_rows) = split_carriers_and_rows(old_rows);
let (mut new_carriers, new_clustering_rows) = split_carriers_and_rows(new_rows);
assert_eq!(
old_clustering_rows, new_clustering_rows,
"streaming path must reconcile the SAME clustering rows, in the \
SAME order, as the whole-partition path"
);
old_carriers.sort_by_key(carrier_sort_key);
new_carriers.sort_by_key(carrier_sort_key);
assert_eq!(
old_carriers, new_carriers,
"streaming path must re-emit the SAME carrier markers (content, \
not necessarily position) as the whole-partition path"
);
}
#[test]
fn step_streaming_matches_step_for_plain_multi_cluster_fixture() {
let schema = test_schema(ClusteringOrder::Asc);
let entries = vec![
live_entry(0, 1, 0, 100),
live_entry(1, 1, 0, 90), live_entry(0, 1, 1, 100),
live_entry(0, 1, 2, 100),
live_entry(0, 2, 0, 100), ];
let mut old_merger = merger_over(entries.clone(), schema.clone());
let old_rows = drain_whole_partition(&mut old_merger);
let mut new_merger = merger_over(entries, schema);
let new_rows = drain_streaming(&mut new_merger);
assert_eq!(old_rows, new_rows);
}
fn ck_value(entry: &MergeEntry) -> i32 {
match entry
.clustering_key
.as_ref()
.and_then(|k| k.columns.first())
{
Some((_, Value::Integer(v))) => *v,
other => panic!("expected a single Integer clustering column, got {other:?}"),
}
}
#[test]
fn step_streaming_matches_step_for_desc_clustering_fixture() {
let schema = test_schema(ClusteringOrder::Desc);
let runs = vec![
vec![live_entry(0, 1, 2, 100)],
vec![live_entry(1, 1, 1, 100)],
vec![live_entry(2, 1, 0, 100)],
];
let expected_order = [2, 1, 0];
let mut old_merger = merger_over_runs(runs.clone(), schema.clone());
let old_rows = drain_whole_partition(&mut old_merger);
let old_order: Vec<i32> = old_rows.iter().map(ck_value).collect();
assert_eq!(
old_order, expected_order,
"step() must emit DESC clustering order"
);
let mut new_merger = merger_over_runs(runs, schema);
let new_rows = drain_streaming(&mut new_merger);
let new_order: Vec<i32> = new_rows.iter().map(ck_value).collect();
assert_eq!(
new_order, expected_order,
"step_streaming() must ALSO emit DESC clustering order, matching step()"
);
}
#[test]
fn step_streaming_matches_step_for_absent_trailing_component_fixture() {
let schema = two_col_schema();
let runs = vec![
vec![live_entry_ck(
0,
1,
ck_multi(&[("ck1", 5), ("ck2", 1)]),
100,
)],
vec![live_entry_ck(1, 1, ck_multi(&[("ck1", 5)]), 100)], vec![live_entry_ck(
2,
1,
ck_multi(&[("ck1", 5), ("ck2", 2)]),
100,
)],
];
fn ck2_marker(entry: &MergeEntry) -> Option<i32> {
entry.clustering_key.as_ref().and_then(|k| {
k.columns
.iter()
.find(|(name, _)| name == "ck2")
.map(|(_, v)| match v {
Value::Integer(v) => *v,
other => panic!("expected Integer ck2, got {other:?}"),
})
})
}
let expected_order = [None, Some(2), Some(1)];
let mut old_merger = merger_over_runs(runs.clone(), schema.clone());
let old_rows = drain_whole_partition(&mut old_merger);
let old_order: Vec<Option<i32>> = old_rows.iter().map(ck2_marker).collect();
assert_eq!(
old_order, expected_order,
"step() must put the absent-ck2 row first (NULL-first), then DESC by ck2"
);
let mut new_merger = merger_over_runs(runs, schema);
let new_rows = drain_streaming(&mut new_merger);
let new_order: Vec<Option<i32>> = new_rows.iter().map(ck2_marker).collect();
assert_eq!(
new_order, expected_order,
"step_streaming() must ALSO put the absent-ck2 row first, then DESC by ck2"
);
}
#[test]
fn range_tombstone_carrier_round_trips_as_cluster_group() {
let carrier = range_carrier(0, 9, 500, 50);
assert!(super::super::carriers::is_range_marker_carrier(&carrier));
let expected = MergeEntry {
run_index: usize::MAX,
..carrier.clone()
};
let mut merger = merger_over(vec![carrier], test_schema(ClusteringOrder::Asc));
let mut stream = StreamingMerger::new(&mut merger);
match stream.step_streaming().unwrap() {
StreamingStep::ClusterGroup { row, .. } => assert_eq!(*row, expected),
other => panic!("expected ClusterGroup, got {other:?}"),
}
}
#[test]
fn static_row_carrier_always_sorts_first_regardless_of_partition_width() {
use crate::schema::Column;
use crate::storage::write_engine::merge::model::CellData;
let mut schema = test_schema(ClusteringOrder::Asc);
schema.columns = vec![Column {
name: "region".to_string(),
data_type: "text".to_string(),
nullable: true,
default: None,
is_static: true,
}];
let mut entries = vec![MergeEntry::new(
0,
key(1),
None,
100,
RowData::Live {
cells: vec![CellData::new(
"region".to_string(),
Value::text("us-east".to_string()),
100,
)],
},
)];
for ck in 0..50 {
entries.push(live_entry(0, 1, ck, 100));
}
let mut merger = merger_over(entries, schema);
let rows = match merger.step().unwrap() {
MergeStep::Partition { rows, .. } => rows,
other => panic!("expected Partition, got {other:?}"),
};
assert_eq!(
rows.len(),
51,
"expected the static row + 50 clustering rows"
);
assert!(
rows[0].clustering_key.is_none(),
"the static row (clustering_key: None) must be FIRST, regardless \
of the 50 regular clustering rows that follow"
);
assert!(
rows[1..].iter().all(|r| r.clustering_key.is_some()),
"every row AFTER the first must be a regular Some(ck) clustering \
row — the None-keyed prefix is exactly one entry wide here"
);
}
}