use doido_jobs::{JobContext, JobRegistry, WorkerEngine};
use std::sync::Arc;
pub async fn run(once: bool) {
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;
}
};
doido_core::tracing::info!(
"starting background worker (backend={:?}, queues={:?}, concurrency={}, once={once})",
config.backend,
config.queues,
config.concurrency,
);
let mut job_ctx = JobContext::new();
if let Some(conn) = doido_model::pool::try_pool() {
job_ctx.insert(conn.clone());
}
let engine = WorkerEngine::with_context(queue, config.engine_config(), job_ctx);
let registry = Arc::new(JobRegistry::from_inventory());
doido_core::tracing::info!("registered jobs: {:?}", registry.names());
let handler = move |job, ctx| {
let registry = Arc::clone(®istry);
async move { registry.dispatch(job, ctx).await }
};
if once {
loop {
match engine.run_once(&handler).await {
Ok(true) => continue,
Ok(false) => break,
Err(e) => {
doido_core::tracing::error!("worker engine error: {e}");
break;
}
}
}
doido_core::tracing::info!("worker drained ready jobs, exiting (once)");
return;
}
let shutdown = async {
let _ = tokio::signal::ctrl_c().await;
doido_core::tracing::info!("shutdown signal received, draining in-flight jobs...");
};
if let Err(e) = engine.run(handler, shutdown).await {
doido_core::tracing::error!("worker engine error: {e}");
}
}