Skip to main content

raft_rust/raft/
message.rs

1// 领导者 Status 中 match_index 用有序 map 保证确定性
2use std::collections::BTreeMap;
3
4// 消息经 bincode 跨节点传输
5use serde::{Deserialize, Serialize};
6
7// 日志条目、索引与节点/任期标识
8use super::{Entry, Index, NodeID, Term};
9// ClientResponse 中携带 Result
10use crate::error::Result;
11// 存储引擎状态嵌入集群 Status
12use crate::storage;
13
14/// 消息信封,标明发送方与接收方。
15// 可克隆以便重试;可序列化走网络
16#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
17// 定义数据结构
18pub struct Envelope {
19    /// 发送方。
20    pub from: NodeID,
21    /// 发送方当前任期。
22    // 接收方用 term 判断是否过期或需要降级为跟随者
23    pub term: Term,
24    /// 接收方。
25    pub to: NodeID,
26    /// 消息本体。
27    pub message: Message,
28// 结束当前作用域
29}
30
31/// Raft 节点之间发送的消息。消息异步发送(非请求/响应模式),可能丢失或乱序。
32///
33/// 实践中它们经 TCP 连接与 crossbeam channel 传递;只要连接保持,通常不会丢失或乱序。
34/// 一条消息及其响应走各自独立的出站 TCP 连接。
35#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
36// 定义枚举
37pub enum Message {
38    /// 预投票请求(Pre-vote):不提升任期、不持久化投票。
39    /// 用于在真正选举前确认能否获得多数,避免分区节点抬升任期打断稳定领导。
40    PreCampaign {
41        /// 候选人最后一条日志的索引。
42        // 投票方用 last_index/last_term 判断日志新旧
43        last_index: Index,
44        /// 候选人最后一条日志的任期。
45        last_term: Term,
46    // 结构项闭合
47    },
48
49    /// 预投票响应。不持久化。
50    PreCampaignResponse {
51        /// 为 true 表示授予预选票。
52        vote: bool,
53    // 结构项闭合
54    },
55
56    /// 候选人向同伴拉票竞选领导者。
57    /// 仅当候选人的日志至少与投票者一样新时才会被授予选票。
58    Campaign {
59        /// 候选人最后一条日志的索引。
60        last_index: Index,
61        /// 候选人最后一条日志的任期。
62        last_term: Term,
63    // 结构项闭合
64    },
65
66    /// 跟随者每个任期只能投一票,且仅当候选人的日志至少与自己一样新时才投票。
67    /// 候选人隐式投票给自己。
68    CampaignResponse {
69        /// 为 true 表示授予选票。false 响应并非必须,但为清晰起见仍会发出。
70        vote: bool,
71    // 结构项闭合
72    },
73
74    /// 领导者发送的周期性心跳,作用包括:
75    ///
76    /// * 告知节点当前领导者,并阻止选举。
77    /// * 检测丢失的 append / read,作为重试机制。
78    /// * 推进跟随者的 commit 索引,以便它们应用条目。
79    ///
80    /// Raft 论文没有独立的心跳消息,而是使用空的 AppendEntries RPC;
81    /// 这里单独定义以便职责更清晰。
82    Heartbeat {
83        /// 领导者最后一条日志的索引。任期即领导者当前任期(当选时会追加 noop)。
84        /// 跟随者据此与本地日志比较,判断是否跟上。
85        last_index: Index,
86        /// 领导者最后已提交日志的索引。跟随者用它推进 commit 索引并应用条目。
87        /// 仅当本地日志在 last_index 处与领导者一致时,提交到此索引才安全。
88        commit_index: Index,
89        /// 领导者在本任期内最新的读序列号。
90        // 兼作线性一致读的 quorum 确认载体
91        read_seq: ReadSequence,
92    // 结构项闭合
93    },
94
95    /// 跟随者在仍认可其为领导者时,对心跳作出响应。
96    HeartbeatResponse {
97        /// 非零表示心跳中的 last_index 与跟随者日志匹配;否则跟随者日志分叉或落后。
98        match_index: Index,
99        /// 心跳中的读序列号。
100        // 回传以便领导者统计该 seq 的多数确认
101        read_seq: ReadSequence,
102    // 结构项闭合
103    },
104
105    /// 领导者在给定 base 条目之后,向跟随者追加日志条目以进行复制。
106    ///
107    /// 若 base 条目与跟随者日志匹配,则两边日志在此之前完全一致(见论文 5.3 节),
108    /// 可以追加(可能替换冲突条目)。否则拒绝追加,领导者需用更早的 base 索引重试,
109    /// 直到找到公共 base。
110    ///
111    /// 空 append(无条目)用于在日志分叉、节点重启或消息丢失时探测公共 match 索引。
112    /// 通常通过递减 base 索引探测,匹配后再发送后续条目。
113    Append {
114        /// 在此日志索引之后追加。
115        base_index: Index,
116        /// base 条目的任期。
117        base_term: Term,
118        /// 要追加的日志条目,必须从 base_index + 1 开始。
119        // 空 Vec 表示仅探测 match 点
120        entries: Vec<Entry>,
121    // 结构项闭合
122    },
123
124    /// 跟随者根据 base 条目是否匹配本地日志,接受或拒绝领导者的 append。
125    AppendResponse {
126        /// 非零表示跟随者已追加到该索引(此前日志与领导者一致)。
127        /// 若未发送条目(探测),则为匹配的 base 索引。
128        match_index: Index,
129        /// 非零表示在该 base 索引处拒绝(base 索引/任期不匹配)。
130        /// 若本地日志短于 base 索引,reject 会降到 last_index+1,避免逐个探测缺失索引。
131        reject_index: Index,
132    // 结构项闭合
133    },
134
135    /// 领导者在提供读服务前需确认自己仍是领导者,以保证线性一致性
136    ///(防止别处已选出新领导者)。读请求在序列号被多数派确认后才执行。
137    // 独立 Read 消息可与心跳并行推进读确认
138    Read { seq: ReadSequence },
139
140    /// 跟随者确认该读序列号下的领导权。
141    ReadResponse { seq: ReadSequence },
142
143    /// 客户端请求。可提交给领导者,或提交给跟随者由后者转发给领导者。
144    /// 若无领导者,或领导者/任期变更,请求会以 `Error::Abort` 的
145    /// ClientResponse 中止,客户端必须重试。
146    ClientRequest {
147        /// 请求 ID。在请求生命周期内必须全局唯一。
148        id: RequestID,
149        /// 请求本体。
150        request: Request,
151    // 结构项闭合
152    },
153
154    /// 客户端响应,通常透传给状态机。
155    ClientResponse {
156        /// 对应原始 ClientRequest 的 ID。
157        id: RequestID,
158        /// 响应,或错误。
159        // 含 Abort/状态机错误等
160        response: Result<Response>,
161    // 结构项闭合
162    },
163
164    /// 安装快照(整包;教学实现不分块)。
165    InstallSnapshot {
166        // 快照覆盖到的最后日志索引
167        last_included_index: Index,
168        // 该索引处日志任期,用于一致性校验
169        last_included_term: Term,
170        // 状态机完整快照字节
171        data: Vec<u8>,
172        // 快照时刻的成员配置,安装后立即生效
173        membership: crate::raft::membership::MembershipEntry,
174    // 结构项闭合
175    },
176
177    /// 快照安装确认。
178    InstallSnapshotResponse {
179        // 回传已安装的 last_included_index,领导者据此推进 match
180        last_included_index: Index,
181    // 结构项闭合
182    },
183// 结束当前作用域
184}
185
186/// 客户端请求 ID。在途期间必须全局唯一。
187///
188/// 为简单起见使用随机 UUIDv4。也可加入节点/进程/MAC 与时间戳以更好避撞
189///(例如 UUIDv6),但在此规模下不必。
190pub type RequestID = uuid::Uuid;
191
192/// 读序列号,用于为线性一致读确认领导权。
193// 领导者本任期内单调递增
194pub type ReadSequence = u64;
195
196/// 客户端请求,通常透传给状态机。
197#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
198// 定义枚举
199pub enum Request {
200    /// 状态机读命令,经 `State::read` 执行。不复制,仅在领导者上求值。
201    Read(Vec<u8>),
202    /// 状态机写命令,经 `State::apply` 执行。复制到所有节点,结果必须确定。
203    Write(Vec<u8>),
204    /// 带客户端 session 的写:日志中保存 (client_id, seq),apply 时去重。
205    WriteSession {
206        // 客户端会话 ID
207        client_id: uuid::Uuid,
208        // 该会话内单调序号
209        seq: u64,
210        // 实际应用层写命令
211        command: Vec<u8>,
212    // 结构项闭合
213    },
214    /// 向领导者查询 Raft 集群状态。
215    Status,
216    /// 变更集群投票成员为目标集合(联合共识)。
217    ChangeMembership {
218        // 目标投票成员集合
219        voters: std::collections::HashSet<NodeID>,
220    // 结构项闭合
221    },
222// 结束当前作用域
223}
224
225/// 客户端响应。外层用 Result 表示错误。
226#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
227// 定义枚举
228pub enum Response {
229    /// 状态机读结果。
230    Read(Vec<u8>),
231    /// 状态机写结果。
232    Write(Vec<u8>),
233    /// 当前 Raft 领导者状态。
234    Status(Status),
235    /// 成员变更已提出(提交索引,可能是 Joint 条目索引)。
236    ChangeMembership { index: Index },
237// 结束当前作用域
238}
239
240/// Raft 集群状态,由领导者生成。
241#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
242// 定义数据结构
243pub struct Status {
244    /// 生成该状态的当前 Raft 领导者。
245    pub leader: NodeID,
246    /// 当前 Raft 任期。
247    pub term: Term,
248    /// 各节点的 match 索引,表示复制进度。使用 BTreeMap 以保证测试确定性。
249    pub match_index: BTreeMap<NodeID, Index>,
250    /// 当前 commit 索引。
251    pub commit_index: Index,
252    /// 当前 applied 索引。
253    pub applied_index: Index,
254    /// 日志存储引擎状态。
255    pub storage: storage::Status,
256    /// 当前生效的投票成员集合。
257    // 缺省空集,兼容旧序列化载荷
258    #[serde(default)]
259    // 业务逻辑步骤
260    pub voters: std::collections::BTreeSet<NodeID>,
261// 结束当前作用域
262}