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};
mod merge_cancel;
#[cfg(not(feature = "tombstones"))]
mod merge_stream_setup;
#[cfg(not(feature = "tombstones"))]
pub(super) use merge_stream_setup::MergeStreamSetupError;
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>,
marker_ts: 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;
let mut has_deleted_data_cell = false;
for cell in cells {
if matches!(cell.value, Value::Tombstone(_)) {
has_deleted_data_cell |= !self.key_columns.contains(&cell.column);
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, marker_ts, max_data_ts, has_deleted_data_cell) {
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 {
let marker_ts = entry.row_liveness.marker_timestamp;
match entry.row_data {
RowData::Live { cells } => {
let Some(surviving) = shadow.filter_live(cover, marker_ts, 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 (_cancel_guard, cancel) = merge_cancel::per_call();
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_cancellable(paths, &schema, cancel.clone())?;
let mut out = Vec::new();
loop {
merge_cancel::check(&cancel)?;
let MergeStep::Partition { key, rows } = merger.step()? else {
break;
};
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::write_engine::merge::{
build_single_partition_merger_from_readers, PointAccessRecording,
};
let admission = reader::scan_stream_windowed::scan_admission::admit().await;
let (_cancel_guard, cancel) = merge_cancel::per_call();
let mut ordered: Vec<Arc<reader::SSTableReader>> = candidates.to_vec();
ordered.sort_by_key(|b| std::cmp::Reverse(b.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,
cancel.clone(), PointAccessRecording::CallerRecords,
)?
else {
return Ok(Vec::new()); };
let mut out = Vec::new();
loop {
merge_cancel::check(&cancel)?;
let MergeStep::Partition { key, rows } = merger.step()? else {
break;
};
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 (_cancel_guard, cancel) = merge_cancel::per_call();
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_cancellable(paths, &merge_schema, cancel.clone())?;
let mut out = Vec::new();
loop {
merge_cancel::check(&cancel)?;
let MergeStep::Partition { key, rows } = merger.step()? else {
break;
};
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 {
let marker_ts = entry.row_liveness.marker_timestamp;
if let RowData::Live { cells } = entry.row_data {
let Some(surviving) = shadow.filter_live(cover, marker_ts, 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,
) -> std::result::Result<reader::RowScanStream, MergeStreamSetupError> {
let mut meter = crate::observability::read_metrics::ReadOpMeter::start(None);
let phase_sink = meter.phase_sink();
let start_key = start_key.cloned();
let end_key = end_key.cloned();
let paths = ordered_generation_paths(reader_list);
let schema = schema.clone();
let fault_scope = crate::storage::producer_fault::FaultScope::capture(|| {
paths.first().cloned().unwrap_or_default()
});
let (ready_tx, ready_rx) = oneshot::channel::<Result<()>>();
let (out_tx, out_rx) = mpsc::channel::<Result<(RowKey, ScanRow)>>(buffer_size.max(1));
let task = tokio::task::spawn_blocking(move || {
let _phases = crate::observability::read_phase::install(phase_sink);
fault_scope.checkpoint(crate::storage::producer_fault::ScanTaskSite::CrossGenerationMerge);
let shadow = ReadShadow::new(&schema, now_epoch_secs());
let constructed = match fault_scope.injected_construction_error() {
Some(injected) => Err(injected),
None => crate::storage::write_engine::merge::merger_deferring_opens(paths, &schema),
};
let mut merger = match constructed {
Ok(m) => {
if ready_tx.send(Ok(())).is_err() {
return;
}
m
}
Err(e) => {
let _ = ready_tx.send(Err(e));
return;
}
};
loop {
if out_tx.is_closed() {
return;
}
let step =
match crate::observability::read_phase::timed_merge_excluding_recv_wait(|| {
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(reader::RowScanStream::new_measured_rows(
out_rx, task, meter,
)),
Ok(Err(e)) => {
let err = MergeStreamSetupError::from_construction_failure(e);
if err.fallback_eligible() {
meter.discard();
} else {
meter.finish();
}
Err(err)
}
Err(_) => {
meter.finish();
Err(MergeStreamSetupError::ProducerDied(
merge_stream_setup::dead_merge_producer_error(task.await),
))
}
}
}
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_key(|b| std::cmp::Reverse(b.generation));
ordered.iter().map(|r| r.file_path()).collect()
}
#[cfg(test)]
mod read_shadow_tests;
#[cfg(all(test, not(feature = "tombstones")))]
pub(super) mod multi_gen_fixture;
#[cfg(all(test, not(feature = "tombstones")))]
mod setup_fail_closed_tests;