use std::collections::BTreeSet;
use nodedb_cluster::{AssignSurrogateRequest, AssignSurrogateResponse, RaftRpc};
use nodedb_types::Surrogate;
use crate::control::server::exchange::resolve::register_peers_from_topology;
use crate::control::state::SharedState;
use crate::types::{DatabaseId, TenantId, TraceId, VShardId};
pub async fn assign_surrogate_routed(
state: &SharedState,
vshard: VShardId,
database_id: DatabaseId,
tenant_id: TenantId,
collection: &str,
pk: &[u8],
trace_id: TraceId,
) -> crate::Result<Surrogate> {
let (Some(transport), Some(routing)) = (
state.cluster_transport.as_ref(),
state.cluster_routing.as_ref(),
) else {
return state
.surrogate_assigner
.assign(database_id, tenant_id, collection, pk);
};
let leader = {
let guard = routing.read().unwrap_or_else(|p| p.into_inner());
guard
.leader_for_vshard(vshard.as_u32())
.map_err(|e| crate::Error::Internal {
detail: format!(
"assign-surrogate: no leader for vshard {} ({collection}): {e}",
vshard.as_u32()
),
})?
};
if leader == 0 {
return Err(crate::Error::Internal {
detail: format!(
"assign-surrogate: no leader elected for home vshard {} ({collection}); \
cannot resolve an authoritative surrogate yet",
vshard.as_u32()
),
});
}
if leader == state.node_id {
return state
.surrogate_assigner
.assign(database_id, tenant_id, collection, pk);
}
let mut targets = BTreeSet::new();
targets.insert(leader);
register_peers_from_topology(state, transport, &targets);
let deadline_remaining_ms = state
.tuning
.network
.default_deadline_secs
.saturating_mul(1000)
.max(1);
let req = AssignSurrogateRequest {
vshard_id: vshard.as_u32(),
database_id: database_id.as_u64(),
tenant_id: tenant_id.as_u64(),
collection: collection.to_string(),
pk: pk.to_vec(),
deadline_remaining_ms,
trace_id: trace_id.0,
};
match transport
.send_rpc(leader, RaftRpc::AssignSurrogateRequest(req))
.await
{
Ok(RaftRpc::AssignSurrogateResponse(AssignSurrogateResponse {
surrogate,
error: None,
})) => Ok(Surrogate::new(surrogate)),
Ok(RaftRpc::AssignSurrogateResponse(AssignSurrogateResponse {
error: Some(e), ..
})) => Err(crate::Error::Internal {
detail: format!("assign-surrogate failed on leader node {leader}: {e:?}"),
}),
Ok(other) => Err(crate::Error::Internal {
detail: format!("assign-surrogate: unexpected reply from node {leader}: {other:?}"),
}),
Err(e) => Err(crate::Error::Internal {
detail: format!("assign-surrogate RPC to node {leader} failed: {e}"),
}),
}
}