1use 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 {
33 #[arg(long)]
37 kind: Option<String>,
38 #[arg(long)]
41 status: Option<String>,
42 #[arg(long, default_value_t = DEFAULT_LIMIT)]
44 limit: i64,
45 #[arg(long, default_value_t = 0)]
47 offset: i64,
48 #[arg(long)]
50 json: bool,
51 },
52 Show {
54 id: String,
56 #[arg(long)]
58 json: bool,
59 },
60 Cancel {
63 id: String,
65 },
66 RunNow {
69 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 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 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}