laminar_db/subscription/
error.rs1use laminar_core::checkpoint::{OutputPartitionId, PartitionSequence};
4
5#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
7pub enum ClusterSubscriptionError {
8 #[error("cluster subscription is unsupported for this plan: {reason}")]
10 UnsupportedPlan {
11 reason: String,
13 },
14 #[error("subscription stream generation does not match the current catalog object")]
16 GenerationMismatch,
17 #[error("subscription epoch {requested} is not committed")]
19 EpochNotCommitted {
20 requested: u64,
22 },
23 #[error("subscription epoch {requested} is no longer retained")]
25 ReplayPruned {
26 requested: u64,
28 },
29 #[error("committed subscription manifest is corrupt: {reason}")]
31 ManifestCorrupt {
32 reason: String,
34 },
35 #[error("committed subscription segment is missing for partition {partition:?} at {first:?}")]
37 SegmentMissing {
38 partition: OutputPartitionId,
40 first: PartitionSequence,
42 },
43 #[error("committed subscription segment is corrupt for partition {partition:?} at {first:?}")]
45 SegmentCorrupt {
46 partition: OutputPartitionId,
48 first: PartitionSequence,
50 },
51 #[error("subscription output schema does not match its distribution certificate")]
53 SchemaMismatch,
54 #[error(
56 "subscription partition {partition:?} expected sequence {expected:?}, found {actual:?}"
57 )]
58 PartitionSequenceGap {
59 partition: OutputPartitionId,
61 expected: PartitionSequence,
63 actual: PartitionSequence,
65 },
66 #[error("subscription frame identity has conflicting immutable content")]
68 ConflictingDuplicateSequence,
69 #[error("subscription output writer is stale")]
71 StaleOutputWriter,
72 #[error("subscription output assignment changed")]
74 AssignmentChanged,
75 #[error("committed subscription backend is unavailable")]
77 BackendUnavailable,
78 #[error("subscription reader exceeded its bounded lag allowance")]
80 SubscriberLagged,
81 #[error("subscription resume token is invalid")]
83 ResumeTokenInvalid,
84 #[error("subscription resume token has expired")]
86 ResumeTokenExpired,
87 #[error("subscription retention no longer covers the required committed range")]
89 RetentionLost,
90 #[error("subscription protocol version {actual} is unsupported")]
92 ProtocolVersion {
93 actual: u16,
95 },
96}
97
98impl ClusterSubscriptionError {
99 #[must_use]
101 pub const fn code(&self) -> &'static str {
102 use laminar_core::error_codes as codes;
103
104 match self {
105 Self::UnsupportedPlan { .. } => codes::SUBSCRIPTION_PLAN_UNSUPPORTED,
106 Self::GenerationMismatch => codes::SUBSCRIPTION_GENERATION_MISMATCH,
107 Self::EpochNotCommitted { .. } => codes::SUBSCRIPTION_EPOCH_NOT_COMMITTED,
108 Self::ReplayPruned { .. } => codes::SUBSCRIPTION_REPLAY_PRUNED,
109 Self::ManifestCorrupt { .. } => codes::SUBSCRIPTION_MANIFEST_CORRUPT,
110 Self::SegmentMissing { .. } => codes::SUBSCRIPTION_SEGMENT_MISSING,
111 Self::SegmentCorrupt { .. } => codes::SUBSCRIPTION_SEGMENT_CORRUPT,
112 Self::SchemaMismatch => codes::SUBSCRIPTION_SCHEMA_MISMATCH,
113 Self::PartitionSequenceGap { .. } => codes::SUBSCRIPTION_SEQUENCE_GAP,
114 Self::ConflictingDuplicateSequence => codes::SUBSCRIPTION_CONFLICTING_DUPLICATE,
115 Self::StaleOutputWriter => codes::SUBSCRIPTION_STALE_WRITER,
116 Self::AssignmentChanged => codes::SUBSCRIPTION_ASSIGNMENT_CHANGED,
117 Self::BackendUnavailable => codes::SUBSCRIPTION_BACKEND_UNAVAILABLE,
118 Self::SubscriberLagged => codes::SUBSCRIPTION_LAGGED,
119 Self::ResumeTokenInvalid => codes::SUBSCRIPTION_RESUME_TOKEN_INVALID,
120 Self::ResumeTokenExpired => codes::SUBSCRIPTION_RESUME_TOKEN_EXPIRED,
121 Self::RetentionLost => codes::SUBSCRIPTION_RETENTION_LOST,
122 Self::ProtocolVersion { .. } => codes::SUBSCRIPTION_PROTOCOL_UNSUPPORTED,
123 }
124 }
125}
126
127#[cfg(test)]
128mod tests {
129 use super::*;
130
131 #[test]
132 fn correctness_failures_have_distinct_stable_codes() {
133 let errors = [
134 ClusterSubscriptionError::GenerationMismatch,
135 ClusterSubscriptionError::ManifestCorrupt {
136 reason: "digest".into(),
137 },
138 ClusterSubscriptionError::ConflictingDuplicateSequence,
139 ClusterSubscriptionError::StaleOutputWriter,
140 ClusterSubscriptionError::SubscriberLagged,
141 ClusterSubscriptionError::ResumeTokenInvalid,
142 ];
143 let codes = errors.map(|error| error.code());
144 assert!(codes.iter().all(|code| code.starts_with("LDB-")));
145 for (index, code) in codes.iter().enumerate() {
146 assert!(!codes[..index].contains(code));
147 }
148 }
149}