use crate::{
types::{RowKey, ScanRow, TableId, TombstoneInfo, TombstoneType, Value},
Result,
};
use std::collections::HashMap;
#[derive(Debug, Clone)]
pub struct EntryMetadata {
pub write_time: i64,
pub generation: u64,
pub ttl: Option<i64>,
}
#[derive(Debug, Clone)]
pub struct GenerationValue {
pub value: ScanRow,
pub metadata: EntryMetadata,
}
impl GenerationValue {
fn tombstone_info(&self) -> Option<&TombstoneInfo> {
match &self.value {
ScanRow::Marker(Value::Tombstone(info)) => Some(info),
_ => None,
}
}
}
pub struct TombstoneMerger {
current_time: i64,
}
fn system_time_to_micros(t: std::time::SystemTime) -> i64 {
t.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_micros() as i64)
.unwrap_or(0)
}
impl TombstoneMerger {
pub fn new() -> Self {
Self {
current_time: system_time_to_micros(std::time::SystemTime::now()),
}
}
pub fn with_time(current_time: i64) -> Self {
Self { current_time }
}
pub fn merge_generations(&self, values: Vec<GenerationValue>) -> Result<Option<ScanRow>> {
if values.is_empty() {
return Ok(None);
}
let mut sorted_values = values;
sorted_values.sort_by(|a, b| {
b.metadata
.generation
.cmp(&a.metadata.generation)
.then_with(|| b.metadata.write_time.cmp(&a.metadata.write_time))
});
let mut latest_tombstone_time: Option<i64> = None;
let mut _latest_tombstone_type: Option<TombstoneType> = None;
for gen_value in &sorted_values {
if let Some(tombstone_info) = gen_value.tombstone_info() {
if !self.is_tombstone_expired(tombstone_info) {
if latest_tombstone_time.is_none_or(|t| tombstone_info.deletion_time > t) {
latest_tombstone_time = Some(tombstone_info.deletion_time);
_latest_tombstone_type = Some(tombstone_info.tombstone_type);
}
}
}
}
for gen_value in sorted_values {
if let Some(tombstone_info) = gen_value.tombstone_info() {
if self.is_tombstone_expired(tombstone_info) {
continue;
}
if let Some(latest_time) = latest_tombstone_time {
if tombstone_info.deletion_time == latest_time {
return Ok(None);
}
}
} else {
if let Some(tombstone_time) = latest_tombstone_time {
if gen_value.metadata.write_time <= tombstone_time {
continue;
}
}
if self.is_value_expired(&gen_value.metadata) {
let expiration_time =
gen_value.metadata.write_time + gen_value.metadata.ttl.unwrap_or(0);
let ttl_tombstone =
Value::ttl_tombstone(expiration_time, gen_value.metadata.ttl.unwrap_or(0));
return Ok(Some(ScanRow::Marker(ttl_tombstone)));
}
return Ok(Some(gen_value.value));
}
}
Ok(None)
}
pub fn merge_row_entries(
&self,
_table_id: &TableId,
_row_key: &RowKey,
entries: Vec<GenerationValue>,
) -> Result<Option<ScanRow>> {
let mut row_tombstone_time: Option<i64> = None;
let mut cell_values = Vec::new();
for entry in entries {
match entry.tombstone_info() {
Some(info) if info.tombstone_type == TombstoneType::RowTombstone => {
if !self.is_tombstone_expired(info) {
if let Some(existing_time) = row_tombstone_time {
if info.deletion_time > existing_time {
row_tombstone_time = Some(info.deletion_time);
}
} else {
row_tombstone_time = Some(info.deletion_time);
}
}
}
_ => {
cell_values.push(entry);
}
}
}
if let Some(tombstone_time) = row_tombstone_time {
cell_values.retain(|entry| entry.metadata.write_time > tombstone_time);
if cell_values.is_empty() {
return Ok(None);
}
}
self.merge_generations(cell_values)
}
fn is_tombstone_expired(&self, tombstone: &TombstoneInfo) -> bool {
if let Some(ttl) = tombstone.ttl {
self.current_time > tombstone.deletion_time + ttl
} else {
false
}
}
fn is_value_expired(&self, metadata: &EntryMetadata) -> bool {
if let Some(ttl) = metadata.ttl {
self.current_time > metadata.write_time + ttl
} else {
false
}
}
pub fn resolve_conflict(&self, values: Vec<GenerationValue>) -> Result<Option<ScanRow>> {
if values.is_empty() {
return Ok(None);
}
let latest = values.into_iter().max_by_key(|v| v.metadata.write_time);
match latest {
Some(gen_value) => {
if self.is_value_expired(&gen_value.metadata) {
Ok(None)
} else if gen_value.tombstone_info().is_some() {
Ok(None)
} else {
Ok(Some(gen_value.value))
}
}
None => Ok(None),
}
}
pub fn merge_cell_tombstones(
&self,
column_values: HashMap<String, Vec<GenerationValue>>,
) -> Result<HashMap<String, Option<ScanRow>>> {
let mut result = HashMap::new();
for (column_name, values) in column_values {
let merged_value = self.merge_generations(values)?;
result.insert(column_name, merged_value);
}
Ok(result)
}
pub fn batch_merge_with_tombstones(
&self,
entries: Vec<(RowKey, Vec<GenerationValue>)>,
batch_size: usize,
) -> Result<Vec<(RowKey, Option<ScanRow>)>> {
let mut results = Vec::with_capacity(entries.len());
for batch in entries.chunks(batch_size) {
for (key, values) in batch {
let merged_value = self.merge_generations(values.clone())?;
results.push((key.clone(), merged_value));
}
}
Ok(results)
}
pub fn identify_garbage_collectible_tombstones(
&self,
tombstones: Vec<GenerationValue>,
gc_grace_seconds: i64,
) -> Result<Vec<GenerationValue>> {
let mut collectible = Vec::new();
let gc_grace_micros = gc_grace_seconds * 1_000_000;
for tombstone_entry in tombstones {
if let Some(tombstone_info) = tombstone_entry.tombstone_info() {
let tombstone_age = self.current_time - tombstone_info.deletion_time;
if tombstone_age > gc_grace_micros {
collectible.push(tombstone_entry);
}
}
}
Ok(collectible)
}
pub fn merge_collection_with_tombstones(
&self,
collection_entries: Vec<GenerationValue>,
) -> Result<Option<ScanRow>> {
let mut sorted_entries = collection_entries;
sorted_entries.sort_by(|a, b| b.metadata.write_time.cmp(&a.metadata.write_time));
for entry in sorted_entries {
if let Some(tombstone_info) = entry.tombstone_info() {
if self.is_tombstone_expired(tombstone_info) {
continue;
}
return Ok(None);
}
if self.is_value_expired(&entry.metadata) {
let ttl_tombstone = Value::ttl_tombstone(
entry.metadata.write_time + entry.metadata.ttl.unwrap_or(0),
entry.metadata.ttl.unwrap_or(0),
);
return Ok(Some(ScanRow::Marker(ttl_tombstone)));
}
return Ok(Some(entry.value));
}
Ok(None)
}
pub fn fast_tombstone_check(&self, value: &Value, _write_time: i64) -> bool {
match value {
Value::Tombstone(info) => {
if info.ttl.is_none() {
true } else {
!self.is_tombstone_expired(info)
}
}
_ => {
false
}
}
}
}
impl Default for TombstoneMerger {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_basic_tombstone_merge() -> Result<()> {
let merger = TombstoneMerger::with_time(5000);
let values = vec![
GenerationValue {
value: ScanRow::Marker(Value::Integer(42)),
metadata: EntryMetadata {
write_time: 1000,
generation: 1,
ttl: None,
},
},
GenerationValue {
value: ScanRow::Marker(Value::row_tombstone(2000)),
metadata: EntryMetadata {
write_time: 2000,
generation: 2,
ttl: None,
},
},
];
let result = merger.merge_generations(values)?;
assert!(result.is_none());
Ok(())
}
#[test]
fn production_new_constructor_merges_without_panicking() -> Result<()> {
let merger = TombstoneMerger::new();
let values = vec![
GenerationValue {
value: ScanRow::Marker(Value::Integer(10)),
metadata: EntryMetadata {
write_time: 1_000,
generation: 1,
ttl: None,
},
},
GenerationValue {
value: ScanRow::Marker(Value::Integer(20)),
metadata: EntryMetadata {
write_time: 2_000,
generation: 2,
ttl: None,
},
},
];
let result = merger.merge_generations(values)?;
assert_eq!(
result,
Some(ScanRow::Marker(Value::Integer(20))),
"the production new() constructor must resolve to the newest live value"
);
Ok(())
}
#[test]
fn system_time_to_micros_falls_back_to_zero_before_epoch() {
let pre_epoch = std::time::UNIX_EPOCH - std::time::Duration::from_secs(1);
assert_eq!(
system_time_to_micros(pre_epoch),
0,
"a SystemTime before UNIX_EPOCH must fall back to 0 micros"
);
let post_epoch = std::time::UNIX_EPOCH + std::time::Duration::from_secs(1);
assert_eq!(system_time_to_micros(post_epoch), 1_000_000);
}
#[test]
fn test_ttl_expiration() -> Result<()> {
let merger = TombstoneMerger::with_time(5000);
let values = vec![GenerationValue {
value: ScanRow::Marker(Value::Integer(42)),
metadata: EntryMetadata {
write_time: 1000,
generation: 1,
ttl: Some(1000), },
}];
let result = merger.merge_generations(values)?;
assert!(matches!(result, Some(ScanRow::Marker(v)) if v.is_tombstone()));
Ok(())
}
#[test]
fn test_row_level_tombstone() {
let merger = TombstoneMerger::with_time(5000);
let table_id = TableId::from("test_table");
let row_key = RowKey::from("test_key");
let entries = vec![
GenerationValue {
value: ScanRow::Marker(Value::Integer(42)),
metadata: EntryMetadata {
write_time: 1000,
generation: 1,
ttl: None,
},
},
GenerationValue {
value: ScanRow::Marker(Value::row_tombstone(2000)),
metadata: EntryMetadata {
write_time: 2000,
generation: 2,
ttl: None,
},
},
GenerationValue {
value: ScanRow::Marker(Value::text("newer_value".to_string())),
metadata: EntryMetadata {
write_time: 3000,
generation: 3,
ttl: None,
},
},
];
let result = merger
.merge_row_entries(&table_id, &row_key, entries)
.unwrap();
assert_eq!(
result,
Some(ScanRow::Marker(Value::text("newer_value".to_string())))
);
}
#[test]
fn test_enhanced_multi_generation_merge() -> Result<()> {
let merger = TombstoneMerger::with_time(10000);
let values = vec![
GenerationValue {
value: ScanRow::Marker(Value::Integer(10)),
metadata: EntryMetadata {
write_time: 1000,
generation: 1,
ttl: None,
},
},
GenerationValue {
value: ScanRow::Marker(Value::cell_tombstone(2000)),
metadata: EntryMetadata {
write_time: 2000,
generation: 2,
ttl: None,
},
},
GenerationValue {
value: ScanRow::Marker(Value::Integer(20)),
metadata: EntryMetadata {
write_time: 1500,
generation: 1,
ttl: None,
},
},
GenerationValue {
value: ScanRow::Marker(Value::Integer(30)),
metadata: EntryMetadata {
write_time: 3000,
generation: 3,
ttl: None,
},
},
];
let result = merger.merge_generations(values)?;
assert_eq!(result, Some(ScanRow::Marker(Value::Integer(30))));
Ok(())
}
#[test]
fn test_batch_processing_performance() -> Result<()> {
let merger = TombstoneMerger::with_time(5000);
let mut entries = Vec::new();
for i in 0..10000 {
let key = RowKey::from(format!("key_{}", i));
let values = vec![GenerationValue {
value: ScanRow::Marker(Value::Integer(i)),
metadata: EntryMetadata {
write_time: 1000 + i as i64,
generation: 1,
ttl: None,
},
}];
entries.push((key, values));
}
let start = std::time::Instant::now();
let result = merger.batch_merge_with_tombstones(entries, 1000)?;
let duration = start.elapsed();
assert_eq!(result.len(), 10000);
assert!(duration.as_millis() < 1000);
Ok(())
}
#[test]
fn test_garbage_collection_identification() {
let merger = TombstoneMerger::with_time(10_000_000);
let tombstones = vec![
GenerationValue {
value: ScanRow::Marker(Value::row_tombstone(1_000_000)), metadata: EntryMetadata {
write_time: 1_000_000,
generation: 1,
ttl: None,
},
},
GenerationValue {
value: ScanRow::Marker(Value::cell_tombstone(8_000_000)), metadata: EntryMetadata {
write_time: 8_000_000,
generation: 2,
ttl: None,
},
},
];
let collectible = merger
.identify_garbage_collectible_tombstones(tombstones, 3)
.unwrap();
assert_eq!(collectible.len(), 1);
assert_eq!(collectible[0].metadata.write_time, 1_000_000);
}
#[test]
fn test_collection_tombstone_handling() {
let merger = TombstoneMerger::with_time(5000);
let list_row = ScanRow::Marker(Value::List(vec![
Value::Integer(1),
Value::Integer(2),
Value::cell_tombstone(2000), Value::Integer(3),
]));
let collection_entries = vec![GenerationValue {
value: list_row.clone(),
metadata: EntryMetadata {
write_time: 3000,
generation: 1,
ttl: None,
},
}];
let result = merger
.merge_collection_with_tombstones(collection_entries)
.unwrap();
assert_eq!(result, Some(list_row));
}
#[test]
fn expired_ttl_tombstone_does_not_delete_older_live_value() -> Result<()> {
let merger = TombstoneMerger::with_time(5000);
let live = ScanRow::Row(vec![(std::sync::Arc::from("v"), Value::Integer(7))]);
let entries = vec![
GenerationValue {
value: ScanRow::Marker(Value::ttl_tombstone(1000, 1000)),
metadata: EntryMetadata {
write_time: 3000,
generation: 2,
ttl: Some(1000),
},
},
GenerationValue {
value: live.clone(),
metadata: EntryMetadata {
write_time: 1000,
generation: 1,
ttl: None,
},
},
];
let result = merger.merge_collection_with_tombstones(entries)?;
assert_eq!(
result,
Some(live),
"expired TTL tombstone must be skipped; older live value survives"
);
Ok(())
}
#[test]
fn live_tombstone_deletes_row_on_collection_merge() -> Result<()> {
let merger = TombstoneMerger::with_time(5000);
let entries = vec![
GenerationValue {
value: ScanRow::Marker(Value::row_tombstone(3000)),
metadata: EntryMetadata {
write_time: 3000,
generation: 2,
ttl: None,
},
},
GenerationValue {
value: ScanRow::Row(vec![(std::sync::Arc::from("v"), Value::Integer(7))]),
metadata: EntryMetadata {
write_time: 1000,
generation: 1,
ttl: None,
},
},
];
assert_eq!(merger.merge_collection_with_tombstones(entries)?, None);
Ok(())
}
#[test]
fn test_fast_tombstone_check_performance() {
let merger = TombstoneMerger::with_time(5000);
let non_tombstone = Value::Integer(42);
let tombstone = Value::row_tombstone(3000);
let ttl_tombstone = Value::ttl_tombstone(2000, 4000);
let start = std::time::Instant::now();
for _ in 0..100000 {
assert!(!merger.fast_tombstone_check(&non_tombstone, 3000));
assert!(merger.fast_tombstone_check(&tombstone, 3000));
assert!(merger.fast_tombstone_check(&ttl_tombstone, 3000));
}
let duration = start.elapsed();
assert!(duration.as_millis() < 100);
}
#[test]
fn test_lww_merge_winner_has_newer_write_time() -> Result<()> {
let merger = TombstoneMerger::with_time(99_999);
let values = vec![
GenerationValue {
value: ScanRow::Marker(Value::Integer(10)), metadata: EntryMetadata {
write_time: 1_000,
generation: 1,
ttl: None,
},
},
GenerationValue {
value: ScanRow::Marker(Value::Integer(30)), metadata: EntryMetadata {
write_time: 3_000,
generation: 2,
ttl: None,
},
},
];
let result = merger.merge_generations(values)?;
assert_eq!(
result,
Some(ScanRow::Marker(Value::Integer(30))),
"LWW merge must pick the newer value (write_time=3000)"
);
Ok(())
}
#[test]
fn test_lww_merge_expired_ttl_returns_tombstone() -> Result<()> {
let now = 100_000_i64;
let merger = TombstoneMerger::with_time(now);
let values = vec![GenerationValue {
value: ScanRow::Marker(Value::text("expiring".to_string())),
metadata: EntryMetadata {
write_time: 5_000,
generation: 1,
ttl: Some(1_000), },
}];
let result = merger.merge_generations(values)?;
assert!(
matches!(result, Some(ScanRow::Marker(v)) if v.is_tombstone()),
"result for expired TTL cell must be a tombstone"
);
Ok(())
}
#[test]
fn test_lww_merge_live_value_newer_than_row_tombstone_survives() -> Result<()> {
let merger = TombstoneMerger::with_time(99_999);
let table_id = TableId::from("ks.tbl");
let row_key = RowKey::from("pk1");
let entries = vec![
GenerationValue {
value: ScanRow::Marker(Value::row_tombstone(1_000)),
metadata: EntryMetadata {
write_time: 1_000,
generation: 1,
ttl: None,
},
},
GenerationValue {
value: ScanRow::Marker(Value::Integer(99)),
metadata: EntryMetadata {
write_time: 5_000,
generation: 2,
ttl: None,
},
},
];
let result = merger.merge_row_entries(&table_id, &row_key, entries)?;
assert_eq!(
result,
Some(ScanRow::Marker(Value::Integer(99))),
"a live value written after the row tombstone must survive the merge"
);
Ok(())
}
}