MCPcopy Create free account
hub / github.com/NodeDB-Lab/nodedb / join_flow

Method join_flow

nodedb-cluster/src/raft_loop/join.rs:86–293  ·  view source on GitHub ↗

Full server-side `JoinRequest` handler. See module docs for the phase-by-phase description.

(&self, req: JoinRequest)

Source from the content-addressed store, hash-verified

84 /// Full server-side `JoinRequest` handler. See module docs for the
85 /// phase-by-phase description.
86 pub(super) async fn join_flow(&self, req: JoinRequest) -> JoinResponse {
87 // 1. Snapshot group-0 leader + clone routing under one lock.
88 let (group0_leader, routing): (u64, RoutingTable) = {
89 let mr = self.multi_raft.lock().unwrap_or_else(|p| p.into_inner());
90 let routing = mr.routing().clone();
91 let leader_id = mr
92 .group_statuses()
93 .into_iter()
94 .find(|s: &GroupStatus| s.group_id == TOPOLOGY_GROUP_ID)
95 .map(|s| s.leader_id)
96 .unwrap_or(0);
97 (leader_id, routing)
98 };
99
100 // Leader check.
101 let leader_addr_hint = if group0_leader != 0 && group0_leader != self.node_id {
102 self.topology
103 .read()
104 .unwrap_or_else(|p| p.into_inner())
105 .get_node(group0_leader)
106 .map(|n| n.addr.clone())
107 } else {
108 None
109 };
110 if let JoinDecision::Redirect { leader_addr } =
111 decide_join(group0_leader, self.node_id, leader_addr_hint)
112 {
113 warn!(
114 joining_node = req.node_id,
115 leader_id = group0_leader,
116 leader_addr = %leader_addr,
117 "JoinRequest received on non-leader; redirecting"
118 );
119 return reject(format!("{LEADER_REDIRECT_PREFIX}{leader_addr}"));
120 }
121
122 // 2. Validate the address.
123 let new_addr: SocketAddr = match req.listen_addr.parse() {
124 Ok(a) => a,
125 Err(e) => {
126 return reject(format!("invalid listen_addr '{}': {e}", req.listen_addr));
127 }
128 };
129
130 // 3. Idempotency / collision check against topology.
131 // `handle_join_request` in step 5 handles the fine-grained
132 // semantics, but we check here first so idempotent re-joins
133 // short-circuit *before* we propose any Raft conf changes.
134 let existing = self
135 .topology
136 .read()
137 .unwrap_or_else(|p| p.into_inner())
138 .get_node(req.node_id)
139 .cloned();
140 if let Some(existing) = existing {
141 if existing.addr != req.listen_addr {
142 return reject(format!(
143 "node_id {} already registered with different address {} (request: {})",

Callers 1

handle_rpcMethod · 0.80

Calls 15

decide_joinFunction · 0.85
handle_join_requestFunction · 0.85
nowFunction · 0.85
broadcast_topologyFunction · 0.85
lockMethod · 0.80
routingMethod · 0.80
get_nodeMethod · 0.80
register_peerMethod · 0.80
load_cluster_idMethod · 0.80
to_stringMethod · 0.80

Tested by

no test coverage detected