use std::collections::BTreeSet;
use nodedb_cluster::{
METADATA_GROUP_ID, RaftRpc, RoutingTable, ShuffleProduceRequest, ShuffleProduceResponse,
};
use crate::control::state::SharedState;
use crate::types::{DatabaseId, VShardId};
pub(super) fn producer_nodes(
routing: &RoutingTable,
database_id: DatabaseId,
collection: &str,
) -> crate::Result<Vec<u64>> {
let vshard = VShardId::from_collection_in_database(database_id, collection).as_u32();
let group = routing
.group_for_vshard(vshard)
.map_err(|e| crate::Error::Internal {
detail: format!("shuffle: no group for vshard {vshard} ({collection}): {e}"),
})?;
let leader = routing
.group_info(group)
.map(|g| g.leader)
.filter(|&l| l != 0)
.ok_or_else(|| crate::Error::Internal {
detail: format!("shuffle: no leader for group {group} ({collection})"),
})?;
Ok(vec![leader])
}
pub(super) fn distinct_data_node_count(routing: &RoutingTable) -> usize {
let mut nodes: BTreeSet<u64> = BTreeSet::new();
for group_id in routing.group_ids() {
if group_id == METADATA_GROUP_ID {
continue;
}
if let Some(info) = routing.group_info(group_id)
&& info.leader != 0
{
nodes.insert(info.leader);
}
}
nodes.len()
}
pub(crate) fn register_peers_from_topology(
state: &SharedState,
transport: &nodedb_cluster::NexarTransport,
nodes: &BTreeSet<u64>,
) {
let Some(topology) = state.cluster_topology.as_ref() else {
return;
};
let topo = topology.read().unwrap_or_else(|p| p.into_inner());
for &node in nodes {
if let Some(info) = topo.get_node(node)
&& let Some(addr) = info.socket_addr()
{
transport.register_peer(node, addr);
}
}
}
pub(super) async fn send_produce(
transport: &nodedb_cluster::NexarTransport,
node: u64,
req: ShuffleProduceRequest,
) -> crate::Result<u64> {
match transport
.send_rpc(node, RaftRpc::ShuffleProduceRequest(req))
.await
{
Ok(RaftRpc::ShuffleProduceResponse(ShuffleProduceResponse {
error: None,
read_version_lsn,
})) => Ok(read_version_lsn),
Ok(RaftRpc::ShuffleProduceResponse(ShuffleProduceResponse { error: Some(e), .. })) => {
Err(crate::Error::Internal {
detail: format!("shuffle produce failed on node {node}: {e:?}"),
})
}
Ok(other) => Err(crate::Error::Internal {
detail: format!("shuffle produce: unexpected reply from node {node}: {other:?}"),
}),
Err(e) => Err(crate::Error::Internal {
detail: format!("shuffle produce RPC to node {node} failed: {e}"),
}),
}
}