use std::sync::{Arc, Mutex};
use boson_core::{ExecutionContextFactory, QueueBackend};
use super::claim::claim_next_job;
use super::config::WorkerSettings;
use super::loop_::WorkerEngine;
use crate::registry::TaskRegistry;
pub struct ManualWorker {
inner: Arc<WorkerEngine>,
in_flight: Mutex<bool>,
}
struct ClearInFlight<'a>(&'a Mutex<bool>);
impl Drop for ClearInFlight<'_> {
fn drop(&mut self) {
if let Ok(mut guard) = self.0.lock() {
*guard = false;
}
}
}
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,
}),
in_flight: Mutex::new(false),
}
}
pub async fn try_run_next(&self) -> bool {
{
let Ok(mut in_flight) = self.in_flight.lock() else {
return false;
};
if *in_flight {
return false;
}
*in_flight = true;
}
let _clear = ClearInFlight(&self.in_flight);
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
}
}