use std::sync::Arc;
use std::time::Duration;
use boson_core::{BosonError, ExecutionContextFactory, Job, JobStatus, QueueBackend, Result};
use tokio::time::sleep;
use crate::registry::TaskRegistry;
pub async fn execute_job(
registry: &TaskRegistry,
identity: &Arc<dyn ExecutionContextFactory>,
backend: &Arc<dyn QueueBackend>,
job: &Job,
) -> Result<()> {
let descriptor = registry.get_or_err(&job.task_name)?;
if job.signature_hash != descriptor.signature_hash {
return Err(BosonError::SignatureMismatch {
expected: job.signature_hash.to_string(),
actual: descriptor.signature_hash.to_string(),
});
}
let ctx = identity
.build(&job.actor_json)
.map_err(|e| BosonError::internal_source("execution context build failed", e))?;
let invoke = (descriptor.invoke)(ctx, job.params_json.clone());
let job_id = job.job_id.clone();
let backend = Arc::clone(backend);
let cancel_watch = async move {
loop {
sleep(Duration::from_millis(50)).await;
match backend.get_job(&job_id).await {
Ok(Some(j)) if j.status == JobStatus::Canceled => return,
Ok(None) | Err(_) => return,
_ => {}
}
}
};
tokio::select! {
result = invoke => result,
() = cancel_watch => Err(BosonError::internal(
"job canceled during execution",
)),
}
}
pub async fn record_run_start(
backend: &Arc<dyn QueueBackend>,
run: &boson_core::Run,
) -> Result<()> {
backend.upsert_run(run).await
}