1use thiserror::Error;
2
3use crate::keys::RunId;
4
5#[derive(Debug, Error)]
7pub enum Error {
8 #[error("step job is missing header `{0}`")]
11 MissingHeader(&'static str),
12
13 #[error("step job has invalid `{header}` header `{value}`")]
16 InvalidStepHeader {
17 header: &'static str,
19 value: String,
21 },
22
23 #[error("submission header `{0}` uses the reserved `workflow.*` prefix")]
27 ReservedHeaderInSubmit(String),
28
29 #[error("invalid run id `{run_id}`: {reason}")]
33 InvalidRunId {
34 run_id: String,
36 reason: &'static str,
38 },
39
40 #[error("run `{0}` exists with a different input; pick a fresh run_id")]
46 InputMismatch(RunId),
47
48 #[error("run `{0}` has inconsistent durable state")]
52 InconsistentRunState(RunId),
53
54 #[error("kv key `{0}` uses the reserved `workflow/` prefix")]
58 ReservedKvKey(String),
59
60 #[error("kv key `{0}` is staged for both a write and a delete")]
65 ConflictingKvEffect(String),
66
67 #[error("the effects handle is sealed; its delivery has returned")]
72 EffectsSealed,
73
74 #[error(transparent)]
76 Queue(#[from] taquba::Error),
77
78 #[error("object store error: {0}")]
80 Store(#[from] taquba::object_store::Error),
81
82 #[error("serialization error: {0}")]
84 Serialization(#[from] rmp_serde::encode::Error),
85
86 #[error("deserialization error: {0}")]
89 Deserialization(#[from] rmp_serde::decode::Error),
90
91 #[error("run `{0}` not found")]
94 RunNotFound(RunId),
95
96 #[error("member `{key}` of group `{group_id}` was not submitted")]
100 MemberNotSubmitted {
101 group_id: RunId,
103 key: String,
105 },
106
107 #[error("duplicate member key `{0}` in group")]
109 DuplicateMemberKey(String),
110
111 #[error("group `{0}` exists with a different member set")]
114 GroupMismatch(RunId),
115
116 #[error("group `{0}` not found")]
118 GroupNotFound(RunId),
119}
120
121impl Error {
122 pub fn is_permanent(&self) -> bool {
128 match self {
129 Self::MissingHeader(_)
130 | Self::InvalidStepHeader { .. }
131 | Self::ReservedHeaderInSubmit(_)
132 | Self::InvalidRunId { .. }
133 | Self::InputMismatch(_)
134 | Self::InconsistentRunState(_)
135 | Self::ReservedKvKey(_)
136 | Self::ConflictingKvEffect(_)
137 | Self::EffectsSealed
138 | Self::Serialization(_)
139 | Self::Deserialization(_)
140 | Self::RunNotFound(_)
141 | Self::MemberNotSubmitted { .. }
142 | Self::DuplicateMemberKey(_)
143 | Self::GroupMismatch(_)
144 | Self::GroupNotFound(_) => true,
145 Self::Queue(e) => e.is_permanent(),
146 Self::Store(_) => false,
147 }
148 }
149}
150
151pub(crate) fn worker_error(err: impl Into<Error>) -> taquba::WorkerError {
155 crate::runner::StepError::from(err.into()).into_worker_error()
156}
157
158pub type Result<T> = std::result::Result<T, Error>;
160
161#[cfg(test)]
162mod tests {
163 use super::*;
164 use crate::test_util::rid;
165
166 struct BadSerialize;
167
168 impl serde::Serialize for BadSerialize {
169 fn serialize<S>(&self, _serializer: S) -> std::result::Result<S::Ok, S::Error>
170 where
171 S: serde::Serializer,
172 {
173 Err(serde::ser::Error::custom("serialization failed"))
174 }
175 }
176
177 #[test]
178 fn is_permanent_classifies_every_variant() {
179 let store_err = taquba::object_store::Error::NotFound {
180 path: "x".into(),
181 source: "missing".into(),
182 };
183 for (error, permanent) in [
184 (Error::MissingHeader("workflow.run_id"), true),
185 (
186 Error::InvalidStepHeader {
187 header: "workflow.step",
188 value: "not-a-u32".into(),
189 },
190 true,
191 ),
192 (Error::ReservedHeaderInSubmit("workflow.foo".into()), true),
193 (
194 Error::InvalidRunId {
195 run_id: String::new(),
196 reason: "run id must not be empty",
197 },
198 true,
199 ),
200 (Error::InputMismatch(rid("run-1")), true),
201 (Error::InconsistentRunState(rid("run-1")), true),
202 (Error::ReservedKvKey("workflow/x".into()), true),
203 (Error::ConflictingKvEffect("k".into()), true),
204 (Error::EffectsSealed, true),
205 (
206 Error::Queue(taquba::Error::JobNotFound("job-1".into())),
207 true,
208 ),
209 (Error::Queue(taquba::Error::InvalidState), true),
210 (
211 Error::Queue(taquba::Error::KvValueTooLarge { size: 10, max: 5 }),
212 true,
213 ),
214 (
215 Error::Queue(taquba::Error::StoreNotInitialized { path: "x".into() }),
216 false,
217 ),
218 (Error::Store(store_err), false),
219 (
220 Error::Serialization(rmp_serde::to_vec_named(&BadSerialize).unwrap_err()),
221 true,
222 ),
223 (
224 Error::Deserialization(rmp_serde::from_slice::<u32>(b"").unwrap_err()),
225 true,
226 ),
227 (Error::DuplicateMemberKey("k".into()), true),
228 (Error::GroupMismatch(rid("b")), true),
229 (Error::GroupNotFound(rid("b")), true),
230 (Error::RunNotFound(rid("run-1")), true),
231 (
232 Error::MemberNotSubmitted {
233 group_id: rid("b"),
234 key: "k".into(),
235 },
236 true,
237 ),
238 ] {
239 assert_eq!(error.is_permanent(), permanent, "{error}");
240 }
241 }
242}