raft_rust/cluster/mod.rs
1//! 进程内 Raft 集群运行时:channel 传输、客户端与故障注入。
2//!
3//! 供 `main`、example 与集成测试复用。不包含真实网络或磁盘持久化。
4
5// 哈希集合:投票人/成员/分区边
6use std::collections::{HashMap, HashSet};
7// 跨线程共享可变传输与客户端状态
8use std::sync::{Arc, Mutex};
9// 节点线程与客户端重试休眠
10use std::thread;
11// 超时与退避间隔
12use std::time::Duration;
13
14// 进程内消息与请求应答通道
15use crossbeam::channel::{self, Receiver, Sender};
16// 丢包注入随机源
17use rand::RngExt as _;
18// 会话/请求唯一 ID
19use uuid::Uuid;
20
21// 统一 Error/Result
22use crate::error::{Error, Result};
23// KV 命令与状态机
24use crate::raft::kv::{self, Command, Kv};
25// Raft 协议核心类型
26use crate::raft::{
27 // 协议类型:信封、日志、消息、节点与请求响应
28 Envelope, Index, Log, Message, Node, NodeID, Options, Request, RequestID, Response, Status,
29 // tick 周期常量
30 TICK_INTERVAL,
31// 当前作用域结束
32};
33// 临时 BitCask 日志引擎
34use crate::storage::BitCask;
35
36// ---------------------------------------------------------------------------
37// 传输层(分区 / 丢包 / 乱序)
38// ---------------------------------------------------------------------------
39
40// 派生 Default,便于 Transport 空状态启动
41#[derive(Default)]
42// 传输可变内核:故障策略与邮箱
43struct TransportInner {
44 /// `(from, to)` 被阻断时,`from` 无法把消息投递给 `to`。
45 partitions: HashSet<(NodeID, NodeID)>,
46 /// 随机丢包率 `[0.0, 1.0]`。
47 drop_rate: f64,
48 /// 简单乱序:缓存一条出站消息,下次再与新消息交换顺序发出。
49 reorder: bool,
50 // 乱序注入时按目标缓存的待交换消息
51 held: HashMap<NodeID, Envelope>,
52 /// 节点是否在线(stop 后为 false)。
53 online: HashMap<NodeID, bool>,
54 /// 各节点入站邮箱。
55 mailboxes: HashMap<NodeID, Sender<Envelope>>,
56// 当前作用域结束
57}
58
59/// 可故障注入的进程内传输。
60#[derive(Clone, Default)]
61// 可故障注入的进程内传输
62pub struct Transport {
63 // 业务:inner: Arc<Mutex<TransportInner>>,
64 inner: Arc<Mutex<TransportInner>>,
65// 当前作用域结束
66}
67
68// 实现该类型的方法
69impl Transport {
70 // 登记节点入站通道并标在线
71 fn register(&self, id: NodeID, tx: Sender<Envelope>) {
72 // 独占传输内核锁
73 let mut g = self.inner.lock().expect("transport lock");
74 // 写入路由/分区/挂起表项
75 g.mailboxes.insert(id, tx);
76 // 写入路由/分区/挂起表项
77 g.online.insert(id, true);
78 // 当前作用域结束
79 }
80
81 // 模拟崩溃/恢复:离线丢弃收发
82 fn set_online(&self, id: NodeID, online: bool) {
83 // TransportInner 字段定义结束
84 let mut g = self.inner.lock().expect("transport lock");
85 // 写入路由/分区/挂起表项
86 g.online.insert(id, online);
87 // 检查节点是否在线
88 if !online {
89 // 移除表项或清理缓存
90 g.held.remove(&id);
91 // 当前作用域结束
92 }
93 // 当前作用域结束
94 }
95
96 // 按分区/丢包/乱序规则投递
97 fn deliver(&self, msg: Envelope) {
98 // Transport 包装结构结束
99 let mut g = self.inner.lock().expect("transport lock");
100 // 信封发送方
101 let from = msg.from;
102 // 信封接收方
103 let to = msg.to;
104
105 // 检查节点是否在线
106 if !g.online.get(&from).copied().unwrap_or(false) {
107 // 完成当前语句
108 return;
109 // 当前作用域结束
110 }
111 // 检查节点是否在线
112 if !g.online.get(&to).copied().unwrap_or(false) {
113 // 完成当前语句
114 return;
115 // 当前作用域结束
116 }
117 // 命中分区则阻断投递
118 if g.partitions.contains(&(from, to)) {
119 // 完成当前语句
120 return;
121 // register:节点入站通道登记完毕
122 }
123 // 按丢包率随机丢弃
124 if g.drop_rate > 0.0 && rand::rng().random::<f64>() < g.drop_rate {
125 // 完成当前语句
126 return;
127 // 当前作用域结束
128 }
129
130 // 进入乱序注入逻辑
131 if g.reorder {
132 // 可选值解构分支
133 if let Some(prev) = g.held.remove(&to) {
134 // 先发当前,再发缓存 → 乱序
135 if let Some(tx) = g.mailboxes.get(&to) {
136 // 发送端通道
137 let _ = tx.try_send(msg);
138 // 发送端通道
139 let _ = tx.try_send(prev);
140 // 当前作用域结束
141 }
142 // 离线时清理 held 缓存的分支结束
143 return;
144 // set_online 方法结束
145 }
146 // 写入路由/分区/挂起表项
147 g.held.insert(to, msg);
148 // 完成当前语句
149 return;
150 // 当前作用域结束
151 }
152
153 // 可选值解构分支
154 if let Some(tx) = g.mailboxes.get(&to) {
155 // 发送端通道
156 let _ = tx.try_send(msg);
157 // 当前作用域结束
158 }
159 // 当前作用域结束
160 }
161
162 /// 双向隔离 `a` 与 `b`。
163 pub fn partition(&self, a: NodeID, b: NodeID) {
164 // 独占传输内核锁
165 let mut g = self.inner.lock().expect("transport lock");
166 // 写入路由/分区/挂起表项
167 g.partitions.insert((a, b));
168 // 写入路由/分区/挂起表项
169 g.partitions.insert((b, a));
170 // 发送方离线则静默丢弃的分支结束
171 }
172
173 /// 单向阻断 `from -> to`。
174 pub fn partition_one_way(&self, from: NodeID, to: NodeID) {
175 // 独占传输内核锁
176 let mut g = self.inner.lock().expect("transport lock");
177 // 接收方离线则静默丢弃的分支结束
178 g.partitions.insert((from, to));
179 // 当前作用域结束
180 }
181
182 /// 按组双向分区:`left` 内节点与 `right` 内节点互不可达。
183 pub fn partition_groups(&self, left: &[NodeID], right: &[NodeID]) {
184 // 命中分区规则阻断的分支结束
185 let mut g = self.inner.lock().expect("transport lock");
186 // 遍历集合或重试轮次
187 for &a in left {
188 // 遍历集合或重试轮次
189 for &b in right {
190 // 写入路由/分区/挂起表项
191 g.partitions.insert((a, b));
192 // 写入路由/分区/挂起表项
193 g.partitions.insert((b, a));
194 // 随机丢包分支结束
195 }
196 // 当前作用域结束
197 }
198 // 当前作用域结束
199 }
200
201 /// 恢复 `a` 与 `b` 之间的连通(双向)。
202 pub fn heal_pair(&self, a: NodeID, b: NodeID) {
203 // 独占传输内核锁
204 let mut g = self.inner.lock().expect("transport lock");
205 // 移除表项或清理缓存
206 g.partitions.remove(&(a, b));
207 // 移除表项或清理缓存
208 g.partitions.remove(&(b, a));
209 // 当前作用域结束
210 }
211
212 /// 清除全部分区,并冲刷乱序缓存。
213 pub fn heal_all(&self) {
214 // 独占传输内核锁
215 let mut g = self.inner.lock().expect("transport lock");
216 // 清空全部分区
217 g.partitions.clear();
218 // 已有 held 时交换发出的 if 结束
219 let held = std::mem::take(&mut g.held);
220 // 遍历集合或重试轮次
221 for (to, msg) in held {
222 // 可选值解构分支
223 if let Some(tx) = g.mailboxes.get(&to) {
224 // 发送端通道
225 let _ = tx.try_send(msg);
226 // 当前作用域结束
227 }
228 // 尚无缓存则 hold 本条的分支结束
229 }
230 // 当前作用域结束
231 }
232
233 // 设置随机丢包率并夹紧到[0,1]
234 pub fn set_drop_rate(&self, rate: f64) {
235 // 独占传输内核锁
236 let mut g = self.inner.lock().expect("transport lock");
237 // 语句/调用结束
238 g.drop_rate = rate.clamp(0.0, 1.0);
239 // 正常路径 try_send 到目标 inbox 结束
240 }
241
242 // 开关乱序;关闭时冲刷 held
243 pub fn set_reorder(&self, on: bool) {
244 // 独占传输内核锁
245 let mut g = self.inner.lock().expect("transport lock");
246 // 完成当前语句
247 g.reorder = on;
248 // 关闭乱序时冲刷 held
249 if !on {
250 // 独占传输内核锁
251 let held = std::mem::take(&mut g.held);
252 // 遍历集合或重试轮次
253 for (to, msg) in held {
254 // 可选值解构分支
255 if let Some(tx) = g.mailboxes.get(&to) {
256 // 发送端通道
257 let _ = tx.try_send(msg);
258 // 当前作用域结束
259 }
260 // 双向 partition 注入结束
261 }
262 // 当前作用域结束
263 }
264 // 当前作用域结束
265 }
266// 当前作用域结束
267}
268
269// ---------------------------------------------------------------------------
270// 客户端
271// ---------------------------------------------------------------------------
272
273// 节点侧请求+应答通道别名
274type RequestTx = Sender<(Request, Sender<Result<Response>>)>;
275
276/// 可向任意本地节点提交请求的客户端(跟随者会转发到领导者)。
277#[derive(Clone)]
278// 可向任意本地节点提交请求的客户端
279pub struct Client {
280 // 各节点请求通道表,可动态增删
281 request_txs: Arc<Mutex<HashMap<NodeID, RequestTx>>>,
282 // 优先尝试的入口节点(常为领导者)
283 preferred: NodeID,
284 // 外层最大重试轮数
285 attempts: u32,
286 // 单次等待节点应答超时
287 per_attempt_timeout: Duration,
288 // 一轮全失败后的选举等待
289 retry_sleep: Duration,
290// 当前作用域结束
291}
292
293// 实现该类型的方法
294impl Client {
295 // 左右组一对节点双向隔离结束
296 fn new(request_txs: HashMap<NodeID, RequestTx>, preferred: NodeID) -> Self {
297 // 左组遍历结束
298 Self {
299 // partition_groups 结束
300 request_txs: Arc::new(Mutex::new(request_txs)),
301 // 优先尝试的入口节点(常为领导者)
302 preferred,
303 // 外层最大重试轮数
304 attempts: 40,
305 // 单次等待节点应答超时
306 per_attempt_timeout: Duration::from_millis(200),
307 // 一轮全失败后的选举等待
308 retry_sleep: Duration::from_millis(50),
309 // 当前作用域结束
310 }
311 // 当前作用域结束
312 }
313
314 /// 注册或更新某节点的请求通道(节点 start 后调用)。
315 pub fn register_node(&self, id: NodeID, tx: RequestTx) {
316 // heal_pair 双向恢复结束
317 self.request_txs.lock().expect("client lock").insert(id, tx);
318 // 当前作用域结束
319 }
320
321 // 节点 start 后注册请求通道
322 pub fn unregister_node(&self, id: NodeID) {
323 // 移除表项或清理缓存
324 self.request_txs.lock().expect("client lock").remove(&id);
325 // 当前作用域结束
326 }
327
328 /// 提示下次优先向该节点发请求。
329 pub fn preferred_hint(&mut self, id: NodeID) {
330 // 更新自身状态字段
331 self.preferred = id;
332 // 当前作用域结束
333 }
334
335 // 按 preferred 优先轮询提交请求
336 pub fn request(&mut self, request: Request) -> Result<Response> {
337 // 协议层请求
338 let txs = self.request_txs.lock().expect("client lock").clone();
339 // 节点尝试顺序(preferred 优先)
340 let mut order: Vec<NodeID> = txs.keys().copied().collect();
341 // 语句/调用结束
342 order.sort();
343 // 冲刷单条 held 到邮箱结束
344 if let Some(pos) = order.iter().position(|&id| id == self.preferred) {
345 // held 补发循环结束
346 let id = order.remove(pos);
347 // heal_all 结束
348 order.insert(0, id);
349 // 当前作用域结束
350 }
351
352 // 可重试的最后错误
353 let mut last_err = Error::Abort;
354 // 遍历集合或重试轮次
355 for _ in 0..self.attempts {
356 // 遍历集合或重试轮次
357 for &node_id in &order {
358 // 请求通道表快照
359 let Some(tx) = txs.get(&node_id) else { continue };
360 // 集群响应
361 let (resp_tx, resp_rx) = channel::bounded(1);
362 // set_drop_rate 夹紧并写回结束
363 if tx.send((request.clone(), resp_tx)).is_err() {
364 // 继续下一轮
365 continue;
366 // 当前作用域结束
367 }
368 // 单次等待节点应答超时
369 match resp_rx.recv_timeout(self.per_attempt_timeout) {
370 // 成功:缓存 preferred 并返回响应
371 Ok(Ok(resp)) => {
372 // 更新自身状态字段
373 self.preferred = node_id;
374 // 返回业务结果或错误
375 return Ok(resp);
376 // 当前作用域结束
377 }
378 // 无主/转发失败,可换节点重试
379 Ok(Err(Error::Abort)) => last_err = Error::Abort,
380 // 业务错误立即返回,不重试
381 Ok(Err(e)) => return Err(e),
382 // 超时或通道错误处理
383 Err(_) => last_err = Error::IO("request timed out".into()),
384 // match 分支结束
385 }
386 // 当前作用域结束
387 }
388 // 一轮全失败后的选举等待
389 thread::sleep(self.retry_sleep);
390 // 当前作用域结束
391 }
392 // 业务:Err(last_err)
393 Err(last_err)
394 // 当前作用域结束
395 }
396
397 /// 只向指定节点发请求(仍在 Abort 时对该节点重试)。
398 pub fn request_on(&mut self, node_id: NodeID, request: Request) -> Result<Response> {
399 // 关闭乱序清理分支结束
400 let txs = self.request_txs.lock().expect("client lock").clone();
401 // set_reorder 结束
402 let Some(tx) = txs.get(&node_id).cloned() else {
403 // Transport 实现块结束
404 return Err(Error::IO(format!("node {node_id} not registered")));
405 // 当前作用域结束
406 };
407 // 可重试的最后错误
408 let mut last_err = Error::Abort;
409 // 遍历集合或重试轮次
410 for _ in 0..self.attempts {
411 // 集群响应
412 let (resp_tx, resp_rx) = channel::bounded(1);
413 // 通道关闭则跳过/失败
414 if tx.send((request.clone(), resp_tx)).is_err() {
415 // 返回业务结果或错误
416 return Err(Error::IO(format!("node {node_id} request channel closed")));
417 // 当前作用域结束
418 }
419 // 单次等待节点应答超时
420 match resp_rx.recv_timeout(self.per_attempt_timeout) {
421 // 成功:缓存 preferred 并返回响应
422 Ok(Ok(resp)) => {
423 // 更新自身状态字段
424 self.preferred = node_id;
425 // 返回业务结果或错误
426 return Ok(resp);
427 // 当前作用域结束
428 }
429 // 无主/转发失败,可换节点重试
430 Ok(Err(Error::Abort)) => last_err = Error::Abort,
431 // 业务错误立即返回,不重试
432 Ok(Err(e)) => return Err(e),
433 // 超时或通道错误处理
434 Err(_) => last_err = Error::IO("request timed out".into()),
435 // match 分支结束
436 }
437 // 一轮全失败后的选举等待
438 thread::sleep(self.retry_sleep);
439 // 当前作用域结束
440 }
441 // 业务:Err(last_err)
442 Err(last_err)
443 // 当前作用域结束
444 }
445
446 // 便捷 Put:编码写并解析提交索引
447 pub fn put(&mut self, key: &str, value: &str) -> Result<Index> {
448 // Client 字段定义结束
449 let req = Request::Write(kv::encode(&Command::Put {
450 // 填充键字段
451 key: key.into(),
452 // 填充值字段
453 value: value.into(),
454 // 语句/调用结束
455 }));
456 // 按响应/结果分支处理
457 match self.request(req)? {
458 // 写回包:解码 KV 层响应
459 Response::Write(bytes) => match kv::decode::<kv::Response>(&bytes)? {
460 // Put 成功,返回提交索引
461 kv::Response::Put(index) => Ok(index),
462 // 未知命令或非预期响应
463 other => Err(Error::InvalidData(format!("unexpected write response: {other:?}"))),
464 // match 分支结束
465 },
466 // 未知命令或非预期响应
467 other => Err(Error::InvalidData(format!("expected Write, got {other:?}"))),
468 // match 分支结束
469 }
470 // 当前作用域结束
471 }
472
473 // 指定节点上的 Put(测转发)
474 pub fn put_on(&mut self, node_id: NodeID, key: &str, value: &str) -> Result<Index> {
475 // 协议层请求
476 let req = Request::Write(kv::encode(&Command::Put {
477 // 填充键字段
478 key: key.into(),
479 // 填充值字段
480 value: value.into(),
481 // 语句/调用结束
482 }));
483 // Client::new 默认超时参数组装结束
484 match self.request_on(node_id, req)? {
485 // Client::new 结束
486 Response::Write(bytes) => match kv::decode::<kv::Response>(&bytes)? {
487 // Put 成功,返回提交索引
488 kv::Response::Put(index) => Ok(index),
489 // 未知命令或非预期响应
490 other => Err(Error::InvalidData(format!("unexpected write response: {other:?}"))),
491 // match 分支结束
492 },
493 // 未知命令或非预期响应
494 other => Err(Error::InvalidData(format!("expected Write, got {other:?}"))),
495 // match 分支结束
496 }
497 // register_node 写入路由表结束
498 }
499
500 // 便捷 Get:线性读路径
501 pub fn get(&mut self, key: &str) -> Result<Option<String>> {
502 // 协议层请求
503 let req = Request::Read(kv::encode(&Command::Get { key: key.into() }));
504 // 按响应/结果分支处理
505 match self.request(req)? {
506 // 读回包:解码 KV 层响应
507 Response::Read(bytes) => match kv::decode::<kv::Response>(&bytes)? {
508 // unregister_node 移除路由结束
509 kv::Response::Get(v) => Ok(v),
510 // 未知命令或非预期响应
511 other => Err(Error::InvalidData(format!("unexpected read response: {other:?}"))),
512 // match 分支结束
513 },
514 // 未知命令或非预期响应
515 other => Err(Error::InvalidData(format!("expected Read, got {other:?}"))),
516 // match 分支结束
517 }
518 // 当前作用域结束
519 }
520
521 // 全量扫描状态机
522 pub fn scan(&mut self) -> Result<std::collections::BTreeMap<String, String>> {
523 // 协议层请求
524 let req = Request::Read(kv::encode(&Command::Scan));
525 // 按响应/结果分支处理
526 match self.request(req)? {
527 // 读回包:解码 KV 层响应
528 Response::Read(bytes) => match kv::decode::<kv::Response>(&bytes)? {
529 // Scan 返回有序 map
530 kv::Response::Scan(map) => Ok(map),
531 // 未知命令或非预期响应
532 other => Err(Error::InvalidData(format!("unexpected scan response: {other:?}"))),
533 // match 分支结束
534 },
535 // 未知命令或非预期响应
536 other => Err(Error::InvalidData(format!("expected Read, got {other:?}"))),
537 // match 分支结束
538 }
539 // 当前作用域结束
540 }
541
542 // 查询集群 Status
543 pub fn status(&mut self) -> Result<Status> {
544 // 按响应/结果分支处理
545 match self.request(Request::Status)? {
546 // 状态回包:展示主从与位点
547 Response::Status(s) => Ok(s),
548 // 未知命令或非预期响应
549 other => Err(Error::InvalidData(format!("expected Status, got {other:?}"))),
550 // 将 preferred 提到队首的调整结束
551 }
552 // 当前作用域结束
553 }
554
555 // 指定节点 Status(观察分区视图)
556 pub fn status_on(&mut self, node_id: NodeID) -> Result<Status> {
557 // 只向指定节点发请求并重试
558 match self.request_on(node_id, Request::Status)? {
559 // 状态回包:展示主从与位点
560 Response::Status(s) => Ok(s),
561 // 未知命令或非预期响应
562 other => Err(Error::InvalidData(format!("expected Status, got {other:?}"))),
563 // match 分支结束
564 }
565 // 当前作用域结束
566 }
567
568 /// 变更集群成员(目标投票人集合)。需协议侧支持 `Request::ChangeMembership`。
569 pub fn change_membership(&mut self, voters: HashSet<NodeID>) -> Result<Index> {
570 // 按响应/结果分支处理
571 match self.request(Request::ChangeMembership { voters })? {
572 // 成员变更已提议,打印日志索引
573 Response::ChangeMembership { index } => Ok(index),
574 // 未知命令或非预期响应
575 other => Err(Error::InvalidData(format!("expected ChangeMembership, got {other:?}"))),
576 // match 分支结束
577 }
578 // 通道已关闭则跳过该节点的分支结束
579 }
580// 当前作用域结束
581}
582
583/// 等待选出领导者。
584pub fn wait_for_leader(client: &mut Client) -> Result<Status> {
585 // 可重试的最后错误
586 let mut last_err = Error::Abort;
587 // 遍历集合或重试轮次
588 for _ in 0..100 {
589 // 按响应/结果分支处理
590 match client.status() {
591 // 成功路径:推进状态或返回
592 Ok(s) => return Ok(s),
593 // 成功响应并缓存 preferred 的分支结束
594 Err(Error::Abort) => {
595 // 完成当前语句
596 last_err = Error::Abort;
597 // 失败后短暂退避,等待选主稳定
598 thread::sleep(Duration::from_millis(50));
599 // 当前作用域结束
600 }
601 // 超时或通道错误处理
602 Err(e) => return Err(e),
603 // match 分支结束
604 }
605 // 当前作用域结束
606 }
607 // 单次 recv_timeout 结果 match 结束
608 Err(last_err)
609// 内层按节点顺序尝试结束
610}
611
612// ---------------------------------------------------------------------------
613// 集群
614// ---------------------------------------------------------------------------
615
616// request 多节点轮询结束
617struct NodeControl {
618 // 通知节点线程退出的发送端
619 stop_tx: Sender<()>,
620// 当前作用域结束
621}
622
623/// 进程内多节点 Raft 集群。
624pub struct Cluster {
625 // 集群统一 Raft 选项快照
626 opts: Options,
627 // 共享故障注入传输
628 transport: Transport,
629 // 运行中节点控制句柄
630 nodes: HashMap<NodeID, NodeControl>,
631 /// 逻辑成员集合(用于新节点 peers 计算);成员变更协议生效前由测试/调用方维护。
632 members: HashSet<NodeID>,
633 // 绑定各节点请求通道的客户端
634 client: Client,
635// 当前作用域结束
636}
637
638// 实现该类型的方法
639impl Cluster {
640 /// 使用默认快速测试选项启动集群。
641 pub fn spawn(node_ids: &[NodeID]) -> Self {
642 // 按 Options 拉起节点组与客户端
643 Self::spawn_with_options(node_ids, test_options())
644 // 当前作用域结束
645 }
646
647 // 集群统一 Raft 选项快照
648 pub fn spawn_with_options(node_ids: &[NodeID], opts: Options) -> Self {
649 // 通道关闭直接失败的分支结束
650 let transport = Transport::default();
651 // 协议层请求
652 let mut request_txs = HashMap::new();
653 // 构造或 tick/step 后的 Node
654 let mut nodes = HashMap::new();
655 // 逻辑成员集合(peers 计算用)
656 let members: HashSet<NodeID> = node_ids.iter().copied().collect();
657
658 // 遍历集合或重试轮次
659 for &id in node_ids {
660 // 启动单节点线程:邮箱、日志、事件循环
661 let (control, req_tx) = spawn_node(id, &members, opts.clone(), transport.clone());
662 // 写入路由/分区/挂起表项
663 request_txs.insert(id, req_tx.clone());
664 // 写入路由/分区/挂起表项
665 nodes.insert(id, control);
666 // request_on 成功返回分支结束
667 }
668
669 // 构造或 tick/step 后的 Node
670 let preferred = node_ids.first().copied().unwrap_or(1);
671 // 协议层请求
672 let client = Client::new(request_txs, preferred);
673
674 // 组装结构体字段
675 Self { opts, transport, nodes, members, client }
676 // 当前作用域结束
677 }
678
679 // 克隆共享客户端句柄
680 pub fn client(&self) -> Client {
681 // 业务:self.client.clone()
682 self.client.clone()
683 // request_on 重试循环结束
684 }
685
686 // 暴露传输以便故障注入
687 pub fn transport(&self) -> Transport {
688 // request_on 方法结束
689 self.transport.clone()
690 // 当前作用域结束
691 }
692
693 // 返回逻辑成员视图副本
694 pub fn members(&self) -> HashSet<NodeID> {
695 // 业务:self.members.clone()
696 self.members.clone()
697 // 当前作用域结束
698 }
699
700 // 只读访问启动选项
701 pub fn options(&self) -> &Options {
702 // 业务:&self.opts
703 &self.opts
704 // 当前作用域结束
705 }
706
707 /// 停止节点(模拟崩溃):不再处理消息/请求,传输层视为离线。
708 pub fn stop(&mut self, id: NodeID) {
709 // 可选值解构分支
710 if let Some(ctrl) = self.nodes.remove(&id) {
711 // 发送端通道
712 let _ = ctrl.stop_tx.send(());
713 // 节点 start 后注册请求通道
714 self.client.unregister_node(id);
715 // 线程退出前标离线,避免继续投递
716 self.transport.set_online(id, false);
717 // 当前作用域结束
718 }
719 // 当前作用域结束
720 }
721
722 /// 以空 BitCask 日志重新拉起节点(用于成员加入;新路径空库)。
723 pub fn start(&mut self, id: NodeID) {
724 // put 协议响应 match 结束
725 if self.nodes.contains_key(&id) {
726 // put 便捷方法结束
727 return;
728 // 当前作用域结束
729 }
730 // 写入路由/分区/挂起表项
731 self.members.insert(id);
732 // 协议层请求
733 let (control, req_tx) =
734 // 启动单节点线程:邮箱、日志、事件循环
735 spawn_node(id, &self.members, self.opts.clone(), self.transport.clone());
736 // 节点 start 后注册请求通道
737 self.client.register_node(id, req_tx);
738 // 写入路由/分区/挂起表项
739 self.nodes.insert(id, control);
740 // 当前作用域结束
741 }
742
743 /// 仅更新本地成员集合视图(协议提交成员变更后由测试调用)。
744 pub fn set_members(&mut self, members: HashSet<NodeID>) {
745 // 更新自身状态字段
746 self.members = members;
747 // 当前作用域结束
748 }
749
750 // 双向隔离两节点
751 pub fn partition(&self, a: NodeID, b: NodeID) {
752 // 语句/调用结束
753 self.transport.partition(a, b);
754 // 当前作用域结束
755 }
756
757 // 左右组双向分区(多数/少数场景)
758 pub fn partition_groups(&self, left: &[NodeID], right: &[NodeID]) {
759 // put_on 的 KV 响应内层 match 结束
760 self.transport.partition_groups(left, right);
761 // 当前作用域结束
762 }
763
764 // put_on 协议响应 match 结束
765 pub fn heal_all(&self) {
766 // put_on 方法结束
767 self.transport.heal_all();
768 // 当前作用域结束
769 }
770
771 // 设置随机丢包率并夹紧到[0,1]
772 pub fn set_drop_rate(&self, rate: f64) {
773 // 设置随机丢包率并夹紧到[0,1]
774 self.transport.set_drop_rate(rate);
775 // 当前作用域结束
776 }
777
778 // 开关乱序;关闭时冲刷 held
779 pub fn set_reorder(&self, on: bool) {
780 // 开关乱序;关闭时冲刷 held
781 self.transport.set_reorder(on);
782 // 当前作用域结束
783 }
784
785 // 节点控制表是否仍登记运行
786 pub fn is_running(&self, id: NodeID) -> bool {
787 // 业务:self.nodes.contains_key(&id)
788 self.nodes.contains_key(&id)
789 // 当前作用域结束
790 }
791// get 的 KV 响应内层 match 结束
792}
793
794/// 测试/演示用较快超时。
795pub fn test_options() -> Options {
796 // get 便捷方法结束
797 Options {
798 // 心跳间隔(tick 数)
799 heartbeat_interval: 2,
800 // 选举超时随机区间
801 election_timeout_range: 5..10,
802 // 单次 AppendEntries 批量上限
803 max_append_entries: 100,
804 // 集成测试默认开启;若 flaky 可在具体用例里覆盖。
805 pre_vote: true,
806 // 领导者定期确认多数派存活
807 check_quorum: true,
808 // 快照阈值(0 表示默认/关闭)
809 snapshot_threshold: 0,
810 // 当前作用域结束
811 }
812// 当前作用域结束
813}
814
815// 启动单节点线程:邮箱、日志、事件循环
816fn spawn_node(
817 // 业务:id: NodeID,
818 id: NodeID,
819 // 逻辑成员集合(peers 计算用)
820 members: &HashSet<NodeID>,
821 // scan 的 KV 响应内层 match 结束
822 opts: Options,
823 // 共享故障注入传输
824 transport: Transport,
825// 进入代码块
826) -> (NodeControl, RequestTx) {
827 // scan 协议响应 match 结束
828 let peers: HashSet<NodeID> = members.iter().copied().filter(|&p| p != id).collect();
829 // scan 便捷方法结束
830 let (inbox_tx, inbox_rx) = channel::unbounded();
831 // 语句/调用结束
832 transport.register(id, inbox_tx);
833
834 // 协议层请求
835 let (request_tx, request_rx) = channel::unbounded();
836 // 发送端通道
837 let (stop_tx, stop_rx) = channel::bounded(1);
838 // 发送端通道
839 let (node_tx, node_rx) = channel::unbounded();
840
841 // 并行测试会同时打开多个 BitCask;路径必须全局唯一以免文件锁冲突。
842 static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
843 // 单调写序号/路径唯一序号
844 let seq = SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
845 // status 响应 match 结束
846 let path = std::env::temp_dir().join(format!(
847 // status 便捷方法结束
848 "raft-cluster-{}-{}-{}-{}.log",
849 // 路径嵌入 pid 避免跨进程冲突
850 std::process::id(),
851 // 业务:id,
852 id,
853 // 业务:seq,
854 seq,
855 // 业务:std::time::SystemTime::now()
856 std::time::SystemTime::now()
857 // 纳秒时间戳进一步保证路径唯一
858 .duration_since(std::time::UNIX_EPOCH)
859 // 时间戳换算(启动期不可失败)
860 .unwrap()
861 // 纳秒时间戳进一步保证路径唯一
862 .as_nanos()
863 // 语句/调用结束
864 ));
865 // 日志或会话文件路径
866 let log = Log::new(Box::new(BitCask::new(path).expect("bitcask"))).expect("log");
867 // status_on 响应 match 结束
868 let node = Node::new(id, peers, log, Kv::new(), node_tx, opts).expect("node");
869
870 // 共享传输
871 let transport_out = transport.clone();
872 // 后台启动节点事件循环
873 thread::spawn(move || {
874 // 单节点主循环:tick/入站/出站/请求
875 run_node(node, inbox_rx, node_rx, request_rx, stop_rx, transport_out);
876 // 语句/调用结束
877 });
878
879 // 业务:(NodeControl { stop_tx }, request_tx)
880 (NodeControl { stop_tx }, request_tx)
881// 当前作用域结束
882}
883
884// 单节点主循环:tick/入站/出站/请求
885fn run_node(
886 // change_membership 响应 match 结束
887 mut node: Node,
888 // change_membership 方法结束
889 peers_rx: Receiver<Envelope>,
890 // Client 实现块结束
891 node_rx: Receiver<Envelope>,
892 // 业务:request_rx: Receiver<(Request, Sender<Re...
893 request_rx: Receiver<(Request, Sender<Result<Response>>)>,
894 // 业务:stop_rx: Receiver<()>,
895 stop_rx: Receiver<()>,
896 // 共享故障注入传输
897 transport: Transport,
898// 进入代码块
899) {
900 // 协议 tick 时钟
901 let ticker = channel::tick(TICK_INTERVAL);
902 // 集群响应
903 let mut response_txs: HashMap<RequestID, Sender<Result<Response>>> = HashMap::new();
904 // 构造或 tick/step 后的 Node
905 let node_id = node.id();
906
907 // 直到 stop 或致命错误
908 loop {
909 // 多路复用 tick/入站/出站/客户端请求
910 crossbeam::select! {
911 // 集群 stop:退出事件循环
912 recv(stop_rx) -> _ => break,
913
914 // 周期 tick:推进超时与心跳
915 recv(ticker) -> _ => {
916 // 推进协议时钟(选举/心跳)
917 node = match node.tick() {
918 // 成功路径:推进状态或返回
919 Ok(n) => n,
920 // 超时或通道错误处理
921 Err(_) => break,
922 // match 分支结束
923 };
924 // 尚无主时休眠再询的分支结束
925 }
926
927 // 传输层投递的对端消息
928 recv(peers_rx) -> msg => {
929 // wait_for_leader 单次 status 匹配结束
930 let Ok(msg) = msg else { break };
931 // 等待选主轮询循环结束
932 node = match node.step(msg) {
933 // 成功路径:推进状态或返回
934 Ok(n) => n,
935 // 超时或通道错误处理
936 Err(_) => break,
937 // wait_for_leader 辅助函数结束
938 };
939 // 当前作用域结束
940 }
941
942 // 本节点协议出站信封
943 recv(node_rx) -> msg => {
944 // 构造/接收的信封
945 let Ok(msg) = msg else { break };
946 // 目标是自己:处理客户端回环
947 if msg.to == node_id {
948 // 本机回环:把结果交还调用方
949 if let Message::ClientResponse { id, response } = msg.message {
950 // 可选值解构分支
951 if let Some(tx) = response_txs.remove(&id) {
952 // 集群响应
953 let _ = tx.send(response);
954 // 当前作用域结束
955 }
956 // NodeControl 仅含 stop 信号发送端
957 }
958 // 继续下一轮
959 continue;
960 // 当前作用域结束
961 }
962 // 经故障注入规则发往对端
963 transport.deliver(msg);
964 // 当前作用域结束
965 }
966
967 // 本地 Client 提交的请求
968 recv(request_rx) -> result => {
969 // 协议层请求
970 let Ok((request, response_tx)) = result else { break };
971 // 请求 ID 或节点 ID
972 let id = Uuid::new_v4();
973 // 构造/接收的信封
974 let msg = Envelope {
975 // 信封来源节点
976 from: node.id(),
977 // 信封目标节点
978 to: node.id(),
979 // 信封携带当前任期
980 term: node.term(),
981 // 封装为自发自收客户端请求
982 message: Message::ClientRequest { id, request },
983 // Cluster 字段定义结束
984 };
985 // 写入路由/分区/挂起表项
986 response_txs.insert(id, response_tx);
987 // 步进处理一封协议/客户端消息
988 node = match node.step(msg) {
989 // 成功路径:推进状态或返回
990 Ok(n) => n,
991 // 超时或通道错误处理
992 Err(_) => break,
993 // match 分支结束
994 };
995 // 当前作用域结束
996 }
997 // 当前作用域结束
998 }
999 // spawn 委托到带 Options 的实现结束
1000 }
1001
1002 // 线程退出前标离线,避免继续投递
1003 transport.set_online(node_id, false);
1004// 当前作用域结束
1005}