use std::future::Future;
use crate::{
controller::{ControllerError, ControllerSpec},
core::deferred_drop::{
DropAdmissionError, DropCapacityError, DropReservation, DropStartError, OWNERSHIP_RESOURCE,
OwnedTask,
},
events::Event,
};
use super::ControllerHandle;
impl ControllerHandle {
async fn reserve_or_closed(
&self,
reservation: impl Future<Output = Result<DropReservation, DropAdmissionError>>,
) -> Result<DropReservation, ControllerError> {
tokio::pin!(reservation);
tokio::select! {
biased;
_ = self.tx.closed() => Err(ControllerError::Closed),
result = &mut reservation => result.map_err(Self::admission_error),
}
}
#[cfg(test)]
async fn reserve_capacity_or_closed(
&self,
reservation: impl Future<Output = Result<DropReservation, DropCapacityError>>,
) -> Result<DropReservation, ControllerError> {
tokio::pin!(reservation);
tokio::select! {
biased;
_ = self.tx.closed() => Err(ControllerError::Closed),
result = &mut reservation => result.map_err(Self::capacity_error),
}
}
fn capacity_error(error: DropCapacityError) -> ControllerError {
match error.limit() {
Some(limit) => ControllerError::ResourceLimit {
resource: OWNERSHIP_RESOURCE,
limit: limit.get(),
},
None => ControllerError::Closed,
}
}
fn start_error(error: DropStartError) -> ControllerError {
ControllerError::ThreadStartFailed {
component: "destructor_isolation",
worker: error.worker(),
kind: error.source_kind(),
raw_os_error: error.raw_os_error(),
}
}
fn admission_error(error: DropAdmissionError) -> ControllerError {
match error {
DropAdmissionError::Start(error) => Self::start_error(error),
DropAdmissionError::Capacity(error) => Self::capacity_error(error),
}
}
pub(super) async fn own(
&self,
spec: ControllerSpec,
) -> Result<OwnedTask<ControllerSpec>, ControllerError> {
#[cfg(test)]
let reservation = match &self.reservation_source {
Some(source) => self.reserve_capacity_or_closed(source.reserve()).await?,
None => self.reserve_or_closed(self.drop_domain.reserve()).await?,
};
#[cfg(not(test))]
let reservation = self.reserve_or_closed(self.drop_domain.reserve()).await?;
let retained = spec.task_spec().task().clone();
let mut owned = OwnedTask::new(spec, retained, reservation);
let bus = self.bus.clone();
owned.cleanup.set_panic_reporter(move |message| {
bus.publish_lazy(|| {
Event::runtime_failure("controller", format!("task_drop_panicked: {message}"))
});
});
Ok(owned)
}
pub(super) fn try_own(
&self,
spec: ControllerSpec,
) -> Result<OwnedTask<ControllerSpec>, ControllerError> {
#[cfg(test)]
let reservation = match &self.reservation_source {
Some(source) => source.try_reserve().map_err(Self::capacity_error),
None => self
.drop_domain
.try_reserve()
.map_err(Self::admission_error),
};
#[cfg(not(test))]
let reservation = self
.drop_domain
.try_reserve()
.map_err(Self::admission_error);
let retained = spec.task_spec().task().clone();
let mut owned = OwnedTask::new(spec, retained, reservation?);
let bus = self.bus.clone();
owned.cleanup.set_panic_reporter(move |message| {
bus.publish_lazy(|| {
Event::runtime_failure("controller", format!("task_drop_panicked: {message}"))
});
});
Ok(owned)
}
}