use std::collections::{HashMap, HashSet};
use std::path::PathBuf;
use std::sync::Arc;
#[cfg(not(feature = "tombstones"))]
use tokio::sync::{mpsc, oneshot};
use super::reader::parsing::row_decoder::now_clock::now_epoch_secs;
use super::reader::parsing::row_decoder::partition_shadow::{
merged_row_shadowed_by_partition, PartitionShadow,
};
#[cfg(not(feature = "tombstones"))]
use super::stream_merge_probe;
use super::{reader, scan_merge};
use crate::storage::write_engine::merge::{CellData, KWayMerger, MergeEntry, MergeStep, RowData};
use crate::types::{CellWriteMetadata, TableId as CqlTableId};
use crate::{Result, RowCells, RowKey, ScanRow, Value};
type MergedMetaRow = (Vec<u8>, ScanRow, Vec<(String, i64)>);
struct ReadShadow {
now_secs: i64,
key_columns: HashSet<String>,
}
impl ReadShadow {
fn new(schema: &crate::schema::TableSchema, now_secs: i64) -> Self {
let mut key_columns =
HashSet::with_capacity(schema.partition_keys.len() + schema.clustering_keys.len());
for k in &schema.partition_keys {
key_columns.insert(k.name.clone());
}
for k in &schema.clustering_keys {
key_columns.insert(k.name.clone());
}
Self {
now_secs,
key_columns,
}
}
fn filter_live(&self, cover: Option<i64>, cells: Vec<CellData>) -> Option<Vec<CellData>> {
let now = self.now_secs;
let mut kept = Vec::with_capacity(cells.len());
let mut max_data_ts: Option<i64> = None;
for cell in cells {
if matches!(cell.value, Value::Tombstone(_)) {
continue;
}
if self.key_columns.contains(&cell.column) {
kept.push(cell);
continue;
}
let eff_ts = Some(cell.timestamp);
let eff_exp = cell_expiry_secs(&cell);
let dropped = PartitionShadow::cell_shadowed_or_expired(cover, now, eff_ts, eff_exp);
if let Some(t) = eff_ts {
max_data_ts = Some(max_data_ts.map_or(t, |m| m.max(t)));
}
if !dropped {
kept.push(cell);
}
}
if merged_row_shadowed_by_partition(cover, max_data_ts) {
return None;
}
Some(kept)
}
}
fn cell_expiry_secs(cell: &CellData) -> Option<i64> {
cell.ttl?;
cell.local_deletion_time.map(|s| (s as u32) as i64)
}
fn partition_cover(rows: &[MergeEntry]) -> Option<i64> {
rows.iter()
.find_map(|e| e.partition_deletion.map(|(mfda, _ldt)| mfda))
}
fn partition_live_rows(
row_key: &RowKey,
rows: Vec<MergeEntry>,
shadow: &ReadShadow,
) -> Vec<(RowKey, ScanRow)> {
let cover = partition_cover(&rows);
let mut out = Vec::new();
for entry in rows {
match entry.row_data {
RowData::Live { cells } => {
let Some(surviving) = shadow.filter_live(cover, cells) else {
continue;
};
let row_cells: RowCells = surviving
.into_iter()
.map(|c| (Arc::from(c.column.as_str()), c.value))
.collect();
if !row_cells.is_empty() {
out.push((row_key.clone(), ScanRow::Row(row_cells)));
}
}
RowData::Tombstone { .. } => {}
}
}
out
}
pub(super) async fn merge_generations_for_read(
reader_list: &[Arc<reader::SSTableReader>],
schema: &crate::schema::TableSchema,
start_key: Option<&RowKey>,
end_key: Option<&RowKey>,
limit: Option<usize>,
target_key: Option<&RowKey>,
) -> Result<Vec<(RowKey, ScanRow)>> {
let admission = reader::scan_stream_windowed::scan_admission::admit().await;
let start_key = start_key.cloned();
let end_key = end_key.cloned();
let target_key = target_key.map(|k| k.as_bytes().to_vec());
let paths = ordered_generation_paths(reader_list);
let schema = schema.clone();
let mut merged = tokio::task::spawn_blocking(move || -> Result<Vec<(RowKey, ScanRow)>> {
let _admission = admission; let shadow = ReadShadow::new(&schema, now_epoch_secs());
let mut merger = KWayMerger::new(paths, &schema)?;
let mut out = Vec::new();
while let MergeStep::Partition { key, rows } = merger.step()? {
let row_key = RowKey::new(key.key.clone());
if let Some(ref target) = target_key {
if row_key.as_bytes() != target.as_slice() {
continue;
}
out.extend(partition_live_rows(&row_key, rows, &shadow));
break;
}
if let Some(ref start) = start_key {
if &row_key < start {
continue;
}
}
if let Some(ref end) = end_key {
if &row_key > end {
continue;
}
}
out.extend(partition_live_rows(&row_key, rows, &shadow));
}
Ok(out)
})
.await
.map_err(|e| crate::Error::Storage(format!("cross-generation read merge task: {e}")))??;
scan_merge::sort_by_token_order(&mut merged, limit, |(k, _)| k);
Ok(merged)
}
#[cfg(all(feature = "write-support", not(feature = "tombstones")))]
pub(super) async fn seek_merge_generations_for_read(
candidates: &[Arc<reader::SSTableReader>],
schema: &crate::schema::TableSchema,
target_key: &RowKey,
) -> Result<Vec<(RowKey, ScanRow)>> {
use crate::storage::scan_cancel::ScanCancel;
use crate::storage::write_engine::merge::build_single_partition_merger_from_readers;
let admission = reader::scan_stream_windowed::scan_admission::admit().await;
let mut ordered: Vec<Arc<reader::SSTableReader>> = candidates.to_vec();
ordered.sort_by(|a, b| b.generation.cmp(&a.generation));
let schema = schema.clone();
let target_bytes = target_key.as_bytes().to_vec();
let mut merged = tokio::task::spawn_blocking(move || -> Result<Vec<(RowKey, ScanRow)>> {
let _admission = admission; let shadow = ReadShadow::new(&schema, now_epoch_secs());
let keys = [target_bytes.clone()];
let Some(mut merger) =
build_single_partition_merger_from_readers(ordered, &keys, &schema, ScanCancel::new())?
else {
return Ok(Vec::new()); };
let mut out = Vec::new();
while let MergeStep::Partition { key, rows } = merger.step()? {
let row_key = RowKey::new(key.key.clone());
if row_key.as_bytes() != target_bytes.as_slice() {
continue;
}
out.extend(partition_live_rows(&row_key, rows, &shadow));
break;
}
Ok(out)
})
.await
.map_err(|e| {
crate::Error::Storage(format!("seeking cross-generation read merge task: {e}"))
})??;
scan_merge::sort_by_token_order(&mut merged, None, |(k, _)| k);
Ok(merged)
}
pub(super) async fn merge_generations_for_read_with_metadata(
reader_list: &[Arc<reader::SSTableReader>],
schema: &crate::schema::TableSchema,
start_key: Option<&RowKey>,
end_key: Option<&RowKey>,
limit: Option<usize>,
target_key: Option<&RowKey>,
) -> Result<Vec<(RowKey, ScanRow, HashMap<String, CellWriteMetadata>)>> {
let admission = reader::scan_stream_windowed::scan_admission::admit().await;
let owned_start = start_key.cloned();
let owned_end = end_key.cloned();
let target_bytes = target_key.map(|k| k.as_bytes().to_vec());
let table_id = CqlTableId::from(format!("{}.{}", schema.keyspace, schema.table).as_str());
let mut ttl_lookup: HashMap<(Vec<u8>, String), CellWriteMetadata> = HashMap::new();
for reader in reader_list {
let (ttl_start, ttl_end): (Option<&RowKey>, Option<&RowKey>) = match target_key {
Some(t) => (Some(t), Some(t)),
None => (owned_start.as_ref(), owned_end.as_ref()),
};
let per_reader = reader
.scan_with_cell_metadata(&table_id, ttl_start, ttl_end, None, Some(schema))
.await?;
for (row_key, _value, meta) in per_reader {
for (column, cell_meta) in meta {
ttl_lookup
.entry((row_key.as_bytes().to_vec(), column))
.and_modify(|existing| {
if cell_meta.write_timestamp_micros > existing.write_timestamp_micros {
*existing = cell_meta.clone();
}
})
.or_insert(cell_meta);
}
}
}
let paths = ordered_generation_paths(reader_list);
let merge_schema = schema.clone();
let start_key = owned_start;
let end_key = owned_end;
let target_for_merge = target_bytes;
let merged_rows = tokio::task::spawn_blocking(move || -> Result<Vec<MergedMetaRow>> {
let _admission = admission; let shadow = ReadShadow::new(&merge_schema, now_epoch_secs());
let mut merger = KWayMerger::new(paths, &merge_schema)?;
let mut out = Vec::new();
while let MergeStep::Partition { key, rows } = merger.step()? {
let row_key = RowKey::new(key.key.clone());
if let Some(ref target) = target_for_merge {
if row_key.as_bytes() != target.as_slice() {
continue;
}
push_metadata_rows(&key.key, rows, &mut out, &shadow);
break;
}
if let Some(ref start) = start_key {
if &row_key < start {
continue;
}
}
if let Some(ref end) = end_key {
if &row_key > end {
continue;
}
}
push_metadata_rows(&key.key, rows, &mut out, &shadow);
}
Ok(out)
})
.await
.map_err(|e| crate::Error::Storage(format!("cross-generation metadata merge task: {e}")))??;
let mut results: Vec<(RowKey, ScanRow, HashMap<String, CellWriteMetadata>)> =
Vec::with_capacity(merged_rows.len());
for (key_bytes, value, timestamps) in merged_rows {
let mut meta_map: HashMap<String, CellWriteMetadata> =
HashMap::with_capacity(timestamps.len());
for (column, write_ts) in timestamps {
let expiration = ttl_lookup
.get(&(key_bytes.clone(), column.clone()))
.filter(|m| m.write_timestamp_micros == write_ts)
.and_then(|m| m.expiration.clone());
meta_map.insert(
column,
CellWriteMetadata {
write_timestamp_micros: write_ts,
expiration,
},
);
}
results.push((RowKey::new(key_bytes), value, meta_map));
}
scan_merge::sort_by_token_order(&mut results, limit, |(k, _, _)| k);
Ok(results)
}
fn push_metadata_rows(
key_bytes: &[u8],
rows: Vec<MergeEntry>,
out: &mut Vec<MergedMetaRow>,
shadow: &ReadShadow,
) {
let cover = partition_cover(&rows);
for entry in rows {
if let RowData::Live { cells } = entry.row_data {
let Some(surviving) = shadow.filter_live(cover, cells) else {
continue;
};
let mut row_cells: RowCells = Vec::with_capacity(surviving.len());
let mut timestamps: Vec<(String, i64)> = Vec::with_capacity(surviving.len());
for c in surviving {
timestamps.push((c.column.clone(), c.timestamp));
row_cells.push((Arc::from(c.column.as_str()), c.value));
}
if !row_cells.is_empty() {
out.push((key_bytes.to_vec(), ScanRow::Row(row_cells), timestamps));
}
}
}
}
#[cfg(not(feature = "tombstones"))]
pub(super) async fn stream_generations_for_read(
reader_list: &[Arc<reader::SSTableReader>],
schema: &crate::schema::TableSchema,
start_key: Option<&RowKey>,
end_key: Option<&RowKey>,
buffer_size: usize,
) -> Result<mpsc::Receiver<Result<(RowKey, ScanRow)>>> {
let start_key = start_key.cloned();
let end_key = end_key.cloned();
let paths = ordered_generation_paths(reader_list);
let schema = schema.clone();
let (ready_tx, ready_rx) = oneshot::channel::<Result<()>>();
let (out_tx, out_rx) = mpsc::channel::<Result<(RowKey, ScanRow)>>(buffer_size.max(1));
tokio::task::spawn_blocking(move || {
let shadow = ReadShadow::new(&schema, now_epoch_secs());
let mut merger = match KWayMerger::new(paths, &schema) {
Ok(m) => {
if ready_tx.send(Ok(())).is_err() {
return;
}
m
}
Err(e) => {
let _ = ready_tx.send(Err(e));
return;
}
};
loop {
let step = match merger.step() {
Ok(s) => s,
Err(e) => {
let _ = out_tx.blocking_send(Err(e));
return;
}
};
let (key, rows) = match step {
MergeStep::Partition { key, rows } => (key, rows),
MergeStep::Complete => return,
};
let row_key = RowKey::new(key.key.clone());
if let Some(ref start) = start_key {
if &row_key < start {
continue;
}
}
if let Some(ref end) = end_key {
if &row_key > end {
continue;
}
}
let live = partition_live_rows(&row_key, rows, &shadow);
stream_merge_probe::record_resident(live.len() as u64);
for entry in live {
if out_tx.blocking_send(Ok(entry)).is_err() {
return; }
}
}
});
match ready_rx.await {
Ok(Ok(())) => Ok(out_rx),
Ok(Err(e)) => Err(e),
Err(_) => Err(crate::Error::Storage(
"cross-generation streaming merge task ended before signalling readiness".to_string(),
)),
}
}
fn ordered_generation_paths(reader_list: &[Arc<reader::SSTableReader>]) -> Vec<PathBuf> {
let mut ordered: Vec<&Arc<reader::SSTableReader>> = reader_list.iter().collect();
ordered.sort_by(|a, b| b.generation.cmp(&a.generation));
ordered.iter().map(|r| r.file_path()).collect()
}
#[cfg(test)]
mod tests {
use super::*;
fn shadow(now_secs: i64) -> ReadShadow {
ReadShadow {
now_secs,
key_columns: HashSet::new(),
}
}
fn shadow_with_keys(now_secs: i64, keys: &[&str]) -> ReadShadow {
ReadShadow {
now_secs,
key_columns: keys.iter().map(|k| k.to_string()).collect(),
}
}
fn live_cell(column: &str, ts: i64) -> CellData {
CellData::new(column.to_string(), Value::Integer(1), ts)
}
fn expiring_cell(column: &str, ts: i64, ldt: i32) -> CellData {
let mut c = CellData::new(column.to_string(), Value::Integer(1), ts);
c.ttl = Some(60);
c.local_deletion_time = Some(ldt);
c
}
#[test]
fn filter_live_drops_ttl_expired_cell_keeps_live() {
let now = 2_000_000i64;
let cells = vec![
live_cell("name", 100),
expiring_cell("token", 100, 1_000), ];
let kept = shadow(now).filter_live(None, cells).expect("row visible");
let names: Vec<&str> = kept.iter().map(|c| c.column.as_str()).collect();
assert_eq!(names, vec!["name"], "expired `token` must be dropped");
}
#[test]
fn filter_live_drops_cell_tombstone() {
let mut tomb = CellData::new("gone".to_string(), Value::Integer(0), 100);
tomb.value = Value::Tombstone(Box::new(crate::types::TombstoneInfo {
deletion_time: 100,
tombstone_type: crate::types::TombstoneType::CellTombstone,
local_deletion_time: 0,
ttl: None,
range_start: None,
range_end: None,
}));
let kept = shadow(0)
.filter_live(None, vec![live_cell("keep", 100), tomb])
.expect("row visible");
assert_eq!(kept.len(), 1);
assert_eq!(kept[0].column, "keep");
}
#[test]
fn filter_live_partition_cover_hides_fully_shadowed_row() {
let cover = Some(2_000i64);
let hidden =
shadow(0).filter_live(cover, vec![live_cell("a", 1_000), live_cell("b", 2_000)]);
assert!(hidden.is_none(), "fully-shadowed row must be hidden");
let kept = shadow(0)
.filter_live(cover, vec![live_cell("a", 1_000), live_cell("b", 3_000)])
.expect("row visible");
let names: Vec<&str> = kept.iter().map(|c| c.column.as_str()).collect();
assert_eq!(names, vec!["b"], "older `a` shadowed, newer `b` survives");
}
#[test]
fn filter_live_keeps_clustering_key_under_partition_cover() {
let cover = Some(2_000i64);
let cells = vec![
live_cell("ck", 1_000),
live_cell("old", 1_000),
live_cell("data", 3_000),
];
let kept = shadow_with_keys(0, &["ck"])
.filter_live(cover, cells)
.expect("row survives via newer `data` cell");
let names: Vec<&str> = kept.iter().map(|c| c.column.as_str()).collect();
assert!(
names.contains(&"ck"),
"clustering-key pseudo-cell must be retained, got {names:?}"
);
assert!(
names.contains(&"data"),
"newer resurrecting data cell must survive, got {names:?}"
);
assert!(
!names.contains(&"old"),
"stale data cell shadowed by the partition tombstone must be dropped, got {names:?}"
);
}
#[test]
fn filter_live_key_cell_does_not_resurrect_shadowed_row() {
let cover = Some(2_000i64);
let cells = vec![live_cell("ck", 1_000), live_cell("old", 1_000)];
let hidden = shadow_with_keys(0, &["ck"]).filter_live(cover, cells);
assert!(
hidden.is_none(),
"a key pseudo-cell alone must not keep a fully-shadowed row visible"
);
}
#[test]
fn cell_expiry_secs_reinterprets_post_2038_unsigned() {
let future: i64 = 2_200_000_000;
let stored = future as u32 as i32; let c = expiring_cell("token", 100, stored);
assert_eq!(
cell_expiry_secs(&c),
Some(future),
"post-2038 LDT must reinterpret unsigned to a future expiry"
);
assert_eq!(cell_expiry_secs(&live_cell("x", 100)), None);
}
}