1use std::path::Path;
17
18use redb::{Database, ReadableTable, TableDefinition};
19use tracing::{debug, info};
20
21use nodedb_raft::message::LogEntry;
22use nodedb_raft::state::HardState;
23use nodedb_raft::storage::LogStorage;
24
25use crate::wire_version::envelope::{decode_versioned, encode_versioned};
26
27const ENTRIES: TableDefinition<&[u8], &[u8]> = TableDefinition::new("raft.entries");
29
30const META: TableDefinition<&str, &[u8]> = TableDefinition::new("raft.meta");
32
33const KEY_HARD_STATE: &str = "hard_state";
34const KEY_SNAPSHOT_INDEX: &str = "snapshot_index";
35const KEY_SNAPSHOT_TERM: &str = "snapshot_term";
36const KEY_APPLIED_INDEX: &str = "applied_index";
37
38pub struct RedbLogStorage {
40 db: Database,
41}
42
43impl RedbLogStorage {
44 pub fn open(path: &Path) -> crate::Result<Self> {
46 if let Some(parent) = path.parent() {
47 std::fs::create_dir_all(parent).map_err(|e| crate::ClusterError::Storage {
48 detail: format!("create raft storage dir: {e}"),
49 })?;
50 }
51
52 let db = Database::create(path).map_err(|e| crate::ClusterError::Storage {
53 detail: format!("open raft storage: {e}"),
54 })?;
55
56 let write_txn = db.begin_write().map_err(|e| crate::ClusterError::Storage {
58 detail: format!("init raft tables: {e}"),
59 })?;
60 {
61 write_txn
62 .open_table(ENTRIES)
63 .map_err(|e| crate::ClusterError::Storage {
64 detail: format!("create entries table: {e}"),
65 })?;
66 write_txn
67 .open_table(META)
68 .map_err(|e| crate::ClusterError::Storage {
69 detail: format!("create meta table: {e}"),
70 })?;
71 }
72 write_txn
73 .commit()
74 .map_err(|e| crate::ClusterError::Storage {
75 detail: format!("commit raft init: {e}"),
76 })?;
77
78 info!(path = %path.display(), "raft log storage opened");
79
80 Ok(Self { db })
81 }
82}
83
84fn index_key(index: u64) -> [u8; 8] {
85 index.to_be_bytes()
86}
87
88impl LogStorage for RedbLogStorage {
89 fn append(&mut self, entries: &[LogEntry]) -> nodedb_raft::error::Result<()> {
90 if entries.is_empty() {
91 return Ok(());
92 }
93
94 let write_txn =
95 self.db
96 .begin_write()
97 .map_err(|e| nodedb_raft::error::RaftError::Storage {
98 detail: format!("write txn: {e}"),
99 })?;
100 {
101 let mut table = write_txn.open_table(ENTRIES).map_err(|e| {
102 nodedb_raft::error::RaftError::Storage {
103 detail: format!("open entries: {e}"),
104 }
105 })?;
106
107 for entry in entries {
108 let key = index_key(entry.index);
109 let value = encode_versioned(entry).map_err(|e| {
110 nodedb_raft::error::RaftError::Storage {
111 detail: format!("serialize entry: {e}"),
112 }
113 })?;
114 table
115 .insert(key.as_slice(), value.as_slice())
116 .map_err(|e| nodedb_raft::error::RaftError::Storage {
117 detail: format!("insert entry: {e}"),
118 })?;
119 }
120 }
121 write_txn
122 .commit()
123 .map_err(|e| nodedb_raft::error::RaftError::Storage {
124 detail: format!("commit append: {e}"),
125 })?;
126
127 debug!(count = entries.len(), "raft log appended");
128 Ok(())
129 }
130
131 fn truncate(&mut self, index: u64) -> nodedb_raft::error::Result<()> {
132 let write_txn =
133 self.db
134 .begin_write()
135 .map_err(|e| nodedb_raft::error::RaftError::Storage {
136 detail: format!("write txn: {e}"),
137 })?;
138 {
139 let mut table = write_txn.open_table(ENTRIES).map_err(|e| {
140 nodedb_raft::error::RaftError::Storage {
141 detail: format!("open entries: {e}"),
142 }
143 })?;
144
145 let start = index_key(index);
147 let keys_to_remove: Vec<[u8; 8]> = table
148 .range(start.as_slice()..)
149 .map_err(|e| nodedb_raft::error::RaftError::Storage {
150 detail: format!("range: {e}"),
151 })?
152 .filter_map(|r| {
153 r.ok().map(|(k, _)| {
154 let mut buf = [0u8; 8];
155 buf.copy_from_slice(k.value());
156 buf
157 })
158 })
159 .collect();
160
161 for key in &keys_to_remove {
162 table.remove(key.as_slice()).map_err(|e| {
163 nodedb_raft::error::RaftError::Storage {
164 detail: format!("remove: {e}"),
165 }
166 })?;
167 }
168 }
169 write_txn
170 .commit()
171 .map_err(|e| nodedb_raft::error::RaftError::Storage {
172 detail: format!("commit truncate: {e}"),
173 })?;
174
175 debug!(from_index = index, "raft log truncated");
176 Ok(())
177 }
178
179 fn load_entries_after(&self, snapshot_index: u64) -> nodedb_raft::error::Result<Vec<LogEntry>> {
180 let read_txn =
181 self.db
182 .begin_read()
183 .map_err(|e| nodedb_raft::error::RaftError::Storage {
184 detail: format!("read txn: {e}"),
185 })?;
186 let table =
187 read_txn
188 .open_table(ENTRIES)
189 .map_err(|e| nodedb_raft::error::RaftError::Storage {
190 detail: format!("open entries: {e}"),
191 })?;
192
193 let start = index_key(snapshot_index + 1);
195 let mut entries = Vec::new();
196
197 for result in
198 table
199 .range(start.as_slice()..)
200 .map_err(|e| nodedb_raft::error::RaftError::Storage {
201 detail: format!("range: {e}"),
202 })?
203 {
204 let (_, value) = result.map_err(|e| nodedb_raft::error::RaftError::Storage {
205 detail: format!("entry read: {e}"),
206 })?;
207 let entry: LogEntry = decode_versioned(value.value()).map_err(|e| {
208 nodedb_raft::error::RaftError::Storage {
209 detail: format!("deserialize entry: {e}"),
210 }
211 })?;
212 entries.push(entry);
213 }
214
215 debug!(
216 count = entries.len(),
217 after = snapshot_index,
218 "raft log loaded"
219 );
220 Ok(entries)
221 }
222
223 fn compact(&mut self, index: u64, term: u64) -> nodedb_raft::error::Result<()> {
224 let write_txn =
225 self.db
226 .begin_write()
227 .map_err(|e| nodedb_raft::error::RaftError::Storage {
228 detail: format!("write txn: {e}"),
229 })?;
230 {
231 let mut table = write_txn.open_table(ENTRIES).map_err(|e| {
233 nodedb_raft::error::RaftError::Storage {
234 detail: format!("open entries: {e}"),
235 }
236 })?;
237
238 let end = index_key(index + 1);
239 let keys_to_remove: Vec<[u8; 8]> = table
240 .range(..end.as_slice())
241 .map_err(|e| nodedb_raft::error::RaftError::Storage {
242 detail: format!("range: {e}"),
243 })?
244 .filter_map(|r| {
245 r.ok().map(|(k, _)| {
246 let mut buf = [0u8; 8];
247 buf.copy_from_slice(k.value());
248 buf
249 })
250 })
251 .collect();
252
253 for key in &keys_to_remove {
254 table.remove(key.as_slice()).map_err(|e| {
255 nodedb_raft::error::RaftError::Storage {
256 detail: format!("remove: {e}"),
257 }
258 })?;
259 }
260
261 let mut meta =
263 write_txn
264 .open_table(META)
265 .map_err(|e| nodedb_raft::error::RaftError::Storage {
266 detail: format!("open meta: {e}"),
267 })?;
268
269 let idx_bytes = zerompk::to_msgpack_vec(&index).map_err(|e| {
270 nodedb_raft::error::RaftError::Storage {
271 detail: format!("serialize: {e}"),
272 }
273 })?;
274 let term_bytes = zerompk::to_msgpack_vec(&term).map_err(|e| {
275 nodedb_raft::error::RaftError::Storage {
276 detail: format!("serialize: {e}"),
277 }
278 })?;
279
280 meta.insert(KEY_SNAPSHOT_INDEX, idx_bytes.as_slice())
281 .map_err(|e| nodedb_raft::error::RaftError::Storage {
282 detail: format!("insert meta: {e}"),
283 })?;
284 meta.insert(KEY_SNAPSHOT_TERM, term_bytes.as_slice())
285 .map_err(|e| nodedb_raft::error::RaftError::Storage {
286 detail: format!("insert meta: {e}"),
287 })?;
288 }
289 write_txn
290 .commit()
291 .map_err(|e| nodedb_raft::error::RaftError::Storage {
292 detail: format!("commit compact: {e}"),
293 })?;
294
295 debug!(index, term, "raft log compacted");
296 Ok(())
297 }
298
299 fn snapshot_metadata(&self) -> (u64, u64) {
300 let Ok(read_txn) = self.db.begin_read() else {
301 return (0, 0);
302 };
303 let Ok(table) = read_txn.open_table(META) else {
304 return (0, 0);
305 };
306
307 let index = table
308 .get(KEY_SNAPSHOT_INDEX)
309 .ok()
310 .flatten()
311 .and_then(|v| zerompk::from_msgpack::<u64>(v.value()).ok())
312 .unwrap_or(0);
313 let term = table
314 .get(KEY_SNAPSHOT_TERM)
315 .ok()
316 .flatten()
317 .and_then(|v| zerompk::from_msgpack::<u64>(v.value()).ok())
318 .unwrap_or(0);
319
320 (index, term)
321 }
322
323 fn save_hard_state(&mut self, state: &HardState) -> nodedb_raft::error::Result<()> {
324 let write_txn =
325 self.db
326 .begin_write()
327 .map_err(|e| nodedb_raft::error::RaftError::Storage {
328 detail: format!("write txn: {e}"),
329 })?;
330 {
331 let mut table =
332 write_txn
333 .open_table(META)
334 .map_err(|e| nodedb_raft::error::RaftError::Storage {
335 detail: format!("open meta: {e}"),
336 })?;
337
338 let bytes = zerompk::to_msgpack_vec(state).map_err(|e| {
339 nodedb_raft::error::RaftError::Storage {
340 detail: format!("serialize: {e}"),
341 }
342 })?;
343 table
344 .insert(KEY_HARD_STATE, bytes.as_slice())
345 .map_err(|e| nodedb_raft::error::RaftError::Storage {
346 detail: format!("insert: {e}"),
347 })?;
348 }
349 write_txn
350 .commit()
351 .map_err(|e| nodedb_raft::error::RaftError::Storage {
352 detail: format!("commit: {e}"),
353 })?;
354
355 debug!(
356 term = state.current_term,
357 voted_for = state.voted_for,
358 "raft hard state saved"
359 );
360 Ok(())
361 }
362
363 fn load_hard_state(&self) -> nodedb_raft::error::Result<HardState> {
364 let read_txn =
365 self.db
366 .begin_read()
367 .map_err(|e| nodedb_raft::error::RaftError::Storage {
368 detail: format!("read txn: {e}"),
369 })?;
370 let table =
371 read_txn
372 .open_table(META)
373 .map_err(|e| nodedb_raft::error::RaftError::Storage {
374 detail: format!("open meta: {e}"),
375 })?;
376
377 match table.get(KEY_HARD_STATE) {
378 Ok(Some(value)) => {
379 let state: HardState = zerompk::from_msgpack(value.value()).map_err(|e| {
380 nodedb_raft::error::RaftError::Storage {
381 detail: format!("deserialize: {e}"),
382 }
383 })?;
384 Ok(state)
385 }
386 Ok(None) => Ok(HardState::default()),
387 Err(e) => Err(nodedb_raft::error::RaftError::Storage {
388 detail: format!("get hard state: {e}"),
389 }),
390 }
391 }
392
393 fn save_applied_index(&mut self, index: u64) -> nodedb_raft::error::Result<()> {
394 let write_txn =
395 self.db
396 .begin_write()
397 .map_err(|e| nodedb_raft::error::RaftError::Storage {
398 detail: format!("write txn: {e}"),
399 })?;
400 {
401 let mut table =
402 write_txn
403 .open_table(META)
404 .map_err(|e| nodedb_raft::error::RaftError::Storage {
405 detail: format!("open meta: {e}"),
406 })?;
407
408 let bytes = zerompk::to_msgpack_vec(&index).map_err(|e| {
409 nodedb_raft::error::RaftError::Storage {
410 detail: format!("serialize: {e}"),
411 }
412 })?;
413 table
414 .insert(KEY_APPLIED_INDEX, bytes.as_slice())
415 .map_err(|e| nodedb_raft::error::RaftError::Storage {
416 detail: format!("insert: {e}"),
417 })?;
418 }
419 write_txn
420 .commit()
421 .map_err(|e| nodedb_raft::error::RaftError::Storage {
422 detail: format!("commit: {e}"),
423 })?;
424
425 debug!(index, "raft applied index saved");
426 Ok(())
427 }
428
429 fn load_applied_index(&self) -> nodedb_raft::error::Result<u64> {
430 let read_txn =
431 self.db
432 .begin_read()
433 .map_err(|e| nodedb_raft::error::RaftError::Storage {
434 detail: format!("read txn: {e}"),
435 })?;
436 let table =
437 read_txn
438 .open_table(META)
439 .map_err(|e| nodedb_raft::error::RaftError::Storage {
440 detail: format!("open meta: {e}"),
441 })?;
442
443 match table.get(KEY_APPLIED_INDEX) {
444 Ok(Some(value)) => {
445 let index: u64 = zerompk::from_msgpack(value.value()).map_err(|e| {
446 nodedb_raft::error::RaftError::Storage {
447 detail: format!("deserialize: {e}"),
448 }
449 })?;
450 Ok(index)
451 }
452 Ok(None) => Ok(0),
457 Err(e) => Err(nodedb_raft::error::RaftError::Storage {
458 detail: format!("get applied index: {e}"),
459 }),
460 }
461 }
462}
463
464#[cfg(test)]
465mod tests {
466 use super::*;
467
468 fn open_temp() -> (RedbLogStorage, tempfile::TempDir) {
469 let dir = tempfile::tempdir().unwrap();
470 let path = dir.path().join("test-raft.redb");
471 let storage = RedbLogStorage::open(&path).unwrap();
472 (storage, dir)
473 }
474
475 #[test]
476 fn append_and_load() {
477 let (mut s, _dir) = open_temp();
478 let entries = vec![
479 LogEntry {
480 term: 1,
481 index: 1,
482 data: b"cmd-a".to_vec(),
483 },
484 LogEntry {
485 term: 1,
486 index: 2,
487 data: b"cmd-b".to_vec(),
488 },
489 LogEntry {
490 term: 2,
491 index: 3,
492 data: b"cmd-c".to_vec(),
493 },
494 ];
495 s.append(&entries).unwrap();
496
497 let loaded = s.load_entries_after(0).unwrap();
498 assert_eq!(loaded.len(), 3);
499 assert_eq!(loaded[0].data, b"cmd-a");
500 assert_eq!(loaded[2].term, 2);
501 }
502
503 #[test]
504 fn truncate_removes_tail() {
505 let (mut s, _dir) = open_temp();
506 for i in 1..=5 {
507 s.append(&[LogEntry {
508 term: 1,
509 index: i,
510 data: vec![],
511 }])
512 .unwrap();
513 }
514 s.truncate(3).unwrap();
515 let loaded = s.load_entries_after(0).unwrap();
516 assert_eq!(loaded.len(), 2);
517 assert_eq!(loaded.last().unwrap().index, 2);
518 }
519
520 #[test]
521 fn compact_removes_prefix() {
522 let (mut s, _dir) = open_temp();
523 for i in 1..=10 {
524 s.append(&[LogEntry {
525 term: 1,
526 index: i,
527 data: vec![],
528 }])
529 .unwrap();
530 }
531 s.compact(5, 1).unwrap();
532 assert_eq!(s.snapshot_metadata(), (5, 1));
533 let loaded = s.load_entries_after(5).unwrap();
534 assert_eq!(loaded.len(), 5);
535 assert_eq!(loaded[0].index, 6);
536 }
537
538 #[test]
539 fn hard_state_roundtrip() {
540 let (mut s, _dir) = open_temp();
541 let hs = HardState {
542 current_term: 7,
543 voted_for: 3,
544 };
545 s.save_hard_state(&hs).unwrap();
546 let loaded = s.load_hard_state().unwrap();
547 assert_eq!(loaded, hs);
548 }
549
550 #[test]
553 fn applied_index_absent_reads_zero() {
554 let (s, _dir) = open_temp();
555 assert_eq!(s.load_applied_index().unwrap(), 0);
556 }
557
558 #[test]
559 fn applied_index_survives_reopen() {
560 let dir = tempfile::tempdir().unwrap();
561 let path = dir.path().join("applied-index-raft.redb");
562
563 {
564 let mut s = RedbLogStorage::open(&path).unwrap();
565 s.save_applied_index(9).unwrap();
566 }
567
568 let s = RedbLogStorage::open(&path).unwrap();
569 assert_eq!(s.load_applied_index().unwrap(), 9);
570 }
571
572 #[test]
574 fn versioned_roundtrip() {
575 let (mut s, _dir) = open_temp();
576 let entry = LogEntry {
577 term: 3,
578 index: 1,
579 data: b"versioned-payload".to_vec(),
580 };
581 s.append(std::slice::from_ref(&entry)).unwrap();
582 let loaded = s.load_entries_after(0).unwrap();
583 assert_eq!(loaded.len(), 1);
584 assert_eq!(loaded[0], entry);
585 }
586
587 #[test]
590 fn version_mismatch_returns_error() {
591 let (s, _dir) = open_temp();
592
593 let inner = zerompk::to_msgpack_vec(&LogEntry {
595 term: 1,
596 index: 77,
597 data: vec![],
598 })
599 .unwrap();
600
601 let mut bad_bytes = Vec::new();
603 bad_bytes.push(0xc1u8); bad_bytes.extend_from_slice(&9999u16.to_be_bytes());
605 bad_bytes.extend_from_slice(&(inner.len() as u32).to_be_bytes());
606 bad_bytes.extend_from_slice(&inner);
607
608 let key = index_key(77);
609 let write_txn = s.db.begin_write().unwrap();
610 {
611 let mut table = write_txn.open_table(ENTRIES).unwrap();
612 table.insert(key.as_slice(), bad_bytes.as_slice()).unwrap();
613 }
614 write_txn.commit().unwrap();
615
616 let err = s.load_entries_after(76).unwrap_err();
617 match err {
618 nodedb_raft::error::RaftError::Storage { detail } => {
619 assert!(
620 detail.contains("deserialize"),
621 "expected 'deserialize' in detail: {detail}"
622 );
623 }
624 other => panic!("expected Storage error, got: {other}"),
625 }
626 }
627
628 #[test]
629 fn survives_reopen() {
630 let dir = tempfile::tempdir().unwrap();
631 let path = dir.path().join("reopen-raft.redb");
632
633 {
634 let mut s = RedbLogStorage::open(&path).unwrap();
635 s.append(&[LogEntry {
636 term: 1,
637 index: 1,
638 data: b"durable".to_vec(),
639 }])
640 .unwrap();
641 s.save_hard_state(&HardState {
642 current_term: 3,
643 voted_for: 1,
644 })
645 .unwrap();
646 }
647
648 let s = RedbLogStorage::open(&path).unwrap();
649 let loaded = s.load_entries_after(0).unwrap();
650 assert_eq!(loaded.len(), 1);
651 assert_eq!(loaded[0].data, b"durable");
652 let hs = s.load_hard_state().unwrap();
653 assert_eq!(hs.current_term, 3);
654 }
655
656 #[test]
664 fn vote_and_log_survive_restart() {
665 use nodedb_raft::node::RaftConfig;
666 use nodedb_raft::{AppendEntriesRequest, RaftNode, RequestVoteRequest};
667 use std::time::Duration;
668
669 let dir = tempfile::tempdir().unwrap();
670 let path = dir.path().join("vote-restart-raft.redb");
671
672 let config = || RaftConfig {
673 node_id: 1,
674 group_id: 7,
675 peers: vec![2, 3],
676 learners: vec![],
677 observers: vec![],
678 starts_as_learner: false,
679 starts_as_observer: false,
680 election_timeout_min: Duration::from_millis(150),
681 election_timeout_max: Duration::from_millis(300),
682 heartbeat_interval: Duration::from_millis(50),
683 log_compaction_threshold: None,
684 };
685
686 const TERM: u64 = 5;
687
688 {
689 let storage = RedbLogStorage::open(&path).unwrap();
690 let mut node = RaftNode::new(config(), storage);
691 node.restore().unwrap();
692
693 let ae = AppendEntriesRequest {
696 term: TERM,
697 leader_id: 2,
698 prev_log_index: 0,
699 prev_log_term: 0,
700 entries: vec![
701 LogEntry {
702 term: TERM,
703 index: 1,
704 data: b"e1".to_vec(),
705 },
706 LogEntry {
707 term: TERM,
708 index: 2,
709 data: b"e2".to_vec(),
710 },
711 ],
712 leader_commit: 0,
713 group_id: 7,
714 };
715 assert!(node.handle_append_entries(&ae).success);
716 node.persist_hard_state_if_dirty().unwrap();
717
718 let rv = RequestVoteRequest {
720 term: TERM,
721 candidate_id: 2,
722 last_log_index: 2,
723 last_log_term: TERM,
724 group_id: 7,
725 };
726 assert!(
727 node.handle_request_vote(&rv).vote_granted,
728 "first vote must be granted"
729 );
730 node.persist_hard_state_if_dirty().unwrap();
731 }
732
733 let storage = RedbLogStorage::open(&path).unwrap();
735 let mut node = RaftNode::new(config(), storage);
736 node.restore().unwrap();
737
738 assert_eq!(node.last_log_index(), 2, "log entries must survive restart");
740 assert_eq!(node.current_term(), TERM);
742
743 let rv2 = RequestVoteRequest {
748 term: TERM,
749 candidate_id: 3,
750 last_log_index: 2,
751 last_log_term: TERM,
752 group_id: 7,
753 };
754 assert!(
755 !node.handle_request_vote(&rv2).vote_granted,
756 "restarted voter must not grant a second vote in the same term"
757 );
758 assert_eq!(node.current_term(), TERM);
759 }
760}