Skip to main content

nodedb_raft/
message.rs

1// SPDX-License-Identifier: BUSL-1.1
2
3/// A single entry in the Raft log.
4///
5/// Each entry carries the term in which it was created and an opaque command
6/// payload. The state machine interprets the payload; Raft only cares about
7/// term and index for consistency.
8#[derive(
9    Debug,
10    Clone,
11    PartialEq,
12    Eq,
13    serde::Serialize,
14    serde::Deserialize,
15    rkyv::Archive,
16    rkyv::Serialize,
17    rkyv::Deserialize,
18    zerompk::ToMessagePack,
19    zerompk::FromMessagePack,
20)]
21pub struct LogEntry {
22    /// The term when this entry was received by the leader.
23    pub term: u64,
24    /// Log index (1-based, monotonically increasing).
25    pub index: u64,
26    /// Opaque command for the state machine. Empty for no-op entries
27    /// (appended by newly elected leaders per Raft paper ยง5.4.2).
28    pub data: Vec<u8>,
29}
30
31/// AppendEntries RPC (Raft paper Figure 2).
32///
33/// Invoked by leader to replicate log entries; also used as heartbeat
34/// (entries is empty).
35#[derive(
36    Debug,
37    Clone,
38    serde::Serialize,
39    serde::Deserialize,
40    rkyv::Archive,
41    rkyv::Serialize,
42    rkyv::Deserialize,
43    zerompk::ToMessagePack,
44    zerompk::FromMessagePack,
45)]
46pub struct AppendEntriesRequest {
47    /// Leader's term.
48    pub term: u64,
49    /// Leader's ID so followers can redirect clients.
50    pub leader_id: u64,
51    /// Index of log entry immediately preceding new ones.
52    pub prev_log_index: u64,
53    /// Term of prev_log_index entry.
54    pub prev_log_term: u64,
55    /// Log entries to store (empty for heartbeat).
56    pub entries: Vec<LogEntry>,
57    /// Leader's commit_index.
58    pub leader_commit: u64,
59    /// Raft group ID for Multi-Raft routing.
60    pub group_id: u64,
61}
62
63#[derive(
64    Debug,
65    Clone,
66    serde::Serialize,
67    serde::Deserialize,
68    rkyv::Archive,
69    rkyv::Serialize,
70    rkyv::Deserialize,
71    zerompk::ToMessagePack,
72    zerompk::FromMessagePack,
73)]
74pub struct AppendEntriesResponse {
75    /// Current term, for leader to update itself.
76    pub term: u64,
77    /// True if follower contained entry matching prev_log_index and prev_log_term.
78    pub success: bool,
79    /// Optimization: on rejection, the follower's last log index.
80    /// Allows leader to skip back faster than decrementing one-by-one.
81    pub last_log_index: u64,
82}
83
84/// RequestVote RPC (Raft paper Figure 2).
85#[derive(
86    Debug,
87    Clone,
88    serde::Serialize,
89    serde::Deserialize,
90    rkyv::Archive,
91    rkyv::Serialize,
92    rkyv::Deserialize,
93    zerompk::ToMessagePack,
94    zerompk::FromMessagePack,
95)]
96pub struct RequestVoteRequest {
97    /// Candidate's term.
98    pub term: u64,
99    /// Candidate requesting vote.
100    pub candidate_id: u64,
101    /// Index of candidate's last log entry.
102    pub last_log_index: u64,
103    /// Term of candidate's last log entry.
104    pub last_log_term: u64,
105    /// Raft group ID for Multi-Raft routing.
106    pub group_id: u64,
107}
108
109#[derive(
110    Debug,
111    Clone,
112    serde::Serialize,
113    serde::Deserialize,
114    rkyv::Archive,
115    rkyv::Serialize,
116    rkyv::Deserialize,
117    zerompk::ToMessagePack,
118    zerompk::FromMessagePack,
119)]
120pub struct RequestVoteResponse {
121    /// Current term, for candidate to update itself.
122    pub term: u64,
123    /// True means candidate received vote.
124    pub vote_granted: bool,
125}
126
127/// TimeoutNow RPC (Raft thesis โ€” leadership transfer).
128///
129/// Sent by a leader to a caught-up target voter to make it immediately start
130/// an election (bypassing its election timeout). It carries no entries and
131/// grants no vote โ€” the recipient runs a normal election at `term + 1`.
132#[derive(
133    Debug,
134    Clone,
135    serde::Serialize,
136    serde::Deserialize,
137    rkyv::Archive,
138    rkyv::Serialize,
139    rkyv::Deserialize,
140    zerompk::ToMessagePack,
141    zerompk::FromMessagePack,
142)]
143pub struct TimeoutNowRequest {
144    /// Leader's term at the time the transfer was initiated.
145    pub term: u64,
146    /// The leader initiating the transfer.
147    pub leader_id: u64,
148    /// Raft group ID for Multi-Raft routing.
149    pub group_id: u64,
150}
151
152/// InstallSnapshot RPC (Raft paper Figure 13).
153///
154/// Used when a follower is too far behind for log-based catch-up.
155#[derive(
156    Debug,
157    Clone,
158    serde::Serialize,
159    serde::Deserialize,
160    rkyv::Archive,
161    rkyv::Serialize,
162    rkyv::Deserialize,
163    zerompk::ToMessagePack,
164    zerompk::FromMessagePack,
165)]
166#[msgpack(map)]
167pub struct InstallSnapshotRequest {
168    /// Leader's term.
169    pub term: u64,
170    /// Leader ID.
171    pub leader_id: u64,
172    /// The snapshot replaces all entries up through and including this index.
173    pub last_included_index: u64,
174    /// Term of last_included_index.
175    pub last_included_term: u64,
176    /// Byte offset where chunk is positioned in the snapshot file.
177    pub offset: u64,
178    /// Raw bytes of the snapshot chunk.
179    pub data: Vec<u8>,
180    /// True if this is the last chunk.
181    pub done: bool,
182    /// Raft group ID for Multi-Raft routing.
183    pub group_id: u64,
184    /// Total snapshot size in bytes. `0` means "unknown" (legacy senders or
185    /// bootstrap stubs). Receivers use this only as an advisory hint; they
186    /// must not reject chunks when `total_size == 0`.
187    #[serde(default)]
188    #[msgpack(default)]
189    pub total_size: u64,
190}
191
192#[derive(
193    Debug,
194    Clone,
195    serde::Serialize,
196    serde::Deserialize,
197    rkyv::Archive,
198    rkyv::Serialize,
199    rkyv::Deserialize,
200    zerompk::ToMessagePack,
201    zerompk::FromMessagePack,
202)]
203pub struct InstallSnapshotResponse {
204    /// Current term, for leader to update itself.
205    pub term: u64,
206}
207
208#[cfg(test)]
209mod tests {
210    use super::*;
211
212    #[test]
213    fn log_entry_serde_roundtrip() {
214        let entry = LogEntry {
215            term: 5,
216            index: 42,
217            data: b"put key=val".to_vec(),
218        };
219        let json = sonic_rs::to_string(&entry).unwrap();
220        let decoded: LogEntry = sonic_rs::from_str(&json).unwrap();
221        assert_eq!(entry, decoded);
222    }
223
224    #[test]
225    fn append_entries_heartbeat() {
226        let req = AppendEntriesRequest {
227            term: 3,
228            leader_id: 1,
229            prev_log_index: 10,
230            prev_log_term: 2,
231            entries: vec![],
232            leader_commit: 8,
233            group_id: 0,
234        };
235        assert!(req.entries.is_empty());
236    }
237
238    #[test]
239    fn request_vote_serde_roundtrip() {
240        let req = RequestVoteRequest {
241            term: 7,
242            candidate_id: 2,
243            last_log_index: 100,
244            last_log_term: 6,
245            group_id: 5,
246        };
247        let json = sonic_rs::to_string(&req).unwrap();
248        let decoded: RequestVoteRequest = sonic_rs::from_str(&json).unwrap();
249        assert_eq!(req.term, decoded.term);
250        assert_eq!(req.candidate_id, decoded.candidate_id);
251    }
252}