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}