use crate::storage::sstable::writer::stats_writer::StatisticsMetadata;
use crate::storage::write_engine::mutation::{CellOperation, Mutation};
pub(crate) fn fold_mutation_stats(stats: &mut StatisticsMetadata, mutation: &Mutation) {
stats.update_timestamp(mutation.timestamp_micros);
if let Some(cell_ts) = &mutation.cell_write_timestamps {
for ts in cell_ts.values() {
stats.update_timestamp(*ts);
}
}
if let Some(ttl) = mutation.ttl_seconds {
stats.update_ttl(ttl as i32);
let now_seconds = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs() as i32)
.unwrap_or(0);
let local_deletion_time = now_seconds.saturating_add(ttl as i32);
stats.update_local_deletion_time(local_deletion_time);
}
for op in &mutation.operations {
match op {
CellOperation::WriteWithTtl {
ttl_seconds,
local_deletion_time,
..
} => {
stats.update_ttl(*ttl_seconds as i32);
let local_deletion_time = match local_deletion_time {
Some(ldt) => *ldt,
None => {
let now_seconds = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs() as i32)
.unwrap_or(0);
now_seconds.saturating_add(*ttl_seconds as i32)
}
};
stats.update_local_deletion_time(local_deletion_time);
}
op @ (CellOperation::Delete { .. } | CellOperation::DeleteRow) => {
let local_deletion_time =
crate::storage::sstable::writer::data_writer::op_cell_local_deletion_time(
op, mutation,
);
stats.update_local_deletion_time(local_deletion_time);
}
CellOperation::ComplexDeletion {
marked_for_delete_at,
local_deletion_time,
..
} => {
stats.update_timestamp(*marked_for_delete_at);
stats.update_local_deletion_time(*local_deletion_time);
}
CellOperation::WriteComplexElement {
timestamp_micros,
ttl_seconds,
local_deletion_time,
is_deleted,
..
} => {
stats.update_timestamp(*timestamp_micros);
if let Some(ttl) = ttl_seconds {
stats.update_ttl(*ttl as i32);
}
if let Some(ldt) = local_deletion_time {
stats.update_local_deletion_time(*ldt);
}
if !*is_deleted && ttl_seconds.is_none() && local_deletion_time.is_none() {
stats.note_live_local_deletion_time();
}
}
CellOperation::Write { value, .. } => {
if mutation.ttl_seconds.is_none() && !matches!(value, crate::types::Value::Null) {
stats.note_live_local_deletion_time();
}
}
}
}
if let Some(pt) = &mutation.partition_tombstone {
stats.update_timestamp(pt.deletion_time);
stats.update_local_deletion_time(pt.local_deletion_time);
stats.mark_partition_level_deletion();
}
for rt in &mutation.range_tombstones {
stats.update_timestamp(rt.deletion_time);
stats.update_local_deletion_time(rt.local_deletion_time);
}
if let Some((deletion_time, ldt)) = mutation.row_tombstone {
stats.update_timestamp(deletion_time);
stats.update_local_deletion_time(ldt);
}
}
pub(crate) fn merge_stats_fold(into: &mut StatisticsMetadata, from: &StatisticsMetadata) {
into.update_timestamp(from.min_timestamp);
into.update_timestamp(from.max_timestamp);
into.update_local_deletion_time(from.min_local_deletion_time);
if from.max_local_deletion_time > i32::MIN {
into.update_local_deletion_time(from.max_local_deletion_time);
}
if from.max_local_deletion_time == i32::MAX {
into.note_live_local_deletion_time();
}
if from.min_ttl != i32::MAX {
into.update_ttl(from.min_ttl);
}
into.update_ttl(from.max_ttl);
if from.has_partition_level_deletions {
into.mark_partition_level_deletion();
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::storage::write_engine::mutation::{
ClusteringBound, ClusteringKey, PartitionKey, PartitionTombstone, RangeTombstone, TableId,
};
use crate::types::Value;
fn table() -> TableId {
TableId::new("ks", "t")
}
fn pk() -> PartitionKey {
PartitionKey::single("id", Value::Integer(1))
}
fn ck(n: i32) -> ClusteringKey {
ClusteringKey::single("ck", Value::Integer(n))
}
fn representative_mutations() -> Vec<Mutation> {
let mut partition_only = Mutation::new(table(), pk(), None, vec![], 500, None);
partition_only.partition_tombstone = Some(PartitionTombstone {
deletion_time: 100,
local_deletion_time: 1_000,
});
let mut range_only = Mutation::new(table(), pk(), None, vec![], 500, None);
range_only.range_tombstones.push(RangeTombstone {
start: ClusteringBound::Inclusive(ck(1)),
end: ClusteringBound::Inclusive(ck(3)),
deletion_time: 200,
local_deletion_time: 2_000,
});
let row_tombstone_row = Mutation::new(table(), pk(), Some(ck(1)), vec![], 300, None)
.with_row_tombstone(50, 3_000);
let ttl_row = Mutation::new(
table(),
pk(),
Some(ck(2)),
vec![CellOperation::WriteWithTtl {
column: "v".to_string(),
value: Value::text("x".to_string()),
ttl_seconds: 60,
local_deletion_time: Some(4_000),
}],
400,
Some(60),
);
let complex_deletion_row = Mutation::new(
table(),
pk(),
Some(ck(3)),
vec![CellOperation::ComplexDeletion {
column: "tags".to_string(),
marked_for_delete_at: 9_000_000,
local_deletion_time: 5_000,
}],
600,
None,
);
let complex_element_row = Mutation::new(
table(),
pk(),
Some(ck(4)),
vec![
CellOperation::WriteComplexElement {
column: "tags".to_string(),
cell_path: b"a".to_vec(),
value: None,
timestamp_micros: 700,
ttl_seconds: None,
local_deletion_time: Some(6_000),
is_deleted: true,
},
CellOperation::WriteComplexElement {
column: "tags".to_string(),
cell_path: b"b".to_vec(),
value: Some(Value::text("b".to_string())),
timestamp_micros: 750,
ttl_seconds: None,
local_deletion_time: None,
is_deleted: false,
},
],
700,
None,
);
let live_write_row = Mutation::new(
table(),
pk(),
Some(ck(5)),
vec![CellOperation::Write {
column: "v".to_string(),
value: Value::text("live".to_string()),
}],
800,
None,
);
let null_write_row = Mutation::new(
table(),
pk(),
Some(ck(6)),
vec![CellOperation::Write {
column: "v".to_string(),
value: Value::Null,
}],
900,
None,
);
vec![
partition_only,
range_only,
row_tombstone_row,
ttl_row,
complex_deletion_row,
complex_element_row,
live_write_row,
null_write_row,
]
}
#[derive(Debug, PartialEq)]
struct FoldSnapshot {
min_timestamp: i64,
max_timestamp: i64,
min_local_deletion_time: i32,
max_local_deletion_time: i32,
min_ttl: i32,
max_ttl: i32,
has_partition_level_deletions: bool,
}
impl From<&StatisticsMetadata> for FoldSnapshot {
fn from(s: &StatisticsMetadata) -> Self {
Self {
min_timestamp: s.min_timestamp,
max_timestamp: s.max_timestamp,
min_local_deletion_time: s.min_local_deletion_time,
max_local_deletion_time: s.max_local_deletion_time,
min_ttl: s.min_ttl,
max_ttl: s.max_ttl,
has_partition_level_deletions: s.has_partition_level_deletions,
}
}
}
#[test]
fn split_and_merge_matches_direct_fold_for_every_mutation_kind() {
let mutations = representative_mutations();
let mut direct = StatisticsMetadata::new();
for m in &mutations {
fold_mutation_stats(&mut direct, m);
}
let mut part_a = StatisticsMetadata::new();
let mut part_b = StatisticsMetadata::new();
let mut part_c = StatisticsMetadata::new();
for (i, m) in mutations.iter().enumerate() {
match i % 3 {
0 => fold_mutation_stats(&mut part_a, m),
1 => fold_mutation_stats(&mut part_b, m),
_ => fold_mutation_stats(&mut part_c, m),
}
}
let mut merged = StatisticsMetadata::new();
merge_stats_fold(&mut merged, &part_a);
merge_stats_fold(&mut merged, &part_b);
merge_stats_fold(&mut merged, &part_c);
assert_eq!(
FoldSnapshot::from(&direct),
FoldSnapshot::from(&merged),
"splitting the fold across partition-scoped sub-accumulators and \
merging must reproduce write_partition's direct fold exactly"
);
assert_eq!(
direct.max_local_deletion_time,
i32::MAX,
"live_write_row must have lifted the live-LDT sentinel"
);
assert!(
direct.has_partition_level_deletions,
"partition_only must have set the partition-level-deletion flag"
);
}
#[test]
fn merging_an_empty_sub_fold_is_a_no_op() {
let mut into = StatisticsMetadata::new();
fold_mutation_stats(
&mut into,
&Mutation::new(
table(),
pk(),
Some(ck(1)),
vec![CellOperation::WriteWithTtl {
column: "v".to_string(),
value: Value::text("x".to_string()),
ttl_seconds: 60,
local_deletion_time: Some(1_000),
}],
100,
Some(60),
),
);
let before = FoldSnapshot::from(&into);
let empty = StatisticsMetadata::new();
merge_stats_fold(&mut into, &empty);
assert_eq!(
FoldSnapshot::from(&into),
before,
"merging a never-folded (default) StatisticsMetadata must not \
change min_local_deletion_time to i32::MIN or max_ttl to i32::MAX"
);
}
}