use tracing::warn;
use super::super::core_loop::CoreLoop;
use super::super::scan_normalize::decoded_col_to_value;
use crate::bridge::envelope::PhysicalPlan;
use crate::types::{DatabaseId, TenantId, VShardId};
use nodedb_physical::physical_plan::{ColumnarInsertIntent, ColumnarOp};
use nodedb_types::columnar::ColumnType;
impl CoreLoop {
pub(super) fn restore_columnar_geometry_indexes(
&mut self,
key: &(DatabaseId, TenantId, String),
engine: &nodedb_columnar::MutationEngine,
segments: &[Vec<u8>],
) -> usize {
let (db_id, tenant_id, collection) = key;
let schema = engine.schema().clone();
if !schema
.columns
.iter()
.any(|c| c.column_type == ColumnType::Geometry)
{
return 0;
}
let mut rows: Vec<Vec<nodedb_types::value::Value>> = Vec::new();
rows.extend(Self::restored_flushed_rows(
engine, segments, &schema, collection,
));
rows.extend(engine.scan_memtable_rows());
let docs: Vec<nodedb_types::Value> = rows
.iter()
.map(|row| {
let mut obj = std::collections::HashMap::with_capacity(schema.columns.len());
for (i, col) in schema.columns.iter().enumerate() {
if let Some(v) = row.get(i) {
obj.insert(col.name.clone(), v.clone());
}
}
nodedb_types::Value::Object(obj)
})
.collect();
let task = Self::replay_task(
*tenant_id,
*db_id,
VShardId::from_collection_in_database(*db_id, collection),
PhysicalPlan::Columnar(ColumnarOp::Insert {
collection: collection.clone(),
payload: Vec::new(),
format: "msgpack".into(),
intent: ColumnarInsertIntent::Insert,
on_conflict_updates: Vec::new(),
surrogates: Vec::new(),
schema_bytes: Vec::new(),
provenance: None,
wal_lsn: None,
}),
None,
);
let indexed = docs.len();
self.index_columnar_geometry_columns(&task, &schema, collection, &docs);
indexed
}
fn restored_flushed_rows(
engine: &nodedb_columnar::MutationEngine,
segments: &[Vec<u8>],
schema: &nodedb_types::columnar::ColumnarSchema,
collection: &str,
) -> Vec<Vec<nodedb_types::value::Value>> {
let mut out = Vec::new();
for (seg_idx, seg_bytes) in segments.iter().enumerate() {
let seg_id = seg_idx as u64 + 1;
let reader = match nodedb_columnar::SegmentReader::open(seg_bytes) {
Ok(r) => r,
Err(e) => {
warn!(
%collection,
seg_id,
error = %e,
"columnar checkpoint restore: flushed segment unreadable; its \
geometry rows are absent from the rebuilt R-tree"
);
continue;
}
};
let mut decoded_cols = Vec::with_capacity(schema.columns.len());
let mut decode_ok = true;
for col_idx in 0..schema.columns.len() {
match reader.read_column(col_idx) {
Ok(dc) => decoded_cols.push(dc),
Err(e) => {
warn!(
%collection,
seg_id,
col_idx,
error = %e,
"columnar checkpoint restore: column decode failed; the \
segment's geometry rows are absent from the rebuilt R-tree"
);
decode_ok = false;
break;
}
}
}
if !decode_ok {
continue;
}
let delete_bm = engine.delete_bitmap(seg_id);
for row_idx in 0..reader.row_count() as usize {
if delete_bm.is_some_and(|bm| bm.is_deleted(row_idx as u32)) {
continue;
}
out.push(
decoded_cols
.iter()
.map(|dc| decoded_col_to_value(dc, row_idx))
.collect(),
);
}
}
out
}
}