use doido_jobs::{build_queue, JobContext, JobPayload, JobsConfig, WorkerEngine};
use std::sync::Arc;
pub async fn run(once: bool) {
let config = JobsConfig::default();
let queue = match build_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 engine = WorkerEngine::with_context(queue, config.engine_config(), JobContext::new());
let handler = |job: JobPayload, _ctx: Arc<JobContext>| async move {
doido_core::tracing::info!("processing job {} on queue {}", job.id, job.queue);
Ok(())
};
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}");
}
}