Full server-side `JoinRequest` handler. See module docs for the phase-by-phase description.
(&self, req: JoinRequest)
| 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: {})", |
no test coverage detected