use std::time::Duration;
use nodedb_types::id::DatabaseId;
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::{TenantId, VShardId};
use nodedb_physical::physical_plan::CrdtOp;
const REISSUE_TIMEOUT: Duration = Duration::from_secs(120);
async fn reissue_crdt_collection(
state: &SharedState,
tenant_id: TenantId,
database_id: DatabaseId,
collection: &str,
bytes: Vec<u8>,
) -> crate::Result<()> {
let vshard = VShardId::from_collection_in_database(database_id, collection);
let plan = PhysicalPlan::Crdt(CrdtOp::ImportSnapshot {
tenant_id: tenant_id.as_u64(),
collection: collection.to_string(),
bytes,
});
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: "restore reissue: crdt import did not map to a replicated write".into(),
})?;
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(())
}
pub(crate) async fn reissue_crdt_snapshots(
state: &SharedState,
crdt_state: Vec<(u64, u64, String, Vec<u8>)>,
) -> crate::Result<usize> {
let mut imported = 0usize;
for (database_id, tid, collection, bytes) in crdt_state {
reissue_crdt_collection(
state,
TenantId::new(tid),
DatabaseId::new(database_id),
&collection,
bytes,
)
.await?;
imported += 1;
}
Ok(imported)
}