Skip to main content

acme_proxy/cli/
jobs.rs

1//! `acme-proxy jobs` — inspect and manage the background queue
2//! (`crates/jobs/src/jobs/` + `crates/store/src/job.rs`), the subsystem whose whole purpose is
3//! surviving the failures an operator gets paged about.
4//!
5//! `list` and `show` are the read half; `cancel` and `run-now` are the two
6//! mutations. Cancelling a `signer_relay_issue` job also abandons the ACME
7//! order it was driving — that coupling lives in [`acme_proxy_admin::admin::ops::cancel_job`],
8//! shared with the runner's own `RelayJob::abandon`.
9
10use std::io::BufRead;
11use std::sync::Arc;
12
13use clap::Subcommand;
14
15use crate::cli::CliError;
16use crate::cli::render;
17use crate::cli::window::{DEFAULT_LIMIT, Window};
18use acme_proxy_admin::admin;
19use acme_proxy_admin::admin::CancelJobOutcome;
20use acme_proxy_admin::admin::RunJobNowOutcome;
21use acme_proxy_core::audit::Actor;
22use acme_proxy_core::audit::ClientContext;
23use acme_proxy_core::palette::Palette;
24use acme_proxy_store::db::Database;
25use acme_proxy_store::job::Job;
26use acme_proxy_store::job::JobQuery;
27use acme_proxy_store::status::JobStatus;
28
29#[derive(Subcommand)]
30pub enum JobsCommand {
31    /// List background jobs, optionally filtered.
32    List {
33        /// A job kind (`signer_relay_issue`, `nonce_sweep`, …). Passed through
34        /// as-is: kinds are an open set, so an unknown one simply matches
35        /// nothing rather than being an error — unlike `--status`.
36        #[arg(long)]
37        kind: Option<String>,
38        /// Only jobs in this state: `ready`, `running`, `done`, `failed` or
39        /// `cancelled`.
40        #[arg(long)]
41        status: Option<String>,
42        /// Rows per page. A value below 1 is read as 1.
43        #[arg(long, default_value_t = DEFAULT_LIMIT)]
44        limit: i64,
45        /// Rows to skip before the page starts.
46        #[arg(long, default_value_t = 0)]
47        offset: i64,
48        /// Print the page as JSON: `{items, total, limit, offset}`.
49        #[arg(long)]
50        json: bool,
51    },
52    /// Show one job, plus the upstream order it drives if it is a relay job.
53    Show {
54        /// The job id.
55        id: String,
56        /// Print it as JSON.
57        #[arg(long)]
58        json: bool,
59    },
60    /// Retire a job. On an in-flight relay issuance this also marks the ACME
61    /// order invalid and abandons the upstream mapping.
62    Cancel {
63        /// The job id.
64        id: String,
65    },
66    /// Make a job eligible to run at the next queue poll. On a failed job this
67    /// grants exactly one more attempt.
68    RunNow {
69        /// The job id.
70        id: String,
71    },
72}
73
74pub async fn run_jobs_command(
75    command: JobsCommand,
76    yes: bool,
77    palette: Palette,
78    reader: &mut impl BufRead,
79    database: Arc<Database>,
80) -> Result<(), CliError> {
81    match command {
82        JobsCommand::List {
83            kind,
84            status,
85            limit,
86            offset,
87            json,
88        } => {
89            // Refused by name rather than passed through: an unknown status
90            // would match no rows, which reads exactly like "nothing is in
91            // that state" (the `order list --status` rule).
92            let status = super::parse_flag::<JobStatus>("--status", status)?;
93            let window = Window::resolve(limit, offset);
94            let query = JobQuery {
95                kind,
96                status,
97                limit: window.limit,
98                offset: window.offset,
99            };
100            let (jobs, total) = Job::search(&query, &database).await?;
101            render::print_page(&jobs, total, window, json, admin::render_job_json, |job| {
102                render::render_job_line(job, palette)
103            });
104        }
105        JobsCommand::Show { id, json } => {
106            let Some(detail) = admin::load_job_detail(&id, database).await? else {
107                return Err(not_found(&id));
108            };
109            if json {
110                println!("{}", admin::render_job_detail_json(&detail));
111            } else {
112                print!("{}", render::render_job_detail_text(&detail, palette));
113            }
114        }
115        JobsCommand::Cancel { id } => {
116            match admin::confirm_cancel_job(
117                &id,
118                yes,
119                reader,
120                Actor::cli(),
121                ClientContext::default(),
122                &acme_proxy_jobs::auditor::Auditor::offline(database.clone()),
123                database,
124            )
125            .await
126            .map_err(|error| CliError::failed(error.to_string()))?
127            {
128                // Not "Cancelled." — on a command named `cancel` that reads
129                // as "the job was cancelled", which is the opposite of what a
130                // declined prompt means. Every other confirm-gated command can
131                // use the idiom; this one cannot.
132                None => println!("Left the job alone."),
133                Some(CancelJobOutcome::NotFound) => return Err(not_found(&id)),
134                Some(CancelJobOutcome::NotCancellable(status)) => {
135                    return Err(CliError::bad_request(format!(
136                        "job {id} is {status}: only ready or failed jobs can be cancelled"
137                    )));
138                }
139                Some(CancelJobOutcome::Cancelled(job)) => {
140                    println!("Cancelled job {id} ({}).", job.kind);
141                }
142                Some(CancelJobOutcome::CancelledAndOrderAbandoned { order_id, .. }) => {
143                    println!(
144                        "Cancelled relay job {id}. Order {order_id} was marked invalid and the \
145                         upstream mapping abandoned; the client will stop polling."
146                    );
147                }
148            }
149        }
150        JobsCommand::RunNow { id } => match admin::run_job_now(
151            &id,
152            Actor::cli(),
153            ClientContext::default(),
154            &acme_proxy_jobs::auditor::Auditor::offline(database.clone()),
155            database.clone(),
156        )
157        .await?
158        {
159            RunJobNowOutcome::NotFound => return Err(not_found(&id)),
160            RunJobNowOutcome::Refused(status) => {
161                return Err(CliError::bad_request(format!(
162                    "job {id} is {status}: run-now applies to ready or failed jobs"
163                )));
164            }
165            RunJobNowOutcome::Nudged(_) => {
166                println!("Job {id} will run at the next queue poll (run_at set to now).");
167            }
168            RunJobNowOutcome::Revived(job) => {
169                println!(
170                    "Job {id} revived: status ready, attempts {}/{} (one more attempt).",
171                    job.attempts, job.max_attempts
172                );
173            }
174        },
175    }
176    Ok(())
177}
178
179fn not_found(id: &str) -> CliError {
180    CliError::bad_request(acme_proxy_admin::admin::subject::Subject::Job.missing(id))
181}
182
183#[cfg(test)]
184mod tests {
185    use super::*;
186    use crate::cli::CliErrorKind;
187    use acme_proxy_store::job::NewJob;
188    use acme_proxy_store::nonce::now_secs;
189    use serde_json::json;
190
191    async fn db() -> Arc<Database> {
192        Arc::new(Database::connect_in_memory().await.unwrap())
193    }
194
195    async fn seed(kind: &str, key: &str, run_at: i64, db: &Database) -> uuid::Uuid {
196        let id = acme_proxy_store::id::mint();
197        Job::enqueue(
198            NewJob {
199                id,
200                kind,
201                dedup_key: key,
202                payload: &json!({}),
203                run_at,
204                deadline: None,
205                max_attempts: 5,
206            },
207            db,
208        )
209        .await
210        .unwrap();
211        id
212    }
213
214    fn palette() -> Palette {
215        Palette::plain()
216    }
217
218    #[tokio::test]
219    async fn every_arm_refuses_an_unknown_job() {
220        let db = db().await;
221        for cmd in [
222            JobsCommand::Show {
223                id: "nope".into(),
224                json: false,
225            },
226            JobsCommand::Cancel { id: "nope".into() },
227            JobsCommand::RunNow { id: "nope".into() },
228        ] {
229            let err = run_jobs_command(cmd, true, palette(), &mut &b""[..], db.clone())
230                .await
231                .unwrap_err();
232            assert!(err.message.contains("no such job"), "{}", err.message);
233            assert_eq!(err.kind(), CliErrorKind::BadRequest);
234        }
235    }
236
237    #[tokio::test]
238    async fn list_renders_both_ways() {
239        let db = db().await;
240        seed("nonce_sweep", "a", now_secs(), &db).await;
241        seed("signer_relay_issue", "b", now_secs(), &db).await;
242        for json in [true, false] {
243            run_jobs_command(
244                JobsCommand::List {
245                    kind: None,
246                    status: None,
247                    limit: 50,
248                    offset: 0,
249                    json,
250                },
251                true,
252                palette(),
253                &mut &b""[..],
254                db.clone(),
255            )
256            .await
257            .unwrap();
258        }
259    }
260
261    #[tokio::test]
262    async fn an_unknown_status_is_refused_by_name() {
263        let db = db().await;
264        let err = run_jobs_command(
265            JobsCommand::List {
266                kind: None,
267                status: Some("halfway".into()),
268                limit: 50,
269                offset: 0,
270                json: false,
271            },
272            true,
273            palette(),
274            &mut &b""[..],
275            db,
276        )
277        .await
278        .unwrap_err();
279        assert!(err.message.contains("--status"), "{}", err.message);
280        assert!(err.message.contains("halfway"), "{}", err.message);
281        assert!(
282            err.message
283                .contains("ready, running, done, failed, cancelled"),
284            "{}",
285            err.message
286        );
287        assert_eq!(err.kind(), CliErrorKind::BadRequest);
288    }
289
290    #[tokio::test]
291    async fn every_job_status_is_accepted_as_a_filter() {
292        let db = db().await;
293        for status in JobStatus::ALL {
294            run_jobs_command(
295                JobsCommand::List {
296                    kind: None,
297                    status: Some(status.as_str().to_string()),
298                    limit: 50,
299                    offset: 0,
300                    json: false,
301                },
302                true,
303                palette(),
304                &mut &b""[..],
305                db.clone(),
306            )
307            .await
308            .unwrap();
309        }
310    }
311
312    #[tokio::test]
313    async fn cancel_is_confirm_gated() {
314        let db = db().await;
315        let id = seed("nonce_sweep", "a", now_secs(), &db).await;
316        run_jobs_command(
317            JobsCommand::Cancel { id: id.to_string() },
318            false,
319            palette(),
320            &mut b"n\n".as_slice(),
321            db.clone(),
322        )
323        .await
324        .unwrap();
325        assert_eq!(
326            Job::find_by_id(id, &db).await.unwrap().unwrap().status,
327            "ready"
328        );
329    }
330
331    #[tokio::test]
332    async fn cancel_a_ready_job_and_refuse_a_running_one() {
333        let db = db().await;
334        let id = seed("nonce_sweep", "a", now_secs(), &db).await;
335        run_jobs_command(
336            JobsCommand::Cancel { id: id.to_string() },
337            true,
338            palette(),
339            &mut &b""[..],
340            db.clone(),
341        )
342        .await
343        .unwrap();
344        assert_eq!(
345            Job::find_by_id(id, &db).await.unwrap().unwrap().status,
346            "cancelled"
347        );
348
349        let running = seed("nonce_sweep", "b", now_secs(), &db).await;
350        sqlx::query("UPDATE jobs SET status = 'running' WHERE id = ?;")
351            .bind(running)
352            .execute(db.raw_pool())
353            .await
354            .unwrap();
355        let err = run_jobs_command(
356            JobsCommand::Cancel {
357                id: running.to_string(),
358            },
359            true,
360            palette(),
361            &mut &b""[..],
362            db,
363        )
364        .await
365        .unwrap_err();
366        assert!(err.message.contains("running"), "{}", err.message);
367        assert_eq!(err.kind(), CliErrorKind::BadRequest);
368    }
369
370    #[tokio::test]
371    async fn run_now_nudges_and_revives() {
372        let db = db().await;
373        let ready = seed("nonce_sweep", "a", now_secs() + 3_600, &db).await;
374        run_jobs_command(
375            JobsCommand::RunNow {
376                id: ready.to_string(),
377            },
378            true,
379            palette(),
380            &mut &b""[..],
381            db.clone(),
382        )
383        .await
384        .unwrap();
385        assert!(Job::find_by_id(ready, &db).await.unwrap().unwrap().run_at <= now_secs() + 1);
386
387        let failed = seed("nonce_sweep", "b", now_secs(), &db).await;
388        sqlx::query("UPDATE jobs SET status = 'failed', attempts = 5 WHERE id = ?;")
389            .bind(failed)
390            .execute(db.raw_pool())
391            .await
392            .unwrap();
393        run_jobs_command(
394            JobsCommand::RunNow {
395                id: failed.to_string(),
396            },
397            true,
398            palette(),
399            &mut &b""[..],
400            db.clone(),
401        )
402        .await
403        .unwrap();
404        let job = Job::find_by_id(failed, &db).await.unwrap().unwrap();
405        assert_eq!(job.status, "ready");
406        assert_eq!(job.attempts, 4);
407    }
408
409    #[tokio::test]
410    async fn run_now_on_a_done_job_is_refused() {
411        let db = db().await;
412        let id = seed("nonce_sweep", "a", now_secs(), &db).await;
413        sqlx::query("UPDATE jobs SET status = 'done' WHERE id = ?;")
414            .bind(id)
415            .execute(db.raw_pool())
416            .await
417            .unwrap();
418        let err = run_jobs_command(
419            JobsCommand::RunNow { id: id.to_string() },
420            true,
421            palette(),
422            &mut &b""[..],
423            db,
424        )
425        .await
426        .unwrap_err();
427        assert!(err.message.contains("done"), "{}", err.message);
428    }
429
430    #[tokio::test]
431    async fn the_window_clamps_a_nonsense_one() {
432        let db = db().await;
433        seed("nonce_sweep", "a", now_secs(), &db).await;
434        run_jobs_command(
435            JobsCommand::List {
436                kind: None,
437                status: None,
438                limit: 0,
439                offset: -5,
440                json: false,
441            },
442            true,
443            palette(),
444            &mut &b""[..],
445            db,
446        )
447        .await
448        .unwrap();
449    }
450}