Skip to main content

laminar_db/checkpoint_coordinator/
protocol.rs

1use std::time::Duration;
2
3#[cfg(feature = "cluster")]
4use laminar_core::checkpoint::{
5    CheckpointAssignmentFence, CheckpointAttempt, CheckpointWatermark, LeaderProof,
6};
7
8#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
9pub enum CheckpointPhase {
10    Idle,
11    PreCommitting,
12    Persisting,
13    Deciding,
14}
15
16impl std::fmt::Display for CheckpointPhase {
17    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
18        match self {
19            Self::Idle => formatter.write_str("Idle"),
20            Self::PreCommitting => formatter.write_str("PreCommitting"),
21            Self::Persisting => formatter.write_str("Persisting"),
22            Self::Deciding => formatter.write_str("Deciding"),
23        }
24    }
25}
26
27#[derive(Debug, Clone, Copy, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
28#[serde(rename_all = "snake_case")]
29pub enum CheckpointFailureDisposition {
30    Retryable,
31    RequiresRecovery,
32}
33
34/// Determines who publishes a prepared successor sink epoch as writable.
35#[derive(Debug, Clone, Copy, PartialEq, Eq)]
36pub(crate) enum SinkEpochPublication {
37    /// Startup and direct coordinator APIs have no callback-owned transition guard.
38    Immediate,
39    /// A spawned pipeline tail publishes only after its terminal result is known successful.
40    DeferredToTail,
41}
42
43#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
44pub struct CheckpointResult {
45    pub success: bool,
46    pub checkpoint_id: u64,
47    pub epoch: u64,
48    pub duration: Duration,
49    pub error: Option<String>,
50    pub failure_disposition: Option<CheckpointFailureDisposition>,
51}
52
53impl CheckpointResult {
54    #[must_use]
55    pub fn continuation_error(&self) -> Option<&str> {
56        self.success.then_some(self.error.as_deref()).flatten()
57    }
58
59    #[must_use]
60    pub fn requires_recovery(&self) -> bool {
61        !self.success
62            && self.failure_disposition == Some(CheckpointFailureDisposition::RequiresRecovery)
63    }
64}
65
66#[cfg(feature = "cluster")]
67pub(crate) type QuorumPeer = laminar_core::cluster::discovery::NodeId;
68
69#[derive(Debug, Clone)]
70pub(crate) enum QuorumStage {
71    RunInline,
72    #[cfg(feature = "cluster")]
73    Captured {
74        cluster_watermark: CheckpointWatermark,
75        participants: Vec<QuorumPeer>,
76        leader_proof: LeaderProof,
77    },
78}
79
80/// Follower-local durability state after immutable capture ownership has been acknowledged.
81#[cfg(feature = "cluster")]
82#[derive(Debug, Clone, Copy, PartialEq, Eq)]
83pub(crate) enum FollowerPrepareOutcome {
84    /// Manifest persistence completed with an acknowledgement.
85    Prepared,
86    /// Manifest Create may be visible even though its acknowledgement was lost. The captured
87    /// phase-one state must remain intact until an authoritative Commit or Abort is observed.
88    InDoubt,
89}
90
91#[cfg(feature = "cluster")]
92pub(crate) struct PrepareQuorum<'a> {
93    pub(super) attempt: CheckpointAttempt,
94    pub(super) local_watermark: CheckpointWatermark,
95    pub(super) assignment_fence: &'a CheckpointAssignmentFence,
96    pub(super) leader_proof: &'a LeaderProof,
97    pub(super) flags: u64,
98}
99
100#[cfg(feature = "cluster")]
101impl<'a> PrepareQuorum<'a> {
102    pub(crate) const fn new(
103        attempt: CheckpointAttempt,
104        local_watermark: CheckpointWatermark,
105        assignment_fence: &'a CheckpointAssignmentFence,
106        leader_proof: &'a LeaderProof,
107        flags: u64,
108    ) -> Self {
109        Self {
110            attempt,
111            local_watermark,
112            assignment_fence,
113            leader_proof,
114            flags,
115        }
116    }
117}