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::effects`] or staged through
55    /// an [`crate::EffectsHandle`] starts with the reserved `workflow/` prefix.
56    /// The runtime owns that prefix. A caller must use another key.
57    #[error("kv key `{0}` uses the reserved `workflow/` prefix")]
58    ReservedKvKey(String),
59
60    /// A key was staged through an [`crate::EffectsHandle`] for both a
61    /// write and a delete within one step. The combination has no defined
62    /// order in the settlement transaction and is rejected when the
63    /// second operation is staged.
64    #[error("kv key `{0}` is staged for both a write and a delete")]
65    ConflictingKvEffect(String),
66
67    /// An effect was staged through an [`crate::EffectsHandle`] or
68    /// [`crate::TerminalEffects`] clone after its delivery returned.
69    /// Effects are collected when the runner or hook returns; an effect
70    /// staged after that point cannot join the settlement.
71    #[error("the effects handle is sealed; its delivery has returned")]
72    EffectsSealed,
73
74    /// Underlying error from a Taquba queue operation.
75    #[error(transparent)]
76    Queue(#[from] taquba::Error),
77
78    /// Reading or writing a blob in object storage failed.
79    #[error("object store error: {0}")]
80    Store(#[from] taquba::object_store::Error),
81
82    /// Serializing a value for workflow storage failed.
83    #[error("serialization error: {0}")]
84    Serialization(#[from] rmp_serde::encode::Error),
85
86    /// Deserializing a stored value, a typed input or a typed output
87    /// failed.
88    #[error("deserialization error: {0}")]
89    Deserialization(#[from] rmp_serde::decode::Error),
90
91    /// A wait named a run the runtime has no record of: never
92    /// submitted, or terminated with no record retained.
93    #[error("run `{0}` not found")]
94    RunNotFound(RunId),
95
96    /// A group operation waited on a member of the manifest that was
97    /// not submitted; [`RunGroup::resume`](crate::RunGroup::resume)
98    /// submits it.
99    #[error("member `{key}` of group `{group_id}` was not submitted")]
100    MemberNotSubmitted {
101        /// The group id.
102        group_id: RunId,
103        /// The member's key.
104        key: String,
105    },
106
107    /// Two members of one group have the same key.
108    #[error("duplicate member key `{0}` in group")]
109    DuplicateMemberKey(String),
110
111    /// A submission to an existing group supplied a different member
112    /// set than the group's manifest.
113    #[error("group `{0}` exists with a different member set")]
114    GroupMismatch(RunId),
115
116    /// A group operation named a group with no manifest.
117    #[error("group `{0}` not found")]
118    GroupNotFound(RunId),
119}
120
121impl Error {
122    /// True if retrying the operation will not change the outcome; callers
123    /// should fast-fail (e.g. dead-letter a step, mark a submission as
124    /// failed) rather than back off and try again.
125    ///
126    /// [`Self::Queue`] delegates to [`taquba::Error::is_permanent`].
127    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
151/// The worker error reporting `err` from a step's delivery: a
152/// [`taquba::PermanentFailure`] for a permanent error, which
153/// dead-letters the step, and a retrying error otherwise.
154pub(crate) fn worker_error(err: impl Into<Error>) -> taquba::WorkerError {
155    crate::runner::StepError::from(err.into()).into_worker_error()
156}
157
158/// Result alias used throughout the crate.
159pub 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}