raft_rust/raft/node.rs
1// 比较工具:截断/推进复制进度与拒绝索引时取上下界
2use std::cmp::{max, min};
3// 集合类型:票集、进度表、转发中请求、有序读队列等
4use std::collections::{BTreeSet, HashMap, HashSet, VecDeque};
5// 选举超时随机区间类型
6use std::ops::Range;
7
8// 出站消息通道:节点经此向网络层投递 Envelope
9use crossbeam::channel::Sender;
10// 排序辅助:中止转发请求时按 id 有序回复,保证可复现
11use itertools::Itertools as _;
12// 协议路径日志:选举、复制、丢弃未知发送方等
13use log::{debug, info, warn};
14// 随机选举超时,降低同时超时引发的选票瓜分
15use rand::RngExt as _;
16
17// Raft 日志条目、索引与持久化日志抽象
18use super::log::{Entry, Index, Log};
19// 成员配置:Simple/Joint 与生效状态机
20use super::membership::{Membership, MembershipEntry, MembershipState};
21// RPC 信封与消息体:选举、复制、读写、快照、成员变更
22use super::message::{Envelope, Message, ReadSequence, Request, RequestID, Response, Status};
23// 状态机接口:apply/read/snapshot/restore
24use super::state::State;
25// 默认心跳间隔、选举超时范围与单次 Append 上限
26use super::{ELECTION_TIMEOUT_RANGE, HEARTBEAT_INTERVAL, MAX_APPEND_ENTRIES};
27// 参数校验错误宏
28use crate::errinput;
29// 统一错误与 Result 类型
30use crate::error::{Error, Result};
31
32/// 节点 ID,在集群内唯一。启动时手动分配。
33pub type NodeID = u8;
34
35/// 领导者任期号。选举时单调递增。
36pub type Term = u64;
37
38/// 逻辑时钟间隔,以 tick 数量表示。
39pub type Ticks = u8;
40
41/// Raft 节点选项。
42#[derive(Clone, Debug, PartialEq)]
43// Raft 运行时可调参数集合
44pub struct Options {
45 /// 领导者心跳之间的 tick 数。
46 pub heartbeat_interval: Ticks,
47 /// 跟随者与候选人的随机选举超时范围。
48 pub election_timeout_range: Range<Ticks>,
49 /// 单条 Append 消息中最多发送的条目数。
50 pub max_append_entries: usize,
51 /// 启用 Pre-vote:真选举前先确认多数,避免分区节点抬升任期。
52 pub pre_vote: bool,
53 /// 启用 CheckQuorum:领导者在无法联系多数时下台。
54 pub check_quorum: bool,
55 /// 距上次快照 apply 了多少条后触发本地快照;0 表示关闭。
56 pub snapshot_threshold: u64,
57// 结束代码块
58}
59
60// 默认配置实现
61impl Default for Options {
62 // 构造与生产默认一致的 Options
63 fn default() -> Self {
64 // 填充各字段默认值
65 Self {
66 // 默认心跳间隔
67 heartbeat_interval: HEARTBEAT_INTERVAL,
68 // 默认选举超时随机区间
69 election_timeout_range: ELECTION_TIMEOUT_RANGE,
70 // 默认单次复制条目上限
71 max_append_entries: MAX_APPEND_ENTRIES,
72 // 默认开启 Pre-vote
73 pre_vote: true,
74 // 默认开启 CheckQuorum
75 check_quorum: true,
76 snapshot_threshold: 0, // 默认关闭,由节点配置开启
77 // 结束代码块
78 }
79 // 结束代码块
80 }
81// 结束代码块
82}
83
84/// 具有动态角色的 Raft 节点。
85pub enum Node {
86 /// 候选人(含 Pre-vote 相位)。
87 Candidate(RawNode<Candidate>),
88 /// 跟随者。
89 Follower(RawNode<Follower>),
90 /// 领导者。
91 Leader(RawNode<Leader>),
92// 结束类型定义
93}
94
95// Node 对外统一入口:创建、step、tick、查询
96impl Node {
97 /// 创建新的 Raft 节点。`peers` 为除自身外的初始投票同伴。
98 pub fn new(
99 // 本节点 ID
100 id: NodeID,
101 // 初始同伴集合(不含自身)
102 peers: HashSet<NodeID>,
103 // 持久化日志
104 log: Log,
105 // 状态机实现
106 state: Box<dyn State>,
107 // 出站消息发送端
108 tx: Sender<Envelope>,
109 // 运行时选项
110 opts: Options,
111 // 可能因参数非法失败
112 ) -> Result<Self> {
113 // 以跟随者形态构造底层 RawNode
114 let node = RawNode::new(id, peers, log, state, tx, opts)?;
115 // 单节点集群无需等待他人选票
116 if node.cluster_size() == 1 {
117 // 单节点:跳过 pre-vote,直接竞选并当选。
118 return Ok(node.into_candidate(false)?.into_leader()?.into());
119 // 结束结构/枚举构造
120 }
121 // 多节点默认以跟随者启动
122 Ok(node.into())
123 // 结束结构/枚举构造
124 }
125
126 // 查询当前节点 ID(与角色无关)
127 pub fn id(&self) -> NodeID {
128 // 按角色分支取出 id
129 match self {
130 // 候选人 id
131 Self::Candidate(node) => node.id,
132 // 跟随者 id
133 Self::Follower(node) => node.id,
134 // 领导者 id
135 Self::Leader(node) => node.id,
136 // 结束 match 分支/块
137 }
138 // 结束 match 分支/块
139 }
140
141 // 查询当前任期
142 pub fn term(&self) -> Term {
143 // 按角色分支取出 term
144 match self {
145 // 候选人任期
146 Self::Candidate(node) => node.term(),
147 // 跟随者任期
148 Self::Follower(node) => node.term(),
149 // 领导者任期
150 Self::Leader(node) => node.term(),
151 // 结束 match 分支/块
152 }
153 // 结束 match 分支/块
154 }
155
156 /// 当前生效的投票成员(含自身)。
157 pub fn voters(&self) -> BTreeSet<NodeID> {
158 // 按角色读取成员配置
159 match self {
160 // 候选人侧投票成员
161 Self::Candidate(node) => node.membership.all_voters(),
162 // 跟随者侧投票成员
163 Self::Follower(node) => node.membership.all_voters(),
164 // 领导者侧投票成员
165 Self::Leader(node) => node.membership.all_voters(),
166 // 结束 match 分支/块
167 }
168 // 结束 match 分支/块
169 }
170
171 // 处理一条入站消息并可能发生角色转换
172 pub fn step(self, msg: Envelope) -> Result<Self> {
173 // 消息目标必须是本节点,防止串包
174 assert_eq!(msg.to, self.id(), "message to other node: {msg:?}");
175
176 // 允许来自当前(可能 joint)配置中的成员或自身;未知发送方丢弃。
177 let known = match &self {
178 // 候选人已知发送方?
179 Self::Candidate(node) => node.is_known_sender(msg.from),
180 // 跟随者已知发送方?
181 Self::Follower(node) => node.is_known_sender(msg.from),
182 // 领导者已知发送方?
183 Self::Leader(node) => node.is_known_sender(msg.from),
184 // 结束 match 分支/块
185 };
186 // 未知发送方:丢弃,防配置外干扰
187 if !known {
188 // 记录丢弃原因
189 warn!("Dropping message from unknown sender: {msg:?}");
190 // 保持当前角色不变
191 return Ok(self);
192 // 结束结构/枚举构造
193 }
194 // 进入角色专属 step
195 debug!("Stepping {msg:?}");
196
197 // 按当前角色分派消息处理
198 match self {
199 // 候选人处理消息(计票/落选等)
200 Self::Candidate(node) => node.step(msg),
201 // 跟随者处理消息(复制/投票/转发)
202 Self::Follower(node) => node.step(msg),
203 // 领导者处理消息(复制确认/客户端/成员变更)
204 Self::Leader(node) => node.step(msg),
205 // 结束 match 分支/块
206 }
207 // 结束 match 分支/块
208 }
209
210 // 逻辑时钟推进一拍,驱动超时与心跳
211 pub fn tick(self) -> Result<Self> {
212 // 按角色分派 tick
213 match self {
214 // 候选人选举计时
215 Self::Candidate(node) => node.tick(),
216 // 跟随者选举超时检测
217 Self::Follower(node) => node.tick(),
218 // 领导者心跳与 CheckQuorum
219 Self::Leader(node) => node.tick(),
220 // 结束 match 分支/块
221 }
222 // 结束 match 分支/块
223 }
224// 结束代码块
225}
226
227// Candidate RawNode 提升为 Node 枚举
228impl From<RawNode<Candidate>> for Node {
229 // 包装为 Candidate 变体
230 fn from(node: RawNode<Candidate>) -> Self {
231 // 构造 Candidate 节点
232 Node::Candidate(node)
233 // 结束函数
234 }
235// 结束函数
236}
237
238// Follower RawNode 提升为 Node 枚举
239impl From<RawNode<Follower>> for Node {
240 // 包装为 Follower 变体
241 fn from(node: RawNode<Follower>) -> Self {
242 // 构造 Follower 节点
243 Node::Follower(node)
244 // 结束函数
245 }
246// 结束函数
247}
248
249// Leader RawNode 提升为 Node 枚举
250impl From<RawNode<Leader>> for Node {
251 // 包装为 Leader 变体
252 fn from(node: RawNode<Leader>) -> Self {
253 // 构造 Leader 节点
254 Node::Leader(node)
255 // 结束函数
256 }
257// 结束函数
258}
259
260// 角色标记 trait:约束 RawNode 的泛型角色
261pub trait Role {}
262
263// 带具体角色状态的 Raft 节点内核
264pub struct RawNode<R: Role> {
265 // 本节点 ID
266 id: NodeID,
267 /// 当前生效的成员配置(日志中最新成员条目,追加后即生效)。
268 membership: MembershipState,
269 // Raft 日志
270 log: Log,
271 // 状态机
272 state: Box<dyn State>,
273 // 出站通道
274 tx: Sender<Envelope>,
275 // 运行选项
276 opts: Options,
277 // 角色专属状态(Follower/Candidate/Leader)
278 role: R,
279// 结束代码块
280}
281
282// 所有角色共享的通用方法
283impl<R: Role> RawNode<R> {
284 // 仅替换角色字段,保留日志/成员/状态机等
285 fn into_role<T: Role>(self, role: T) -> RawNode<T> {
286 // 重建同配置不同角色的 RawNode
287 RawNode {
288 // 保留节点 ID
289 id: self.id,
290 // 保留成员配置状态
291 membership: self.membership,
292 // 保留日志
293 log: self.log,
294 // 保留状态机
295 state: self.state,
296 // 保留发送通道
297 tx: self.tx,
298 // 保留选项
299 opts: self.opts,
300 // 换上新角色状态
301 role,
302 // 结束代码块
303 }
304 // 结束代码块
305 }
306
307 // 从持久化 term/vote 读取当前任期
308 fn term(&self) -> Term {
309 // term 存于日志元数据 (term, voted_for)
310 self.log.get_term_vote().0
311 // 结束函数
312 }
313
314 // 除自身外的同伴集合(用于广播)
315 fn peers(&self) -> BTreeSet<NodeID> {
316 // 由成员配置推导 peers
317 self.membership.peers_of(self.id)
318 // 结束函数
319 }
320
321 // 是否接受该发送方消息:自身或当前投票成员
322 fn is_known_sender(&self, from: NodeID) -> bool {
323 // 自身回环或配置内成员
324 from == self.id || self.membership.all_voters().contains(&from)
325 // 结束函数
326 }
327
328 /// 配置中的投票节点数(joint 时为并集大小,仅用于进度 map 等)。
329 fn cluster_size(&self) -> usize {
330 // 统计 all_voters 大小
331 self.membership.all_voters().len()
332 // 结束函数
333 }
334
335 // 在配置区间内随机选举超时,打散同时超时
336 fn random_election_timeout(&self) -> Ticks {
337 // 均匀采样 [start, end)
338 rand::rng().random_range(self.opts.election_timeout_range.clone())
339 // 结束函数
340 }
341
342 /// 选举超时上界,用于 check-quorum 窗口。
343 fn election_timeout_max(&self) -> Ticks {
344 // 取 range.end-1,至少为 1
345 self.opts.election_timeout_range.end.saturating_sub(1).max(1)
346 // 结束函数
347 }
348
349 // 向指定节点发送带本任期的消息
350 fn send(&self, to: NodeID, message: Message) -> Result<()> {
351 // 封装 Envelope 并经通道投递
352 Self::send_via(&self.tx, Envelope { from: self.id, to, term: self.term(), message })
353 // 结束函数
354 }
355
356 // 底层发送:写通道并记录调试日志
357 fn send_via(tx: &Sender<Envelope>, msg: Envelope) -> Result<()> {
358 // 记录出站消息
359 debug!("Sending {msg:?}");
360 // 通道错误上浮为 Result
361 Ok(tx.send(msg)?)
362 // 结束结构/枚举构造
363 }
364
365 // 向所有同伴广播同一消息(如心跳/拉票)
366 fn broadcast(&self, message: Message) -> Result<()> {
367 // 遍历当前 peers
368 for id in self.peers() {
369 // 逐个发送克隆后的消息体
370 self.send(id, message.clone())?;
371 // 结束循环
372 }
373 // 广播完成
374 Ok(())
375 // 结束结构/枚举构造
376 }
377
378 /// 日志追加后立即应用其中的成员配置。
379 fn maybe_apply_membership_from_entries(&mut self, entries: &[Entry]) {
380 // 扫描刚追加/拼接的条目
381 for e in entries {
382 // 仅处理带 membership 的配置条目
383 if let Some(ref m) = e.membership {
384 // 记录配置切换
385 info!("Node {} adopting membership from log index {}: {m:?}", self.id, e.index);
386 // 更新本地 MembershipState
387 self.membership.apply_entry(m);
388 // 结束条件分支
389 }
390 // 结束条件分支
391 }
392 // 结束条件分支
393 }
394
395 /// 从日志恢复最新成员配置(启动时)。
396 fn restore_membership_from_log(&mut self) -> Result<()> {
397 // 读取日志中最新 membership 条目
398 if let Some((_idx, m)) = self.log.latest_membership()? {
399 // 应用到内存配置
400 self.membership.apply_entry(&m);
401 // 若该条目已提交且为 Simple,清除 pending。
402 let (commit, _) = self.log.get_commit_index();
403 // 最新配置为已提交的 Simple 时清除 pending
404 if let Some((idx, MembershipEntry::Simple(_))) = self.log.latest_membership()?
405 // 条目索引不超过 commit 才算已提交
406 && idx <= commit
407 // 进入代码块
408 {
409 // 无在途成员变更
410 self.membership.change_pending = false;
411 // 结束条件分支
412 }
413 // 结束条件分支
414 }
415 // 恢复完成
416 Ok(())
417 // 结束结构/枚举构造
418 }
419// 结束结构/枚举构造
420}
421
422// =============================================================================
423// Follower
424// =============================================================================
425
426// 跟随者角色状态
427pub struct Follower {
428 // 当前已知领导者;None 表示无主
429 leader: Option<NodeID>,
430 // 距上次收到领导者消息的 tick 计数
431 leader_seen: Ticks,
432 // 本任期随机选举超时阈值
433 election_timeout: Ticks,
434 // 已转发至领导者、等待回包的客户端请求 ID
435 forwarded: HashSet<RequestID>,
436// 结束代码块
437}
438
439// 跟随者构造辅助
440impl Follower {
441 // 创建跟随者状态:未见领导者、清空转发集
442 fn new(leader: Option<NodeID>, election_timeout: Ticks) -> Self {
443 // leader_seen 从 0 起算
444 Self { leader, leader_seen: 0, election_timeout, forwarded: HashSet::new() }
445 // 结束结构/枚举构造
446 }
447// 结束结构/枚举构造
448}
449
450// Follower 实现 Role 标记
451impl Role for Follower {}
452
453// 跟随者节点行为:选举超时、复制、投票、客户端代理
454impl RawNode<Follower> {
455 // 以跟随者身份初始化集群节点
456 fn new(
457 // 节点 ID
458 id: NodeID,
459 // 初始同伴
460 peers: HashSet<NodeID>,
461 // 日志
462 log: Log,
463 // 状态机
464 state: Box<dyn State>,
465 // 发送通道
466 tx: Sender<Envelope>,
467 // 选项
468 opts: Options,
469 // 初始化可能失败
470 ) -> Result<Self> {
471 // ID 不得出现在 peers 中
472 if peers.contains(&id) {
473 // 参数错误
474 return errinput!("node ID {id} can't be in peers");
475 // 结束条件分支
476 }
477 // 用自身与 peers 引导初始 Simple 配置
478 let membership = MembershipState::bootstrap(id, peers.iter().copied());
479 // 占位角色,稍后填随机超时
480 let role = Follower::new(None, 0);
481 // 组装 RawNode
482 let mut node = Self { id, membership, log, state, tx, opts, role };
483 // 为跟随者抽一次随机选举超时
484 node.role.election_timeout = node.random_election_timeout();
485 // 从持久化日志恢复成员配置
486 node.restore_membership_from_log()?;
487 // 补齐已提交但未 apply 的条目
488 node.maybe_apply()?;
489 // 返回就绪跟随者
490 Ok(node)
491 // 结束结构/枚举构造
492 }
493
494 /// `use_prevote`: true 时先 Pre-vote;false 时直接真选举(单节点或关闭 pre_vote)。
495 fn into_candidate(mut self, use_prevote: bool) -> Result<RawNode<Candidate>> {
496 // 角色切换前中止未完成的转发请求
497 self.abort_forwarded()?;
498 // 切换前尽量应用已提交日志
499 self.maybe_apply()?;
500 // 新角色使用新的随机选举超时
501 let election_timeout = self.random_election_timeout();
502 // 根据参数与配置决定是否走 Pre-vote
503 let phase = if use_prevote && self.opts.pre_vote {
504 // Pre-vote 相位:不抬升本地任期
505 ElectionPhase::PreVote
506 // 不走 Pre-vote 时直接真选举
507 } else {
508 // 真选举相位:立即 term+1 并持久化自投
509 ElectionPhase::Election
510 // 结束条件分支
511 };
512 // 切换为候选人角色
513 let mut node = self.into_role(Candidate::new(election_timeout, phase));
514 // 按相位发起预选或真选
515 match node.role.phase {
516 // 广播 PreCampaign
517 ElectionPhase::PreVote => node.pre_campaign()?,
518 // 广播 Campaign 并自投
519 ElectionPhase::Election => node.campaign()?,
520 // 结束 match 分支/块
521 }
522 // 返回候选人节点
523 Ok(node)
524 // 结束结构/枚举构造
525 }
526
527 // 转为跟随者:跟随已知领导或发现更高任期
528 fn into_follower(mut self, term: Term, leader: Option<NodeID>) -> Result<RawNode<Follower>> {
529 // term 0 非法,协议从 1 起
530 assert_ne!(term, 0, "can't become follower in term 0");
531 // 切换前中止转发中的客户端请求
532 self.abort_forwarded()?;
533
534 // 已知领导者:同任期追随
535 if let Some(leader) = leader {
536 // 领导者必须在当前投票配置内
537 assert!(self.membership.all_voters().contains(&leader), "leader is not a voter");
538 // 同任期不应重复设置领导
539 assert_eq!(self.role.leader, None, "already have leader in term");
540 // 追随的领导任期须与本地一致
541 assert_eq!(term, self.term(), "can't follow leader in different term");
542 // 记录开始跟随
543 info!("Following leader {leader} in term {term}");
544 // 写入领导 ID,保留原选举超时
545 self.role = Follower::new(Some(leader), self.role.election_timeout);
546 // 无已知领导时的分支
547 } else {
548 // 无主跟随者:必须来自更高任期
549 assert_ne!(term, self.term(), "can't become leaderless follower in current term");
550 // 记录发现新任期
551 info!("Discovered new term {term}");
552 // 持久化新任期并清空投票
553 self.log.set_term_vote(term, None)?;
554 // 无主并重置随机选举超时
555 self.role = Follower::new(None, self.random_election_timeout());
556 // 结束代码块
557 }
558 // 返回跟随者
559 Ok(self)
560 // 结束结构/枚举构造
561 }
562
563 // 跟随者消息处理主循环
564 fn step(mut self, msg: Envelope) -> Result<Node> {
565 // Pre-vote 消息:不因更高 term 而切换任期。
566 if matches!(msg.message, Message::PreCampaign { .. } | Message::PreCampaignResponse { .. }) {
567 // 委托 step_prevote
568 return self.step_prevote(msg);
569 // 结束条件分支
570 }
571
572 // 过期任期消息直接丢弃
573 if msg.term < self.term() {
574 // 调试记录
575 debug!("Dropping message from past term: {msg:?}");
576 // 保持跟随者
577 return Ok(self.into());
578 // 结束结构/枚举构造
579 }
580 // 更高任期:先无主降级再重放该消息
581 if msg.term > self.term() {
582 // 递归 step 以在新任期处理
583 return self.into_follower(msg.term, None)?.step(msg);
584 // 结束条件分支
585 }
586
587 // 来自当前领导者的任意消息重置见领导计时
588 if Some(msg.from) == self.role.leader {
589 // leader_seen 清零,推迟选举超时
590 self.role.leader_seen = 0;
591 // 结束条件分支
592 }
593
594 // 按消息类型分支
595 match msg.message {
596 // 心跳:同步 commit 并应答 match/read_seq
597 Message::Heartbeat { last_index, commit_index, read_seq } => {
598 // commit 不得越过领导者 last_index
599 assert!(commit_index <= last_index, "commit_index after last_index");
600 // 校验心跳来源是否为当前领导
601 match self.role.leader {
602 // 非当前领导的心跳忽略
603 Some(leader) if msg.from != leader => {
604 // 告警双重领导迹象
605 warn!(
606 // 说明当前领导
607 "node {} ignoring Heartbeat from {} (current leader {})",
608 // 附带冲突方
609 self.id, msg.from, leader
610 // 业务逻辑
611 );
612 // 不处理冲突心跳
613 return Ok(self.into());
614 // 结束结构/枚举构造
615 }
616 // 已有正确领导:继续
617 Some(_) => {}
618 // 无主时通过心跳确立领导者
619 None => self = self.into_follower(msg.term, Some(msg.from))?,
620 // 结束结构/枚举构造
621 }
622 // 日志匹配则回报 match,否则 0 触发探测
623 let match_index = if self.log.has(last_index, msg.term)? { last_index } else { 0 };
624 // 心跳应答携带 match_index 与读序列
625 self.send(msg.from, Message::HeartbeatResponse { match_index, read_seq })?;
626 // 仅当日志匹配时才可安全推进本地 commit
627 if match_index != 0 && commit_index > self.log.get_commit_index().0 {
628 // 提交到领导者给出的 commit_index
629 self.log.commit(commit_index)?;
630 // 提交后应用到状态机
631 self.maybe_apply()?;
632 // 结束条件分支
633 }
634 // 结束条件分支
635 }
636
637 // 日志复制 AppendEntries
638 Message::Append { base_index, base_term, entries } => {
639 // 一致性检查:base 应对齐首条前一索引
640 if let Some(first) = entries.first() {
641 // base_index 等于 first.index - 1
642 assert_eq!(base_index, first.index - 1, "base index mismatch");
643 // 结束条件分支
644 }
645 // 校验 Append 是否来自当前领导
646 match self.role.leader {
647 // 非领导 Append 忽略
648 Some(leader) if msg.from != leader => {
649 // 告警
650 warn!(
651 // 冲突来源
652 "node {} ignoring Append from {} (current leader {})",
653 // 当前领导
654 self.id, msg.from, leader
655 // 业务逻辑
656 );
657 // 丢弃
658 return Ok(self.into());
659 // 结束结构/枚举构造
660 }
661 // 已有领导
662 Some(_) => {}
663 // 无主时确立领导
664 None => self = self.into_follower(msg.term, Some(msg.from))?,
665 // 结束结构/枚举构造
666 }
667 // prevLog 匹配(或 base=0)则接受拼接
668 if base_index == 0 || self.log.has(base_index, base_term)? {
669 // 匹配点为最后一条新条目或 base
670 let match_index = entries.last().map(|e| e.index).unwrap_or(base_index);
671 // splice:截断冲突后缀并追加
672 self.log.splice(entries.clone())?;
673 // 配置在日志中出现后立即生效(论文联合共识)。
674 self.maybe_apply_membership_from_entries(&entries);
675 // 成功应答 match_index
676 self.send(msg.from, Message::AppendResponse { match_index, reject_index: 0 })?;
677 // 普通业务条目(非成员变更)分支
678 } else {
679 // 拒绝:回报可回退的 reject_index
680 let reject_index = min(base_index, self.log.get_last_index().0 + 1);
681 // match_index=0 表示拒绝
682 self.send(msg.from, Message::AppendResponse { reject_index, match_index: 0 })?;
683 // 结束结构/枚举构造
684 }
685 // 结束结构/枚举构造
686 }
687
688 // 线性一致读探测:跟随者仅回显 seq
689 Message::Read { seq } => {
690 // 校验来源为当前领导
691 match self.role.leader {
692 // 非领导 Read 忽略
693 Some(leader) if msg.from != leader => {
694 // 告警
695 warn!(
696 // 冲突说明
697 "node {} ignoring Read from {} (current leader {})",
698 // 当前领导
699 self.id, msg.from, leader
700 // 业务逻辑
701 );
702 // 丢弃
703 return Ok(self.into());
704 // 结束结构/枚举构造
705 }
706 // 已有领导
707 Some(_) => {}
708 // 无主时确立领导
709 None => self = self.into_follower(msg.term, Some(msg.from))?,
710 // 结束结构/枚举构造
711 }
712 // 确认本节点存活,供领导者统计读多数
713 self.send(msg.from, Message::ReadResponse { seq })?;
714 // 结束结构/枚举构造
715 }
716
717 // 真选举投票请求 RequestVote
718 Message::Campaign { last_index, last_term } => {
719 // 本任期已投票给他人则拒绝
720 if let (_, Some(vote)) = self.log.get_term_vote()
721 // 仅允许投给已记录的候选人
722 && msg.from != vote
723 // 进入代码块
724 {
725 // 拒绝票
726 self.send(msg.from, Message::CampaignResponse { vote: false })?;
727 // 结束处理
728 return Ok(self.into());
729 // 结束结构/枚举构造
730 }
731 // 比较候选人日志新旧(term 优先,再 index)
732 let (log_index, log_term) = self.log.get_last_index();
733 // 本地日志更新则拒票,保证选主带最新日志
734 if log_term > last_term || log_term == last_term && log_index > last_index {
735 // 拒绝
736 self.send(msg.from, Message::CampaignResponse { vote: false })?;
737 // 返回
738 return Ok(self.into());
739 // 结束结构/枚举构造
740 }
741 // 授予选票
742 info!("Voting for {} in term {} election", msg.from, msg.term);
743 // 持久化 term 与 voted_for
744 self.log.set_term_vote(msg.term, Some(msg.from))?;
745 // 赞成票
746 self.send(msg.from, Message::CampaignResponse { vote: true })?;
747 // 结束结构/枚举构造
748 }
749
750 // 客户端请求:跟随者代理转发领导者
751 Message::ClientRequest { id, request: _ } => {
752 // 仅接受本节点注入的客户端请求;其它 from 直接 Abort,避免 panic。
753 if msg.from != self.id {
754 // 外来 ClientRequest 拒绝
755 warn!(
756 // 说明节点与来源
757 "node {} rejecting ClientRequest from foreign sender {}",
758 // 外来 from
759 self.id, msg.from
760 // 业务逻辑
761 );
762 // 回复 Abort
763 self.send(
764 // 目标为外来发送方
765 msg.from,
766 // 错误响应
767 Message::ClientResponse { id, response: Err(Error::Abort) },
768 // 完成发送并向上传播错误
769 )?;
770 // 结束
771 return Ok(self.into());
772 // 结束结构/枚举构造
773 }
774 // 已知领导则转发
775 if let Some(leader) = self.role.leader {
776 // 记录转发
777 debug!("Forwarding request to leader {leader}: {msg:?}");
778 // 跟踪请求 id 以便回包
779 self.role.forwarded.insert(id);
780 // 原样转给领导者
781 self.send(leader, msg.message)?;
782 // 无已知领导时的分支
783 } else {
784 // 无主:无法服务,Abort
785 self.send(msg.from, Message::ClientResponse { id, response: Err(Error::Abort) })?;
786 // 结束结构/枚举构造
787 }
788 // 结束结构/枚举构造
789 }
790
791 // 领导者对转发请求的响应
792 Message::ClientResponse { id, response } => {
793 // 必须来自当前领导
794 if Some(msg.from) != self.role.leader {
795 // 非领导响应忽略
796 warn!(
797 // 告警
798 "node {} ignoring ClientResponse from non-leader {}",
799 // 来源
800 self.id, msg.from
801 // 业务逻辑
802 );
803 // 丢弃
804 return Ok(self.into());
805 // 结束结构/枚举构造
806 }
807 // 仅当仍在转发集合中才回传客户端
808 if self.role.forwarded.remove(&id) {
809 // 回环给本节点客户端出口
810 self.send(self.id, Message::ClientResponse { id, response })?;
811 // 结束结构/枚举构造
812 }
813 // 结束结构/枚举构造
814 }
815
816 // 跟随者不收集选票响应
817 Message::CampaignResponse { .. } => {}
818
819 // Pre-vote 已在入口处理
820 Message::PreCampaign { .. } | Message::PreCampaignResponse { .. } => {
821 // 不可达分支
822 unreachable!("handled above")
823 // 结束结构/枚举构造
824 }
825
826 // 安装快照:落后跟随者用快照追赶
827 Message::InstallSnapshot {
828 // 快照覆盖到的最后索引
829 last_included_index,
830 // 对应任期
831 last_included_term,
832 // 状态机快照字节
833 data,
834 // 快照中的成员配置
835 membership,
836 // 匹配该模式后进入处理体
837 } => {
838 // 校验来源为当前领导
839 match self.role.leader {
840 // 非领导快照忽略
841 Some(leader) if msg.from != leader => {
842 // 告警
843 warn!(
844 // 冲突
845 "node {} ignoring InstallSnapshot from {} (current leader {})",
846 // 当前领导
847 self.id, msg.from, leader
848 // 业务逻辑
849 );
850 // 丢弃
851 return Ok(self.into());
852 // 结束结构/枚举构造
853 }
854 // 已有领导
855 Some(_) => {}
856 // 无主时确立领导
857 None => self = self.into_follower(msg.term, Some(msg.from))?,
858 // 结束结构/枚举构造
859 }
860 // 忽略过期 / 重复快照,避免把状态机回滚到更旧点。
861 let (snap_idx, _) = self.log.get_snapshot_meta();
862 // 状态机已应用点
863 let applied = self.state.get_applied_index();
864 // 只接受不落后于本地快照/已应用的快照
865 if last_included_index >= snap_idx && last_included_index >= applied {
866 // 恢复状态机到快照点
867 self.state.restore(&data, last_included_index)?;
868 // 持久化快照字节便于重启
869 self.log.engine.set(&super::log::Key::SnapshotData.encode(), data.clone())?;
870 // 日志丢弃快照点之前条目并更新元数据
871 self.log.reset_with_snapshot(last_included_index, last_included_term)?;
872 // 采用快照内成员配置
873 self.membership.apply_entry(&membership);
874 // 快照中的配置视为已提交生效。
875 if matches!(membership, MembershipEntry::Simple(_)) {
876 // 无在途变更
877 self.membership.change_pending = false;
878 // 结束条件分支
879 }
880 // 结束条件分支
881 }
882 // 无论是否应用都应答,避免领导者卡住
883 self.send(
884 // 回给领导者
885 msg.from,
886 // 携带 last_included_index 推进 match
887 Message::InstallSnapshotResponse { last_included_index },
888 // 完成发送并向上传播错误
889 )?;
890 // 结束结构/枚举构造
891 }
892
893 // 跟随者不处理快照应答
894 Message::InstallSnapshotResponse { .. } => {}
895
896 // 领导者专属响应不应到达跟随者
897 Message::HeartbeatResponse { .. }
898 // Append 应答
899 | Message::AppendResponse { .. }
900 // 读确认应答
901 | Message::ReadResponse { .. } => {
902 // 协议错误:panic 暴露 bug
903 panic!("follower received unexpected message {msg:?}")
904 // 结束结构/枚举构造
905 }
906 // 结束结构/枚举构造
907 };
908 // 跟随者 step 完成,包装回 Node
909 Ok(self.into())
910 // 结束结构/枚举构造
911 }
912
913 // 处理 Pre-vote:不修改本地 term/vote
914 fn step_prevote(self, msg: Envelope) -> Result<Node> {
915 // 按 Pre-vote 消息类型
916 match msg.message {
917 // 预选票请求
918 Message::PreCampaign { last_index, last_term } => {
919 // 仍能联系当前领导则拒绝预选票。
920 if self.role.leader.is_some() && self.role.leader_seen < self.role.election_timeout {
921 // 拒预选票
922 self.send(msg.from, Message::PreCampaignResponse { vote: false })?;
923 // 返回
924 return Ok(self.into());
925 // 结束结构/枚举构造
926 }
927 // 日志新旧检查,规则同真投票
928 let (log_index, log_term) = self.log.get_last_index();
929 // log_ok:候选人不落后
930 let log_ok =
931 // 本地更新则 false
932 !(log_term > last_term || log_term == last_term && log_index > last_index);
933 // 预选票请求的 envelope.term 为 intended term(current+1),不更新本地 term。
934 let term_ok = msg.term > self.term();
935 // 同时满足日志与任期条件才授预选票
936 let vote = log_ok && term_ok;
937 // 调试预投票结果
938 debug!(
939 // 候选人与 intended term
940 "Pre-vote for {} (term {}): vote={vote} log_ok={log_ok} term_ok={term_ok}",
941 // 来源与任期
942 msg.from, msg.term
943 // 业务逻辑
944 );
945 // 发送预选票结果(不持久化)
946 self.send(msg.from, Message::PreCampaignResponse { vote })?;
947 // 结束结构/枚举构造
948 }
949 // 跟随者不累计预选票
950 Message::PreCampaignResponse { .. } => {
951 // 忽略
952 // 跟随者不收集预选票。
953 }
954 // 其它消息忽略
955 _ => {}
956 // 结束结构/枚举构造
957 }
958 // Pre-vote 处理结束
959 Ok(self.into())
960 // 结束结构/枚举构造
961 }
962
963 // 跟随者时钟:累计未见领导时间
964 fn tick(mut self) -> Result<Node> {
965 // 每个 tick 增加 leader_seen
966 self.role.leader_seen += 1;
967 // 达到选举超时则发起竞选(优先 Pre-vote)
968 if self.role.leader_seen >= self.role.election_timeout {
969 // 转为候选人并进入选举流程
970 return Ok(self.into_candidate(true)?.into());
971 // 结束结构/枚举构造
972 }
973 // 未超时:保持跟随者
974 Ok(self.into())
975 // 结束结构/枚举构造
976 }
977
978 // 角色切换时中止所有在途转发请求
979 fn abort_forwarded(&mut self) -> Result<()> {
980 // 取出并按 id 排序,保证回复顺序稳定
981 for id in std::mem::take(&mut self.role.forwarded).into_iter().sorted() {
982 // 调试中止
983 debug!("Aborting forwarded request {id}");
984 // 向本节点客户端回 Abort
985 self.send(self.id, Message::ClientResponse { id, response: Err(Error::Abort) })?;
986 // 结束结构/枚举构造
987 }
988 // 清理完成
989 Ok(())
990 // 结束结构/枚举构造
991 }
992
993 // 将已提交未应用的日志应用到状态机
994 fn maybe_apply(&mut self) -> Result<()> {
995 // 从 applied_index 之后扫描
996 let mut iter = self.log.scan_apply(self.state.get_applied_index());
997 // 逐条取出
998 while let Some(entry) = iter.next().transpose()? {
999 // 调试 apply
1000 debug!("Applying {entry:?}");
1001 // 成员条目提交时更新 pending 等
1002 if let Some(ref m) = entry.membership {
1003 // on_commit 处理配置提交语义
1004 self.membership.on_commit(m);
1005 // 结束条件分支
1006 }
1007 // 成员变更条目:状态机收到 noop 式 command=None。
1008 _ = self.state.apply(entry);
1009 // 结束条件分支
1010 }
1011 // apply 完成
1012 Ok(())
1013 // 结束结构/枚举构造
1014 }
1015// 结束结构/枚举构造
1016}
1017
1018// =============================================================================
1019// Candidate (+ Pre-vote phase)
1020// =============================================================================
1021
1022// 选举相位为轻量可拷贝枚举,便于匹配与日志
1023#[derive(Clone, Copy, Debug, PartialEq, Eq)]
1024// 选举相位:预选或真选
1025enum ElectionPhase {
1026 // Pre-vote:试探多数,不抬升任期
1027 PreVote,
1028 // 真选举:term+1 并持久化自投
1029 Election,
1030// 结束类型定义
1031}
1032
1033// 候选人角色状态
1034pub struct Candidate {
1035 // 已获(预)选票的节点集合
1036 votes: HashSet<NodeID>,
1037 // 本轮选举已进行的 tick
1038 election_duration: Ticks,
1039 // 本轮选举超时阈值
1040 election_timeout: Ticks,
1041 // 当前处于 PreVote 还是 Election
1042 phase: ElectionPhase,
1043// 结束代码块
1044}
1045
1046// 构造候选人状态
1047impl Candidate {
1048 // 空票集、计时归零
1049 fn new(election_timeout: Ticks, phase: ElectionPhase) -> Self {
1050 // 初始化各字段
1051 Self { votes: HashSet::new(), election_duration: 0, election_timeout, phase }
1052 // 结束结构/枚举构造
1053 }
1054// 结束结构/枚举构造
1055}
1056
1057// Candidate 实现 Role
1058impl Role for Candidate {}
1059
1060// 候选人行为:预选、拉票、计票、当选/落选
1061impl RawNode<Candidate> {
1062 // 落选或发现更高任期,转为跟随者
1063 fn into_follower(mut self, term: Term, leader: Option<NodeID>) -> Result<RawNode<Follower>> {
1064 // 新跟随者使用新随机超时
1065 let election_timeout = self.random_election_timeout();
1066 // 跟随已知领导(同任期)
1067 if let Some(leader) = leader {
1068 // 任期必须一致
1069 assert_eq!(term, self.term(), "can't follow leader in different term");
1070 // 记录落选跟随
1071 info!("Lost election, following leader {leader} in term {term}");
1072 // 切换角色并带上领导 ID
1073 Ok(self.into_role(Follower::new(Some(leader), election_timeout)))
1074 // 无已知领导时的分支
1075 } else {
1076 // 无主:必须更高任期
1077 assert_ne!(term, self.term(), "can't become leaderless follower in current term");
1078 // 记录新任期
1079 info!("Discovered new term {term}");
1080 // 持久化新任期、清空投票
1081 self.log.set_term_vote(term, None)?;
1082 // 无主跟随者
1083 Ok(self.into_role(Follower::new(None, election_timeout)))
1084 // 结束结构/枚举构造
1085 }
1086 // 结束结构/枚举构造
1087 }
1088
1089 // 获得多数真选票后转领导者
1090 fn into_leader(self) -> Result<RawNode<Leader>> {
1091 // 校验持久化投票状态
1092 let (term, vote) = self.log.get_term_vote();
1093 // 领导任期不可为 0
1094 assert_ne!(term, 0, "leaders can't have term 0");
1095 // 当选者必须已自投
1096 assert_eq!(vote, Some(self.id), "leader did not vote for self");
1097 // 仅真选举胜出可当领导(Pre-vote 不够)
1098 assert_eq!(self.role.phase, ElectionPhase::Election, "must win real election");
1099
1100 // 记录当选
1101 info!("Won election for term {term}, becoming leader");
1102 // 当前同伴集
1103 let peers = self.peers();
1104 // 日志末端作为初始 next_index 基准
1105 let (last_index, _) = self.log.get_last_index();
1106 // 构造领导者:各 peer 进度从 last+1 起
1107 let mut node = self.into_role(Leader::new(peers, last_index));
1108 // 空条目(no-op)确立本任期领导权并推进提交
1109 node.propose(None)?;
1110 // 尝试提交并应用(单节点立即提交)
1111 let _ = node.maybe_commit_and_apply()?;
1112 // 立即广播心跳宣告领导
1113 node.heartbeat()?;
1114 // 返回领导者
1115 Ok(node)
1116 // 结束结构/枚举构造
1117 }
1118
1119 // 候选人消息处理
1120 fn step(mut self, msg: Envelope) -> Result<Node> {
1121 // Pre-vote:不因更高 term 切换;可向日志足够新的节点授予预选票。
1122 if let Message::PreCampaign { last_index, last_term } = msg.message {
1123 // 本地日志末端
1124 let (log_index, log_term) = self.log.get_last_index();
1125 // 日志是否不比候选人新
1126 let log_ok =
1127 // term/index 比较
1128 !(log_term > last_term || log_term == last_term && log_index > last_index);
1129 // intended term 应大于本地 term;同 term 的预选也允许(大家都在 term T 抢 T+1)。
1130 let term_ok = msg.term > self.term();
1131 // 日志与任期都通过才投预选票
1132 let vote = log_ok && term_ok;
1133 // 回复预选票
1134 self.send(msg.from, Message::PreCampaignResponse { vote })?;
1135 // 保持候选人
1136 return Ok(self.into());
1137 // 结束结构/枚举构造
1138 }
1139 // 收集预选票响应
1140 if let Message::PreCampaignResponse { vote } = msg.message {
1141 // 仅 PreVote 相位且赞成票有效
1142 if self.role.phase == ElectionPhase::PreVote && vote {
1143 // 发送方须在投票配置内
1144 if self.membership.all_voters().contains(&msg.from) {
1145 // 记入票集
1146 self.role.votes.insert(msg.from);
1147 // 结束条件分支
1148 }
1149 // 预选票达多数则启动真选举
1150 if self.quorum_reached(&self.role.votes) {
1151 // 记录 Pre-vote 成功
1152 info!("Pre-vote won, starting real election");
1153 // 抬升任期、自投、广播 Campaign
1154 self.campaign()?;
1155 // 单节点或已够多数时立即当选。
1156 if self.quorum_reached(&self.role.votes)
1157 // 相位已切到 Election
1158 && self.role.phase == ElectionPhase::Election
1159 // 进入代码块
1160 {
1161 // 转为领导者
1162 return Ok(self.into_leader()?.into());
1163 // 结束结构/枚举构造
1164 }
1165 // 结束结构/枚举构造
1166 }
1167 // 结束结构/枚举构造
1168 }
1169 // 预选票未达多数:继续等待
1170 return Ok(self.into());
1171 // 结束结构/枚举构造
1172 }
1173
1174 // 过期任期消息丢弃
1175 if msg.term < self.term() {
1176 // 调试
1177 debug!("Dropping message from past term: {msg:?}");
1178 // 保持
1179 return Ok(self.into());
1180 // 结束结构/枚举构造
1181 }
1182 // 更高任期:降为无主跟随者并重放
1183 if msg.term > self.term() {
1184 // 递归处理
1185 return self.into_follower(msg.term, None)?.step(msg);
1186 // 结束条件分支
1187 }
1188
1189 // 真选举相关消息
1190 match msg.message {
1191 // 获得赞成票
1192 Message::CampaignResponse { vote: true } => {
1193 // 仅在真选举相位计票
1194 if self.role.phase == ElectionPhase::Election {
1195 // 投票者须在配置内
1196 if self.membership.all_voters().contains(&msg.from) {
1197 // 记票
1198 self.role.votes.insert(msg.from);
1199 // 结束条件分支
1200 }
1201 // 达多数则当选
1202 if self.quorum_reached(&self.role.votes) {
1203 // 进入领导者
1204 return Ok(self.into_leader()?.into());
1205 // 结束结构/枚举构造
1206 }
1207 // 结束结构/枚举构造
1208 }
1209 // 结束结构/枚举构造
1210 }
1211 // 反对票:忽略,等超时再选
1212 Message::CampaignResponse { vote: false } => {}
1213 // 同任期他人拉票:候选人已自投,拒票
1214 Message::Campaign { .. } => {
1215 // 明确拒绝
1216 self.send(msg.from, Message::CampaignResponse { vote: false })?;
1217 // 结束结构/枚举构造
1218 }
1219 // 收到领导类消息:说明已有领导,立即追随
1220 Message::Heartbeat { .. }
1221 // Append
1222 | Message::Append { .. }
1223 // Read 探测
1224 | Message::Read { .. }
1225 // 快照
1226 | Message::InstallSnapshot { .. } => {
1227 // 降为该发送方的跟随者并重放消息
1228 return self.into_follower(msg.term, Some(msg.from))?.step(msg);
1229 // 结束结构/枚举构造
1230 }
1231 // 选举期间拒绝客户端请求
1232 Message::ClientRequest { id, request: _ } => {
1233 // Abort
1234 self.send(msg.from, Message::ClientResponse { id, response: Err(Error::Abort) })?;
1235 // 结束结构/枚举构造
1236 }
1237 // 忽略快照应答
1238 Message::InstallSnapshotResponse { .. } => {}
1239 // 不应出现的领导侧响应
1240 Message::HeartbeatResponse { .. }
1241 // Append 应答
1242 | Message::AppendResponse { .. }
1243 // 读应答
1244 | Message::ReadResponse { .. }
1245 // 客户端应答
1246 | Message::ClientResponse { .. }
1247 // PreCampaign 已处理
1248 | Message::PreCampaign { .. }
1249 // PreCampaignResponse 已处理
1250 | Message::PreCampaignResponse { .. } => {
1251 // 异常消息 panic
1252 panic!("unexpected message {msg:?}")
1253 // 结束结构/枚举构造
1254 }
1255 // 结束结构/枚举构造
1256 }
1257 // 候选人 step 结束
1258 Ok(self.into())
1259 // 结束结构/枚举构造
1260 }
1261
1262 // 候选人时钟:选举超时则重启一轮
1263 fn tick(mut self) -> Result<Node> {
1264 // 累计本轮时长
1265 self.role.election_duration += 1;
1266 // 超时
1267 if self.role.election_duration >= self.role.election_timeout {
1268 // 超时后重新从 Pre-vote 开始(若启用),避免 term 风暴。
1269 if self.opts.pre_vote {
1270 // 重置为 PreVote 并广播
1271 self.pre_campaign()?;
1272 // 不走 Pre-vote 时直接真选举
1273 } else {
1274 // 关闭 Pre-vote:直接真选举
1275 self.campaign()?;
1276 // 结束条件分支
1277 }
1278 // 结束条件分支
1279 }
1280 // 未超时或重启后继续候选
1281 Ok(self.into())
1282 // 结束结构/枚举构造
1283 }
1284
1285 /// 当前票集是否达到(联合)法定人数。候选人总是含自己。
1286 fn quorum_reached(&self, votes: &HashSet<NodeID>) -> bool {
1287 // 转为有序集合供 has_quorum
1288 let matched: BTreeSet<NodeID> = votes.iter().copied().collect();
1289 // 委托成员配置的多数判断(joint 需双多数)
1290 self.membership.has_quorum(&matched)
1291 // 结束函数
1292 }
1293
1294 // 发起 Pre-vote 轮次
1295 fn pre_campaign(&mut self) -> Result<()> {
1296 // 新随机超时
1297 let timeout = self.random_election_timeout();
1298 // 重置角色为 PreVote 空票集
1299 self.role = Candidate::new(timeout, ElectionPhase::PreVote);
1300 // 预选票先计入自己
1301 self.role.votes.insert(self.id);
1302 // 单节点:直接进入真选举。
1303 if self.cluster_size() == 1 {
1304 // campaign
1305 return self.campaign();
1306 // 结束条件分支
1307 }
1308 // 带上日志末端供他人比较新旧
1309 let (last_index, last_term) = self.log.get_last_index();
1310 // Pre-vote 使用 intended term = current+1,接收方不持久化。
1311 let intended = self.term() + 1;
1312 // 记录预选开始
1313 info!("Starting pre-vote for intended term {intended}");
1314 // 向每个同伴发 PreCampaign
1315 for id in self.peers() {
1316 // 直接构造 Envelope 以写入 intended term
1317 Self::send_via(
1318 // 使用本节点通道
1319 &self.tx,
1320 // 信封
1321 Envelope {
1322 // 发送方
1323 from: self.id,
1324 // 目标同伴
1325 to: id,
1326 // intended term 而非本地 term
1327 term: intended,
1328 // 预选消息体
1329 message: Message::PreCampaign { last_index, last_term },
1330 // 业务逻辑
1331 },
1332 // 完成发送并向上传播错误
1333 )?;
1334 // 结束结构/枚举构造
1335 }
1336 // 已有自己一票;若自己已构成多数(不应发生在多节点)则直接竞选。
1337 if self.quorum_reached(&self.role.votes) {
1338 // 进入 campaign
1339 self.campaign()?;
1340 // 结束条件分支
1341 }
1342 // 预选发起完成
1343 Ok(())
1344 // 结束结构/枚举构造
1345 }
1346
1347 // 发起真选举:term+1、自投、广播 RequestVote
1348 fn campaign(&mut self) -> Result<()> {
1349 // 新任期
1350 let term = self.term() + 1;
1351 // 记录
1352 info!("Starting new election for term {term}");
1353 // 新超时
1354 let timeout = self.random_election_timeout();
1355 // 重置为 Election 相位
1356 self.role = Candidate::new(timeout, ElectionPhase::Election);
1357 // 自投票
1358 self.role.votes.insert(self.id);
1359 // 持久化新任期与 voted_for=self
1360 self.log.set_term_vote(term, Some(self.id))?;
1361 // 日志位置用于投票比较
1362 let (last_index, last_term) = self.log.get_last_index();
1363 // 广播 Campaign
1364 self.broadcast(Message::Campaign { last_index, last_term })?;
1365 // 单节点等已达多数:由调用方 into_leader
1366 if self.quorum_reached(&self.role.votes) {
1367 // 单节点等情况:立即当选由调用方 into_leader;此处仅标记。
1368 // 实际 into_leader 在 step/tick 路径;对单节点 Node::new 会链式调用。
1369 }
1370 // 真选举发起完成
1371 Ok(())
1372 // 结束结构/枚举构造
1373 }
1374// 结束结构/枚举构造
1375}
1376
1377// =============================================================================
1378// Leader
1379// =============================================================================
1380
1381// 领导者角色状态:复制进度、读写队列、心跳与活跃性
1382pub struct Leader {
1383 // 每个跟随者的 match/next/read_seq 进度
1384 progress: HashMap<NodeID, Progress>,
1385 // 在途写请求:日志索引到客户端
1386 writes: HashMap<Index, Write>,
1387 /// 在途成员变更客户端请求(joint 条目索引)。
1388 membership_writes: HashMap<Index, Write>,
1389 // 有序读请求队列(按 read_seq)
1390 reads: VecDeque<Read>,
1391 // 领导者发出的读序列号,单调递增
1392 read_seq: ReadSequence,
1393 // 距上次心跳的 tick
1394 since_heartbeat: Ticks,
1395 /// 自上次收到该 peer 有效响应以来的 tick;自身不计入。
1396 peer_seen: HashMap<NodeID, Ticks>,
1397// 结束代码块
1398}
1399
1400// 单跟随者复制与读确认进度
1401struct Progress {
1402 // 已确认匹配的最高日志索引
1403 match_index: Index,
1404 // 下一条待发送的日志索引
1405 next_index: Index,
1406 // 该跟随者已确认的最大读序列
1407 read_seq: ReadSequence,
1408// 结束类型定义
1409}
1410
1411// 进度推进辅助
1412impl Progress {
1413 // 跟随者日志匹配点前进
1414 fn advance(&mut self, match_index: Index) -> bool {
1415 // 不前进则忽略陈旧应答
1416 if match_index <= self.match_index {
1417 // 返回 false
1418 return false;
1419 // 结束条件分支
1420 }
1421 // 更新 match_index
1422 self.match_index = match_index;
1423 // next 至少为 match+1
1424 self.next_index = max(self.next_index, match_index + 1);
1425 // 发生了有效推进
1426 true
1427 // 结束代码块
1428 }
1429
1430 // 读序列确认前进
1431 fn advance_read(&mut self, read_seq: ReadSequence) -> bool {
1432 // 陈旧 read_seq 忽略
1433 if read_seq <= self.read_seq {
1434 // false
1435 return false;
1436 // 结束条件分支
1437 }
1438 // 更新
1439 self.read_seq = read_seq;
1440 // 有效推进
1441 true
1442 // 结束条件分支
1443 }
1444
1445 // Append 被拒后回退 next_index
1446 fn regress_next(&mut self, next_index: Index) -> bool {
1447 // 不能回退到已匹配之后或无效位置
1448 if next_index >= self.next_index || self.next_index <= self.match_index + 1 {
1449 // 无需回退
1450 return false;
1451 // 结束条件分支
1452 }
1453 // next 不低于 match+1
1454 self.next_index = max(next_index, self.match_index + 1);
1455 // 发生回退
1456 true
1457 // 结束条件分支
1458 }
1459// 结束代码块
1460}
1461
1462// 在途写:记录客户端来源与请求 id
1463struct Write {
1464 // 客户端节点(常为本节点)
1465 from: NodeID,
1466 // 请求关联 id
1467 id: RequestID,
1468// 结束类型定义
1469}
1470
1471// 在途线性一致读
1472struct Read {
1473 // 对应 read_seq
1474 seq: ReadSequence,
1475 // 客户端来源
1476 from: NodeID,
1477 // 请求 id
1478 id: RequestID,
1479 // 只读命令字节
1480 command: Vec<u8>,
1481// 结束代码块
1482}
1483
1484// 领导者状态构造
1485impl Leader {
1486 // 按同伴初始化进度与活跃计时
1487 fn new(peers: BTreeSet<NodeID>, last_index: Index) -> Self {
1488 // 乐观假设同伴缺最后一条,从 last+1 开始
1489 let next_index = last_index + 1;
1490 // 为每个 peer 建 Progress
1491 let progress = peers
1492 // 遍历同伴
1493 .iter()
1494 // 复制
1495 .copied()
1496 // match=0, read_seq=0, next=last+1
1497 .map(|p| (p, Progress { next_index, match_index: 0, read_seq: 0 }))
1498 // 聚合为最终集合/映射
1499 .collect();
1500 // peer_seen 初始 0 表示刚当选视为刚联系过
1501 let peer_seen = peers.iter().copied().map(|p| (p, 0)).collect();
1502 // 组装 Leader 状态
1503 Self {
1504 // 进度表
1505 progress,
1506 // 无在途写
1507 writes: HashMap::new(),
1508 // 无在途成员变更写
1509 membership_writes: HashMap::new(),
1510 // 空读队列
1511 reads: VecDeque::new(),
1512 // 读序列从 0 起
1513 read_seq: 0,
1514 // 立即允许发心跳
1515 since_heartbeat: 0,
1516 // 活跃性表
1517 peer_seen,
1518 // 结束代码块
1519 }
1520 // 结束代码块
1521 }
1522// 结束代码块
1523}
1524
1525// Leader 实现 Role
1526impl Role for Leader {}
1527
1528// 领导者行为:心跳、复制、提交、读写、成员变更、CheckQuorum
1529impl RawNode<Leader> {
1530 /// 因更高任期或 check-quorum 下台。
1531 fn into_follower(mut self, term: Term) -> Result<RawNode<Follower>> {
1532 // check-quorum 可能同任期下台。
1533 if term > self.term() {
1534 // 记录
1535 info!("Discovered new term {term}");
1536 // set_term_vote
1537 self.log.set_term_vote(term, None)?;
1538 // 同任期下台或无更高任期时的分支
1539 } else {
1540 // 同任期 step-down(CheckQuorum/移出配置)
1541 info!("Leader stepping down in term {}", self.term());
1542 // 结束条件分支
1543 }
1544
1545 // 中止所有在途写请求
1546 for write in std::mem::take(&mut self.role.writes).into_values().sorted_by_key(|w| w.id) {
1547 // 客户端 Abort
1548 self.send(write.from, Message::ClientResponse { id: write.id, response: Err(Error::Abort) })?;
1549 // 结束结构/枚举构造
1550 }
1551 // 中止在途成员变更请求
1552 for write in
1553 // 按 id 排序
1554 std::mem::take(&mut self.role.membership_writes).into_values().sorted_by_key(|w| w.id)
1555 // 遍历
1556 {
1557 // Abort
1558 self.send(write.from, Message::ClientResponse { id: write.id, response: Err(Error::Abort) })?;
1559 // 结束结构/枚举构造
1560 }
1561 // 中止在途读请求
1562 for read in std::mem::take(&mut self.role.reads).into_iter().sorted_by_key(|r| r.id) {
1563 // Abort
1564 self.send(read.from, Message::ClientResponse { id: read.id, response: Err(Error::Abort) })?;
1565 // 结束结构/枚举构造
1566 }
1567
1568 // 下台后使用新随机选举超时
1569 let election_timeout = self.random_election_timeout();
1570 // 无主跟随者
1571 Ok(self.into_role(Follower::new(None, election_timeout)))
1572 // 结束结构/枚举构造
1573 }
1574
1575 // 记录 peer 刚有有效响应,供 CheckQuorum
1576 fn note_peer_response(&mut self, from: NodeID) {
1577 // 存在则清零 seen
1578 if let Some(seen) = self.role.peer_seen.get_mut(&from) {
1579 // 重置活跃计时
1580 *seen = 0;
1581 // 结束条件分支
1582 }
1583 // 结束条件分支
1584 }
1585
1586 // 成员变更后同步 progress/peer_seen 与 peers 集合
1587 fn sync_progress_with_membership(&mut self) {
1588 // 当前同伴
1589 let peers = self.peers();
1590 // 日志末端
1591 let (last_index, _) = self.log.get_last_index();
1592 // 新同伴从 last+1 开始追赶
1593 let next_index = last_index + 1;
1594 // 添加新同伴。
1595 for p in &peers {
1596 // or_insert 避免覆盖已有 match
1597 self.role.progress.entry(*p).or_insert(Progress {
1598 // 初始 next
1599 next_index,
1600 // 尚未匹配
1601 match_index: 0,
1602 // 读序列 0
1603 read_seq: 0,
1604 // 结束映射/闭包构造
1605 });
1606 // 新同伴活跃计时
1607 self.role.peer_seen.entry(*p).or_insert(0);
1608 // 结束代码块
1609 }
1610 // 移除旧同伴。
1611 self.role.progress.retain(|id, _| peers.contains(id));
1612 // 同步 peer_seen
1613 self.role.peer_seen.retain(|id, _| peers.contains(id));
1614 // 结束代码块
1615 }
1616
1617 // 领导者消息处理主循环
1618 fn step(mut self, msg: Envelope) -> Result<Node> {
1619 // 成员变更提交后可能需要 step_down
1620 let mut step_down = false;
1621 // 领导者拒绝预选票:表明自己仍存活
1622 if matches!(msg.message, Message::PreCampaign { .. }) {
1623 // 领导者拒绝预选票(自己仍活着)。
1624 self.send(msg.from, Message::PreCampaignResponse { vote: false })?;
1625 // 结束
1626 return Ok(self.into());
1627 // 结束结构/枚举构造
1628 }
1629 // 领导者不收集预选票
1630 if matches!(msg.message, Message::PreCampaignResponse { .. }) {
1631 // 忽略
1632 return Ok(self.into());
1633 // 结束结构/枚举构造
1634 }
1635
1636 // 过期任期丢弃
1637 if msg.term < self.term() {
1638 // 调试
1639 debug!("Dropping message from past term: {msg:?}");
1640 // 保持领导
1641 return Ok(self.into());
1642 // 结束结构/枚举构造
1643 }
1644 // 更高任期:下台并重放
1645 if msg.term > self.term() {
1646 // into_follower 后 step
1647 return self.into_follower(msg.term)?.step(msg);
1648 // 结束条件分支
1649 }
1650
1651 // 按消息类型处理
1652 match msg.message {
1653 // 心跳应答:更新读确认、探测落后、推进提交
1654 Message::HeartbeatResponse { match_index, read_seq } => {
1655 // 已移除的 peer 忽略
1656 if !self.role.progress.contains_key(&msg.from) {
1657 // 返回
1658 return Ok(self.into());
1659 // 结束结构/枚举构造
1660 }
1661 // 标记 peer 活跃
1662 self.note_peer_response(msg.from);
1663 // 本地日志末端
1664 let (last_index, _) = self.log.get_last_index();
1665 // match 不能超过本地 last
1666 assert!(match_index <= last_index, "future match index");
1667 // read_seq 不能超过领导者当前值
1668 assert!(read_seq <= self.role.read_seq, "future read sequence number");
1669
1670 // 读序列前进则尝试放行读
1671 if self.progress(msg.from).advance_read(read_seq) {
1672 // maybe_read
1673 self.maybe_read()?;
1674 // 结束条件分支
1675 }
1676 // match=0 表示日志不匹配,回退并探测
1677 if match_index == 0 {
1678 // 将 next 拉向 last 以便 probe
1679 self.progress(msg.from).regress_next(last_index);
1680 // 强制 probe Append
1681 self.maybe_send_append(msg.from, true)?;
1682 // 结束条件分支
1683 }
1684 // match 前进则尝试提交
1685 if self.progress(msg.from).advance(match_index) {
1686 // 可能因移出配置需下台
1687 step_down |= self.maybe_commit_and_apply()?.1;
1688 // 结束条件分支
1689 }
1690 // 结束条件分支
1691 }
1692
1693 // 成功的 Append 应答
1694 Message::AppendResponse { match_index, reject_index: 0 } if match_index > 0 => {
1695 // 未知 peer 忽略
1696 if !self.role.progress.contains_key(&msg.from) {
1697 // 返回
1698 return Ok(self.into());
1699 // 结束结构/枚举构造
1700 }
1701 // 活跃
1702 self.note_peer_response(msg.from);
1703 // 末端
1704 let (last_index, _) = self.log.get_last_index();
1705 // 校验
1706 assert!(match_index <= last_index, "future match index");
1707 // 推进 match
1708 if self.progress(msg.from).advance(match_index) {
1709 // 尝试提交/应用
1710 step_down |= self.maybe_commit_and_apply()?.1;
1711 // 结束条件分支
1712 }
1713 // 流水线继续发送后续条目
1714 self.maybe_send_append(msg.from, false)?;
1715 // 结束条件分支
1716 }
1717
1718 // 独立读确认应答
1719 Message::ReadResponse { seq } => {
1720 // 未知 peer
1721 if !self.role.progress.contains_key(&msg.from) {
1722 // 返回
1723 return Ok(self.into());
1724 // 结束结构/枚举构造
1725 }
1726 // 活跃
1727 self.note_peer_response(msg.from);
1728 // 推进 read_seq
1729 if self.progress(msg.from).advance_read(seq) {
1730 // 尝试完成读
1731 self.maybe_read()?;
1732 // 结束条件分支
1733 }
1734 // 结束条件分支
1735 }
1736
1737 // Append 被拒绝:回退 next 并重发
1738 Message::AppendResponse { reject_index, match_index: 0 } if reject_index > 0 => {
1739 // 未知 peer
1740 if !self.role.progress.contains_key(&msg.from) {
1741 // 返回
1742 return Ok(self.into());
1743 // 结束结构/枚举构造
1744 }
1745 // 活跃
1746 self.note_peer_response(msg.from);
1747 // 末端
1748 let (last_index, _) = self.log.get_last_index();
1749 // reject 不得超 last
1750 assert!(reject_index <= last_index, "future reject index");
1751 // 拒绝点不大于已 match:陈旧应答
1752 if reject_index <= self.progress(msg.from).match_index {
1753 // 忽略
1754 return Ok(self.into());
1755 // 结束结构/枚举构造
1756 }
1757 // 回退 next_index
1758 if self.progress(msg.from).regress_next(reject_index) {
1759 // probe 重发
1760 self.maybe_send_append(msg.from, true)?;
1761 // 结束条件分支
1762 }
1763 // 结束条件分支
1764 }
1765
1766 // 非法 AppendResponse 形态
1767 Message::AppendResponse { .. } => panic!("invalid message {msg:?}"),
1768
1769 // 客户端写:追加日志并登记回调
1770 Message::ClientRequest { id, request: Request::Write(command) } => {
1771 // propose 复制 command
1772 let index = self.propose(Some(command))?;
1773 // 记录写等待提交
1774 self.role.writes.insert(index, Write { from: msg.from, id });
1775 // 单节点立即提交
1776 if self.cluster_size() == 1 {
1777 // 提交/应用,检查 step_down
1778 step_down |= self.maybe_commit_and_apply()?.1;
1779 // 结束条件分支
1780 }
1781 // 结束条件分支
1782 }
1783
1784 // 带会话的幂等写
1785 Message::ClientRequest {
1786 // 请求 id
1787 id,
1788 // 会话字段
1789 request: Request::WriteSession { client_id, seq, command },
1790 // 处理
1791 } => {
1792 // 编码为状态机可识别的会话包装
1793 let wrapped = super::session::encode_session(client_id, seq, command);
1794 // 作为普通写提出
1795 let index = self.propose(Some(wrapped))?;
1796 // 登记回调
1797 self.role.writes.insert(index, Write { from: msg.from, id });
1798 // 单节点快路径
1799 if self.cluster_size() == 1 {
1800 // 提交
1801 step_down |= self.maybe_commit_and_apply()?.1;
1802 // 结束条件分支
1803 }
1804 // 结束条件分支
1805 }
1806
1807 // 线性一致读:分配 read_seq 并广播确认
1808 Message::ClientRequest { id, request: Request::Read(command) } => {
1809 // 递增全局读序列
1810 self.role.read_seq += 1;
1811 // 入队读请求
1812 let read = Read { seq: self.role.read_seq, from: msg.from, id, command };
1813 // 保持 FIFO
1814 self.role.reads.push_back(read);
1815 // 广播 Read 让跟随者确认领导仍在
1816 self.broadcast(Message::Read { seq: self.role.read_seq })?;
1817 // 单节点无需多数确认
1818 if self.cluster_size() == 1 {
1819 // 立即 maybe_read
1820 self.maybe_read()?;
1821 // 结束条件分支
1822 }
1823 // 结束条件分支
1824 }
1825
1826 // 集群状态查询:本地即可回答
1827 Message::ClientRequest { id, request: Request::Status } => {
1828 // 组装 Status
1829 let response = self.status().map(Response::Status);
1830 // 直接响应客户端
1831 self.send(msg.from, Message::ClientResponse { id, response })?;
1832 // 结束结构/枚举构造
1833 }
1834
1835 // 成员变更请求
1836 Message::ClientRequest {
1837 // 请求 id
1838 id,
1839 // 新投票集合
1840 request: Request::ChangeMembership { voters },
1841 // 处理
1842 } => {
1843 // 提出 joint 配置
1844 let result = self.propose_membership_change(voters);
1845 // 成功或立即失败
1846 match result {
1847 // 成功:登记在途成员变更写
1848 Ok(index) => {
1849 // 等待 joint 提交后回复
1850 self.role.membership_writes.insert(index, Write { from: msg.from, id });
1851 // 单节点提交
1852 if self.cluster_size() == 1 {
1853 // 检查 step_down
1854 step_down |= self.maybe_commit_and_apply()?.1;
1855 // 结束条件分支
1856 }
1857 // 结束条件分支
1858 }
1859 // 失败:立即错误响应(如已有在途变更)
1860 Err(e) => {
1861 // 发送
1862 self.send(
1863 // 客户端
1864 msg.from,
1865 // 错误
1866 Message::ClientResponse { id, response: Err(e) },
1867 // 完成发送并向上传播错误
1868 )?;
1869 // 结束结构/枚举构造
1870 }
1871 // 结束结构/枚举构造
1872 }
1873 // 结束结构/枚举构造
1874 }
1875
1876 // 同任期竞选:领导者拒票
1877 Message::Campaign { .. } => {
1878 // vote=false
1879 self.send(msg.from, Message::CampaignResponse { vote: false })?;
1880 // 结束结构/枚举构造
1881 }
1882 // 忽略选票响应
1883 Message::CampaignResponse { .. } => {}
1884 // 同任期不应出现另一领导的复制/读消息
1885 Message::Heartbeat { .. } | Message::Append { .. } | Message::Read { .. } => {
1886 // 同任期不应出现另一领导;丢弃以免陈旧/异常消息拖垮进程。
1887 warn!(
1888 // 说明
1889 "leader {} ignoring peer-leader message from {} in term {}",
1890 // 来源与任期
1891 self.id, msg.from, msg.term
1892 // 业务逻辑
1893 );
1894 // 结束结构/枚举构造
1895 }
1896 // 快照安装成功应答:推进 match
1897 Message::InstallSnapshotResponse { last_included_index } => {
1898 // 未知 peer
1899 if !self.role.progress.contains_key(&msg.from) {
1900 // 返回
1901 return Ok(self.into());
1902 // 结束结构/枚举构造
1903 }
1904 // 活跃
1905 self.note_peer_response(msg.from);
1906 // 用 last_included 作为 match
1907 if self.progress(msg.from).advance(last_included_index) {
1908 // 尝试提交
1909 step_down |= self.maybe_commit_and_apply()?.1;
1910 // 结束条件分支
1911 }
1912 // 继续追加后续日志
1913 self.maybe_send_append(msg.from, false)?;
1914 // 结束条件分支
1915 }
1916 // 领导者不应安装他人快照
1917 Message::InstallSnapshot { .. } => {
1918 // 告警忽略
1919 warn!("leader {} ignoring InstallSnapshot from {}", self.id, msg.from);
1920 // 结束结构/枚举构造
1921 }
1922 // 异常消息
1923 Message::ClientResponse { .. }
1924 // PreCampaign 已处理
1925 | Message::PreCampaign { .. }
1926 // panic
1927 | Message::PreCampaignResponse { .. } => panic!("unexpected message {msg:?}"),
1928 // 结束结构/枚举构造
1929 }
1930
1931 // 本步处理中触发领导转移/移出配置
1932 if step_down {
1933 // 当前任期下台
1934 let term = self.term();
1935 // 转为无主跟随者
1936 return Ok(self.into_follower(term)?.into());
1937 // 结束结构/枚举构造
1938 }
1939 // 保持领导者
1940 Ok(self.into())
1941 // 结束结构/枚举构造
1942 }
1943
1944 // 领导者时钟:心跳与 CheckQuorum
1945 fn tick(mut self) -> Result<Node> {
1946 // 推进 peer_seen。
1947 for seen in self.role.peer_seen.values_mut() {
1948 // 饱和递增防溢出
1949 *seen = seen.saturating_add(1);
1950 // 结束循环
1951 }
1952
1953 // 心跳间隔计时
1954 self.role.since_heartbeat += 1;
1955 // 到达心跳周期
1956 if self.role.since_heartbeat >= self.opts.heartbeat_interval {
1957 // 广播 Heartbeat(含 commit 与 read_seq)
1958 self.heartbeat()?;
1959 // 结束条件分支
1960 }
1961
1962 // CheckQuorum:窗口内活跃节点(含自己)是否构成多数。
1963 if self.opts.check_quorum && self.cluster_size() > 1 {
1964 // 活跃窗口约等于选举超时上界
1965 let window = self.election_timeout_max();
1966 // 活跃集合
1967 let mut matched: BTreeSet<NodeID> = BTreeSet::new();
1968 // 领导者自身算活跃
1969 matched.insert(self.id);
1970 // 检查每个 peer
1971 for (peer, seen) in &self.role.peer_seen {
1972 // 窗口内有过响应
1973 if *seen < window {
1974 // 计入活跃
1975 matched.insert(*peer);
1976 // 结束条件分支
1977 }
1978 // 结束条件分支
1979 }
1980 // 不构成(联合)多数则主动下台
1981 if !self.membership.has_quorum(&matched) {
1982 // 告警
1983 warn!(
1984 // 失联多数
1985 "CheckQuorum: leader {} lost majority contact, stepping down",
1986 // 本节点
1987 self.id
1988 // 业务逻辑
1989 );
1990 // 同任期下台
1991 let term = self.term();
1992 // 转为跟随者
1993 return Ok(self.into_follower(term)?.into());
1994 // 结束结构/枚举构造
1995 }
1996 // 结束结构/枚举构造
1997 }
1998
1999 // 通过 CheckQuorum,保持领导
2000 Ok(self.into())
2001 // 结束结构/枚举构造
2002 }
2003
2004 // 广播心跳:宣告存活并捎带 commit/read_seq
2005 fn heartbeat(&mut self) -> Result<()> {
2006 // 领导者日志末端
2007 let (last_index, last_term) = self.log.get_last_index();
2008 // 当前提交点
2009 let (commit_index, _) = self.log.get_commit_index();
2010 // 当前读序列
2011 let read_seq = self.role.read_seq;
2012 // 领导者 last 条目必须属本任期(no-op 保证)
2013 assert_eq!(last_term, self.term(), "leader's last_term not in current term");
2014 // 重置心跳计时
2015 self.role.since_heartbeat = 0;
2016 // 向所有同伴发 Heartbeat
2017 self.broadcast(Message::Heartbeat { last_index, commit_index, read_seq })
2018 // 结束结构/枚举构造
2019 }
2020
2021 // 领导者提出新日志条目并触发复制
2022 fn propose(&mut self, command: Option<Vec<u8>>) -> Result<Index> {
2023 // 追加到本地日志,返回索引
2024 let index = self.log.append(command)?;
2025 // 对进度刚好指向该索引的 peer 立即发送
2026 for peer in self.peers() {
2027 // 避免无谓空转
2028 if index == self.progress(peer).next_index {
2029 // 发送 Append
2030 self.maybe_send_append(peer, false)?;
2031 // 结束条件分支
2032 }
2033 // 结束条件分支
2034 }
2035 // 返回新条目索引
2036 Ok(index)
2037 // 结束结构/枚举构造
2038 }
2039
2040 // 提出成员变更:写入 Joint(old,new) 并即时生效
2041 fn propose_membership_change(&mut self, voters: HashSet<NodeID>) -> Result<Index> {
2042 // 禁止重叠的成员变更
2043 if self.membership.change_pending {
2044 // 错误
2045 return Err(Error::InvalidInput("membership change already in progress".into()));
2046 // 结束条件分支
2047 }
2048 // 空投票集非法
2049 if voters.is_empty() {
2050 // 错误
2051 return Err(Error::InvalidInput("voters must not be empty".into()));
2052 // 结束条件分支
2053 }
2054 // 允许移除当前领导:Simple 配置提交后领导 step down(领导转移)。
2055 let old = match &self.membership.active {
2056 // 取出旧配置
2057 MembershipEntry::Simple(m) => m.clone(),
2058 // 已在 joint 中禁止嵌套
2059 MembershipEntry::Joint { .. } => {
2060 // 错误
2061 return Err(Error::InvalidInput("already in joint consensus".into()));
2062 // 结束 match 分支/块
2063 }
2064 // 结束 match 分支/块
2065 };
2066 // 由目标 voters 构造新 Membership
2067 let new = Membership::from_iter(voters);
2068 // 无变化则拒绝
2069 if old == new {
2070 // 错误
2071 return Err(Error::InvalidInput("membership unchanged".into()));
2072 // 结束条件分支
2073 }
2074
2075 // 构造 Joint 条目
2076 let joint = MembershipEntry::Joint { old, new: new.clone() };
2077 // 记录
2078 info!("Proposing joint membership: {joint:?}");
2079 // 追加成员日志
2080 let index = self.log.append_membership(joint.clone())?;
2081 // 立即采用 joint(两侧都需多数提交)
2082 self.membership.apply_entry(&joint);
2083 // 为新节点建 progress
2084 self.sync_progress_with_membership();
2085
2086 // 向需要的 peer 推送该配置条目
2087 for peer in self.peers() {
2088 // 进度对齐时发送
2089 if index == self.progress(peer).next_index {
2090 // Append
2091 self.maybe_send_append(peer, false)?;
2092 // 结束条件分支
2093 }
2094 // 结束条件分支
2095 }
2096 // 返回 joint 日志索引
2097 Ok(index)
2098 // 结束结构/枚举构造
2099 }
2100
2101 /// 提交 joint 后自动提出 Simple(C_new)。
2102 fn maybe_propose_simple_after_joint(&mut self, committed: &MembershipEntry) -> Result<()> {
2103 // 仅对 Joint 生效
2104 if let MembershipEntry::Joint { new, .. } = committed {
2105 // 目标配置的 Simple 形态
2106 let simple = MembershipEntry::Simple(new.clone());
2107 // 记录
2108 info!("Joint committed, proposing simple membership: {simple:?}");
2109 // 追加 Simple 成员条目
2110 let index = self.log.append_membership(simple.clone())?;
2111 // 立即采用 Simple
2112 self.membership.apply_entry(&simple);
2113 // 同步进度(可能移除旧节点)
2114 self.sync_progress_with_membership();
2115 // 推送新配置
2116 for peer in self.peers() {
2117 // 对齐则发
2118 if index == self.progress(peer).next_index {
2119 // Append
2120 self.maybe_send_append(peer, false)?;
2121 // 结束条件分支
2122 }
2123 // 结束条件分支
2124 }
2125 // 结束条件分支
2126 }
2127 // 非 Joint 无操作
2128 Ok(())
2129 // 结束结构/枚举构造
2130 }
2131
2132 // 根据多数 match 推进 commit,并 apply、回复客户端
2133 fn maybe_commit_and_apply(&mut self) -> Result<(Index, bool)> {
2134 // 日志末端
2135 let (last_index, _) = self.log.get_last_index();
2136
2137 // 基于当前(可能 joint)配置计算可提交索引:
2138 // 从 last_index 向下找第一个被多数复制的本任期索引。
2139 let mut commit_index = self.log.get_commit_index().0;
2140 // 自高向低找第一个可提交索引
2141 for idx in (commit_index + 1..=last_index).rev() {
2142 // 读取条目
2143 let Some(entry) = self.log.get(idx)? else { continue };
2144 // 只提交本任期条目(Raft 安全规则)
2145 if entry.term != self.term() {
2146 // 跳过其它任期
2147 continue;
2148 // 结束条件分支
2149 }
2150 // 统计复制到该索引的节点
2151 let mut matched: BTreeSet<NodeID> = BTreeSet::new();
2152 matched.insert(self.id); // 领导者已拥有
2153 // 遍历 progress
2154 for (peer, p) in &self.role.progress {
2155 // match 达到 idx
2156 if p.match_index >= idx {
2157 // 计入
2158 matched.insert(*peer);
2159 // 结束条件分支
2160 }
2161 // 结束条件分支
2162 }
2163 // 达到(联合)多数则选定该 commit
2164 if self.membership.has_quorum(&matched) {
2165 // 更新候选 commit_index
2166 commit_index = idx;
2167 // 找到最高可提交即停
2168 break;
2169 // 结束条件分支
2170 }
2171 // 结束条件分支
2172 }
2173
2174 // 读取旧提交点
2175 let (old_index, old_term) = self.log.get_commit_index();
2176 // 无前进则返回
2177 if commit_index <= old_index {
2178 // step_down=false
2179 return Ok((old_index, false));
2180 // 结束结构/枚举构造
2181 }
2182
2183 // 再次确认提交点属本任期
2184 match self.log.get(commit_index)? {
2185 // 本任期 OK
2186 Some(entry) if entry.term == self.term() => {}
2187 // 非本任期不可提交
2188 Some(_) => return Ok((old_index, false)),
2189 // 缺失条目严重错误
2190 None => panic!("commit index {commit_index} missing"),
2191 // 结束结构/枚举构造
2192 }
2193
2194 // 持久化新 commit_index
2195 self.log.commit(commit_index)?;
2196
2197 // 应用时发送响应需带当前 term
2198 let term = self.term();
2199 // 从已应用点扫描到 commit
2200 let mut iter = self.log.scan_apply(self.state.get_applied_index());
2201 // 收集本批提交的成员条目,稍后处理 joint 到 simple
2202 let mut committed_membership = Vec::new();
2203 // 逐条 apply
2204 while let Some(entry) = iter.next().transpose()? {
2205 // 调试
2206 debug!("Applying {entry:?}");
2207 // 成员条目
2208 if let Some(ref m) = entry.membership {
2209 // 提交语义:如清除 pending
2210 self.membership.on_commit(m);
2211 // 记录以便后续 propose simple
2212 committed_membership.push(m.clone());
2213 // 结束条件分支
2214 }
2215
2216 // 取出对应在途写
2217 let write = self.role.writes.remove(&entry.index);
2218 // 取出在途成员变更写
2219 let mwrite = self.role.membership_writes.remove(&entry.index);
2220 // 条目索引供响应
2221 let entry_index = entry.index;
2222
2223 // 成员变更条目:状态机 noop
2224 if entry.membership.is_some() {
2225 // 成员变更:状态机按 noop 应用,推进 applied_index。
2226 let _ = self.state.apply(Entry {
2227 // 索引
2228 index: entry.index,
2229 // 任期
2230 term: entry.term,
2231 // 无业务命令
2232 command: None,
2233 // membership 不传给状态机
2234 membership: None,
2235 // 结束映射/闭包构造
2236 });
2237 // 若有客户端等待成员变更结果则回复
2238 if let Some(Write { id, from: to }) = mwrite {
2239 // 经通道发送
2240 Self::send_via(
2241 // tx
2242 &self.tx,
2243 // 信封
2244 Envelope {
2245 // 领导者 id
2246 from: self.id,
2247 // 任期
2248 term,
2249 // 客户端
2250 to,
2251 // 成功响应
2252 message: Message::ClientResponse {
2253 // 原请求 id
2254 id,
2255 // 返回配置条目索引
2256 response: Ok(Response::ChangeMembership { index: entry_index }),
2257 // 业务逻辑
2258 },
2259 // 业务逻辑
2260 },
2261 // 完成发送并向上传播错误
2262 )?;
2263 // 结束结构/枚举构造
2264 }
2265 // 普通业务条目(非成员变更)分支
2266 } else {
2267 // 普通业务条目:apply 并回复写结果
2268 let result = self.state.apply(entry);
2269 // 有等待的写则回包
2270 if let Some(Write { id, from: to }) = write {
2271 // 发送
2272 Self::send_via(
2273 // tx
2274 &self.tx,
2275 // 信封
2276 Envelope {
2277 // from
2278 from: self.id,
2279 // term
2280 term,
2281 // to 客户端
2282 to,
2283 // 写响应
2284 message: Message::ClientResponse {
2285 // id
2286 id,
2287 // 状态机结果映射为 Write
2288 response: result.map(Response::Write),
2289 // 业务逻辑
2290 },
2291 // 业务逻辑
2292 },
2293 // 完成发送并向上传播错误
2294 )?;
2295 // 结束结构/枚举构造
2296 }
2297 // 结束代码块
2298 }
2299 // 结束代码块
2300 }
2301 // 释放日志扫描迭代器,避免与后续日志操作冲突
2302 drop(iter);
2303
2304 // 是否因移出配置而下台
2305 let mut step_down = false;
2306 // 处理本批提交的成员配置
2307 for m in committed_membership {
2308 // joint 提交则提出 Simple
2309 self.maybe_propose_simple_after_joint(&m)?;
2310 // 配置变化后同步进度
2311 self.sync_progress_with_membership();
2312 // Simple 配置已生效且自身不在投票集合 → 领导转移完成,下台。
2313 if matches!(self.membership.active, MembershipEntry::Simple(_))
2314 // 检查是否仍为 voter
2315 && !self.membership.all_voters().contains(&self.id)
2316 // 进入代码块
2317 {
2318 // 记录即将下台
2319 info!(
2320 // 节点 id
2321 "Leader {} removed from membership, will step down",
2322 // 说明
2323 self.id
2324 // 业务逻辑
2325 );
2326 // 标记 step_down
2327 step_down = true;
2328 // 结束代码块
2329 }
2330 // 结束代码块
2331 }
2332
2333 // 本任期首次提交后可放行读(领导权确认)
2334 if old_term != self.term() {
2335 // maybe_read
2336 self.maybe_read()?;
2337 // 结束条件分支
2338 }
2339
2340 // 不下台时考虑本地快照压缩
2341 if !step_down {
2342 // maybe_snapshot
2343 self.maybe_snapshot()?;
2344 // 结束条件分支
2345 }
2346
2347 // 返回新 commit 与是否下台
2348 Ok((commit_index, step_down))
2349 // 结束结构/枚举构造
2350 }
2351
2352 // 达到阈值则对已应用前缀做本地快照并压缩日志
2353 fn maybe_snapshot(&mut self) -> Result<()> {
2354 // 读取阈值
2355 let threshold = self.opts.snapshot_threshold;
2356 // 0 表示关闭
2357 if threshold == 0 {
2358 // 直接返回
2359 return Ok(());
2360 // 结束结构/枚举构造
2361 }
2362 // 状态机已应用索引
2363 let applied = self.state.get_applied_index();
2364 // 当前快照点
2365 let (snap_idx, _) = self.log.get_snapshot_meta();
2366 // 增量不足阈值则跳过
2367 if applied <= snap_idx || applied - snap_idx < threshold {
2368 // 返回
2369 return Ok(());
2370 // 结束结构/枚举构造
2371 }
2372 // 提交点
2373 let (commit, commit_term) = self.log.get_commit_index();
2374 // 未应用到 commit 前不做快照
2375 if applied < commit {
2376 // 返回
2377 return Ok(());
2378 // 结束结构/枚举构造
2379 }
2380 // 确定 applied 条目的 term
2381 let term = match self.log.get(applied)? {
2382 // 正常从日志取
2383 Some(e) => e.term,
2384 // 恰为快照点则无需
2385 None if applied == snap_idx => return Ok(()),
2386 // 缺失时用 commit_term 兜底
2387 None => commit_term,
2388 // 结束条件分支
2389 };
2390 // 状态机导出快照字节
2391 let data = self.state.snapshot()?;
2392 // 将快照字节存入引擎,便于重启恢复状态机。
2393 self.log.engine.set(&super::log::Key::SnapshotData.encode(), data)?;
2394 // 记录压缩
2395 info!("Compacting log through index {applied} term {term}");
2396 // 丢弃 applied 之前日志
2397 self.log.compact_to(applied, term)?;
2398 // 完成
2399 Ok(())
2400 // 结束结构/枚举构造
2401 }
2402
2403 // 在领导权与读多数确认后执行只读并回复
2404 fn maybe_read(&mut self) -> Result<()> {
2405 // 无在途读
2406 if self.role.reads.is_empty() {
2407 // 返回
2408 return Ok(());
2409 // 结束结构/枚举构造
2410 }
2411 // 当前提交点与任期
2412 let (commit_index, commit_term) = self.log.get_commit_index();
2413 // 已应用点
2414 let applied_index = self.state.get_applied_index();
2415 // 须本任期已提交且状态机追上 commit 才安全读
2416 if commit_term < self.term() || applied_index < commit_index {
2417 // 等待
2418 return Ok(());
2419 // 结束结构/枚举构造
2420 }
2421
2422 // 读确认:read_seq 被多数确认。
2423 // 构造每个 voter 的 read_seq(自己为 role.read_seq)。
2424 let mut matched_seq: Vec<(NodeID, ReadSequence)> = Vec::new();
2425 // 领导者自身视为已确认当前 read_seq
2426 matched_seq.push((self.id, self.role.read_seq));
2427 // 跟随者进度中的 read_seq
2428 for (peer, p) in &self.role.progress {
2429 // 加入列表
2430 matched_seq.push((*peer, p.read_seq));
2431 // 结束循环
2432 }
2433
2434 // 找最大的 seq,使得拥有 >= seq 的节点构成多数。
2435 let mut quorum_read_seq = 0;
2436 // 所有出现过的 seq 作为候选
2437 let candidates: BTreeSet<ReadSequence> = matched_seq.iter().map(|(_, s)| *s).collect();
2438 // 从大到小找第一个获多数的 seq
2439 for seq in candidates.into_iter().rev() {
2440 // 拥有 >= seq 的节点集
2441 let matched: BTreeSet<NodeID> =
2442 // 过滤
2443 matched_seq.iter().filter(|(_, s)| *s >= seq).map(|(id, _)| *id).collect();
2444 // 多数确认
2445 if self.membership.has_quorum(&matched) {
2446 // 记录 quorum_read_seq
2447 quorum_read_seq = seq;
2448 // 找到即停
2449 break;
2450 // 结束条件分支
2451 }
2452 // 结束条件分支
2453 }
2454
2455 // 按队列顺序放行 seq 已确认的读
2456 while let Some(read) = self.role.reads.front() {
2457 // 后续读尚未获多数
2458 if read.seq > quorum_read_seq {
2459 // 停止
2460 break;
2461 // 结束条件分支
2462 }
2463 // 出队
2464 let read = self.role.reads.pop_front().unwrap();
2465 // 状态机只读
2466 let response = self.state.read(read.command).map(Response::Read);
2467 // 回复客户端
2468 self.send(read.from, Message::ClientResponse { id: read.id, response })?;
2469 // 结束结构/枚举构造
2470 }
2471 // 读处理完成
2472 Ok(())
2473 // 结束结构/枚举构造
2474 }
2475
2476 // 向指定跟随者发送 Append 或快照以追赶日志
2477 fn maybe_send_append(&mut self, peer: NodeID, mut probe: bool) -> Result<()> {
2478 // peer 已不在进度表则跳过
2479 if !self.role.progress.contains_key(&peer) {
2480 // 返回
2481 return Ok(());
2482 // 结束结构/枚举构造
2483 }
2484 // 本地 last
2485 let (last_index, _) = self.log.get_last_index();
2486 // 取可变进度
2487 let progress = self.role.progress.get_mut(&peer).unwrap();
2488 // next 不得为 0
2489 assert_ne!(progress.next_index, 0, "invalid next_index");
2490 // next 必须在 match 之后
2491 assert!(progress.next_index > progress.match_index, "invalid next_index <= match_index");
2492 // match 不超过 last
2493 assert!(progress.match_index <= last_index, "invalid match_index > last_index");
2494 // next 最多 last+1
2495 assert!(progress.next_index <= last_index + 1, "invalid next_index > last_index + 1");
2496
2497 // 日志中仍保留的最早索引(快照之后)
2498 let first_index = self.log.get_first_index();
2499 // next 落在快照之前:只能发 InstallSnapshot
2500 if progress.next_index < first_index {
2501 // 快照元数据
2502 let (last_included_index, last_included_term) = self.log.get_snapshot_meta();
2503 // 当前状态机快照
2504 let data = self.state.snapshot()?;
2505 // 当前生效成员配置随快照发送
2506 let membership = self.membership.active.clone();
2507 // 发送快照消息
2508 return self.send(
2509 // 目标 peer
2510 peer,
2511 // InstallSnapshot
2512 Message::InstallSnapshot {
2513 // 覆盖索引
2514 last_included_index,
2515 // 覆盖任期
2516 last_included_term,
2517 // 数据
2518 data,
2519 // 配置
2520 membership,
2521 // 业务逻辑
2522 },
2523 // 业务逻辑
2524 );
2525 // 结束代码块
2526 }
2527
2528 // 已完全追平则无需发送
2529 if progress.match_index == last_index {
2530 // 返回
2531 return Ok(());
2532 // 结束结构/枚举构造
2533 }
2534 // probe 模式:next 与 match 有空洞时发空 Append 探测
2535 probe = probe && progress.next_index > progress.match_index + 1;
2536 // next 已超 last 且非 probe:无数据可发
2537 if progress.next_index > last_index && !probe {
2538 // 返回
2539 return Ok(());
2540 // 结束结构/枚举构造
2541 }
2542
2543 // 确定 prevLogIndex/prevLogTerm
2544 let (base_index, base_term) = match progress.next_index {
2545 // next=0 非法
2546 0 => panic!("next_index=0 for node {peer}"),
2547 // 从日志起点复制:base=0
2548 1 => (0, 0),
2549 // 否则取 next-1 条目的 index/term
2550 next => {
2551 // 缺失 base 条目则严重错误
2552 self.log.get(next - 1)?.map(|e| (e.index, e.term)).expect("missing base entry")
2553 // 结束循环
2554 }
2555 // 结束循环
2556 };
2557 // probe 发空,否则批量取最多 max_append_entries 条
2558 let entries = match probe {
2559 // 正常复制
2560 false => self
2561 // 从 next 扫描
2562 .log
2563 // 范围
2564 .scan(progress.next_index..)
2565 // 限制批量大小
2566 .take(self.opts.max_append_entries)
2567 // 收集为 Vec
2568 .try_collect()?,
2569 // 探测:空 entries
2570 true => Vec::new(),
2571 // 结束代码块
2572 };
2573
2574 // 乐观推进 next,等待应答确认 match
2575 if let Some(last) = entries.last() {
2576 // 发到 last+1
2577 progress.next_index = last.index + 1;
2578 // 结束条件分支
2579 }
2580
2581 // 调试复制规模
2582 debug!("Replicating {} entries with base {base_index} to {peer}", entries.len());
2583 // 发送 AppendEntries
2584 self.send(peer, Message::Append { base_index, base_term, entries })
2585 // 结束结构/枚举构造
2586 }
2587
2588 // 汇总集群状态供客户端 Status 请求
2589 fn status(&mut self) -> Result<Status> {
2590 // 构造 Status
2591 Ok(Status {
2592 // 当前领导为自己
2593 leader: self.id,
2594 // 当前任期
2595 term: self.term(),
2596 // 各节点 match_index 视图
2597 match_index: self
2598 // 从 progress
2599 .role
2600 // 迭代
2601 .progress
2602 // peer 与 match
2603 .iter()
2604 // 映射
2605 .map(|(id, p)| (*id, p.match_index))
2606 // 领导者自身 match=last_index
2607 .chain(std::iter::once((self.id, self.log.get_last_index().0)))
2608 // 收集
2609 .collect(),
2610 // 提交点
2611 commit_index: self.log.get_commit_index().0,
2612 // 状态机应用点
2613 applied_index: self.state.get_applied_index(),
2614 // 存储引擎状态
2615 storage: self.log.status()?,
2616 // 当前投票成员
2617 voters: self.membership.all_voters(),
2618 // 结束 Status/结构构造表达式
2619 })
2620 // 结束代码块
2621 }
2622
2623 // 按 id 取可变复制进度,缺失则协议 bug
2624 fn progress(&mut self, id: NodeID) -> &mut Progress {
2625 // expect 未知节点
2626 self.role.progress.get_mut(&id).expect("unknown node")
2627 // 结束函数
2628 }
2629// 结束函数
2630}