acme_proxy/jobs/mod.rs
1//! The durable job runner: background work that survives the process.
2//!
3//! ## Why this is not a `tokio::spawn`
4//!
5//! A spawned task is a promise the process makes to itself, and a restart breaks
6//! it. That was tolerable while exactly one subsystem needed one — the `relay`
7//! signer backend, whose state lived on its own `upstream_orders` row and whose
8//! `resume` sweep re-created the task at startup — but it bought only *recovery*,
9//! never *retry*: with nowhere to record that an attempt had failed and should
10//! happen again, every failure had to be terminal. A five-second upstream blip
11//! therefore invalidated an order the client then had to place again.
12//!
13//! Here the schedule is a row. An attempt that fails writes its next `run_at`
14//! and goes back in the queue; a process that dies leaves a lease that expires
15//! and a row another runner takes. Both are the same table and the same loop, so
16//! the next subsystem that wants a queue — expiry reminders, a CRL sweep, an
17//! OCSP refresh — registers a handler rather than writing this again.
18//!
19//! ## The three things a handler must know
20//!
21//! 1. **`run` may be called again from scratch.** The runner promises nothing
22//! about what a previous attempt did. `signer::relay::flow` is the worked
23//! example: it re-reads the upstream order and skips every authorization that
24//! is no longer `pending`, so re-entering at any point is safe by
25//! construction rather than by checkpointing.
26//! 2. **[`JobOutcome::Retry`] and [`JobOutcome::Failed`] are a real
27//! distinction.** `Retry` is "the network, the nameserver, a 503"; `Failed` is
28//! "the CA stated a reason" or "this payload can never parse". Classifying a
29//! permanent failure as retryable wastes the budget and delays the client's
30//! real answer; the reverse throws away the retry this module exists for.
31//! 3. **[`JobHandler::abandon`] fires exactly once**, when the job is retired for
32//! good — and never between retries. It is where a handler tells its *subject*
33//! what happened, so a client polling through a transient blip keeps seeing
34//! work in progress rather than a terminal failure it must start over from.
35//!
36//! ## Wake-up is in-process
37//!
38//! [`JobQueue::enqueue`] notifies the runner directly, so a job queued by a
39//! request starts within microseconds rather than waiting for
40//! `jobs.poll_interval_ms`. That matters because the ACME client that triggered
41//! it is already polling its order. The notification is a
42//! [`tokio::sync::Notify`], which is process-local: with two processes over one
43//! database, the second sees the first's enqueue at its next tick. That is a
44//! deliberate limit and not a bug — the *claim* is race-free across processes,
45//! only the wake-up is not.
46
47use std::sync::Arc;
48use std::sync::atomic::{AtomicU32, Ordering};
49use std::time::Duration;
50
51use async_trait::async_trait;
52use serde_json::{Value, json};
53use tokio::sync::Notify;
54use tracing::error;
55
56use crate::config::JobsConfig;
57use crate::sqlite::db::Database;
58use crate::sqlite::job::Job;
59use crate::sqlite::nonce::now_secs;
60
61pub mod registry;
62pub mod runner;
63pub mod sweep;
64
65pub use registry::JobRegistry;
66pub use runner::{spawn_runner, spawn_runner_watching};
67pub use sweep::SweepJob;
68
69/// How one attempt ended.
70#[derive(Debug)]
71pub enum JobOutcome {
72 /// Finished. Terminal, and nothing else runs.
73 Done,
74 /// This attempt failed for a reason that may not recur: schedule another
75 /// after the backoff, unless the attempts or the deadline have run out.
76 Retry(String),
77 /// This will never succeed. Terminal, and [`JobHandler::abandon`] is called.
78 Failed(String),
79 /// This *occurrence* is finished; run again after the given delay.
80 ///
81 /// The row stays live, so its identity stays taken and nothing can queue a
82 /// second copy — which is what makes a periodic job one row rather than a
83 /// growing pile of them, with no cron table to keep in step.
84 Reschedule(Duration),
85}
86
87/// What to queue.
88#[derive(Debug, Clone)]
89pub struct JobSpec {
90 /// Which handler runs it.
91 pub kind: &'static str,
92 /// The identity of the work within its kind. A live job holds it; a settled
93 /// one releases it.
94 pub key: String,
95 /// The *identity* of the subject, never a snapshot of its state — a snapshot
96 /// goes stale across a retry.
97 pub payload: Value,
98 /// Not before this instant, epoch seconds. `now_secs()` for immediate.
99 pub run_at: i64,
100 /// The outer bound on retrying, epoch seconds. `None` for no deadline.
101 pub deadline: Option<i64>,
102 /// Overrides `jobs.max_attempts` for this one job.
103 pub max_attempts: Option<u32>,
104}
105
106impl JobSpec {
107 /// A job to run as soon as the runner can take it.
108 #[must_use]
109 pub fn now(kind: &'static str, key: impl Into<String>) -> Self {
110 Self {
111 kind,
112 key: key.into(),
113 payload: json!({}),
114 run_at: now_secs(),
115 deadline: None,
116 max_attempts: None,
117 }
118 }
119
120 #[must_use]
121 pub fn with_payload(mut self, payload: Value) -> Self {
122 self.payload = payload;
123 self
124 }
125
126 #[must_use]
127 pub fn with_deadline(mut self, deadline: Option<i64>) -> Self {
128 self.deadline = deadline;
129 self
130 }
131
132 #[must_use]
133 pub fn with_delay(mut self, delay: Duration) -> Self {
134 self.run_at = now_secs().saturating_add(seconds(delay));
135 self
136 }
137}
138
139/// One kind of background work.
140#[async_trait]
141pub trait JobHandler: Send + Sync {
142 /// The `jobs.kind` this handler answers for. One handler per kind — the
143 /// registry refuses a second, since two would each get half the rows.
144 fn kind(&self) -> &'static str;
145
146 /// Runs one attempt.
147 ///
148 /// Must be safe to run again from scratch: see the module documentation.
149 async fn run(&self, job: &Job) -> JobOutcome;
150
151 /// How long one attempt may take, and therefore how long the lease is held.
152 ///
153 /// `None` takes `jobs.lease_seconds`. The runner enforces it with a timeout
154 /// and writes a lease slightly longer, so the in-process deadline always
155 /// fires before another runner could steal the row.
156 fn lease(&self) -> Option<Duration> {
157 None
158 }
159
160 /// Called exactly once when the runner retires this job permanently — the
161 /// handler said [`JobOutcome::Failed`], the attempts ran out, or the deadline
162 /// passed.
163 ///
164 /// This is where a handler tells its subject what happened. Deliberately not
165 /// called between retries: a client polling an order must not be told it
166 /// failed while the server is still trying.
167 async fn abandon(&self, _job: &Job, _reason: &str) {}
168
169 /// Re-derives work a previous process left unfinished. Run once, at startup,
170 /// before the loop begins.
171 ///
172 /// The generic replacement for the old `SignerBackend::resume`: an enqueue
173 /// rather than a spawn, and safely repeatable because the identity index
174 /// refuses a duplicate.
175 async fn recover(&self, _queue: &JobQueue) {}
176}
177
178/// The enqueue side of the queue, plus the runner's wake-up.
179///
180/// Cloneable and cheap to hold: everything inside is an `Arc` or a scalar. A
181/// subsystem that queues work holds one of these and never sees the runner.
182#[derive(Clone)]
183pub struct JobQueue {
184 database: Arc<Database>,
185 notify: Arc<Notify>,
186 /// Shared rather than copied, because a reload has to reach the clones.
187 /// This queue is cloned into [`crate::Assembly`], every `RelaySigner` and
188 /// every `NotifyDispatcher` at startup and none of them is ever rebuilt, so
189 /// a plain `u32` field would leave `jobs.max_attempts` readable only where
190 /// the reload happened to be holding a handle. See
191 /// [`set_max_attempts`](JobQueue::set_max_attempts).
192 default_max_attempts: Arc<AtomicU32>,
193}
194
195impl JobQueue {
196 #[must_use]
197 pub fn new(database: Arc<Database>, config: &JobsConfig) -> Self {
198 Self {
199 database,
200 notify: Arc::new(Notify::new()),
201 default_max_attempts: Arc::new(AtomicU32::new(config.max_attempts)),
202 }
203 }
204
205 /// Republishes `jobs.max_attempts`, for every clone of this queue at once.
206 ///
207 /// Synchronous and infallible, which is what lets a reload call it from
208 /// `cli::apply_reload`'s publishing run beside the `watch` sends.
209 ///
210 /// It sets the budget for work queued from **now on** and does not touch the
211 /// backlog: `max_attempts` is frozen onto each row at enqueue, so a job
212 /// already waiting keeps the budget it was queued under. Raising this to
213 /// rescue rows that are about to give up is therefore not what it does —
214 /// that would be an `UPDATE` over pending rows, and a deliberately different
215 /// promise from the one `crate::sqlite::job` makes.
216 pub fn set_max_attempts(&self, max_attempts: u32) {
217 self.default_max_attempts
218 .store(max_attempts, Ordering::Relaxed);
219 }
220
221 /// Queues one job.
222 ///
223 /// `Ok(false)` means a live job already holds this `(kind, key)` — not an
224 /// error, but the answer two racing callers must both survive. The caller
225 /// reads it the way `RelaySigner::issue` reads `UpstreamOrder::create`'s
226 /// `Ok(None)`: somebody else is already on it.
227 pub async fn enqueue(&self, spec: JobSpec) -> Result<bool, sqlx::Error> {
228 let max_attempts = i64::from(
229 spec.max_attempts
230 .unwrap_or_else(|| self.default_max_attempts.load(Ordering::Relaxed)),
231 );
232 let queued = Job::enqueue(
233 crate::sqlite::job::NewJob {
234 id: &uuid::Uuid::new_v4().to_string(),
235 kind: spec.kind,
236 dedup_key: &spec.key,
237 payload: &spec.payload,
238 run_at: spec.run_at,
239 deadline: spec.deadline,
240 max_attempts,
241 },
242 &self.database,
243 )
244 .await?;
245
246 if queued {
247 // Wakes the runner now rather than at its next tick: the request
248 // that queued this is often one an ACME client is already polling.
249 self.notify.notify_one();
250 }
251 Ok(queued)
252 }
253
254 /// [`JobQueue::enqueue`], logging rather than returning a database failure.
255 ///
256 /// For the callers with nowhere to report one — `recover`, which the trait
257 /// defines as best-effort, and any path where refusing to queue is worse
258 /// than the error it would surface.
259 pub async fn enqueue_or_log(&self, spec: JobSpec) -> bool {
260 let kind = spec.kind;
261 let key = spec.key.clone();
262 match self.enqueue(spec).await {
263 Ok(queued) => queued,
264 Err(error) => {
265 error!(
266 event = "job_enqueue_failed",
267 outcome = "failure",
268 job_kind = %kind,
269 dedup_key = %key,
270 error = %error,
271 );
272 false
273 }
274 }
275 }
276
277 /// The database this queue writes to, for a handler that needs one and would
278 /// otherwise have to be handed a second copy.
279 #[must_use]
280 pub fn database(&self) -> &Arc<Database> {
281 &self.database
282 }
283}
284
285/// A `Duration` as whole seconds, saturating.
286///
287/// Job schedules are epoch seconds — the column type the whole schema uses — so
288/// every `Duration` crossing that boundary goes through here rather than through
289/// a bare `as` cast at four call sites.
290pub(crate) fn seconds(duration: Duration) -> i64 {
291 i64::try_from(duration.as_secs()).unwrap_or(i64::MAX)
292}
293
294#[cfg(test)]
295mod tests {
296 use super::*;
297
298 async fn queue() -> JobQueue {
299 let database = Arc::new(Database::connect_in_memory().await.unwrap());
300 JobQueue::new(database, &JobsConfig::default())
301 }
302
303 #[tokio::test]
304 async fn enqueue_reports_whether_it_took_the_identity() {
305 let queue = queue().await;
306 assert!(queue.enqueue(JobSpec::now("test", "k")).await.unwrap());
307 assert!(
308 !queue.enqueue(JobSpec::now("test", "k")).await.unwrap(),
309 "a second live job for one key is refused, not duplicated"
310 );
311 assert!(queue.enqueue(JobSpec::now("test", "other")).await.unwrap());
312 }
313
314 #[tokio::test]
315 async fn a_spec_carries_its_payload_deadline_and_attempt_budget() {
316 let queue = queue().await;
317 let deadline = now_secs() + 60;
318 let spec = JobSpec {
319 max_attempts: Some(9),
320 ..JobSpec::now("test", "k")
321 .with_payload(json!({"order_id": "ord-1"}))
322 .with_deadline(Some(deadline))
323 };
324 assert!(queue.enqueue(spec).await.unwrap());
325
326 let job = Job::find_live("test", "k", queue.database())
327 .await
328 .unwrap()
329 .unwrap();
330 assert_eq!(job.payload, json!({"order_id": "ord-1"}));
331 assert_eq!(job.deadline, Some(deadline));
332 assert_eq!(job.max_attempts, 9);
333 }
334
335 #[tokio::test]
336 async fn a_spec_with_no_budget_of_its_own_takes_the_configured_one() {
337 let queue = queue().await;
338 assert!(queue.enqueue(JobSpec::now("test", "k")).await.unwrap());
339 let job = Job::find_live("test", "k", queue.database())
340 .await
341 .unwrap()
342 .unwrap();
343 assert_eq!(
344 job.max_attempts,
345 i64::from(JobsConfig::default().max_attempts)
346 );
347 }
348
349 /// A reloaded `jobs.max_attempts` has to reach the clones, because the
350 /// clones are all there is: this queue is handed to `Assembly`, to every
351 /// relay signer and to every notify dispatcher at startup, and none of them
352 /// is ever rebuilt. A copied `u32` would have left the new value visible only
353 /// wherever the reload happened to be holding a handle.
354 ///
355 /// And it reaches *future* work only. A row already waiting keeps the budget
356 /// frozen onto it at enqueue, which is the promise `crate::sqlite::job` makes
357 /// and the reason raising this is not a way to rescue a backlog.
358 #[tokio::test]
359 async fn a_reloaded_max_attempts_reaches_the_clones_but_not_the_backlog() {
360 let queue = queue().await;
361 let held = queue.clone();
362
363 queue.enqueue(JobSpec::now("test", "before")).await.unwrap();
364 queue.set_max_attempts(11);
365 held.enqueue(JobSpec::now("test", "after")).await.unwrap();
366
367 let queued = |key: &'static str| {
368 let database = queue.database().clone();
369 async move {
370 Job::find_live("test", key, &database)
371 .await
372 .unwrap()
373 .unwrap()
374 .max_attempts
375 }
376 };
377
378 assert_eq!(
379 queued("before").await,
380 i64::from(JobsConfig::default().max_attempts),
381 "a row already queued keeps what it was queued under",
382 );
383 assert_eq!(
384 queued("after").await,
385 11,
386 "a clone taken before the change still enqueues under the new value",
387 );
388 }
389
390 #[tokio::test]
391 async fn with_delay_pushes_the_run_time_out() {
392 let queue = queue().await;
393 let spec = JobSpec::now("test", "k").with_delay(Duration::from_secs(600));
394 assert!(spec.run_at >= now_secs() + 599);
395 assert!(queue.enqueue(spec).await.unwrap());
396 }
397
398 /// The failure path a `recover` implementation relies on: a closed pool is
399 /// logged and reported as "not queued", never propagated into a caller with
400 /// nowhere to put it.
401 #[tokio::test]
402 async fn enqueue_or_log_swallows_a_database_failure() {
403 let queue = queue().await;
404 queue.database().pool.close().await;
405 assert!(!queue.enqueue_or_log(JobSpec::now("test", "k")).await);
406 }
407
408 #[test]
409 fn seconds_saturates_rather_than_wrapping() {
410 assert_eq!(seconds(Duration::from_secs(90)), 90);
411 assert_eq!(seconds(Duration::MAX), i64::MAX);
412 }
413}