use std::collections::HashSet;
use nodedb_types::Surrogate;
use nodedb_types::columnar::ColumnarSchema;
use nodedb_types::value::Value;
use crate::bridge::expr_eval::ComputedColumn;
use crate::bridge::scan_filter::ScanFilter;
use crate::data::executor::core_loop::CoreLoop;
use crate::data::executor::handlers::columnar_read::convert::row_to_projected_json;
use crate::data::executor::handlers::columnar_read::filter::row_matches_filters;
use crate::data::executor::handlers::transaction::overlay::Staged;
use crate::types::{DatabaseId, TenantId, TxnId};
pub(in crate::data::executor) type ColumnarMatchedRow =
(Option<Surrogate>, Vec<Value>, serde_json::Value);
pub(in crate::data::executor) struct ColumnarOverlayMergeParams<'a> {
pub txn_id: TxnId,
pub coll_key: &'a (DatabaseId, TenantId, String),
pub schema: &'a ColumnarSchema,
pub projection: &'a [String],
pub filter_predicates: &'a [ScanFilter],
pub computed_cols: &'a [ComputedColumn],
pub all_versions: bool,
}
fn decode_staged_row(body: &[u8]) -> Option<Vec<Value>> {
match nodedb_types::value_from_msgpack(body) {
Ok(Value::Array(values)) => Some(values),
_ => None,
}
}
impl CoreLoop {
pub(in crate::data::executor) fn merge_overlay_into_columnar_scan(
&self,
params: ColumnarOverlayMergeParams<'_>,
matched: &mut Vec<ColumnarMatchedRow>,
) {
let ColumnarOverlayMergeParams {
txn_id,
coll_key,
schema,
projection,
filter_predicates,
computed_cols,
all_versions,
} = params;
self.touch_overlay(txn_id);
let Some(overlay) = self.txn_overlays.get(&txn_id) else {
return;
};
let predicate = |row: &[Value]| -> bool {
filter_predicates.is_empty() || row_matches_filters(row, schema, filter_predicates)
};
let mut seen: HashSet<u32> = matched
.iter()
.filter_map(|(surrogate, _, _)| surrogate.map(|s| s.0))
.collect();
matched.retain_mut(|(surrogate, row, json)| {
let Some(s) = surrogate else {
return true;
};
match overlay.get(coll_key, s.0) {
Some(Staged::Tombstone) => false,
Some(Staged::Put(body)) => match decode_staged_row(body) {
Some(new_row) => {
if !predicate(&new_row) {
return false;
}
*json = row_to_projected_json(
&new_row,
schema,
projection,
computed_cols,
all_versions,
);
*row = new_row;
true
}
None => false,
},
None => true,
}
});
for (surrogate, staged) in overlay.iter_for_collection(coll_key) {
if seen.contains(&surrogate) {
continue;
}
let Staged::Put(body) = staged else {
continue;
};
let Some(new_row) = decode_staged_row(body) else {
continue;
};
if !predicate(&new_row) {
continue;
}
let json =
row_to_projected_json(&new_row, schema, projection, computed_cols, all_versions);
matched.push((Some(Surrogate::new(surrogate)), new_row, json));
seen.insert(surrogate);
}
}
}