use std::sync::Arc;
use crate::{
core::{OutcomeTx, TaskOutcome, deferred_drop::OwnedTask},
events::{Event, EventKind, RejectionKind},
identity::TaskId,
};
use super::super::{Controller, state::PendingSubmission};
pub(super) struct AdmissionWatcher<'a> {
controller: &'a Controller,
id: TaskId,
event_task: Option<Arc<str>>,
state: AdmissionWatcherState,
owned: Option<OwnedTask<crate::ControllerSpec>>,
}
enum AdmissionWatcherState {
Local(Option<OutcomeTx>),
Parked,
Committed,
}
impl<'a> AdmissionWatcher<'a> {
pub(super) fn new(
controller: &'a Controller,
id: TaskId,
owned: OwnedTask<crate::ControllerSpec>,
done: Option<OutcomeTx>,
event_task: Option<Arc<str>>,
) -> Self {
Self {
controller,
id,
event_task,
state: AdmissionWatcherState::Local(done),
owned: Some(owned),
}
}
pub(super) fn take_pending(&mut self, id: TaskId, task_name: Arc<str>) -> PendingSubmission {
let owned = self
.owned
.take()
.expect("controller ownership is transferred once")
.map(crate::ControllerSpec::into_task_spec);
PendingSubmission::new(id, task_name, owned)
}
fn dispose_owned(&mut self, terminal: Option<TaskOutcome>) {
let Some(owned) = self.owned.take() else {
return;
};
let (spec, mut cleanup) = owned.into_parts();
drop(spec);
if let Some(terminal) = terminal {
cleanup.attach_outcome(terminal);
}
cleanup.submit();
}
pub(super) fn park(&mut self) {
let state = std::mem::replace(&mut self.state, AdmissionWatcherState::Committed);
match state {
AdmissionWatcherState::Local(Some(tx)) => {
self.controller.state().watchers.insert(self.id, tx);
self.state = AdmissionWatcherState::Parked;
}
AdmissionWatcherState::Local(None) => {
self.state = AdmissionWatcherState::Parked;
}
AdmissionWatcherState::Committed => {}
AdmissionWatcherState::Parked => {
self.state = AdmissionWatcherState::Parked;
}
}
}
pub(super) fn commit(&mut self) {
debug_assert!(
!matches!(self.state, AdmissionWatcherState::Local(Some(_))),
"a watched admission must be parked before commit"
);
self.state = AdmissionWatcherState::Committed;
}
fn reject(&mut self, kind: RejectionKind, reason: &str) -> Option<TaskOutcome> {
let state = std::mem::replace(&mut self.state, AdmissionWatcherState::Committed);
let undelivered = match state {
AdmissionWatcherState::Local(Some(tx)) => tx
.send(TaskOutcome::Rejected {
kind,
reason: Arc::from(reason),
})
.err(),
AdmissionWatcherState::Parked => {
self.controller.finalize_rejected(self.id, kind, reason)
}
AdmissionWatcherState::Local(None) | AdmissionWatcherState::Committed => None,
};
if self.owned.is_some() {
self.dispose_owned(undelivered);
None
} else {
undelivered
}
}
pub(super) fn reject_with_event(
&mut self,
kind: RejectionKind,
reason: &str,
) -> Option<TaskOutcome> {
if matches!(self.state, AdmissionWatcherState::Committed) {
return None;
}
self.controller.bus.publish_lazy(|| {
let mut event = Event::new(EventKind::ControllerRejected)
.with_id(self.id)
.with_rejection_kind(kind)
.with_reason(reason);
if let Some(task) = &self.event_task {
event = event.with_task(Arc::clone(task));
}
event
});
self.reject(kind, reason)
}
}
impl Drop for AdmissionWatcher<'_> {
fn drop(&mut self) {
if matches!(self.state, AdmissionWatcherState::Committed) {
return;
}
drop(self.reject_with_event(
RejectionKind::AdmissionFailed,
crate::reasons::CONTROLLER_ADMISSION_INTERRUPTED,
));
}
}