use clap::Subcommand;
use doido_jobs::{JobPayload, JobQueue, JobsConfig};
#[derive(Subcommand)]
pub enum JobsCommand {
Failed,
Retry,
Discard,
}
pub async fn run(cmd: JobsCommand) {
let config = doido_jobs::config::load();
let queue = match doido_jobs::config::build_configured_queue(&config).await {
Ok(q) => q,
Err(e) => {
doido_core::tracing::error!("failed to build jobs backend: {e}");
return;
}
};
match cmd {
JobsCommand::Failed => list_failed(queue.as_ref(), &config).await,
JobsCommand::Retry => retry_failed(queue.as_ref(), &config).await,
JobsCommand::Discard => discard_failed(queue.as_ref(), &config).await,
}
}
async fn list_failed(queue: &dyn JobQueue, config: &JobsConfig) {
let mut total = 0;
for q in &config.queues {
match queue.dead_jobs(q).await {
Ok(jobs) => {
for job in &jobs {
total += 1;
doido_core::tracing::info!(
"[{}] {} attempts={} error={}",
job.queue,
job.id,
job.attempts,
job.error.as_deref().unwrap_or("-"),
);
}
}
Err(e) => doido_core::tracing::error!("failed to read dead jobs for {q}: {e}"),
}
}
if total == 0 {
doido_core::tracing::info!("failed jobs: (none)");
} else {
doido_core::tracing::info!("{total} failed job(s)");
}
}
async fn retry_failed(queue: &dyn JobQueue, config: &JobsConfig) {
let mut retried = 0;
for q in &config.queues {
let jobs = match queue.dead_jobs(q).await {
Ok(j) => j,
Err(e) => {
doido_core::tracing::error!("failed to read dead jobs for {q}: {e}");
continue;
}
};
for job in jobs {
let fresh = JobPayload::new(job.queue.clone(), job.payload.clone(), job.max_retries)
.with_priority(job.priority)
.with_backoff(job.backoff, job.backoff_base)
.with_timeout(job.timeout);
match queue.enqueue(fresh).await {
Ok(_) => retried += 1,
Err(e) => doido_core::tracing::error!("re-enqueue failed: {e}"),
}
}
if let Err(e) = queue.discard_dead(q).await {
doido_core::tracing::error!("failed to clear dead store for {q}: {e}");
}
}
doido_core::tracing::info!("retried {retried} failed job(s)");
}
async fn discard_failed(queue: &dyn JobQueue, config: &JobsConfig) {
let mut discarded = 0;
for q in &config.queues {
match queue.discard_dead(q).await {
Ok(n) => discarded += n,
Err(e) => doido_core::tracing::error!("failed to discard dead jobs for {q}: {e}"),
}
}
doido_core::tracing::info!("discarded {discarded} failed job(s)");
}