Skip to main content

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}