use std::collections::HashMap;
use std::time::Duration;
use nodedb_columnar::{ColumnarEngineSnapshot, MutationEngine, materialize_segment_live_rows};
use nodedb_types::surrogate::Surrogate;
use nodedb_types::value::Value;
use crate::Error;
use crate::bridge::envelope::PhysicalPlan;
use crate::control::server::shared::ddl::sync_dispatch;
use crate::control::server::wal_dispatch::wal_append_if_write;
use crate::control::state::SharedState;
use crate::types::{DatabaseId, TenantId, VShardId};
use nodedb_physical::physical_plan::{ColumnarInsertIntent, ColumnarOp};
pub struct DecodedColumnarRows {
pub rows: Vec<Value>,
pub surrogates: Vec<Surrogate>,
pub schema_bytes: Vec<u8>,
}
pub fn decode_snapshot_live_rows(
collection: &str,
snap: ColumnarEngineSnapshot,
kek: Option<&nodedb_wal::crypto::WalEncryptionKey>,
) -> crate::Result<DecodedColumnarRows> {
let schema = snap.schema.clone();
let schema_bytes = zerompk::to_msgpack_vec(&schema).map_err(|e| Error::Serialization {
format: "msgpack".into(),
detail: format!("restore reissue: encode schema for '{collection}': {e}"),
})?;
let column_names: Vec<String> = schema.columns.iter().map(|c| c.name.clone()).collect();
let (engine, flushed_segments, flushed_surrogates): (MutationEngine, Vec<Vec<u8>>, _) =
MutationEngine::from_snapshot(snap).map_err(|e| Error::Storage {
engine: "columnar".into(),
detail: format!("restore reissue: from_snapshot for '{collection}': {e}"),
})?;
let mut rows: Vec<Value> = Vec::new();
let mut surrogates: Vec<Surrogate> = Vec::new();
for (surrogate, values) in engine.scan_memtable_rows_with_surrogates() {
let surrogate = surrogate.ok_or_else(|| Error::Storage {
engine: "columnar".into(),
detail: format!(
"restore reissue: memtable row for '{collection}' has no surrogate; \
cannot preserve cross-engine identity"
),
})?;
rows.push(row_values_to_object(&column_names, values, collection)?);
surrogates.push(surrogate);
}
for (idx, blob) in flushed_segments.iter().enumerate() {
let segment_id = idx as u64 + 1;
let empty_deletes = nodedb_columnar::DeleteBitmap::new();
let deletes = engine.delete_bitmap(segment_id).unwrap_or(&empty_deletes);
let seg_surrogates: &[Option<Surrogate>] = flushed_surrogates
.get(idx)
.map(|v| v.as_slice())
.unwrap_or(&[]);
let live = materialize_segment_live_rows(blob, kek, &schema, deletes, seg_surrogates)
.map_err(|e| Error::Storage {
engine: "columnar".into(),
detail: format!(
"restore reissue: materialize segment {segment_id} for '{collection}': {e}"
),
})?;
for (row, surrogate) in live {
let surrogate = surrogate.ok_or_else(|| Error::Storage {
engine: "columnar".into(),
detail: format!(
"restore reissue: flushed row in segment {segment_id} for '{collection}' \
has no surrogate; cannot preserve cross-engine identity"
),
})?;
rows.push(row);
surrogates.push(surrogate);
}
}
Ok(DecodedColumnarRows {
rows,
surrogates,
schema_bytes,
})
}
pub fn build_columnar_insert_plan(
collection: &str,
decoded: DecodedColumnarRows,
) -> crate::Result<PhysicalPlan> {
let payload = nodedb_types::value_to_msgpack(&Value::Array(decoded.rows)).map_err(|e| {
Error::Serialization {
format: "msgpack".into(),
detail: format!("restore reissue: encode rows for '{collection}': {e}"),
}
})?;
Ok(PhysicalPlan::Columnar(ColumnarOp::Insert {
collection: collection.to_string(),
payload,
format: "msgpack".into(),
intent: ColumnarInsertIntent::Insert,
on_conflict_updates: Vec::new(),
surrogates: decoded.surrogates,
schema_bytes: decoded.schema_bytes,
provenance: None,
wal_lsn: None,
}))
}
pub async fn reissue_columnar_durably(
state: &SharedState,
tenant_id: TenantId,
database_id: DatabaseId,
collection: &str,
plan: PhysicalPlan,
) -> crate::Result<()> {
let vshard = VShardId::from_collection_in_database(database_id, collection);
if let Some(proposer) = state.async_raft_proposer.get() {
let entry = crate::control::wal_replication::to_replicated_entry(
tenant_id,
database_id,
vshard,
&plan,
)
.ok_or_else(|| Error::Internal {
detail: format!(
"restore reissue: columnar plan for '{collection}' did not map to a \
replicated write"
),
})?;
crate::control::wal_replication::propose_replicated_entry(state, proposer, entry).await?;
return Ok(());
}
wal_append_if_write(&state.wal, tenant_id, vshard, database_id, &plan)?;
sync_dispatch::dispatch_async(
state,
tenant_id,
database_id,
collection,
plan,
REISSUE_TIMEOUT,
)
.await?;
Ok(())
}
const REISSUE_TIMEOUT: Duration = Duration::from_secs(120);
fn row_values_to_object(
column_names: &[String],
values: Vec<Value>,
collection: &str,
) -> crate::Result<Value> {
if values.len() != column_names.len() {
return Err(Error::Storage {
engine: "columnar".into(),
detail: format!(
"restore reissue: row arity {} != schema column count {} for '{collection}'",
values.len(),
column_names.len()
),
});
}
let mut map: HashMap<String, Value> = HashMap::with_capacity(values.len());
for (name, value) in column_names.iter().zip(values) {
map.insert(name.clone(), value);
}
Ok(Value::Object(map))
}