use std::sync::Arc;
use boson_core::ExecutionContextFactory;
use boson_core::QueueBackend;
use tokio::sync::Mutex;
use super::claim::claim_next_job;
use super::config::WorkerSettings;
use super::loop_::WorkerEngine;
use crate::registry::TaskRegistry;
pub struct ManualWorker {
inner: Arc<WorkerEngine>,
lock: Mutex<()>,
}
impl ManualWorker {
pub fn new(
backend: Arc<dyn QueueBackend>,
registry: Arc<TaskRegistry>,
identity: Arc<dyn ExecutionContextFactory>,
worker: WorkerSettings,
) -> Self {
Self {
inner: Arc::new(WorkerEngine {
backend,
registry,
identity,
worker,
}),
lock: Mutex::new(()),
}
}
pub async fn try_run_next(&self) -> bool {
let _guard = self.lock.lock().await;
let discovered = self
.inner
.backend
.distinct_pools_queued()
.await
.unwrap_or_default();
let pools = self.inner.worker.pools_to_poll(discovered);
for pool in pools {
if let Ok(Some((job, lease_id))) = claim_next_job(
&self.inner.backend,
&pool,
&self.inner.worker.worker_id,
self.inner.worker.lease_ttl_secs,
)
.await
{
self.inner.drive_run(job, lease_id).await;
return true;
}
}
false
}
}