use std::sync::{Arc, atomic::AtomicBool};
use tokio::sync::watch;
use super::{
Registry,
completion::{OutcomeTx, RemovalCompletion},
scheduler::{ActorHandle, ActorRegistration, ScheduledActor},
};
use crate::{
core::{
actor::{ActorExitReason, TaskActor, TaskActorParams, TaskActorResources},
deferred_drop::{DropBundle, OwnedTask},
outcome::TaskOutcome,
},
identity::TaskId,
tasks::TaskSpec,
};
mod batch;
mod single;
struct PreparedRegistration {
id: TaskId,
label: Arc<str>,
join: ActorHandle,
cancel: tokio_util::sync::CancellationToken,
done: Option<OutcomeTx>,
completion: RemovalCompletion,
scheduled: ScheduledActor,
cleanup: DropBundle,
activity: Arc<AtomicBool>,
}
fn deliver_or_attach_rejection(
done: Option<OutcomeTx>,
outcome: TaskOutcome,
cleanup: &mut DropBundle,
) {
match done {
Some(done) => {
if let Err(undelivered) = done.send(outcome) {
cleanup.attach_outcome(undelivered);
}
}
None => cleanup.attach_outcome(outcome),
}
}
impl Registry {
fn prepare_registration(
&self,
id: TaskId,
label: Arc<str>,
owned: OwnedTask<TaskSpec>,
done: Option<OutcomeTx>,
completion: Option<RemovalCompletion>,
mut start: watch::Receiver<bool>,
) -> PreparedRegistration {
let task_token = self.runtime_token.child_token();
let (spec, cleanup) = owned.into_parts();
let spec = spec.resolve(&self.task_defaults);
let task = spec.task().clone();
let activity = Arc::new(AtomicBool::new(false));
let cleanup_poisoned = Arc::new(AtomicBool::new(false));
let completion = completion.unwrap_or_else(RemovalCompletion::new);
let actor = TaskActor::new(
self.bus.clone(),
Arc::clone(&label),
task,
TaskActorParams {
restart: spec.restart(),
backoff: spec.backoff(),
timeout: spec.timeout(),
max_retries: spec.max_retries(),
},
TaskActorResources {
semaphore: self.semaphore.clone(),
activity: Arc::clone(&activity),
cleanup_poisoned: Arc::clone(&cleanup_poisoned),
},
id,
);
let task_token_clone = task_token.clone();
let actor_future = async move {
loop {
if *start.borrow_and_update() {
break;
}
if start.changed().await.is_err() {
return ActorExitReason::Canceled;
}
}
actor.run(task_token_clone).await
};
let (scheduled, join) = ScheduledActor::new(
ActorRegistration {
id,
label: Arc::clone(&label),
activity: Arc::clone(&activity),
cleanup_poisoned,
physical_release: completion.clone(),
reaper: self.actors.attempt_reaper(),
completion_tx: self.listener.completion_tx.clone(),
},
actor_future,
);
PreparedRegistration {
id,
label,
join,
cancel: task_token,
done,
completion,
scheduled,
cleanup,
activity,
}
}
fn registered_limit_exceeded(&self, current: usize, incoming: usize) -> Option<usize> {
let limit = self.max_registered_tasks?.get();
let exceeds = current
.checked_add(self.actors.reaping_attempts())
.and_then(|used| used.checked_add(incoming))
.is_none_or(|total| total > limit);
exceeds.then_some(limit)
}
}