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")]
59 ReservedKvKey(String),
60
61 #[error("kv key `{0}` is staged for both a write and a delete")]
66 ConflictingKvEffect(String),
67
68 #[error("the effects handle is sealed; its delivery has returned")]
73 EffectsSealed,
74
75 #[error(transparent)]
77 Queue(#[from] taquba::Error),
78
79 #[error("object store error: {0}")]
81 Store(#[from] taquba::object_store::Error),
82
83 #[error("serialization error: {0}")]
85 Serialization(#[from] rmp_serde::encode::Error),
86
87 #[error("deserialization error: {0}")]
90 Deserialization(#[from] rmp_serde::decode::Error),
91
92 #[error("run `{0}` not found")]
95 RunNotFound(RunId),
96
97 #[error("member `{key}` of group `{group_id}` was not submitted")]
101 MemberNotSubmitted {
102 group_id: RunId,
104 key: String,
106 },
107
108 #[error("duplicate member key `{0}` in group")]
110 DuplicateMemberKey(String),
111
112 #[error("group `{0}` exists with a different member set")]
115 GroupMismatch(RunId),
116
117 #[error("group `{0}` not found")]
119 GroupNotFound(RunId),
120}
121
122impl Error {
123 pub fn is_permanent(&self) -> bool {
129 match self {
130 Self::MissingHeader(_)
131 | Self::InvalidStepHeader { .. }
132 | Self::ReservedHeaderInSubmit(_)
133 | Self::InvalidRunId { .. }
134 | Self::InputMismatch(_)
135 | Self::InconsistentRunState(_)
136 | Self::ReservedKvKey(_)
137 | Self::ConflictingKvEffect(_)
138 | Self::EffectsSealed
139 | Self::Serialization(_)
140 | Self::Deserialization(_)
141 | Self::RunNotFound(_)
142 | Self::MemberNotSubmitted { .. }
143 | Self::DuplicateMemberKey(_)
144 | Self::GroupMismatch(_)
145 | Self::GroupNotFound(_) => true,
146 Self::Queue(e) => e.is_permanent(),
147 Self::Store(_) => false,
148 }
149 }
150}
151
152pub(crate) fn worker_error(err: impl Into<Error>) -> taquba::WorkerError {
156 crate::runner::StepError::from(err.into()).into_worker_error()
157}
158
159pub type Result<T> = std::result::Result<T, Error>;
161
162#[cfg(test)]
163mod tests {
164 use super::*;
165 use crate::test_util::rid;
166
167 struct BadSerialize;
168
169 impl serde::Serialize for BadSerialize {
170 fn serialize<S>(&self, _serializer: S) -> std::result::Result<S::Ok, S::Error>
171 where
172 S: serde::Serializer,
173 {
174 Err(serde::ser::Error::custom("serialization failed"))
175 }
176 }
177
178 #[test]
179 fn is_permanent_classifies_every_variant() {
180 let store_err = taquba::object_store::Error::NotFound {
181 path: "x".into(),
182 source: "missing".into(),
183 };
184 for (error, permanent) in [
185 (Error::MissingHeader("workflow.run_id"), true),
186 (
187 Error::InvalidStepHeader {
188 header: "workflow.step",
189 value: "not-a-u32".into(),
190 },
191 true,
192 ),
193 (Error::ReservedHeaderInSubmit("workflow.foo".into()), true),
194 (
195 Error::InvalidRunId {
196 run_id: String::new(),
197 reason: "run id must not be empty",
198 },
199 true,
200 ),
201 (Error::InputMismatch(rid("run-1")), true),
202 (Error::InconsistentRunState(rid("run-1")), true),
203 (Error::ReservedKvKey("workflow/x".into()), true),
204 (Error::ConflictingKvEffect("k".into()), true),
205 (Error::EffectsSealed, true),
206 (
207 Error::Queue(taquba::Error::JobNotFound("job-1".into())),
208 true,
209 ),
210 (Error::Queue(taquba::Error::InvalidState), true),
211 (
212 Error::Queue(taquba::Error::KvValueTooLarge { size: 10, max: 5 }),
213 true,
214 ),
215 (
216 Error::Queue(taquba::Error::StoreNotInitialized { path: "x".into() }),
217 false,
218 ),
219 (Error::Store(store_err), false),
220 (
221 Error::Serialization(rmp_serde::to_vec_named(&BadSerialize).unwrap_err()),
222 true,
223 ),
224 (
225 Error::Deserialization(rmp_serde::from_slice::<u32>(b"").unwrap_err()),
226 true,
227 ),
228 (Error::DuplicateMemberKey("k".into()), true),
229 (Error::GroupMismatch(rid("b")), true),
230 (Error::GroupNotFound(rid("b")), true),
231 (Error::RunNotFound(rid("run-1")), true),
232 (
233 Error::MemberNotSubmitted {
234 group_id: rid("b"),
235 key: "k".into(),
236 },
237 true,
238 ),
239 ] {
240 assert_eq!(error.is_permanent(), permanent, "{error}");
241 }
242 }
243}