1use crate::hash::{FNV_OFFSET, fnv64};
2use serde::{Deserialize, Serialize};
3
4const COMMIT_MARKER: u64 = 0x4943_4D45_4D43_4F4D;
5
6#[derive(Clone, Copy, Debug, Eq, PartialEq)]
7enum CommitSlotIndex {
8 Slot0,
9 Slot1,
10}
11
12impl CommitSlotIndex {
13 const fn opposite(self) -> Self {
14 match self {
15 Self::Slot0 => Self::Slot1,
16 Self::Slot1 => Self::Slot0,
17 }
18 }
19}
20
21#[derive(Clone, Copy, Debug, Eq, PartialEq)]
22struct AuthoritativeSlot<'slot> {
23 index: CommitSlotIndex,
24 record: &'slot CommittedGenerationBytes,
25}
26
27fn select_authoritative_slot<'slot>(
28 slot0: Option<&'slot CommittedGenerationBytes>,
29 slot1: Option<&'slot CommittedGenerationBytes>,
30) -> Result<AuthoritativeSlot<'slot>, CommitRecoveryError> {
31 let slot0_invalid = slot0.is_some_and(|slot| !slot.validates());
32 let slot1_invalid = slot1.is_some_and(|slot| !slot.validates());
33 if slot0_invalid || slot1_invalid {
34 return Err(CommitRecoveryError::InvalidCommitSlots {
35 slot0_invalid,
36 slot1_invalid,
37 });
38 }
39
40 let slot0 = slot0.map(|record| AuthoritativeSlot {
41 index: CommitSlotIndex::Slot0,
42 record,
43 });
44 let slot1 = slot1.map(|record| AuthoritativeSlot {
45 index: CommitSlotIndex::Slot1,
46 record,
47 });
48
49 match (slot0, slot1) {
50 (Some(left), Some(right))
51 if left.record.generation() == right.record.generation()
52 && left.record != right.record =>
53 {
54 Err(CommitRecoveryError::AmbiguousGeneration {
55 generation: left.record.generation(),
56 })
57 }
58 (Some(left), Some(right)) if right.record.generation() > left.record.generation() => {
59 Ok(right)
60 }
61 (Some(left), Some(_) | None) => Ok(left),
62 (None, Some(right)) => Ok(right),
63 (None, None) => Err(CommitRecoveryError::NoValidGeneration),
64 }
65}
66
67#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
81#[serde(deny_unknown_fields)]
82pub struct CommittedGenerationBytes {
83 pub(crate) generation: u64,
85 pub(crate) commit_marker: u64,
87 pub(crate) checksum: u64,
89 #[serde(with = "opaque_payload")]
91 pub(crate) payload: Vec<u8>,
92}
93
94mod opaque_payload {
97 use crate::constants::MAX_COMMITTED_PAYLOAD_BYTES;
98 use serde::{Deserializer, Serialize, Serializer, de::Visitor};
99
100 pub fn serialize<S: Serializer>(bytes: &[u8], serializer: S) -> Result<S::Ok, S::Error> {
101 if bytes.len() > MAX_COMMITTED_PAYLOAD_BYTES {
102 return Err(serde::ser::Error::custom(
103 "ledger byte string exceeds payload bound",
104 ));
105 }
106 if serializer.is_human_readable() {
107 bytes.serialize(serializer)
108 } else {
109 serializer.serialize_bytes(bytes)
110 }
111 }
112
113 pub fn deserialize<'de, D: Deserializer<'de>>(deserializer: D) -> Result<Vec<u8>, D::Error> {
114 struct Bytes;
115 impl Visitor<'_> for Bytes {
116 type Value = Vec<u8>;
117
118 fn expecting(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
119 f.write_str("a bounded opaque ledger byte string")
120 }
121
122 fn visit_bytes<E: serde::de::Error>(self, bytes: &[u8]) -> Result<Self::Value, E> {
123 if bytes.len() > MAX_COMMITTED_PAYLOAD_BYTES {
124 return Err(E::custom("ledger byte string exceeds payload bound"));
125 }
126 Ok(bytes.to_vec())
127 }
128
129 fn visit_byte_buf<E: serde::de::Error>(self, bytes: Vec<u8>) -> Result<Self::Value, E> {
130 if bytes.len() > MAX_COMMITTED_PAYLOAD_BYTES {
131 return Err(E::custom("ledger byte string exceeds payload bound"));
132 }
133 Ok(bytes)
134 }
135 }
136 if deserializer.is_human_readable() {
137 return crate::cbor::deserialize_bounded_vec::<D, u8, MAX_COMMITTED_PAYLOAD_BYTES>(
138 deserializer,
139 );
140 }
141 deserializer.deserialize_byte_buf(Bytes)
142 }
143}
144
145impl CommittedGenerationBytes {
146 #[must_use]
148 pub fn new(generation: u64, payload: Vec<u8>) -> Self {
149 let mut record = Self {
150 generation,
151 commit_marker: COMMIT_MARKER,
152 checksum: 0,
153 payload,
154 };
155 record.checksum = generation_checksum(&record);
156 record
157 }
158
159 #[must_use]
161 pub const fn generation(&self) -> u64 {
162 self.generation
163 }
164
165 #[must_use]
171 pub const fn commit_marker(&self) -> u64 {
172 self.commit_marker
173 }
174
175 #[must_use]
180 pub const fn checksum(&self) -> u64 {
181 self.checksum
182 }
183
184 #[must_use]
186 pub fn payload(&self) -> &[u8] {
187 &self.payload
188 }
189
190 #[must_use]
192 pub fn validates(&self) -> bool {
193 self.commit_marker == COMMIT_MARKER && self.checksum == generation_checksum(self)
194 }
195}
196
197#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
220#[serde(deny_unknown_fields)]
221pub struct DualCommitStore {
222 #[serde(deserialize_with = "crate::cbor::deserialize_present_option")]
224 pub(crate) slot0: Option<CommittedGenerationBytes>,
225 #[serde(deserialize_with = "crate::cbor::deserialize_present_option")]
227 pub(crate) slot1: Option<CommittedGenerationBytes>,
228}
229
230impl DualCommitStore {
231 #[must_use]
233 pub const fn is_uninitialized(&self) -> bool {
234 self.slot0.is_none() && self.slot1.is_none()
235 }
236
237 #[must_use]
242 pub const fn slot0(&self) -> Option<&CommittedGenerationBytes> {
243 self.slot0.as_ref()
244 }
245
246 #[must_use]
251 pub const fn slot1(&self) -> Option<&CommittedGenerationBytes> {
252 self.slot1.as_ref()
253 }
254
255 fn authoritative_slot(&self) -> Result<AuthoritativeSlot<'_>, CommitRecoveryError> {
256 select_authoritative_slot(self.slot0(), self.slot1())
257 }
258
259 #[cfg(test)]
260 fn inactive_slot_index(&self) -> CommitSlotIndex {
261 match self.authoritative_slot() {
262 Ok(authoritative) => authoritative.index.opposite(),
263 Err(_) if self.slot0.is_none() => CommitSlotIndex::Slot0,
264 Err(_) => CommitSlotIndex::Slot1,
265 }
266 }
267
268 pub fn authoritative(&self) -> Result<&CommittedGenerationBytes, CommitRecoveryError> {
270 self.authoritative_slot()
271 .map(|authoritative| authoritative.record)
272 }
273
274 #[must_use]
276 pub fn diagnostic(&self) -> CommitStoreDiagnostic {
277 CommitStoreDiagnostic::from_store(self)
278 }
279
280 pub fn commit_payload(
285 &mut self,
286 payload: Vec<u8>,
287 ) -> Result<&CommittedGenerationBytes, CommitRecoveryError> {
288 self.commit_payload_with_generation(None, payload)
289 }
290
291 pub fn commit_payload_at_generation(
302 &mut self,
303 generation: u64,
304 payload: Vec<u8>,
305 ) -> Result<&CommittedGenerationBytes, CommitRecoveryError> {
306 self.commit_payload_with_generation(Some(generation), payload)
307 }
308
309 fn commit_payload_with_generation(
310 &mut self,
311 requested: Option<u64>,
312 payload: Vec<u8>,
313 ) -> Result<&CommittedGenerationBytes, CommitRecoveryError> {
314 let (index, generation) = match self.authoritative_slot() {
316 Ok(authoritative) => {
317 let expected = authoritative.record.generation.checked_add(1).ok_or(
318 CommitRecoveryError::GenerationOverflow {
319 generation: authoritative.record.generation,
320 },
321 )?;
322 let generation = requested.unwrap_or(expected);
323 if generation != expected {
324 return Err(CommitRecoveryError::UnexpectedGeneration {
325 expected,
326 actual: generation,
327 });
328 }
329 (authoritative.index.opposite(), generation)
330 }
331 Err(CommitRecoveryError::NoValidGeneration) if self.is_uninitialized() => {
332 (CommitSlotIndex::Slot0, requested.unwrap_or(0))
333 }
334 Err(err) => return Err(err),
335 };
336 let slot = match index {
337 CommitSlotIndex::Slot0 => &mut self.slot0,
338 CommitSlotIndex::Slot1 => &mut self.slot1,
339 };
340 Ok(slot.insert(CommittedGenerationBytes::new(generation, payload)))
343 }
344
345 #[cfg(test)]
350 pub fn write_corrupt_inactive_slot(&mut self, generation: u64, payload: Vec<u8>) {
351 let mut corrupt = CommittedGenerationBytes::new(generation, payload);
352 corrupt.checksum = corrupt.checksum.wrapping_add(1);
353
354 if self.inactive_slot_index() == CommitSlotIndex::Slot0 {
355 self.slot0 = Some(corrupt);
356 } else {
357 self.slot1 = Some(corrupt);
358 }
359 }
360}
361
362#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
369#[serde(deny_unknown_fields)]
370pub struct CommitStoreDiagnostic {
371 pub slot0: CommitSlotDiagnostic,
373 pub slot1: CommitSlotDiagnostic,
375 pub recovery: Result<u64, CommitRecoveryError>,
377}
378
379impl CommitStoreDiagnostic {
380 #[must_use]
382 pub fn from_store(store: &DualCommitStore) -> Self {
383 Self {
384 slot0: CommitSlotDiagnostic::from_slot(store.slot0()),
385 slot1: CommitSlotDiagnostic::from_slot(store.slot1()),
386 recovery: store
387 .authoritative_slot()
388 .map(|slot| slot.record.generation()),
389 }
390 }
391}
392
393#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
400#[serde(deny_unknown_fields)]
401pub enum CommitSlotDiagnostic {
402 Empty,
404 Valid {
406 generation: u64,
408 },
409 Invalid {
411 generation: u64,
413 },
414}
415
416impl CommitSlotDiagnostic {
417 fn from_slot(slot: Option<&CommittedGenerationBytes>) -> Self {
418 match slot {
419 Some(record) if record.validates() => Self::Valid {
420 generation: record.generation(),
421 },
422 Some(record) => Self::Invalid {
423 generation: record.generation(),
424 },
425 None => Self::Empty,
426 }
427 }
428}
429
430#[non_exhaustive]
437#[derive(Clone, Copy, Debug, Deserialize, Eq, thiserror::Error, PartialEq, Serialize)]
438pub enum CommitRecoveryError {
439 #[error("no committed ledger generation is present")]
441 NoValidGeneration,
442 #[error(
444 "present commit slot validation failed (slot0_invalid={slot0_invalid}, slot1_invalid={slot1_invalid})"
445 )]
446 InvalidCommitSlots {
447 slot0_invalid: bool,
449 slot1_invalid: bool,
451 },
452 #[error("ambiguous committed ledger generation {generation}")]
454 AmbiguousGeneration {
455 generation: u64,
457 },
458 #[error("committed ledger generation {generation} cannot be advanced without overflow")]
460 GenerationOverflow {
461 generation: u64,
463 },
464 #[error("expected committed ledger generation {expected}, got {actual}")]
466 UnexpectedGeneration {
467 expected: u64,
469 actual: u64,
471 },
472}
473
474fn generation_checksum(generation: &CommittedGenerationBytes) -> u64 {
475 let header = [
476 generation.generation.to_le_bytes(),
477 generation.commit_marker.to_le_bytes(),
478 (generation.payload.len() as u64).to_le_bytes(),
479 ];
480 let hash = fnv64(FNV_OFFSET, header.as_flattened());
481 fnv64(hash, &generation.payload)
482}
483
484#[cfg(test)]
485mod tests {
486 use super::*;
487
488 #[test]
489 fn automatic_and_explicit_commits_preserve_rejected_predecessors() {
490 let mut store = DualCommitStore::default();
491 assert_eq!(
493 store
494 .commit_payload_at_generation(7, payload(1))
495 .unwrap()
496 .generation(),
497 7
498 );
499 let before = store.clone();
500 assert_eq!(
501 store.commit_payload_at_generation(9, payload(2)),
502 Err(CommitRecoveryError::UnexpectedGeneration {
503 expected: 8,
504 actual: 9
505 })
506 );
507 assert_eq!(store, before);
508 assert_eq!(store.commit_payload(payload(2)).unwrap().generation(), 8);
509 assert_eq!(store.slot0().unwrap().generation(), 7);
510 assert_eq!(store.slot1().unwrap().generation(), 8);
511 assert_eq!(
512 store
513 .commit_payload_at_generation(9, payload(3))
514 .unwrap()
515 .generation(),
516 9
517 );
518 assert_eq!(store.slot0().unwrap().generation(), 9);
519 assert_eq!(store.slot1().unwrap().generation(), 8);
520
521 store.slot1.as_mut().unwrap().checksum ^= 1;
522 let corrupt = store.clone();
523 for result in [
524 store
525 .commit_payload(payload(4))
526 .map(CommittedGenerationBytes::generation),
527 store
528 .commit_payload_at_generation(10, payload(4))
529 .map(CommittedGenerationBytes::generation),
530 ] {
531 assert_eq!(
532 result,
533 Err(CommitRecoveryError::InvalidCommitSlots {
534 slot0_invalid: false,
535 slot1_invalid: true,
536 })
537 );
538 }
539 assert_eq!(store, corrupt);
540 }
541
542 #[test]
543 fn payload_uses_binary_bytes_and_human_readable_arrays() {
544 let record = CommittedGenerationBytes::new(7, vec![0, 24, 255]);
545 let bytes = crate::test_cbor::to_vec(&record).unwrap();
546 let value: ciborium::Value = crate::cbor::from_slice_exact(&bytes).unwrap();
547 let ciborium::Value::Map(mut fields) = value else {
548 panic!("record map")
549 };
550 let (_, payload) = fields
551 .iter_mut()
552 .find(|(key, _)| key.as_text() == Some("payload"))
553 .unwrap();
554 assert_eq!(*payload, ciborium::Value::Bytes(vec![0, 24, 255]));
555 *payload = ciborium::Value::Array(vec![0.into(), 24.into(), 255.into()]);
557 let removed = crate::test_cbor::to_vec(&ciborium::Value::Map(fields)).unwrap();
558 assert!(crate::cbor::from_slice_exact::<CommittedGenerationBytes>(&removed).is_err());
559 assert_eq!(
560 crate::cbor::from_slice_exact::<CommittedGenerationBytes>(&bytes).unwrap(),
561 record
562 );
563 let json = serde_json::to_value(&record).unwrap();
564 assert_eq!(json["payload"], serde_json::json!([0, 24, 255]));
565 assert_eq!(
566 serde_json::from_value::<CommittedGenerationBytes>(json).unwrap(),
567 record
568 );
569 }
570
571 #[test]
572 fn maximum_payload_writer_reader_agree() {
573 let record = CommittedGenerationBytes::new(
574 0,
575 vec![0; crate::constants::MAX_COMMITTED_PAYLOAD_BYTES],
576 );
577 let store = DualCommitStore {
578 slot0: Some(record.clone()),
579 slot1: Some(record),
580 };
581 let bytes = crate::test_cbor::to_vec(&store).unwrap();
582 assert!(bytes.len() <= crate::constants::MAX_LEDGER_RECORD_BYTES);
583 let decoded: DualCommitStore = crate::cbor::from_slice_exact(&bytes).unwrap();
584 assert_eq!(decoded, store);
585 assert!(decoded.authoritative().is_ok());
586 let oversized = CommittedGenerationBytes::new(
587 0,
588 vec![0; crate::constants::MAX_COMMITTED_PAYLOAD_BYTES + 1],
589 );
590 assert!(crate::test_cbor::to_vec(&oversized).is_err());
591 }
592
593 fn payload(value: u8) -> Vec<u8> {
594 vec![value; 4]
595 }
596
597 #[test]
598 fn committed_generation_validates_marker_and_checksum() {
599 let mut generation = CommittedGenerationBytes::new(7, payload(1));
600 assert!(generation.validates());
601
602 generation.checksum = generation.checksum.wrapping_add(1);
603 assert!(!generation.validates());
604 }
605
606 #[test]
607 fn physical_commit_accessors_expose_read_only_state() {
608 let mut store = DualCommitStore::default();
609 store.commit_payload(payload(1)).expect("first commit");
610
611 let slot = store.slot0().expect("first slot");
612
613 assert_eq!(slot.generation(), 0);
614 assert_eq!(slot.payload(), payload(1).as_slice());
615 assert_eq!(slot.commit_marker(), COMMIT_MARKER);
616 assert_eq!(slot.checksum(), generation_checksum(slot));
617 assert!(store.slot1().is_none());
618 }
619
620 #[test]
621 fn authoritative_selects_highest_valid_generation() {
622 let mut store = DualCommitStore::default();
623 store.commit_payload(payload(1)).expect("first commit");
624 store.commit_payload(payload(2)).expect("second commit");
625
626 let authoritative = store.authoritative().expect("authoritative");
627 let authoritative_slot =
628 select_authoritative_slot(store.slot0.as_ref(), store.slot1.as_ref())
629 .expect("authoritative slot");
630
631 assert_eq!(authoritative.generation, 1);
632 assert_eq!(authoritative.payload, payload(2));
633 assert_eq!(authoritative_slot.index, CommitSlotIndex::Slot1);
634 assert_eq!(authoritative_slot.record.payload, payload(2));
635 }
636
637 #[test]
638 fn corrupt_newer_slot_fails_closed() {
639 let mut store = DualCommitStore::default();
640 store.commit_payload(payload(1)).expect("first commit");
641 store.write_corrupt_inactive_slot(1, payload(2));
642
643 let err = store.authoritative().expect_err("corrupt slot");
644
645 assert_eq!(
646 err,
647 CommitRecoveryError::InvalidCommitSlots {
648 slot0_invalid: false,
649 slot1_invalid: true,
650 }
651 );
652 }
653
654 #[test]
655 fn two_invalid_commit_slots_fail_closed() {
656 let mut store = DualCommitStore::default();
657 store.write_corrupt_inactive_slot(0, payload(1));
658 store.write_corrupt_inactive_slot(1, payload(2));
659
660 let err = store.authoritative().expect_err("invalid slots");
661
662 assert_eq!(
663 err,
664 CommitRecoveryError::InvalidCommitSlots {
665 slot0_invalid: true,
666 slot1_invalid: true,
667 }
668 );
669 }
670
671 #[test]
672 fn same_generation_identical_slots_recover_deterministically() {
673 let committed = CommittedGenerationBytes::new(7, payload(1));
674 let store = DualCommitStore {
675 slot0: Some(committed.clone()),
676 slot1: Some(committed),
677 };
678
679 let authoritative = store.authoritative_slot().expect("authoritative");
680
681 assert_eq!(authoritative.index, CommitSlotIndex::Slot0);
682 assert_eq!(authoritative.record.generation, 7);
683 }
684
685 #[test]
686 fn same_generation_divergent_slots_fail_closed() {
687 let store = DualCommitStore {
688 slot0: Some(CommittedGenerationBytes::new(7, payload(1))),
689 slot1: Some(CommittedGenerationBytes::new(7, payload(2))),
690 };
691
692 let err = store.authoritative().expect_err("ambiguous generation");
693
694 assert_eq!(
695 err,
696 CommitRecoveryError::AmbiguousGeneration { generation: 7 }
697 );
698 }
699
700 #[test]
701 fn physical_generation_overflow_fails_closed() {
702 let mut store = DualCommitStore {
703 slot0: Some(CommittedGenerationBytes::new(u64::MAX, payload(1))),
704 slot1: None,
705 };
706
707 let err = store
708 .commit_payload(payload(2))
709 .expect_err("overflow must fail");
710
711 assert_eq!(
712 err,
713 CommitRecoveryError::GenerationOverflow {
714 generation: u64::MAX
715 }
716 );
717 }
718
719 #[test]
720 fn diagnostic_reports_corrupt_slots_without_an_authoritative_generation() {
721 let mut store = DualCommitStore::default();
722 store.commit_payload(payload(1)).expect("first commit");
723 store.write_corrupt_inactive_slot(1, payload(2));
724
725 let diagnostic = store.diagnostic();
726
727 assert_eq!(
728 diagnostic.recovery,
729 Err(CommitRecoveryError::InvalidCommitSlots {
730 slot0_invalid: false,
731 slot1_invalid: true,
732 })
733 );
734 assert_eq!(
735 diagnostic.slot0,
736 CommitSlotDiagnostic::Valid { generation: 0 }
737 );
738 assert_eq!(
739 diagnostic.slot1,
740 CommitSlotDiagnostic::Invalid { generation: 1 }
741 );
742 let bytes = crate::test_cbor::to_vec(&diagnostic).expect("diagnostic bytes");
743 let decoded: CommitStoreDiagnostic =
744 crate::test_cbor::from_slice(&bytes).expect("diagnostic round trip");
745 assert_eq!(decoded, diagnostic);
746 }
747
748 #[test]
749 fn diagnostic_reports_no_valid_generation_for_empty_store() {
750 let diagnostic = DualCommitStore::default().diagnostic();
751
752 assert_eq!(
753 diagnostic.recovery,
754 Err(CommitRecoveryError::NoValidGeneration)
755 );
756 assert_eq!(diagnostic.slot0, CommitSlotDiagnostic::Empty);
757 assert_eq!(diagnostic.slot1, CommitSlotDiagnostic::Empty);
758 }
759
760 #[test]
761 fn uninitialized_distinguishes_empty_from_corrupt() {
762 let mut store = DualCommitStore::default();
763 assert!(store.is_uninitialized());
764
765 store.write_corrupt_inactive_slot(0, payload(1));
766
767 assert!(!store.is_uninitialized());
768 }
769
770 #[test]
771 fn commit_after_corrupt_slot_fails_closed() {
772 let mut store = DualCommitStore::default();
773 store.commit_payload(payload(1)).expect("first commit");
774 store.write_corrupt_inactive_slot(1, payload(2));
775
776 let err = store
777 .commit_payload(payload(3))
778 .expect_err("corrupt history must not be overwritten");
779
780 assert_eq!(
781 err,
782 CommitRecoveryError::InvalidCommitSlots {
783 slot0_invalid: false,
784 slot1_invalid: true,
785 }
786 );
787 }
788}