Skip to main content

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}