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}