acme-proxy 0.5.0

An ACME (RFC 8555) server that issues from a local CA, relays to an upstream CA, or delegates to a script
Documentation
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
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
//! The durable job runner: background work that survives the process.
//!
//! ## Why this is not a `tokio::spawn`
//!
//! A spawned task is a promise the process makes to itself, and a restart breaks
//! it. That was tolerable while exactly one subsystem needed one — the `relay`
//! signer backend, whose state lived on its own `upstream_orders` row and whose
//! `resume` sweep re-created the task at startup — but it bought only *recovery*,
//! never *retry*: with nowhere to record that an attempt had failed and should
//! happen again, every failure had to be terminal. A five-second upstream blip
//! therefore invalidated an order the client then had to place again.
//!
//! Here the schedule is a row. An attempt that fails writes its next `run_at`
//! and goes back in the queue; a process that dies leaves a lease that expires
//! and a row another runner takes. Both are the same table and the same loop, so
//! the next subsystem that wants a queue — expiry reminders, a CRL sweep, an
//! OCSP refresh — registers a handler rather than writing this again.
//!
//! ## The three things a handler must know
//!
//! 1. **`run` may be called again from scratch.** The runner promises nothing
//!    about what a previous attempt did. `signer::relay::flow` is the worked
//!    example: it re-reads the upstream order and skips every authorization that
//!    is no longer `pending`, so re-entering at any point is safe by
//!    construction rather than by checkpointing.
//! 2. **[`JobOutcome::Retry`] and [`JobOutcome::Failed`] are a real
//!    distinction.** `Retry` is "the network, the nameserver, a 503"; `Failed` is
//!    "the CA stated a reason" or "this payload can never parse". Classifying a
//!    permanent failure as retryable wastes the budget and delays the client's
//!    real answer; the reverse throws away the retry this module exists for.
//! 3. **[`JobHandler::abandon`] fires exactly once**, when the job is retired for
//!    good — and never between retries. It is where a handler tells its *subject*
//!    what happened, so a client polling through a transient blip keeps seeing
//!    work in progress rather than a terminal failure it must start over from.
//!
//! ## Wake-up is in-process
//!
//! [`JobQueue::enqueue`] notifies the runner directly, so a job queued by a
//! request starts within microseconds rather than waiting for
//! `jobs.poll_interval_ms`. That matters because the ACME client that triggered
//! it is already polling its order. The notification is a
//! [`tokio::sync::Notify`], which is process-local: with two processes over one
//! database, the second sees the first's enqueue at its next tick. That is a
//! deliberate limit and not a bug — the *claim* is race-free across processes,
//! only the wake-up is not.

use std::sync::Arc;
use std::sync::atomic::{AtomicU32, Ordering};
use std::time::Duration;

use async_trait::async_trait;
use serde_json::{Value, json};
use tokio::sync::Notify;
use tracing::error;

use crate::config::JobsConfig;
use crate::sqlite::db::Database;
use crate::sqlite::job::Job;
use crate::sqlite::nonce::now_secs;

pub mod registry;
pub mod runner;
pub mod sweep;

pub use registry::JobRegistry;
pub use runner::{spawn_runner, spawn_runner_watching};
pub use sweep::SweepJob;

/// How one attempt ended.
#[derive(Debug)]
pub enum JobOutcome {
    /// Finished. Terminal, and nothing else runs.
    Done,
    /// This attempt failed for a reason that may not recur: schedule another
    /// after the backoff, unless the attempts or the deadline have run out.
    Retry(String),
    /// This will never succeed. Terminal, and [`JobHandler::abandon`] is called.
    Failed(String),
    /// This *occurrence* is finished; run again after the given delay.
    ///
    /// The row stays live, so its identity stays taken and nothing can queue a
    /// second copy — which is what makes a periodic job one row rather than a
    /// growing pile of them, with no cron table to keep in step.
    Reschedule(Duration),
}

/// What to queue.
#[derive(Debug, Clone)]
pub struct JobSpec {
    /// Which handler runs it.
    pub kind: &'static str,
    /// The identity of the work within its kind. A live job holds it; a settled
    /// one releases it.
    pub key: String,
    /// The *identity* of the subject, never a snapshot of its state — a snapshot
    /// goes stale across a retry.
    pub payload: Value,
    /// Not before this instant, epoch seconds. `now_secs()` for immediate.
    pub run_at: i64,
    /// The outer bound on retrying, epoch seconds. `None` for no deadline.
    pub deadline: Option<i64>,
    /// Overrides `jobs.max_attempts` for this one job.
    pub max_attempts: Option<u32>,
}

impl JobSpec {
    /// A job to run as soon as the runner can take it.
    #[must_use]
    pub fn now(kind: &'static str, key: impl Into<String>) -> Self {
        Self {
            kind,
            key: key.into(),
            payload: json!({}),
            run_at: now_secs(),
            deadline: None,
            max_attempts: None,
        }
    }

    #[must_use]
    pub fn with_payload(mut self, payload: Value) -> Self {
        self.payload = payload;
        self
    }

    #[must_use]
    pub fn with_deadline(mut self, deadline: Option<i64>) -> Self {
        self.deadline = deadline;
        self
    }

    #[must_use]
    pub fn with_delay(mut self, delay: Duration) -> Self {
        self.run_at = now_secs().saturating_add(seconds(delay));
        self
    }
}

/// One kind of background work.
#[async_trait]
pub trait JobHandler: Send + Sync {
    /// The `jobs.kind` this handler answers for. One handler per kind — the
    /// registry refuses a second, since two would each get half the rows.
    fn kind(&self) -> &'static str;

    /// Runs one attempt.
    ///
    /// Must be safe to run again from scratch: see the module documentation.
    async fn run(&self, job: &Job) -> JobOutcome;

    /// How long one attempt at *this row* may take, and therefore how long the
    /// lease is held.
    ///
    /// `None` takes `jobs.lease_seconds`. The runner enforces it with a timeout
    /// and writes a lease slightly longer, so the in-process deadline always
    /// fires before another runner could steal the row.
    ///
    /// Takes the job because a handler may serve several configurations of one
    /// subsystem — `signer::relay::flow::RelayJob` is one handler over every
    /// relay profile in the process, and `signer.relay.poll_timeout_secs` is
    /// per profile. The runner has the claimed row in hand when it asks, so
    /// this costs nothing; a handler with one budget ignores the argument.
    fn lease(&self, _job: &Job) -> Option<Duration> {
        None
    }

    /// Called exactly once when the runner retires this job permanently — the
    /// handler said [`JobOutcome::Failed`], the attempts ran out, or the deadline
    /// passed.
    ///
    /// This is where a handler tells its subject what happened. Deliberately not
    /// called between retries: a client polling an order must not be told it
    /// failed while the server is still trying.
    async fn abandon(&self, _job: &Job, _reason: &str) {}

    /// Re-derives work a previous process left unfinished. Run once, at startup,
    /// before the loop begins.
    ///
    /// The generic replacement for the old `SignerBackend::resume`: an enqueue
    /// rather than a spawn, and safely repeatable because the identity index
    /// refuses a duplicate.
    async fn recover(&self, _queue: &JobQueue) {}
}

/// The enqueue side of the queue, plus the runner's wake-up.
///
/// Cloneable and cheap to hold: everything inside is an `Arc` or a scalar. A
/// subsystem that queues work holds one of these and never sees the runner.
#[derive(Clone)]
pub struct JobQueue {
    database: Arc<Database>,
    notify: Arc<Notify>,
    /// Shared rather than copied, because a reload has to reach the clones.
    /// This queue is cloned into [`crate::Assembly`], every `RelaySigner` and
    /// every `NotifyDispatcher` at startup and none of them is ever rebuilt, so
    /// a plain `u32` field would leave `jobs.max_attempts` readable only where
    /// the reload happened to be holding a handle. See
    /// [`set_max_attempts`](JobQueue::set_max_attempts).
    default_max_attempts: Arc<AtomicU32>,
}

impl JobQueue {
    #[must_use]
    pub fn new(database: Arc<Database>, config: &JobsConfig) -> Self {
        Self {
            database,
            notify: Arc::new(Notify::new()),
            default_max_attempts: Arc::new(AtomicU32::new(config.max_attempts)),
        }
    }

    /// Republishes `jobs.max_attempts`, for every clone of this queue at once.
    ///
    /// Synchronous and infallible, which is what lets a reload call it from
    /// `cli::apply_reload`'s publishing run beside the `watch` sends.
    ///
    /// It sets the budget for work queued from **now on** and does not touch the
    /// backlog: `max_attempts` is frozen onto each row at enqueue, so a job
    /// already waiting keeps the budget it was queued under. Raising this to
    /// rescue rows that are about to give up is therefore not what it does —
    /// that would be an `UPDATE` over pending rows, and a deliberately different
    /// promise from the one `crate::sqlite::job` makes.
    pub fn set_max_attempts(&self, max_attempts: u32) {
        self.default_max_attempts
            .store(max_attempts, Ordering::Relaxed);
    }

    /// Queues one job.
    ///
    /// `Ok(false)` means a live job already holds this `(kind, key)` — not an
    /// error, but the answer two racing callers must both survive. The caller
    /// reads it the way `RelaySigner::issue` reads `UpstreamOrder::create`'s
    /// `Ok(None)`: somebody else is already on it.
    pub async fn enqueue(&self, spec: JobSpec) -> Result<bool, sqlx::Error> {
        let max_attempts = i64::from(
            spec.max_attempts
                .unwrap_or_else(|| self.default_max_attempts.load(Ordering::Relaxed)),
        );
        let queued = Job::enqueue(
            crate::sqlite::job::NewJob {
                id: crate::sqlite::id::mint(),
                kind: spec.kind,
                dedup_key: &spec.key,
                payload: &spec.payload,
                run_at: spec.run_at,
                deadline: spec.deadline,
                max_attempts,
            },
            &self.database,
        )
        .await?;

        if queued {
            // Wakes the runner now rather than at its next tick: the request
            // that queued this is often one an ACME client is already polling.
            self.notify.notify_one();
        }
        Ok(queued)
    }

    /// [`JobQueue::enqueue`], logging rather than returning a database failure.
    ///
    /// For the callers with nowhere to report one — `recover`, which the trait
    /// defines as best-effort, and any path where refusing to queue is worse
    /// than the error it would surface.
    pub async fn enqueue_or_log(&self, spec: JobSpec) -> bool {
        let kind = spec.kind;
        let key = spec.key.clone();
        match self.enqueue(spec).await {
            Ok(queued) => queued,
            Err(error) => {
                error!(
                    event = "job_enqueue_failed",
                    outcome = "failure",
                    job_kind = %kind,
                    dedup_key = %key,
                    error = %error,
                );
                false
            }
        }
    }

    /// The database this queue writes to, for a handler that needs one and would
    /// otherwise have to be handed a second copy.
    #[must_use]
    pub fn database(&self) -> &Arc<Database> {
        &self.database
    }
}

/// A `Duration` as whole seconds, saturating.
///
/// Job schedules are epoch seconds — the column type the whole schema uses — so
/// every `Duration` crossing that boundary goes through here rather than through
/// a bare `as` cast at four call sites.
pub(crate) fn seconds(duration: Duration) -> i64 {
    i64::try_from(duration.as_secs()).unwrap_or(i64::MAX)
}

#[cfg(test)]
mod tests {
    use super::*;

    async fn queue() -> JobQueue {
        let database = Arc::new(Database::connect_in_memory().await.unwrap());
        JobQueue::new(database, &JobsConfig::default())
    }

    #[tokio::test]
    async fn enqueue_reports_whether_it_took_the_identity() {
        let queue = queue().await;
        assert!(queue.enqueue(JobSpec::now("test", "k")).await.unwrap());
        assert!(
            !queue.enqueue(JobSpec::now("test", "k")).await.unwrap(),
            "a second live job for one key is refused, not duplicated"
        );
        assert!(queue.enqueue(JobSpec::now("test", "other")).await.unwrap());
    }

    #[tokio::test]
    async fn a_spec_carries_its_payload_deadline_and_attempt_budget() {
        let queue = queue().await;
        let deadline = now_secs() + 60;
        let spec = JobSpec {
            max_attempts: Some(9),
            ..JobSpec::now("test", "k")
                .with_payload(json!({"order_id": "ord-1"}))
                .with_deadline(Some(deadline))
        };
        assert!(queue.enqueue(spec).await.unwrap());

        let job = Job::find_live("test", "k", queue.database())
            .await
            .unwrap()
            .unwrap();
        assert_eq!(job.payload, json!({"order_id": "ord-1"}));
        assert_eq!(job.deadline, Some(deadline));
        assert_eq!(job.max_attempts, 9);
    }

    #[tokio::test]
    async fn a_spec_with_no_budget_of_its_own_takes_the_configured_one() {
        let queue = queue().await;
        assert!(queue.enqueue(JobSpec::now("test", "k")).await.unwrap());
        let job = Job::find_live("test", "k", queue.database())
            .await
            .unwrap()
            .unwrap();
        assert_eq!(
            job.max_attempts,
            i64::from(JobsConfig::default().max_attempts)
        );
    }

    /// A reloaded `jobs.max_attempts` has to reach the clones, because the
    /// clones are all there is: this queue is handed to `Assembly`, to every
    /// relay signer and to every notify dispatcher at startup, and none of them
    /// is ever rebuilt. A copied `u32` would have left the new value visible only
    /// wherever the reload happened to be holding a handle.
    ///
    /// And it reaches *future* work only. A row already waiting keeps the budget
    /// frozen onto it at enqueue, which is the promise `crate::sqlite::job` makes
    /// and the reason raising this is not a way to rescue a backlog.
    #[tokio::test]
    async fn a_reloaded_max_attempts_reaches_the_clones_but_not_the_backlog() {
        let queue = queue().await;
        let held = queue.clone();

        queue.enqueue(JobSpec::now("test", "before")).await.unwrap();
        queue.set_max_attempts(11);
        held.enqueue(JobSpec::now("test", "after")).await.unwrap();

        let queued = |key: &'static str| {
            let database = queue.database().clone();
            async move {
                Job::find_live("test", key, &database)
                    .await
                    .unwrap()
                    .unwrap()
                    .max_attempts
            }
        };

        assert_eq!(
            queued("before").await,
            i64::from(JobsConfig::default().max_attempts),
            "a row already queued keeps what it was queued under",
        );
        assert_eq!(
            queued("after").await,
            11,
            "a clone taken before the change still enqueues under the new value",
        );
    }

    #[tokio::test]
    async fn with_delay_pushes_the_run_time_out() {
        let queue = queue().await;
        let spec = JobSpec::now("test", "k").with_delay(Duration::from_secs(600));
        assert!(spec.run_at >= now_secs() + 599);
        assert!(queue.enqueue(spec).await.unwrap());
    }

    /// The failure path a `recover` implementation relies on: a closed pool is
    /// logged and reported as "not queued", never propagated into a caller with
    /// nowhere to put it.
    #[tokio::test]
    async fn enqueue_or_log_swallows_a_database_failure() {
        let queue = queue().await;
        queue.database().pool.close().await;
        assert!(!queue.enqueue_or_log(JobSpec::now("test", "k")).await);
    }

    #[test]
    fn seconds_saturates_rather_than_wrapping() {
        assert_eq!(seconds(Duration::from_secs(90)), 90);
        assert_eq!(seconds(Duration::MAX), i64::MAX);
    }
}