use std::collections::HashSet;
use nodedb_physical::physical_plan::SpatialPredicate;
use nodedb_types::Surrogate;
use nodedb_types::columnar::ColumnarSchema;
use nodedb_types::geometry::Geometry;
use nodedb_types::value::Value;
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::spatial_refine::{
apply_predicate, extract_geometry, project_doc,
};
use crate::data::executor::handlers::transaction::overlay::Staged;
use crate::engine::document::store::surrogate_to_doc_id;
use crate::types::{DatabaseId, TenantId, TxnId};
pub(in crate::data::executor) struct SpatialOverlayMergeParams<'a> {
pub txn_id: TxnId,
pub coll_key: &'a (DatabaseId, TenantId, String),
pub field: &'a str,
pub predicate: &'a SpatialPredicate,
pub query_geom: &'a Geometry,
pub distance_meters: f64,
pub projection: &'a [String],
pub attr_filters: &'a [ScanFilter],
pub row_level_filters: &'a [ScanFilter],
}
fn decode_staged_spatial_row(body: &[u8], schema: Option<&ColumnarSchema>) -> Option<Value> {
match nodedb_types::value_from_msgpack(body).ok()? {
Value::Array(row) => {
let schema = schema?;
let json = row_to_projected_json(&row, schema, &[], &[], false);
Some(Value::from(json))
}
obj @ Value::Object(_) => Some(obj),
_ => None,
}
}
fn row_surrogate(row: &Value) -> Option<u32> {
match row {
Value::Object(map) => match map.get("id") {
Some(Value::String(s)) => u32::from_str_radix(s, 16).ok(),
_ => None,
},
_ => None,
}
}
impl CoreLoop {
pub(in crate::data::executor) fn merge_overlay_into_spatial_scan(
&self,
params: SpatialOverlayMergeParams<'_>,
results: &mut Vec<Value>,
) {
let SpatialOverlayMergeParams {
txn_id,
coll_key,
field,
predicate,
query_geom,
distance_meters,
projection,
attr_filters,
row_level_filters,
} = params;
self.touch_overlay(txn_id);
let Some(overlay) = self.txn_overlays.get(&txn_id) else {
return;
};
let schema = self.columnar_engines.get(coll_key).map(|e| e.schema());
let row_matches = |doc: &Value| -> bool {
let Some(doc_geom) = extract_geometry(doc, field) else {
return false;
};
if !apply_predicate(predicate, query_geom, &doc_geom, distance_meters) {
return false;
}
attr_filters.iter().all(|f| f.matches_value(doc))
&& row_level_filters.iter().all(|f| f.matches_value(doc))
};
let mut seen: HashSet<u32> = results.iter().filter_map(row_surrogate).collect();
results.retain_mut(|row| {
let Some(raw) = row_surrogate(row) else {
return true;
};
match overlay.get(coll_key, raw) {
Some(Staged::Tombstone) => false,
Some(Staged::Put(body)) => match decode_staged_spatial_row(body, schema) {
Some(doc) => {
if !row_matches(&doc) {
return false;
}
let doc_id = surrogate_to_doc_id(Surrogate(raw));
*row = project_doc(&doc, &doc_id, projection);
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(doc) = decode_staged_spatial_row(body, schema) else {
continue;
};
if !row_matches(&doc) {
continue;
}
let doc_id = surrogate_to_doc_id(Surrogate(surrogate));
results.push(project_doc(&doc, &doc_id, projection));
seen.insert(surrogate);
}
}
}