use std::collections::{HashMap, HashSet};
use crate::storage::write_engine::mutation::{ClusteringKey, DecoratedKey, RangeTombstone};
use crate::storage::write_engine::reconcile_rules;
use super::{CellData, ComplexDeletion, KWayMerger, MergeEntry, PurgeCounts, RowData};
type CellKey = (String, Option<Vec<u8>>);
pub(super) struct ReconcileState {
clustering_key: Option<ClusteringKey>,
key: Option<DecoratedKey>,
run_index: usize,
row_del: Option<i64>,
row_del_ldt: i32,
order: Vec<CellKey>,
winners: HashMap<CellKey, CellData>,
complex_deletions: Vec<ComplexDeletion>,
range_deletion: Option<RangeTombstone>,
after_row_del: Vec<CellData>,
surviving: Vec<CellData>,
had_data_before: bool,
}
impl ReconcileState {
pub(super) fn new(clustering_key: Option<ClusteringKey>) -> Self {
Self {
clustering_key,
key: None,
run_index: usize::MAX,
row_del: None,
row_del_ldt: 0,
order: Vec::new(),
winners: HashMap::new(),
complex_deletions: Vec::new(),
range_deletion: None,
after_row_del: Vec::new(),
surviving: Vec::new(),
had_data_before: false,
}
}
pub(super) fn has_key(&self) -> bool {
self.key.is_some()
}
pub(super) fn fold_row_deletions(&mut self, cluster_rows: &[MergeEntry]) {
for entry in cluster_rows {
if self.key.is_none() {
self.key = Some(entry.key.clone());
}
self.run_index = self.run_index.min(entry.run_index);
for cd in &entry.complex_deletions {
if !self.complex_deletions.contains(cd) {
self.complex_deletions.push(cd.clone());
}
}
if let Some(rd) = &entry.range_deletion {
let replace = match &self.range_deletion {
None => true,
Some(current) => rd.deletion_time > current.deletion_time,
};
if replace {
self.range_deletion = Some(rd.clone());
}
}
if let Some((del_ts, del_ldt)) = entry.row_deletion {
if self.row_del.is_none_or(|d| del_ts > d) {
self.row_del = Some(del_ts);
self.row_del_ldt = del_ldt;
}
}
if let RowData::Tombstone {
deletion_time,
local_deletion_time,
} = &entry.row_data
{
if self.row_del.is_none_or(|d| *deletion_time > d) {
self.row_del = Some(*deletion_time);
self.row_del_ldt = *local_deletion_time;
}
}
}
}
pub(super) fn resolve_cell_winners(&mut self, cluster_rows: &[MergeEntry]) {
for entry in cluster_rows {
if let RowData::Live { cells } = &entry.row_data {
for cell in cells {
let cell_key: CellKey = (cell.column.clone(), cell.cell_path.clone());
match self.winners.get(&cell_key) {
None => {
self.order.push(cell_key.clone());
self.winners.insert(cell_key, cell.clone());
}
Some(existing) => {
if reconcile_rules::cell_wins(cell, existing) {
self.winners.insert(cell_key, cell.clone());
}
}
}
}
}
}
}
pub(super) fn apply_complex_deletions(&mut self) {
if !self.complex_deletions.is_empty() {
let mut active: HashMap<String, ComplexDeletion> = HashMap::new();
let mut active_order: Vec<String> = Vec::new();
for cd in self.complex_deletions.drain(..) {
match active.get_mut(&cd.column) {
None => {
active_order.push(cd.column.clone());
active.insert(cd.column.clone(), cd);
}
Some(existing)
if reconcile_rules::complex_deletion_supersedes(
cd.marked_for_delete_at,
existing.marked_for_delete_at,
) =>
{
*existing = cd;
}
Some(_) => {}
}
}
for column in &active_order {
if let Some(cd) = active.get(column) {
let mfda = cd.marked_for_delete_at;
let winners = &mut self.winners;
self.order.retain(|cell_key| {
let (cell_column, cell_path) = cell_key;
if cell_column != column || cell_path.is_none() {
return true;
}
match winners.get(cell_key) {
Some(cell)
if reconcile_rules::element_survives_complex_deletion(
cell.timestamp,
mfda,
) =>
{
true
}
Some(_) => {
winners.remove(cell_key);
false
}
None => false,
}
});
}
}
self.complex_deletions = active_order
.into_iter()
.filter_map(|column| active.remove(&column))
.collect();
}
}
pub(super) fn shadow_by_row_deletion(&mut self) {
let row_del = self.row_del;
let winners = &mut self.winners;
self.after_row_del = std::mem::take(&mut self.order)
.into_iter()
.filter_map(|cell_key| winners.remove(&cell_key))
.filter(|cell| match row_del {
Some(d) => cell.timestamp > d,
None => true,
})
.collect();
}
pub(super) fn filter_dropped_columns(&mut self, dropped_columns: &HashMap<String, i64>) {
self.surviving = self
.after_row_del
.iter()
.filter(|cell| match dropped_columns.get(&cell.column) {
Some(drop_time) => cell.timestamp > *drop_time,
None => true,
})
.cloned()
.collect();
let ck_names: HashSet<&str> = self
.clustering_key
.as_ref()
.map(|ck| ck.columns.iter().map(|(n, _)| n.as_str()).collect())
.unwrap_or_default();
let is_data_cell = |cell: &CellData| !ck_names.contains(cell.column.as_str());
self.had_data_before = self.after_row_del.iter().any(is_data_cell);
}
pub(super) fn expire_ttl_cells(&mut self, now_secs: Option<i64>) {
let Some(now) = now_secs else {
return; };
for cell in &mut self.surviving {
if cell.is_complex_element || KWayMerger::is_cell_tombstone(cell) {
continue;
}
let (Some(_ttl), Some(ldt)) = (cell.ttl, cell.local_deletion_time) else {
continue; };
if i64::from(ldt as u32) < now {
cell.value = crate::types::Value::Tombstone(crate::types::TombstoneInfo {
deletion_time: cell.timestamp,
tombstone_type: crate::types::TombstoneType::CellTombstone,
local_deletion_time: i64::from(ldt as u32),
ttl: None,
range_start: None,
range_end: None,
});
cell.ttl = None;
}
}
}
pub(super) fn purge_gc_grace(
&mut self,
gc_before_secs: Option<i64>,
max_purgeable_timestamp: i64,
purges: &mut PurgeCounts,
) {
if let Some(gc_before) = gc_before_secs {
self.surviving.retain(|cell| {
if KWayMerger::is_cell_tombstone(cell) || cell.is_deleted {
let gc_purgeable = match KWayMerger::cell_effective_ldt(cell) {
Some(ldt) => i64::from(ldt as u32) < gc_before,
None => false,
};
let overlap_purgeable = cell.timestamp < max_purgeable_timestamp;
let keep = !(gc_purgeable && overlap_purgeable);
if !keep {
purges.cell_tombstones += 1;
}
keep
} else {
true
}
});
if self.row_del_ldt != 0
&& i64::from(self.row_del_ldt as u32) < gc_before
&& self.row_del.is_some_and(|d| d < max_purgeable_timestamp)
{
self.row_del = None;
purges.row_tombstones += 1;
}
self.complex_deletions.retain(|cd| {
let gc_purgeable = i64::from(cd.local_deletion_time as u32) < gc_before;
let overlap_purgeable = cd.marked_for_delete_at < max_purgeable_timestamp;
let keep = !(gc_purgeable && overlap_purgeable);
if !keep {
purges.complex_deletions += 1;
}
keep
});
}
}
pub(super) fn build(self) -> Option<MergeEntry> {
let ReconcileState {
clustering_key,
key,
run_index,
row_del,
row_del_ldt,
complex_deletions,
range_deletion,
surviving,
had_data_before,
..
} = self;
let key = key?;
let ck_names: HashSet<&str> = clustering_key
.as_ref()
.map(|ck| ck.columns.iter().map(|(n, _)| n.as_str()).collect())
.unwrap_or_default();
let is_data_cell = |cell: &CellData| !ck_names.contains(cell.column.as_str());
let has_data_after = surviving.iter().any(is_data_cell);
let purged_to_empty = had_data_before && !has_data_after;
drop(ck_names);
let has_carried_metadata = !complex_deletions.is_empty() || range_deletion.is_some();
let built = if !surviving.is_empty() && !purged_to_empty {
let row_ts = surviving.iter().map(|c| c.timestamp).max().unwrap_or(0);
let live = MergeEntry::new(
run_index,
key,
clustering_key,
row_ts,
RowData::Live { cells: surviving },
);
Some(match row_del {
Some(deletion_time) => live.with_row_deletion(deletion_time, row_del_ldt),
None => live,
})
} else if let Some(deletion_time) = row_del {
Some(MergeEntry::new(
run_index,
key,
clustering_key,
deletion_time,
RowData::Tombstone {
deletion_time,
local_deletion_time: row_del_ldt,
},
))
} else if has_carried_metadata {
Some(MergeEntry::new(
run_index,
key,
clustering_key,
0,
RowData::Live { cells: vec![] },
))
} else {
None
};
built.map(|entry| {
let entry = if complex_deletions.is_empty() {
entry
} else {
entry.with_complex_deletions(complex_deletions)
};
match range_deletion {
Some(rd) => entry.with_range_deletion(rd),
None => entry,
}
})
}
}