Skip to main content

raft_rust/raft/
session.rs

1//! 客户端 session 去重:把 (client_id, seq) 编入日志命令,apply 时幂等。
2
3// 每个 client 缓存最后一次成功写的 seq 与响应
4use std::collections::HashMap;
5
6// SessionCommand 写入日志需可序列化
7use serde::{Deserialize, Serialize};
8// 客户端会话标识
9use uuid::Uuid;
10
11// 包装底层 State
12use super::{Entry, Index, State};
13// apply/read 错误类型
14use crate::error::Result;
15
16// session 编解码与快照使用统一 bincode 配置
17const BINCODE: bincode::config::Configuration = bincode::config::standard();
18
19/// 写入日志的 session 包装命令。
20#[derive(Clone, Debug, Serialize, Deserialize)]
21// 定义数据结构
22pub struct SessionCommand {
23    // 客户端会话 ID
24    pub client_id: Uuid,
25    // 会话内单调序号
26    pub seq: u64,
27    // 真正交给内层状态机的命令
28    pub payload: Vec<u8>,
29// 结束当前作用域
30}
31
32// 将 session 元数据 + 应用命令编码为日志命令字节
33pub fn encode_session(client_id: Uuid, seq: u64, payload: Vec<u8>) -> Vec<u8> {
34    // 魔数前缀,避免与裸应用命令混淆。
35    let mut out = b"SESS".to_vec();
36    // 序列化 SessionCommand 体
37    let body = bincode::serde::encode_to_vec(
38        // 借用参数
39        &SessionCommand { client_id, seq, payload },
40        // 列表或字段续项
41        BINCODE,
42    // 结束调用参数列表
43    )
44    // 失败则说明原因
45    .expect("session encode");
46    // 拼在魔数之后
47    out.extend_from_slice(&body);
48    // 返回完整日志命令
49    out
50// 结束当前作用域
51}
52
53// 尝试从日志命令字节解析 session 包装;非 session 返回 None
54pub fn decode_session(bytes: &[u8]) -> Option<SessionCommand> {
55    // 长度不足或魔数不匹配 → 非 session 命令
56    if bytes.len() < 4 || &bytes[..4] != b"SESS" {
57        // 提前返回
58        return None;
59    // 结束当前作用域
60    }
61    // 解码 body;失败视为非 session(或损坏,上层当普通命令处理可能再失败)
62    bincode::serde::borrow_decode_from_slice(&bytes[4..], BINCODE)
63        // 方法链调用
64        .ok()
65        // 链式变换结果
66        .map(|(c, _)| c)
67// 结束当前作用域
68}
69
70/// 包装任意 State,对带 session 头的写命令做去重。
71pub struct SessionState {
72    // 被包装的真实应用状态机
73    inner: Box<dyn State>,
74    /// client_id -> (last_seq, cached response)
75    sessions: HashMap<Uuid, (u64, Vec<u8>)>,
76// 结束当前作用域
77}
78
79// 为类型实现方法
80impl SessionState {
81    // 用内层状态机构造带 session 去重的包装
82    pub fn new(inner: Box<dyn State>) -> Box<Self> {
83        // 初始无会话缓存
84        Box::new(Self { inner, sessions: HashMap::new() })
85    // 结束当前作用域
86    }
87
88    // 拆出内层状态机(测试或迁移用)
89    pub fn into_inner(self) -> Box<dyn State> {
90        // 业务逻辑步骤
91        self.inner
92    // 结束当前作用域
93    }
94// 结束当前作用域
95}
96
97// 实现 State:写路径做幂等,读/快照委托并附加 session 表
98impl State for SessionState {
99    // 委托内层 applied 索引
100    fn get_applied_index(&self) -> Index {
101        // 业务逻辑步骤
102        self.inner.get_applied_index()
103    // 结束当前作用域
104    }
105
106    // apply:识别 session 头并按 seq 幂等
107    fn apply(&mut self, entry: Entry) -> Result<Vec<u8>> {
108        // 成员变更 / noop
109        // 无 command:直接交给内层(noop 推进 applied)
110        let Some(cmd) = entry.command.as_ref() else {
111            // 提前返回
112            return self.inner.apply(entry);
113        // 结束当前作用域
114        };
115        // 若是 session 包装命令
116        if let Some(sess) = decode_session(cmd) {
117            // 幂等规则(只缓存每个 client 最后一次成功写):
118            // - seq == last:同一请求重试 / 日志里重复提出 → 返回缓存,不二次执行
119            // - seq <  last:过期序号 → noop,不返回「更新的那次」缓存(避免串结果)
120            // - seq >  last:正常执行并更新缓存
121            if let Some((last, cached)) = self.sessions.get(&sess.client_id).cloned() {
122                // 重试同一 seq:推进 applied 但不改业务状态
123                if sess.seq == last {
124                    // 构造 noop 条目推进内层 applied_index
125                    let noop = Entry {
126                        // 日志条目字段
127                        index: entry.index,
128                        // 携带当前任期
129                        term: entry.term,
130                        // 日志条目字段
131                        command: None,
132                        // 日志条目字段
133                        membership: None,
134                    // 结束当前作用域
135                    };
136                    // 忽略 noop 返回值
137                    let _ = self.inner.apply(noop)?;
138                    // 返回上次缓存的成功响应
139                    return Ok(cached);
140                // 结束当前作用域
141                }
142                // 过期 seq:仍推进 applied,返回空
143                if sess.seq < last {
144                    // 绑定中间结果
145                    let noop = Entry {
146                        // 日志条目字段
147                        index: entry.index,
148                        // 携带当前任期
149                        term: entry.term,
150                        // 日志条目字段
151                        command: None,
152                        // 日志条目字段
153                        membership: None,
154                    // 结束当前作用域
155                    };
156                    // 绑定中间结果
157                    let _ = self.inner.apply(noop)?;
158                    // 不返回旧缓存,避免客户端拿到「更新的那次」结果
159                    return Ok(Vec::new());
160                // 结束当前作用域
161                }
162            // 结束当前作用域
163            }
164            // 新 seq:解开 payload 交给内层
165            let inner_entry = Entry {
166                // 日志条目字段
167                index: entry.index,
168                // 携带当前任期
169                term: entry.term,
170                // 仅应用层命令
171                command: Some(sess.payload),
172                // session 写不应携带 membership
173                membership: None,
174            // 结束当前作用域
175            };
176            // 真正执行业务写
177            let resp = self.inner.apply(inner_entry)?;
178            // 缓存 (seq, 响应) 供重试
179            self.sessions.insert(sess.client_id, (sess.seq, resp.clone()));
180            // 提前返回
181            return Ok(resp);
182        // 结束当前作用域
183        }
184        // 非 session 命令:原样委托
185        self.inner.apply(entry)
186    // 结束当前作用域
187    }
188
189    // 读不涉及 session,直接委托
190    fn read(&self, command: Vec<u8>) -> Result<Vec<u8>> {
191        // 业务逻辑步骤
192        self.inner.read(command)
193    // 结束当前作用域
194    }
195
196    // 快照 = 内层快照 + 全部 session 缓存
197    fn snapshot(&self) -> Result<Vec<u8>> {
198        // 先导出内层状态
199        let inner_snap = self.inner.snapshot()?;
200        // 将会话表展平为可序列化列表
201        let sessions: Vec<(Uuid, u64, Vec<u8>)> = self
202            // 业务逻辑步骤
203            .sessions
204            // 方法链调用
205            .iter()
206            // 链式变换结果
207            .map(|(id, (seq, resp))| (*id, *seq, resp.clone()))
208            // 方法链调用
209            .collect();
210        // 打包编码
211        Ok(bincode::serde::encode_to_vec(&(inner_snap, sessions), BINCODE)
212            // 失败则说明原因
213            .expect("session snapshot"))
214    // 结束当前作用域
215    }
216
217    // 从组合快照恢复内层与 session 表
218    fn restore(&mut self, data: &[u8], index: Index) -> Result<()> {
219        // 解码 (内层快照, sessions)
220        let (inner_snap, sessions): (Vec<u8>, Vec<(Uuid, u64, Vec<u8>)>) =
221            // 业务逻辑步骤
222            bincode::serde::borrow_decode_from_slice(data, BINCODE)
223                // 链式变换结果
224                .map_err(|e| crate::error::Error::InvalidData(e.to_string()))?
225                // 语法续行
226                .0;
227        // 恢复内层状态机
228        self.inner.restore(&inner_snap, index)?;
229        // 重建 session 哈希表
230        self.sessions = sessions.into_iter().map(|(id, seq, resp)| (id, (seq, resp))).collect();
231        // 成功返回
232        Ok(())
233    // 结束当前作用域
234    }
235// 结束当前作用域
236}