use std::any::Any;
use std::panic::{self, AssertUnwindSafe};
use std::ptr::NonNull;
use super::latch::Latch;
use crate::deque::Element;
use crate::sync::UnsafeCell;
#[repr(C)]
pub(super) struct JobHeader {
execute: unsafe fn(NonNull<JobHeader>),
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub(super) struct JobRef(NonNull<JobHeader>);
unsafe impl Send for JobRef {}
unsafe impl Element for JobRef {
fn into_raw(self) -> NonNull<()> {
self.0.cast()
}
unsafe fn from_raw(ptr: NonNull<()>) -> Self {
Self(ptr.cast())
}
}
impl JobRef {
pub(super) unsafe fn execute(self) {
unsafe { (self.0.as_ref().execute)(self.0) }
}
}
#[repr(C)]
pub(super) struct HeapJob<F> {
header: JobHeader,
func: F,
}
impl<F: FnOnce() + Send + 'static> HeapJob<F> {
pub(super) fn new_ref(func: F) -> JobRef {
let job = Box::new(Self {
header: JobHeader {
execute: Self::execute,
},
func,
});
JobRef(NonNull::from(Box::leak(job)).cast())
}
unsafe fn execute(this: NonNull<JobHeader>) {
let job = unsafe { Box::from_raw(this.cast::<Self>().as_ptr()) };
if panic::catch_unwind(AssertUnwindSafe(job.func)).is_err() {
eprintln!("parkring: a spawned job panicked; aborting");
std::process::abort();
}
}
}
pub(super) enum JobResult<R> {
None,
Ok(R),
Panic(Box<dyn Any + Send>),
}
#[repr(C)]
pub(super) struct StackJob<F, R, L> {
header: JobHeader,
func: UnsafeCell<Option<F>>,
result: UnsafeCell<JobResult<R>>,
pub(super) latch: L,
}
impl<F, R, L> StackJob<F, R, L>
where
F: FnOnce() -> R + Send,
R: Send,
L: Latch,
{
pub(super) fn new(func: F, latch: L) -> Self {
Self {
header: JobHeader {
execute: Self::execute,
},
func: UnsafeCell::new(Some(func)),
result: UnsafeCell::new(JobResult::None),
latch,
}
}
pub(super) unsafe fn as_job_ref(&self) -> JobRef {
JobRef(NonNull::from(self).cast())
}
unsafe fn execute(this: NonNull<JobHeader>) {
let job: *const Self = this.cast::<Self>().as_ptr();
let func = unsafe { (*job).func.with_mut(|f| (*f).take()) }.expect("job ran twice");
let result = match panic::catch_unwind(AssertUnwindSafe(func)) {
Ok(r) => JobResult::Ok(r),
Err(payload) => JobResult::Panic(payload),
};
unsafe { (*job).result.with_mut(|r| *r = result) };
unsafe { L::set(&raw const (*job).latch) };
}
pub(super) unsafe fn run_inline(&self) -> std::thread::Result<R> {
let func = self
.func
.with_mut(|f| unsafe { (*f).take() })
.expect("job ran twice");
panic::catch_unwind(AssertUnwindSafe(func))
}
pub(super) fn into_result(self) -> std::thread::Result<R> {
let result = self.result.with_mut(|r| {
unsafe { std::mem::replace(&mut *r, JobResult::None) }
});
match result {
JobResult::Ok(r) => Ok(r),
JobResult::Panic(payload) => Err(payload),
JobResult::None => unreachable!("latch set without a result"),
}
}
}