Skip to main content

taquba_workflow/
error.rs

1use thiserror::Error;
2
3use crate::keys::RunId;
4
5/// Errors returned by the runtime's submission and worker paths.
6#[derive(Debug, Error)]
7pub enum Error {
8    /// A step job is missing the [`crate::HEADER_RUN_ID`] header.
9    /// Permanent: a misconfigured job will not become valid on retry.
10    #[error("step job is missing header `{0}`")]
11    MissingHeader(&'static str),
12
13    /// A step job's [`crate::HEADER_STEP`] header is not a valid `u32`.
14    /// Permanent: header value won't change across retries.
15    #[error("step job has invalid `{header}` header `{value}`")]
16    InvalidStepHeader {
17        /// Header name.
18        header: &'static str,
19        /// Offending value.
20        value: String,
21    },
22
23    /// A submission included a user header starting with the reserved
24    /// `workflow.*` prefix. The runtime owns that prefix; submitters must use
25    /// any other key.
26    #[error("submission header `{0}` uses the reserved `workflow.*` prefix")]
27    ReservedHeaderInSubmit(String),
28
29    /// A run id or group id is empty, longer than
30    /// [`crate::MAX_RUN_ID_LEN`] bytes or contains a character outside
31    /// `[A-Za-z0-9_-]`. See [`crate::RunId`].
32    #[error("invalid run id `{run_id}`: {reason}")]
33    InvalidRunId {
34        /// The rejected run id.
35        run_id: String,
36        /// Which rule the run id broke.
37        reason: &'static str,
38    },
39
40    /// A re-submission of `run_id` carried `spec.input` bytes that differ
41    /// from the original submission's: the run is active, or it is a
42    /// typed job whose run result record is retained. Reusing a `run_id`
43    /// with new input is treated as a programmer error: pick a fresh
44    /// `run_id` for a new run.
45    #[error("run `{0}` exists with a different input; pick a fresh run_id")]
46    InputMismatch(RunId),
47
48    /// The durable records of a run disagree with one another. The runtime
49    /// writes and deletes them together, so the error reports a store the
50    /// runtime did not write.
51    #[error("run `{0}` has inconsistent durable state")]
52    InconsistentRunState(RunId),
53
54    /// A caller KV key passed via [`crate::RunSpec::kv_writes`] or staged
55    /// through an [`crate::EffectsHandle`] starts with the reserved
56    /// `workflow/` prefix. The runtime owns that prefix; callers must use
57    /// any other key.
58    #[error("kv key `{0}` uses the reserved `workflow/` prefix")]
59    ReservedKvKey(String),
60
61    /// A key was staged through an [`crate::EffectsHandle`] for both a
62    /// write and a delete within one step. The combination has no defined
63    /// order in the settlement transaction and is rejected when the
64    /// second operation is staged.
65    #[error("kv key `{0}` is staged for both a write and a delete")]
66    ConflictingKvEffect(String),
67
68    /// An effect was staged through an [`crate::EffectsHandle`] or
69    /// [`crate::TerminalEffects`] clone after its delivery returned.
70    /// Effects are collected when the runner or hook returns; an effect
71    /// staged after that point cannot join the settlement.
72    #[error("the effects handle is sealed; its delivery has returned")]
73    EffectsSealed,
74
75    /// Underlying error from a Taquba queue operation.
76    #[error(transparent)]
77    Queue(#[from] taquba::Error),
78
79    /// Reading or writing a blob in object storage failed.
80    #[error("object store error: {0}")]
81    Store(#[from] taquba::object_store::Error),
82
83    /// Serializing a value for workflow storage failed.
84    #[error("serialization error: {0}")]
85    Serialization(#[from] rmp_serde::encode::Error),
86
87    /// Deserializing a stored value, a typed input or a typed output
88    /// failed.
89    #[error("deserialization error: {0}")]
90    Deserialization(#[from] rmp_serde::decode::Error),
91
92    /// A wait named a run the runtime has no record of: never
93    /// submitted, or terminated with no record retained.
94    #[error("run `{0}` not found")]
95    RunNotFound(RunId),
96
97    /// A group operation waited on a member of the manifest that was
98    /// not submitted; [`RunGroup::resume`](crate::RunGroup::resume)
99    /// submits it.
100    #[error("member `{key}` of group `{group_id}` was not submitted")]
101    MemberNotSubmitted {
102        /// The group id.
103        group_id: RunId,
104        /// The member's key.
105        key: String,
106    },
107
108    /// Two members of one group have the same key.
109    #[error("duplicate member key `{0}` in group")]
110    DuplicateMemberKey(String),
111
112    /// A submission to an existing group supplied a different member
113    /// set than the group's manifest.
114    #[error("group `{0}` exists with a different member set")]
115    GroupMismatch(RunId),
116
117    /// A group operation named a group with no manifest.
118    #[error("group `{0}` not found")]
119    GroupNotFound(RunId),
120}
121
122impl Error {
123    /// True if retrying the operation will not change the outcome; callers
124    /// should fast-fail (e.g. dead-letter a step, mark a submission as
125    /// failed) rather than back off and try again.
126    ///
127    /// [`Self::Queue`] delegates to [`taquba::Error::is_permanent`].
128    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
152/// The worker error reporting `err` from a step's delivery: a
153/// [`taquba::PermanentFailure`] for a permanent error, which
154/// dead-letters the step, and a retrying error otherwise.
155pub(crate) fn worker_error(err: impl Into<Error>) -> taquba::WorkerError {
156    crate::runner::StepError::from(err.into()).into_worker_error()
157}
158
159/// Result alias used throughout the crate.
160pub 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}