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 at *this row* may take, and therefore how long the
152    /// lease is held.
153    ///
154    /// `None` takes `jobs.lease_seconds`. The runner enforces it with a timeout
155    /// and writes a lease slightly longer, so the in-process deadline always
156    /// fires before another runner could steal the row.
157    ///
158    /// Takes the job because a handler may serve several configurations of one
159    /// subsystem — `signer::relay::flow::RelayJob` is one handler over every
160    /// relay profile in the process, and `signer.relay.poll_timeout_secs` is
161    /// per profile. The runner has the claimed row in hand when it asks, so
162    /// this costs nothing; a handler with one budget ignores the argument.
163    fn lease(&self, _job: &Job) -> Option<Duration> {
164        None
165    }
166
167    /// Called exactly once when the runner retires this job permanently — the
168    /// handler said [`JobOutcome::Failed`], the attempts ran out, or the deadline
169    /// passed.
170    ///
171    /// This is where a handler tells its subject what happened. Deliberately not
172    /// called between retries: a client polling an order must not be told it
173    /// failed while the server is still trying.
174    async fn abandon(&self, _job: &Job, _reason: &str) {}
175
176    /// Re-derives work a previous process left unfinished. Run once, at startup,
177    /// before the loop begins.
178    ///
179    /// The generic replacement for the old `SignerBackend::resume`: an enqueue
180    /// rather than a spawn, and safely repeatable because the identity index
181    /// refuses a duplicate.
182    async fn recover(&self, _queue: &JobQueue) {}
183}
184
185/// The enqueue side of the queue, plus the runner's wake-up.
186///
187/// Cloneable and cheap to hold: everything inside is an `Arc` or a scalar. A
188/// subsystem that queues work holds one of these and never sees the runner.
189#[derive(Clone)]
190pub struct JobQueue {
191    database: Arc<Database>,
192    notify: Arc<Notify>,
193    /// Shared rather than copied, because a reload has to reach the clones.
194    /// This queue is cloned into [`crate::Assembly`], every `RelaySigner` and
195    /// every `NotifyDispatcher` at startup and none of them is ever rebuilt, so
196    /// a plain `u32` field would leave `jobs.max_attempts` readable only where
197    /// the reload happened to be holding a handle. See
198    /// [`set_max_attempts`](JobQueue::set_max_attempts).
199    default_max_attempts: Arc<AtomicU32>,
200}
201
202impl JobQueue {
203    #[must_use]
204    pub fn new(database: Arc<Database>, config: &JobsConfig) -> Self {
205        Self {
206            database,
207            notify: Arc::new(Notify::new()),
208            default_max_attempts: Arc::new(AtomicU32::new(config.max_attempts)),
209        }
210    }
211
212    /// Republishes `jobs.max_attempts`, for every clone of this queue at once.
213    ///
214    /// Synchronous and infallible, which is what lets a reload call it from
215    /// `cli::apply_reload`'s publishing run beside the `watch` sends.
216    ///
217    /// It sets the budget for work queued from **now on** and does not touch the
218    /// backlog: `max_attempts` is frozen onto each row at enqueue, so a job
219    /// already waiting keeps the budget it was queued under. Raising this to
220    /// rescue rows that are about to give up is therefore not what it does —
221    /// that would be an `UPDATE` over pending rows, and a deliberately different
222    /// promise from the one `crate::sqlite::job` makes.
223    pub fn set_max_attempts(&self, max_attempts: u32) {
224        self.default_max_attempts
225            .store(max_attempts, Ordering::Relaxed);
226    }
227
228    /// Queues one job.
229    ///
230    /// `Ok(false)` means a live job already holds this `(kind, key)` — not an
231    /// error, but the answer two racing callers must both survive. The caller
232    /// reads it the way `RelaySigner::issue` reads `UpstreamOrder::create`'s
233    /// `Ok(None)`: somebody else is already on it.
234    pub async fn enqueue(&self, spec: JobSpec) -> Result<bool, sqlx::Error> {
235        let max_attempts = i64::from(
236            spec.max_attempts
237                .unwrap_or_else(|| self.default_max_attempts.load(Ordering::Relaxed)),
238        );
239        let queued = Job::enqueue(
240            crate::sqlite::job::NewJob {
241                id: crate::sqlite::id::mint(),
242                kind: spec.kind,
243                dedup_key: &spec.key,
244                payload: &spec.payload,
245                run_at: spec.run_at,
246                deadline: spec.deadline,
247                max_attempts,
248            },
249            &self.database,
250        )
251        .await?;
252
253        if queued {
254            // Wakes the runner now rather than at its next tick: the request
255            // that queued this is often one an ACME client is already polling.
256            self.notify.notify_one();
257        }
258        Ok(queued)
259    }
260
261    /// [`JobQueue::enqueue`], logging rather than returning a database failure.
262    ///
263    /// For the callers with nowhere to report one — `recover`, which the trait
264    /// defines as best-effort, and any path where refusing to queue is worse
265    /// than the error it would surface.
266    pub async fn enqueue_or_log(&self, spec: JobSpec) -> bool {
267        let kind = spec.kind;
268        let key = spec.key.clone();
269        match self.enqueue(spec).await {
270            Ok(queued) => queued,
271            Err(error) => {
272                error!(
273                    event = "job_enqueue_failed",
274                    outcome = "failure",
275                    job_kind = %kind,
276                    dedup_key = %key,
277                    error = %error,
278                );
279                false
280            }
281        }
282    }
283
284    /// The database this queue writes to, for a handler that needs one and would
285    /// otherwise have to be handed a second copy.
286    #[must_use]
287    pub fn database(&self) -> &Arc<Database> {
288        &self.database
289    }
290}
291
292/// A `Duration` as whole seconds, saturating.
293///
294/// Job schedules are epoch seconds — the column type the whole schema uses — so
295/// every `Duration` crossing that boundary goes through here rather than through
296/// a bare `as` cast at four call sites.
297pub(crate) fn seconds(duration: Duration) -> i64 {
298    i64::try_from(duration.as_secs()).unwrap_or(i64::MAX)
299}
300
301#[cfg(test)]
302mod tests {
303    use super::*;
304
305    async fn queue() -> JobQueue {
306        let database = Arc::new(Database::connect_in_memory().await.unwrap());
307        JobQueue::new(database, &JobsConfig::default())
308    }
309
310    #[tokio::test]
311    async fn enqueue_reports_whether_it_took_the_identity() {
312        let queue = queue().await;
313        assert!(queue.enqueue(JobSpec::now("test", "k")).await.unwrap());
314        assert!(
315            !queue.enqueue(JobSpec::now("test", "k")).await.unwrap(),
316            "a second live job for one key is refused, not duplicated"
317        );
318        assert!(queue.enqueue(JobSpec::now("test", "other")).await.unwrap());
319    }
320
321    #[tokio::test]
322    async fn a_spec_carries_its_payload_deadline_and_attempt_budget() {
323        let queue = queue().await;
324        let deadline = now_secs() + 60;
325        let spec = JobSpec {
326            max_attempts: Some(9),
327            ..JobSpec::now("test", "k")
328                .with_payload(json!({"order_id": "ord-1"}))
329                .with_deadline(Some(deadline))
330        };
331        assert!(queue.enqueue(spec).await.unwrap());
332
333        let job = Job::find_live("test", "k", queue.database())
334            .await
335            .unwrap()
336            .unwrap();
337        assert_eq!(job.payload, json!({"order_id": "ord-1"}));
338        assert_eq!(job.deadline, Some(deadline));
339        assert_eq!(job.max_attempts, 9);
340    }
341
342    #[tokio::test]
343    async fn a_spec_with_no_budget_of_its_own_takes_the_configured_one() {
344        let queue = queue().await;
345        assert!(queue.enqueue(JobSpec::now("test", "k")).await.unwrap());
346        let job = Job::find_live("test", "k", queue.database())
347            .await
348            .unwrap()
349            .unwrap();
350        assert_eq!(
351            job.max_attempts,
352            i64::from(JobsConfig::default().max_attempts)
353        );
354    }
355
356    /// A reloaded `jobs.max_attempts` has to reach the clones, because the
357    /// clones are all there is: this queue is handed to `Assembly`, to every
358    /// relay signer and to every notify dispatcher at startup, and none of them
359    /// is ever rebuilt. A copied `u32` would have left the new value visible only
360    /// wherever the reload happened to be holding a handle.
361    ///
362    /// And it reaches *future* work only. A row already waiting keeps the budget
363    /// frozen onto it at enqueue, which is the promise `crate::sqlite::job` makes
364    /// and the reason raising this is not a way to rescue a backlog.
365    #[tokio::test]
366    async fn a_reloaded_max_attempts_reaches_the_clones_but_not_the_backlog() {
367        let queue = queue().await;
368        let held = queue.clone();
369
370        queue.enqueue(JobSpec::now("test", "before")).await.unwrap();
371        queue.set_max_attempts(11);
372        held.enqueue(JobSpec::now("test", "after")).await.unwrap();
373
374        let queued = |key: &'static str| {
375            let database = queue.database().clone();
376            async move {
377                Job::find_live("test", key, &database)
378                    .await
379                    .unwrap()
380                    .unwrap()
381                    .max_attempts
382            }
383        };
384
385        assert_eq!(
386            queued("before").await,
387            i64::from(JobsConfig::default().max_attempts),
388            "a row already queued keeps what it was queued under",
389        );
390        assert_eq!(
391            queued("after").await,
392            11,
393            "a clone taken before the change still enqueues under the new value",
394        );
395    }
396
397    #[tokio::test]
398    async fn with_delay_pushes_the_run_time_out() {
399        let queue = queue().await;
400        let spec = JobSpec::now("test", "k").with_delay(Duration::from_secs(600));
401        assert!(spec.run_at >= now_secs() + 599);
402        assert!(queue.enqueue(spec).await.unwrap());
403    }
404
405    /// The failure path a `recover` implementation relies on: a closed pool is
406    /// logged and reported as "not queued", never propagated into a caller with
407    /// nowhere to put it.
408    #[tokio::test]
409    async fn enqueue_or_log_swallows_a_database_failure() {
410        let queue = queue().await;
411        queue.database().pool.close().await;
412        assert!(!queue.enqueue_or_log(JobSpec::now("test", "k")).await);
413    }
414
415    #[test]
416    fn seconds_saturates_rather_than_wrapping() {
417        assert_eq!(seconds(Duration::from_secs(90)), 90);
418        assert_eq!(seconds(Duration::MAX), i64::MAX);
419    }
420}