1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
use Future;
use Serialize;
use DeserializeOwned;
use crateResult;
use crateStepErrorKind;
use crateJobContext;
use cratehex_sha256;
/// A unit of durable background work.
///
/// Implement this trait to define a job type: the struct fields are the
/// typed input, [`Output`](Self::Output) is the typed result, and
/// [`run`](Self::run) is the work. The remaining methods have defaults that
/// work without configuration; override them to customize idempotency, retry
/// limits, and error classification.
///
/// A `Job` must round-trip through [`serde`]: the runner serializes the
/// input to enqueue it and the output to persist it. It must also be
/// `Send + Sync + 'static` so the runner can dispatch it across worker tasks.
///
/// # Example
///
/// ```
/// use serde::{Serialize, Deserialize};
/// use taquba_workflow::StepErrorKind;
/// use taquba_workflow::jobs::{Job, JobContext};
///
/// #[derive(Serialize, Deserialize)]
/// struct ResizeImage {
/// bucket: String,
/// key: String,
/// }
///
/// #[derive(Debug, thiserror::Error)]
/// #[error("resize failed: {0}")]
/// struct ResizeError(String);
///
/// impl Job for ResizeImage {
/// const NAME: &'static str = "media.resize-image";
/// type Output = u64; // bytes written
/// type Error = ResizeError;
///
/// async fn run(&self, _ctx: JobContext<'_>) -> Result<u64, ResizeError> {
/// // ... do the work ...
/// Ok(4096)
/// }
///
/// fn idempotency_key(&self) -> Option<String> {
/// Some(format!("resize:{}:{}", self.bucket, self.key))
/// }
///
/// fn classify(&self, _err: &ResizeError) -> StepErrorKind {
/// StepErrorKind::Transient
/// }
/// }
/// ```
/// Derive a stable idempotency key by hashing a job's serialized form.
///
/// A convenience for jobs that want hash-based deduplication without
/// hand-writing a key: return this from
/// [`Job::idempotency_key`]. Two submissions of an identical job value
/// collapse onto a single execution.
///
/// This is opt-in by design: collapsing identical submissions silently
/// discards intentional duplicate work, so it is never the default.
///
/// ```
/// # use serde::{Serialize, Deserialize};
/// # use taquba_workflow::jobs::{Job, JobContext, payload_idempotency_key};
/// # #[derive(Serialize, Deserialize)]
/// # struct SendDigest { user_id: u64 }
/// # #[derive(Debug, thiserror::Error)]
/// # #[error("err")]
/// # struct E;
/// impl Job for SendDigest {
/// const NAME: &'static str = "email.send-digest";
/// type Output = ();
/// type Error = E;
/// async fn run(&self, _ctx: JobContext<'_>) -> Result<(), E> { Ok(()) }
/// fn idempotency_key(&self) -> Option<String> {
/// payload_idempotency_key(self).ok()
/// }
/// }
/// ```