use nodedb_types::value::Value;
use crate::bridge::scan_filter::ScanFilter;
use crate::data::executor::core_loop::CoreLoop;
use crate::data::executor::handlers::transaction::overlay::Staged;
use crate::types::{DatabaseId, TenantId, TxnId};
pub(in crate::data::executor) struct TimeseriesOverlayMergeParams<'a> {
pub txn_id: TxnId,
pub coll_key: &'a (DatabaseId, TenantId, String),
pub time_range: (i64, i64),
pub filter_predicates: &'a [ScanFilter],
pub has_filters: bool,
pub limit: usize,
}
fn decode_staged_row(body: &[u8]) -> Option<Value> {
match nodedb_types::value_from_msgpack(body) {
Ok(v @ Value::Object(_)) => Some(v),
_ => None,
}
}
fn row_timestamp_ms(row: &Value) -> Option<i64> {
let Value::Object(map) = row else {
return None;
};
for key in ["timestamp", "ts", "time"] {
match map.get(key) {
Some(Value::Integer(ms)) => return Some(*ms),
Some(Value::Float(ms)) => return Some(*ms as i64),
_ => {}
}
}
None
}
fn staged_row_to_rmpv(row: &Value) -> rmpv::Value {
let Value::Object(map) = row else {
return rmpv::Value::Nil;
};
let fields: Vec<(rmpv::Value, rmpv::Value)> = map
.iter()
.map(|(k, v)| (rmpv::Value::String(k.as_str().into()), scalar_to_rmpv(v)))
.collect();
rmpv::Value::Map(fields)
}
fn scalar_to_rmpv(v: &Value) -> rmpv::Value {
match v {
Value::Integer(n) => rmpv::Value::Integer((*n).into()),
Value::Float(f) => rmpv::Value::F64(*f),
Value::String(s) => rmpv::Value::String(s.as_str().into()),
Value::Bool(b) => rmpv::Value::Boolean(*b),
_ => rmpv::Value::Nil,
}
}
impl CoreLoop {
pub(in crate::data::executor) fn merge_overlay_into_timeseries_scan(
&self,
params: TimeseriesOverlayMergeParams<'_>,
results: &mut Vec<rmpv::Value>,
) {
let TimeseriesOverlayMergeParams {
txn_id,
coll_key,
time_range,
filter_predicates,
has_filters,
limit,
} = params;
self.touch_overlay(txn_id);
let Some(overlay) = self.txn_overlays.get(&txn_id) else {
return;
};
for (_surrogate, staged) in overlay.iter_for_collection(coll_key) {
if results.len() >= limit {
break;
}
let Staged::Put(body) = staged else {
continue;
};
let Some(row) = decode_staged_row(body) else {
continue;
};
match row_timestamp_ms(&row) {
Some(ts) if ts >= time_range.0 && ts <= time_range.1 => {}
_ => continue,
}
if has_filters && !filter_predicates.iter().all(|f| f.matches_binary(body)) {
continue;
}
results.push(staged_row_to_rmpv(&row));
}
}
}