use std::{
any::Any,
sync::{Arc, Mutex},
};
use crate::{core::outcome::TaskOutcome, tasks::TaskRef};
use super::{
capacity::OwnershipPermit,
executor::{DropBatch, DropExecutor},
};
pub(super) type DropJob = Box<dyn FnOnce() + Send + 'static>;
pub(super) type PanicReporter = Box<dyn FnOnce(String) + Send + 'static>;
pub(crate) struct DropReservation {
pub(super) executor: Arc<DropExecutor>,
pub(super) permit: Option<OwnershipPermit>,
}
impl DropReservation {
pub(super) fn new(executor: Arc<DropExecutor>, permit: OwnershipPermit) -> Self {
Self {
executor,
permit: Some(permit),
}
}
pub(crate) fn bundle<T>(self, retained: T) -> DropBundle
where
T: Send + 'static,
{
DropBundle::new(self, Box::new(move || drop(retained)))
}
fn submit(
mut self,
retained: DropJob,
undelivered_outcome: Option<DropJob>,
auxiliary: Option<DropJob>,
panic_reporter: Option<PanicReporter>,
poisoned: bool,
) {
let permit = self
.permit
.take()
.expect("one ownership reservation submits at most one bundle");
self.executor.submit(DropBatch::new(
permit,
retained,
undelivered_outcome,
auxiliary,
panic_reporter,
poisoned,
));
}
}
struct DropBundleInner {
reservation: DropReservation,
retained: DropJob,
undelivered_outcome: Option<DropJob>,
auxiliary: Option<DropJob>,
panic_reporter: Option<PanicReporter>,
poisoned: bool,
}
pub(crate) struct DropBundle {
inner: Mutex<Option<DropBundleInner>>,
}
impl DropBundle {
fn new(reservation: DropReservation, retained: DropJob) -> Self {
Self {
inner: Mutex::new(Some(DropBundleInner {
reservation,
retained,
undelivered_outcome: None,
auxiliary: None,
panic_reporter: None,
poisoned: false,
})),
}
}
fn inner_mut(&mut self) -> Option<&mut DropBundleInner> {
self.inner
.get_mut()
.unwrap_or_else(|error| error.into_inner())
.as_mut()
}
pub(crate) fn attach_outcome(&mut self, outcome: TaskOutcome) {
let Some(inner) = self.inner_mut() else {
std::mem::forget(outcome);
return;
};
if inner.undelivered_outcome.is_some() {
inner.poisoned = true;
std::mem::forget(outcome);
return;
}
inner.undelivered_outcome = Some(Box::new(move || drop(outcome)));
}
pub(crate) fn attach_panic_payload(&mut self, payload: Box<dyn Any + Send>) {
self.attach_auxiliary(payload);
}
pub(crate) fn set_panic_reporter<F>(&mut self, reporter: F)
where
F: FnOnce(String) + Send + 'static,
{
let reporter: PanicReporter = Box::new(reporter);
let Some(inner) = self.inner_mut() else {
std::mem::forget(reporter);
return;
};
if inner.panic_reporter.is_some() {
inner.poisoned = true;
std::mem::forget(reporter);
return;
}
inner.panic_reporter = Some(reporter);
}
pub(crate) fn attach_physical<T>(&mut self, value: T)
where
T: Send + 'static,
{
self.attach_auxiliary(value);
}
fn attach_auxiliary<T>(&mut self, value: T)
where
T: Send + 'static,
{
let Some(inner) = self.inner_mut() else {
std::mem::forget(value);
return;
};
if inner.auxiliary.is_some() {
inner.poisoned = true;
std::mem::forget(value);
return;
}
inner.auxiliary = Some(Box::new(move || drop(value)));
}
pub(crate) fn poison(&mut self) {
if let Some(inner) = self.inner_mut() {
inner.poisoned = true;
}
}
pub(crate) fn submit(self) {
self.submit_inner();
}
fn submit_inner(&self) {
let Some(inner) = self
.inner
.lock()
.unwrap_or_else(|error| error.into_inner())
.take()
else {
return;
};
inner.reservation.submit(
inner.retained,
inner.undelivered_outcome,
inner.auxiliary,
inner.panic_reporter,
inner.poisoned,
);
}
}
impl Drop for DropBundle {
fn drop(&mut self) {
self.submit_inner();
}
}
pub(crate) struct OwnedTask<T> {
pub(crate) value: T,
pub(crate) cleanup: DropBundle,
}
impl<T> OwnedTask<T> {
pub(crate) fn new(value: T, retained: TaskRef, reservation: DropReservation) -> Self {
Self {
value,
cleanup: reservation.bundle(retained),
}
}
#[cfg(feature = "controller")]
pub(crate) fn map<U>(self, map: impl FnOnce(T) -> U) -> OwnedTask<U> {
let Self { value, cleanup } = self;
OwnedTask {
value: map(value),
cleanup,
}
}
pub(crate) fn into_parts(self) -> (T, DropBundle) {
let Self { value, cleanup } = self;
(value, cleanup)
}
}