use std::collections::HashMap;
use nodedb_types::columnar::schema::TS_SYSTEM;
use crate::engine::timeseries::columnar_memtable::{ColumnData, ColumnType};
pub(super) fn rmpv_system_time(row: &rmpv::Value) -> i64 {
let rmpv::Value::Map(entries) = row else {
return i64::MIN;
};
for (k, v) in entries {
if let rmpv::Value::String(s) = k
&& s.as_str() == Some(TS_SYSTEM)
&& let rmpv::Value::Integer(i) = v
{
return i.as_i64().unwrap_or(i64::MIN);
}
}
i64::MIN
}
pub(super) fn emit_memtable_row(
mt: &crate::engine::timeseries::columnar_memtable::ColumnarMemtable,
columns: &[(usize, &String, &ColumnType, &ColumnData)],
idx: usize,
) -> rmpv::Value {
let mut buf = Vec::with_capacity(columns.len() * 32);
nodedb_query::msgpack_scan::write_map_header(&mut buf, columns.len());
for (col_idx, col_name, col_type, col_data) in columns {
nodedb_query::msgpack_scan::write_str(&mut buf, col_name);
crate::data::executor::handlers::columnar_read::emit_column_value(
&mut buf, mt, *col_idx, col_type, col_data, idx,
);
}
rmpv::decode::read_value(&mut buf.as_slice()).unwrap_or(rmpv::Value::Nil)
}
pub(super) fn emit_partition_row(
schema: &[(String, ColumnType)],
col_data: &[Option<ColumnData>],
sym_dicts: &HashMap<usize, nodedb_types::timeseries::SymbolDictionary>,
idx: usize,
) -> rmpv::Value {
let mut fields: Vec<(rmpv::Value, rmpv::Value)> = Vec::with_capacity(schema.len());
for (col_i, (col_name, col_type)) in schema.iter().enumerate() {
let Some(data) = &col_data[col_i] else {
continue;
};
let val = match col_type {
ColumnType::Timestamp => rmpv::Value::Integer(data.as_timestamps()[idx].into()),
ColumnType::Float64 => {
let v = data.as_f64()[idx];
if v.is_nan() {
rmpv::Value::Nil
} else {
rmpv::Value::F64(v)
}
}
ColumnType::Int64 => {
if let ColumnData::Int64(vals) = data {
rmpv::Value::Integer(vals[idx].into())
} else {
rmpv::Value::Nil
}
}
ColumnType::Symbol => {
if let ColumnData::Symbol(ids) = data {
sym_dicts
.get(&col_i)
.and_then(|dict| dict.get(ids[idx]))
.map(|s| rmpv::Value::String(s.into()))
.unwrap_or(rmpv::Value::Nil)
} else {
rmpv::Value::Nil
}
}
};
fields.push((rmpv::Value::String(col_name.as_str().into()), val));
}
rmpv::Value::Map(fields)
}
pub(super) fn extract_timestamp(row: &rmpv::Value) -> i64 {
if let rmpv::Value::Map(fields) = row {
for (_, v) in fields {
if let rmpv::Value::Integer(n) = v {
return n.as_i64().unwrap_or(0);
}
}
}
0
}
pub(super) fn apply_computed_columns_rmpv(
row: rmpv::Value,
computed_cols: &[crate::bridge::expr_eval::ComputedColumn],
) -> rmpv::Value {
let doc = rmpv_to_nodedb_value(&row);
let mut fields: Vec<(rmpv::Value, rmpv::Value)> = Vec::with_capacity(computed_cols.len());
for cc in computed_cols {
let result = cc.expr.eval(&doc);
fields.push((
rmpv::Value::String(cc.alias.as_str().into()),
nodedb_value_to_rmpv(&result),
));
}
rmpv::Value::Map(fields)
}
pub(super) fn rmpv_to_nodedb_value(row: &rmpv::Value) -> nodedb_types::Value {
match row {
rmpv::Value::Map(fields) => {
let mut map = std::collections::HashMap::new();
for (k, v) in fields {
let key = match k {
rmpv::Value::String(s) => s.as_str().unwrap_or("").to_string(),
_ => continue,
};
let val = match v {
rmpv::Value::Integer(n) => {
nodedb_types::Value::Integer(n.as_i64().unwrap_or(0))
}
rmpv::Value::F64(f) => nodedb_types::Value::Float(*f),
rmpv::Value::String(s) => {
nodedb_types::Value::String(s.as_str().unwrap_or("").to_string())
}
rmpv::Value::Nil => nodedb_types::Value::Null,
rmpv::Value::Boolean(b) => nodedb_types::Value::Bool(*b),
_ => nodedb_types::Value::Null,
};
map.insert(key, val);
}
nodedb_types::Value::Object(map)
}
_ => nodedb_types::Value::Null,
}
}
pub(super) fn nodedb_value_to_rmpv(v: &nodedb_types::Value) -> rmpv::Value {
match v {
nodedb_types::Value::Integer(n) => rmpv::Value::Integer((*n).into()),
nodedb_types::Value::Float(f) => rmpv::Value::F64(*f),
nodedb_types::Value::String(s) => rmpv::Value::String(s.as_str().into()),
nodedb_types::Value::Bool(b) => rmpv::Value::Boolean(*b),
nodedb_types::Value::Null => rmpv::Value::Nil,
_ => rmpv::Value::Nil,
}
}