use crate::error::Result;
use crate::forward::PlanExecutor;
use crate::health;
use crate::rpc_codec::{PingRequest, RaftRpc, TopologyUpdate};
use super::super::loop_core::{CommitApplier, RaftLoop};
pub(in crate::raft_loop) const TOPOLOGY_GROUP_ID: u64 = 0;
#[derive(Debug, PartialEq, Eq)]
pub(in crate::raft_loop) enum JoinDecision {
Admit,
Redirect { leader_addr: String },
}
pub(in crate::raft_loop) fn decide_join(
group0_leader: u64,
self_node_id: u64,
leader_addr: Option<String>,
) -> JoinDecision {
if group0_leader == 0 || group0_leader == self_node_id {
JoinDecision::Admit
} else {
JoinDecision::Redirect {
leader_addr: leader_addr.unwrap_or_default(),
}
}
}
impl<A: CommitApplier, P: PlanExecutor> RaftLoop<A, P> {
pub(super) fn handle_ping_rpc(&self, req: PingRequest) -> Result<RaftRpc> {
let topo_version = {
let topo = self.topology.read().unwrap_or_else(|p| p.into_inner());
topo.version()
};
Ok(health::handle_ping(self.node_id, topo_version, &req))
}
pub(super) fn handle_topology_update_rpc(&self, update: TopologyUpdate) -> Result<RaftRpc> {
let (updated, ack) = health::handle_topology_update(self.node_id, &self.topology, &update);
if updated {
for node in &update.nodes {
if node.node_id == self.node_id {
continue;
}
match node.addr.parse::<std::net::SocketAddr>() {
Ok(addr) => self.transport.register_peer(node.node_id, addr),
Err(e) => tracing::warn!(
node_id = node.node_id,
addr = %node.addr,
error = %e,
"topology update contains unparseable peer address; skipping register_peer"
),
}
}
if let Some(catalog) = self.catalog.as_ref() {
let snap = self
.topology
.read()
.unwrap_or_else(|p| p.into_inner())
.clone();
if let Err(e) = catalog.save_topology(&snap) {
tracing::warn!(error = %e, "failed to persist topology update to catalog");
}
}
}
Ok(ack)
}
}