use std::sync::Arc;
use nodedb_types::surrogate::Surrogate;
use crate::Error;
use crate::control::state::SharedState;
use crate::engine::vector::index_config::IndexConfig;
use crate::types::TenantId;
use crate::control::backup::snapshot_keys::extract_db_scoped_collection;
pub(super) async fn reissue_timeseries_snapshots(
state: &Arc<SharedState>,
tenant_id: u64,
memtables: Vec<(String, Vec<u8>)>,
flushed: Vec<crate::types::TsFlushedCollectionBlob>,
) -> Result<usize, Error> {
let kek = state.wal.encryption_key().cloned();
let database_id = crate::types::DatabaseId::DEFAULT;
let mut memtable_by_key: std::collections::HashMap<String, Vec<u8>> =
memtables.into_iter().collect();
let mut keys_in_order: Vec<String> = Vec::new();
let mut flushed_by_key: std::collections::HashMap<
String,
crate::types::TsFlushedCollectionBlob,
> = std::collections::HashMap::new();
for blob in flushed {
keys_in_order.push(blob.collection_key.clone());
flushed_by_key.insert(blob.collection_key.clone(), blob);
}
for key in memtable_by_key.keys() {
if !flushed_by_key.contains_key(key) {
keys_in_order.push(key.clone());
}
}
let empty_flushed = crate::types::TsFlushedCollectionBlob::default();
let mut reissued = 0usize;
for key in keys_in_order {
let Some(collection) = extract_db_scoped_collection(&key, tenant_id) else {
return Err(Error::Internal {
detail: format!("restore reissue: malformed timeseries snapshot key '{key}'"),
});
};
let collection = collection.to_owned();
let memtable_bytes = memtable_by_key.remove(&key);
let flushed_blob = flushed_by_key.get(&key).unwrap_or(&empty_flushed);
let rows = super::super::timeseries_reissue::decode_timeseries_live_rows(
&collection,
memtable_bytes.as_deref(),
flushed_blob,
kek.as_ref(),
)?;
if rows.is_empty() {
continue;
}
let plan =
super::super::timeseries_reissue::build_timeseries_ingest_plan(&collection, rows)?;
super::super::timeseries_reissue::reissue_timeseries_durably(
state,
TenantId::new(tenant_id),
database_id,
&collection,
plan,
)
.await?;
reissued += 1;
}
Ok(reissued)
}
pub(super) async fn reissue_columnar_snapshots(
state: &Arc<SharedState>,
tenant_id: u64,
entries: Vec<(String, Vec<u8>)>,
) -> Result<usize, Error> {
let kek = state.wal.encryption_key().cloned();
let database_id = crate::types::DatabaseId::DEFAULT;
let mut reissued = 0usize;
for (key, bytes) in entries {
let Some(collection) = extract_db_scoped_collection(&key, tenant_id) else {
return Err(Error::Internal {
detail: format!("restore reissue: malformed columnar snapshot key '{key}'"),
});
};
let collection = collection.to_owned();
let snap: nodedb_columnar::ColumnarEngineSnapshot =
zerompk::from_msgpack(&bytes).map_err(|e| Error::Serialization {
format: "msgpack".into(),
detail: format!(
"restore reissue: deserialize ColumnarEngineSnapshot for '{collection}': {e}"
),
})?;
let decoded = super::super::columnar_reissue::decode_snapshot_live_rows(
&collection,
snap,
kek.as_ref(),
)?;
if decoded.rows.is_empty() {
continue;
}
let plan =
super::super::columnar_reissue::build_columnar_insert_plan(&collection, decoded)?;
super::super::columnar_reissue::reissue_columnar_durably(
state,
TenantId::new(tenant_id),
database_id,
&collection,
plan,
)
.await?;
reissued += 1;
}
Ok(reissued)
}
pub(super) async fn reissue_vector_snapshots(
state: &Arc<SharedState>,
tenant_id: u64,
entries: Vec<(String, Vec<u8>)>,
) -> Result<usize, Error> {
let database_id = crate::types::DatabaseId::DEFAULT;
let mut reissued = 0usize;
for (key, bytes) in entries {
let Some(coll_key) = extract_db_scoped_collection(&key, tenant_id) else {
return Err(Error::Internal {
detail: format!("restore reissue: malformed vector snapshot key '{key}'"),
});
};
let (collection, field_name) =
super::super::vector_reissue::split_vector_coll_key(coll_key);
let collection = collection.to_owned();
let field_name = field_name.to_owned();
let vectors: Vec<(u32, Vec<f32>, Option<Surrogate>)> = zerompk::from_msgpack(&bytes)
.map_err(|e| Error::Serialization {
format: "msgpack".into(),
detail: format!(
"restore reissue: deserialize vector snapshot for '{collection}': {e}"
),
})?;
if vectors.is_empty() {
continue;
}
for (_node_id, vector, surrogate) in vectors {
let surrogate = surrogate.unwrap_or(Surrogate::ZERO);
let plan = super::super::vector_reissue::build_vector_insert_plan(
&collection,
&field_name,
vector,
surrogate,
);
super::super::vector_reissue::reissue_vector_durably(
state,
TenantId::new(tenant_id),
database_id,
&collection,
plan,
)
.await?;
reissued += 1;
}
}
Ok(reissued)
}
pub(super) async fn reissue_vector_params(
state: &Arc<SharedState>,
tenant_id: u64,
params: Vec<(String, Vec<u8>)>,
index_configs: Vec<(String, Vec<u8>)>,
) -> Result<usize, Error> {
let database_id = crate::types::DatabaseId::DEFAULT;
let mut resolved: std::collections::HashMap<String, IndexConfig> =
std::collections::HashMap::new();
let mut keys_in_order: Vec<String> = Vec::new();
for (key, bytes) in index_configs {
let Some(coll_key) = extract_db_scoped_collection(&key, tenant_id) else {
return Err(Error::Internal {
detail: format!("restore reissue: malformed index_configs snapshot key '{key}'"),
});
};
let coll_key = coll_key.to_owned();
let cfg: IndexConfig = zerompk::from_msgpack(&bytes).map_err(|e| Error::Serialization {
format: "msgpack".into(),
detail: format!("restore reissue: deserialize IndexConfig for '{coll_key}': {e}"),
})?;
keys_in_order.push(coll_key.clone());
resolved.insert(coll_key, cfg);
}
for (key, bytes) in params {
let Some(coll_key) = extract_db_scoped_collection(&key, tenant_id) else {
return Err(Error::Internal {
detail: format!("restore reissue: malformed vector_params snapshot key '{key}'"),
});
};
let coll_key = coll_key.to_owned();
if resolved.contains_key(&coll_key) {
continue;
}
let hnsw: nodedb_types::hnsw::HnswParams =
zerompk::from_msgpack(&bytes).map_err(|e| Error::Serialization {
format: "msgpack".into(),
detail: format!("restore reissue: deserialize HnswParams for '{coll_key}': {e}"),
})?;
keys_in_order.push(coll_key.clone());
resolved.insert(
coll_key,
IndexConfig {
hnsw,
..IndexConfig::default()
},
);
}
let mut reissued = 0usize;
for coll_key in keys_in_order {
let Some(config) = resolved.remove(&coll_key) else {
continue;
};
let (collection, field_name) =
super::super::vector_reissue::split_vector_coll_key(&coll_key);
let collection = collection.to_owned();
let field_name = field_name.to_owned();
let plan = super::super::vector_reissue::build_vector_set_params_plan(
&collection,
&field_name,
&config,
);
super::super::vector_reissue::reissue_vector_durably(
state,
TenantId::new(tenant_id),
database_id,
&collection,
plan,
)
.await?;
reissued += 1;
}
Ok(reissued)
}