use nodedb_types::columnar::schema::{TS_SYSTEM, TS_VALID_FROM, TS_VALID_UNTIL};
use nodedb_types::columnar::{ColumnDef, ColumnType, ColumnarSchema};
use nodedb_types::value::Value;
use crate::data::executor::core_loop::CoreLoop;
use crate::types::{DatabaseId, TenantId};
impl CoreLoop {
pub(in crate::data::executor) fn ensure_columnar_engine_schema(
&mut self,
engine_key: &(DatabaseId, TenantId, String),
collection: &str,
bitemporal: bool,
first_row: &Value,
schema_bytes: &[u8],
) -> ColumnarSchema {
if let Some(engine) = self.columnar_engines.get(engine_key) {
return engine.schema().clone();
}
let flush_threshold = self.query_tuning.columnar_flush_threshold;
let engine = self
.columnar_engines
.entry(engine_key.clone())
.or_insert_with(|| {
let base_schema = if !schema_bytes.is_empty() {
zerompk::from_msgpack::<ColumnarSchema>(schema_bytes)
.unwrap_or_else(|_| infer_schema_from_value(first_row))
} else {
infer_schema_from_value(first_row)
};
let schema = if bitemporal {
prepend_bitemporal_columns(base_schema)
} else {
base_schema
};
nodedb_columnar::MutationEngine::with_flush_threshold(
collection.to_string(),
schema,
flush_threshold,
)
});
engine.schema().clone()
}
}
pub(super) fn row_values_to_object(schema: &ColumnarSchema, row: &[Value]) -> nodedb_types::Value {
let mut map = std::collections::HashMap::with_capacity(schema.columns.len());
for (col, val) in schema.columns.iter().zip(row.iter()) {
map.insert(col.name.clone(), val.clone());
}
nodedb_types::Value::Object(map)
}
pub(in crate::data::executor) fn ndb_field_to_value(
val: Option<&Value>,
col_type: &ColumnType,
) -> crate::Result<Value> {
let Some(val) = val else {
return Ok(Value::Null);
};
let v = match (col_type, val) {
(_, Value::Null) => Value::Null,
(ColumnType::Int64, Value::Integer(_)) => val.clone(),
(ColumnType::Int64, Value::Float(f)) => Value::Integer(*f as i64),
(ColumnType::Int64, Value::String(s)) => {
s.parse::<i64>().map(Value::Integer).unwrap_or(Value::Null)
}
(ColumnType::Float64, Value::Float(_)) => val.clone(),
(ColumnType::Float64, Value::Integer(n)) => Value::Float(*n as f64),
(ColumnType::Float64, Value::String(s)) => {
s.parse::<f64>().map(Value::Float).unwrap_or(Value::Null)
}
(ColumnType::Bool, Value::Bool(_)) => val.clone(),
(ColumnType::String, Value::String(_)) => val.clone(),
(ColumnType::Timestamp, Value::Integer(n)) => {
Value::NaiveDateTime(nodedb_types::NdbDateTime::from_millis(*n).map_err(|e| {
crate::Error::BadRequest {
detail: format!("timestamp coercion: {e}"),
}
})?)
}
(ColumnType::Timestamp, Value::Float(f)) => {
Value::NaiveDateTime(nodedb_types::NdbDateTime::from_millis(*f as i64).map_err(
|e| crate::Error::BadRequest {
detail: format!("timestamp coercion: {e}"),
},
)?)
}
(ColumnType::Timestamp, Value::String(s)) => nodedb_types::datetime::NdbDateTime::parse(s)
.map(Value::NaiveDateTime)
.unwrap_or_else(|| Value::String(s.clone())),
(ColumnType::Timestamptz, Value::Integer(n)) => {
Value::DateTime(nodedb_types::NdbDateTime::from_millis(*n).map_err(|e| {
crate::Error::BadRequest {
detail: format!("timestamptz coercion: {e}"),
}
})?)
}
(ColumnType::Timestamptz, Value::Float(f)) => Value::DateTime(
nodedb_types::NdbDateTime::from_millis(*f as i64).map_err(|e| {
crate::Error::BadRequest {
detail: format!("timestamptz coercion: {e}"),
}
})?,
),
(ColumnType::Timestamptz, Value::String(s)) => {
nodedb_types::datetime::NdbDateTime::parse(s)
.map(Value::DateTime)
.unwrap_or_else(|| Value::String(s.clone()))
}
(ColumnType::Uuid, Value::String(_)) => val.clone(),
(ColumnType::Float64, _) => Value::Null,
(ColumnType::Int64, _) => Value::Null,
_ => val.clone(),
};
Ok(v)
}
pub(in crate::data::executor) fn infer_schema_from_value(row: &Value) -> ColumnarSchema {
let obj = match row {
Value::Object(m) => m,
_ => {
return ColumnarSchema::new(vec![ColumnDef::required("value", ColumnType::Float64)])
.expect("single-column schema");
}
};
let mut columns = Vec::new();
for (key, val) in obj {
let lower = key.to_lowercase();
let col_type = if lower == "timestamp" || lower == "ts" || lower == "time" {
ColumnType::Timestamp
} else {
match val {
Value::Float(_) => ColumnType::Float64,
Value::Integer(_) => ColumnType::Int64,
Value::Bool(_) => ColumnType::Bool,
_ => ColumnType::String,
}
};
if lower == "id" {
columns.push(ColumnDef::required(key.clone(), col_type).with_primary_key());
} else {
columns.push(ColumnDef::nullable(key.clone(), col_type));
}
}
if columns.is_empty() {
columns.push(ColumnDef::required("value", ColumnType::Float64));
}
ColumnarSchema::new(columns).expect("inferred schema must be valid")
}
pub(in crate::data::executor) fn prepend_bitemporal_columns(
base: ColumnarSchema,
) -> ColumnarSchema {
let mut cols = Vec::with_capacity(3 + base.columns.len());
cols.push(ColumnDef::required(TS_SYSTEM, ColumnType::Int64));
cols.push(ColumnDef::required(TS_VALID_FROM, ColumnType::Int64));
cols.push(ColumnDef::required(TS_VALID_UNTIL, ColumnType::Int64));
cols.extend(base.columns);
ColumnarSchema::new(cols).expect("bitemporal columnar schema must be valid")
}
pub(super) fn infer_schema_from_json(row: &serde_json::Value) -> ColumnarSchema {
let ndb: Value = row.clone().into();
infer_schema_from_value(&ndb)
}