use super::super::result::{cql_type_to_data_type, ColumnInfo, QueryRow};
use crate::{
parser::complex_types::ComplexTypeParser,
schema::CqlType,
types::{RowKey, ScanRow, Value},
TableId,
};
use std::collections::HashMap;
use std::sync::Arc;
pub(super) fn parse_table_id(table_id: &TableId) -> (Option<String>, String) {
let table_str = table_id.name();
match table_str.rfind('.') {
Some(dot) => (
Some(table_str[..dot].to_string()),
table_str[dot + 1..].to_string(),
),
None => (None, table_str.to_string()),
}
}
pub(super) fn parse_cql_type_str(type_str: &str) -> Option<CqlType> {
let parser = ComplexTypeParser::new();
parser
.parse_type(type_str)
.ok()
.map(|parsed| parsed.cql_type)
}
pub(super) fn column_info_from_type_str(
name: String,
type_str: &str,
position: usize,
table_name: Option<String>,
) -> ColumnInfo {
let cql_type_opt = parse_cql_type_str(type_str);
let data_type = cql_type_opt
.as_ref()
.map(cql_type_to_data_type)
.unwrap_or(crate::types::DataType::Text);
let mut col_info = ColumnInfo {
name,
data_type,
nullable: true,
position,
table_name,
cql_type: None,
};
if let Some(cql_type) = cql_type_opt {
col_info = col_info.with_cql_type(cql_type);
}
col_info
}
type DecodedPartitionKey = (Arc<[u8]>, u64, Vec<(Arc<str>, Value)>);
fn pk_schema_fingerprint(schema: &crate::schema::TableSchema) -> u64 {
use std::hash::{Hash, Hasher};
let mut hasher = std::collections::hash_map::DefaultHasher::new();
schema.partition_keys.len().hash(&mut hasher);
for col in &schema.partition_keys {
col.name.hash(&mut hasher);
col.data_type.hash(&mut hasher);
col.position.hash(&mut hasher);
}
hasher.finish()
}
#[derive(Default)]
pub struct PartitionKeyCache {
decoded: Option<DecodedPartitionKey>,
}
impl PartitionKeyCache {
fn columns_for<'a>(
&'a mut self,
key_bytes: &Arc<[u8]>,
schema: &crate::schema::TableSchema,
) -> &'a [(Arc<str>, Value)] {
let fingerprint = pk_schema_fingerprint(schema);
let hit = self.decoded.as_ref().is_some_and(|(cached, fp, _)| {
*fp == fingerprint && cached.as_ref() == key_bytes.as_ref()
});
if !hit {
#[cfg(test)]
super::PARTITION_KEY_DECODES.with(|c| c.set(c.get() + 1));
let cols = match crate::storage::partition_key_codec::decode_partition_key_columns(
key_bytes, schema,
) {
Ok(pk_columns) => pk_columns
.into_iter()
.map(|(name, value)| (Arc::<str>::from(name), value))
.collect(),
Err(e) => {
tracing::warn!(
"Failed to reconstruct partition-key columns from row key \
(len={} bytes) for {}.{}: {}",
key_bytes.len(),
schema.keyspace,
schema.table,
e
);
Vec::new()
}
};
self.decoded = Some((Arc::clone(key_bytes), fingerprint, cols));
}
match &self.decoded {
Some((_, _, cols)) => cols,
None => &[],
}
}
}
pub fn build_row_from_scan(
key: RowKey,
row: ScanRow,
projection: &[String],
schema: Option<&crate::schema::TableSchema>,
) -> Option<QueryRow> {
let mut pk_cache = PartitionKeyCache::default();
build_row_from_scan_cached(key, row, projection, schema, &mut pk_cache)
}
pub fn build_row_from_scan_cached(
key: RowKey,
row: ScanRow,
projection: &[String],
schema: Option<&crate::schema::TableSchema>,
pk_cache: &mut PartitionKeyCache,
) -> Option<QueryRow> {
let cells = row.into_cells()?;
let pk_hint = schema.map(|s| s.partition_keys.len()).unwrap_or(0);
let mut row_values: HashMap<Arc<str>, Value> = HashMap::with_capacity(cells.len() + pk_hint);
let project = |name: &str| projection.is_empty() || projection.iter().any(|p| p == name);
for (name, col_value) in cells {
if project(&name) {
row_values.insert(name, col_value.into_owned());
}
}
if let Some(schema) = schema {
for (name, value) in pk_cache.columns_for(&key.0, schema) {
if project(name) {
row_values.insert(Arc::clone(name), value.clone());
}
}
}
Some(QueryRow {
values: row_values,
key,
metadata: Default::default(),
cell_metadata: None,
})
}
#[cfg(test)]
mod tests {
use super::super::predicate::evaluate_predicates;
use super::super::test_support::single_pk_schema;
use super::*;
#[test]
fn build_row_from_scan_materialises_single_text_pk() {
let key = RowKey::new(b"k0000000000000000".to_vec());
let value = ScanRow::Row(vec![(Arc::from("name"), Value::text("name-0".to_string()))]);
let schema = single_pk_schema("id", "text");
let row = build_row_from_scan(key, value, &[], Some(&schema))
.expect("row must be built (not tombstoned)");
assert_eq!(
row.values.get("id"),
Some(&Value::text("k0000000000000000".to_string())),
"Issue #586: single TEXT PK column must be reconstructed from the raw row key"
);
assert_eq!(
row.values.get("name"),
Some(&Value::text("name-0".to_string()))
);
}
#[test]
fn scan_built_row_matches_text_pk_equality_predicate() {
use super::super::super::select_optimizer::{SSTableFilterOp, SSTablePredicate};
let key = RowKey::new(b"k0000000000000000".to_vec());
let value = ScanRow::Row(vec![(Arc::from("age"), Value::Integer(0))]);
let schema = single_pk_schema("id", "text");
let row = build_row_from_scan(key, value, &[], Some(&schema)).unwrap();
let predicate = SSTablePredicate::column(
"id",
SSTableFilterOp::Equal,
vec![Value::text("k0000000000000000".to_string())],
);
assert!(
evaluate_predicates(&row, std::slice::from_ref(&predicate)).unwrap(),
"Issue #586: WHERE id = '<literal>' must match the reconstructed PK column"
);
}
#[test]
fn build_row_from_scan_multi_column_row_has_no_data_fallback() {
let key = RowKey::new(b"k0000000000000000".to_vec());
let value = ScanRow::Row(vec![
(Arc::from("name"), Value::text("alice".to_string())),
(Arc::from("score"), Value::Integer(42)),
]);
let row = build_row_from_scan(key, value, &[], None)
.expect("a live row must build (not tombstoned)");
assert_eq!(
row.values.get("name"),
Some(&Value::text("alice".to_string())),
"real text column value must survive the row-carrier disassembly"
);
assert_eq!(
row.values.get("score"),
Some(&Value::Integer(42)),
"real int column value must survive the row-carrier disassembly"
);
assert!(
!row.values.contains_key("data"),
"roborev H2: column values must NOT collapse into a synthetic 'data' fallback"
);
assert_eq!(
row.values.len(),
2,
"exactly the two real columns, no extras"
);
}
#[test]
fn build_row_from_scan_presizes_value_map() {
let cells: Vec<(Arc<str>, Value)> = (0..8)
.map(|i| (Arc::from(format!("c{i}").as_str()), Value::Integer(i)))
.collect();
let key = RowKey::new(b"k".to_vec());
let row = build_row_from_scan(key, ScanRow::Row(cells), &["c0".to_string()], None)
.expect("a live row must build");
assert_eq!(row.values.len(), 1, "projection keeps exactly one column");
assert!(
row.values.capacity() >= 8,
"issue #1584: value map must be pre-sized to the decoded cell count \
(>= 8), not grown from empty to the projected size; got capacity {}",
row.values.capacity()
);
}
#[test]
fn pk_decode_is_once_per_partition_not_per_row() {
use super::super::PARTITION_KEY_DECODES;
let schema = single_pk_schema("id", "text");
let key_bytes = b"partition-A".to_vec();
const N: usize = 50;
PARTITION_KEY_DECODES.with(|c| c.set(0));
let mut pk_cache = PartitionKeyCache::default();
let mut rows = Vec::with_capacity(N);
for i in 0..N {
let cells = ScanRow::Row(vec![(Arc::from("v"), Value::Integer(i as i32))]);
let row = build_row_from_scan_cached(
RowKey::new(key_bytes.clone()),
cells,
&[],
Some(&schema),
&mut pk_cache,
)
.expect("a live row must build");
rows.push(row);
}
let decodes = PARTITION_KEY_DECODES.with(|c| c.get());
assert_eq!(
decodes, 1,
"issue #1817: {N} rows of ONE partition must decode the partition key \
ONCE (O(partitions)), not once per row (would be {N})"
);
assert_eq!(rows.len(), N);
for (i, row) in rows.iter().enumerate() {
assert_eq!(
row.values.get("id"),
Some(&Value::text("partition-A".to_string())),
"each row must carry the reconstructed PK column value"
);
assert_eq!(row.values.get("v"), Some(&Value::Integer(i as i32)));
}
}
#[test]
fn pk_decode_counts_partitions_not_rows() {
use super::super::PARTITION_KEY_DECODES;
let schema = single_pk_schema("id", "text");
const P: usize = 4;
const ROWS_PER_PART: usize = 10;
PARTITION_KEY_DECODES.with(|c| c.set(0));
let mut pk_cache = PartitionKeyCache::default();
let mut total_rows = 0;
for p in 0..P {
let key_bytes = format!("partition-{p}").into_bytes();
for r in 0..ROWS_PER_PART {
let cells = ScanRow::Row(vec![(Arc::from("v"), Value::Integer(r as i32))]);
let row = build_row_from_scan_cached(
RowKey::new(key_bytes.clone()),
cells,
&[],
Some(&schema),
&mut pk_cache,
)
.expect("a live row must build");
assert_eq!(
row.values.get("id"),
Some(&Value::text(format!("partition-{p}"))),
"PK column reconstructed per partition"
);
total_rows += 1;
}
}
let decodes = PARTITION_KEY_DECODES.with(|c| c.get());
assert_eq!(total_rows, P * ROWS_PER_PART);
assert_eq!(
decodes,
P,
"issue #1817: decode count must be O(partitions) = {P}, not O(rows) = {}",
P * ROWS_PER_PART
);
}
#[test]
fn build_row_from_scan_wrapper_decodes_per_call() {
use super::super::PARTITION_KEY_DECODES;
let schema = single_pk_schema("id", "text");
let key_bytes = b"partition-A".to_vec();
const N: usize = 20;
PARTITION_KEY_DECODES.with(|c| c.set(0));
for i in 0..N {
let cells = ScanRow::Row(vec![(Arc::from("v"), Value::Integer(i as i32))]);
let row =
build_row_from_scan(RowKey::new(key_bytes.clone()), cells, &[], Some(&schema))
.expect("a live row must build");
assert_eq!(
row.values.get("id"),
Some(&Value::text("partition-A".to_string()))
);
}
let decodes = PARTITION_KEY_DECODES.with(|c| c.get());
assert_eq!(
decodes, N,
"the single-use wrapper decodes once per call ({N}); the shared-cache \
`build_row_from_scan_cached` is what makes a loop O(partitions)"
);
}
#[test]
fn pk_cache_reuse_across_schemas_does_not_leak_columns() {
let schema_a = single_pk_schema("id_a", "text");
let schema_b = single_pk_schema("id_b", "text");
let key_bytes: Arc<[u8]> = Arc::from(b"shared-key".as_slice());
let mut cache = PartitionKeyCache::default();
let cols_a: Vec<(Arc<str>, Value)> = cache.columns_for(&key_bytes, &schema_a).to_vec();
assert_eq!(cols_a.len(), 1);
assert_eq!(cols_a[0].0.as_ref(), "id_a");
assert_eq!(cols_a[0].1, Value::text("shared-key".to_string()));
let cols_b: Vec<(Arc<str>, Value)> = cache.columns_for(&key_bytes, &schema_b).to_vec();
assert_eq!(cols_b.len(), 1);
assert_eq!(
cols_b[0].0.as_ref(),
"id_b",
"roborev: a cache reused across schemas must NOT leak the first \
schema's partition-key column names — differing schema is a MISS"
);
assert_eq!(cols_b[0].1, Value::text("shared-key".to_string()));
}
#[test]
fn build_row_from_scan_marker_is_suppressed() {
let key = RowKey::new(b"k".to_vec());
assert!(
build_row_from_scan(key, ScanRow::Marker(Value::Null), &[], None).is_none(),
"a marker (tombstone/null) row must be suppressed from user output"
);
}
#[test]
fn live_value_dropped_as_marker_surfaces_as_row() {
let key = RowKey::new(b"k".to_vec());
let live = Value::text("synthetic-fallback".to_string());
assert!(
build_row_from_scan(key.clone(), ScanRow::Marker(live.clone()), &[], None).is_none(),
"a LIVE value mis-wrapped as Marker is dropped — the bug the producers must avoid"
);
let row = build_row_from_scan(
key,
ScanRow::Row(vec![(Arc::from("data"), live.clone())]),
&[],
None,
)
.expect("a live Row must surface");
assert_eq!(row.values.get("data"), Some(&live));
}
}