use core::ffi::c_void;
use std::io;
use std::marker::PhantomData;
use std::os::windows::io::BorrowedHandle;
use std::ptr;
use std::sync::Mutex;
use std::time::{Duration, SystemTime};
use windows_sys::Win32::Foundation::{FALSE, TRUE};
use windows_sys::Win32::System::Threading::{
CloseThreadpoolCleanupGroup, CloseThreadpoolCleanupGroupMembers, CreateThreadpoolCleanupGroup,
IsThreadpoolTimerSet, PTP_CLEANUP_GROUP, PTP_TIMER, PTP_WAIT, PTP_WORK, SubmitThreadpoolWork,
WaitForThreadpoolTimerCallbacks, WaitForThreadpoolWaitCallbacks,
WaitForThreadpoolWorkCallbacks,
};
use crate::callback_env::CallbackEnviron;
use crate::timer::{
PeriodicTick, ThreadpoolPeriodicTimer, ThreadpoolTimer, TimerFiring, absolute_filetime,
arm_raw, disarm_raw, millis_u32, relative_filetime,
};
use crate::wait::{ThreadpoolWait, WaitActivation, WaitTarget, WaitableHandle};
use crate::work::ThreadpoolWork;
struct OwnedResource {
ptr: *mut c_void,
prepare_shutdown: unsafe fn(*mut c_void),
free: unsafe fn(*mut c_void),
}
unsafe impl Send for OwnedResource {}
unsafe fn free_boxed<T>(ptr: *mut c_void) {
drop(unsafe { Box::from_raw(ptr.cast::<T>()) });
}
fn prepare_shutdown_noop(_ptr: *mut c_void) {}
pub struct CleanupGroup {
group: PTP_CLEANUP_GROUP,
resources: Mutex<Vec<OwnedResource>>,
}
unsafe impl Send for CleanupGroup {}
unsafe impl Sync for CleanupGroup {}
impl CleanupGroup {
pub fn new() -> io::Result<Self> {
let group = unsafe { CreateThreadpoolCleanupGroup() };
if group == 0 {
return Err(io::Error::last_os_error());
}
Ok(Self {
group,
resources: Mutex::new(Vec::new()),
})
}
fn member_environment(&self, env: Option<&CallbackEnviron<'_>>) -> CallbackEnviron<'_> {
let mut member_env = match env {
Some(env) => CallbackEnviron::from_inner(*env.as_inner()),
None => CallbackEnviron::new(),
};
unsafe { member_env.set_cleanup_group(self.group, None) };
member_env
}
fn adopt(&self, resource: OwnedResource) {
self.resources
.lock()
.unwrap_or_else(|poison| poison.into_inner())
.push(resource);
}
pub fn create_work<F>(
&self,
callback: F,
env: Option<&CallbackEnviron>,
) -> io::Result<WorkMember<'_>>
where
F: Fn() + Send + Sync + 'static,
{
let mut member_env = self.member_environment(env);
let work = ThreadpoolWork::new(callback, Some(&mut member_env))?;
let (handle, context) = work.into_parts();
self.adopt(OwnedResource {
ptr: context,
prepare_shutdown: prepare_shutdown_noop,
free: ThreadpoolWork::drop_context,
});
Ok(WorkMember {
handle,
_group: PhantomData,
})
}
pub fn create_timer<F>(
&self,
callback: F,
env: Option<&CallbackEnviron>,
) -> io::Result<TimerMember<'_>>
where
F: Fn(&TimerFiring<'_>) + Send + Sync + 'static,
{
let mut member_env = self.member_environment(env);
let timer = ThreadpoolTimer::new(callback, Some(&mut member_env))?;
let (handle, context) = timer.into_parts();
self.adopt(OwnedResource {
ptr: context,
prepare_shutdown: ThreadpoolTimer::prepare_shutdown,
free: ThreadpoolTimer::drop_context,
});
Ok(TimerMember {
handle,
_group: PhantomData,
})
}
pub fn create_periodic_timer<F>(
&self,
period: Duration,
callback: F,
env: Option<&CallbackEnviron>,
) -> io::Result<PeriodicTimerMember<'_>>
where
F: Fn(&PeriodicTick<'_>) + Send + Sync + 'static,
{
let mut member_env = self.member_environment(env);
let timer = ThreadpoolPeriodicTimer::new(period, callback, Some(&mut member_env))?;
let (handle, context, period) = timer.into_parts();
self.adopt(OwnedResource {
ptr: context,
prepare_shutdown: prepare_shutdown_noop,
free: ThreadpoolPeriodicTimer::drop_context,
});
Ok(PeriodicTimerMember {
handle,
period,
_group: PhantomData,
})
}
pub fn create_wait<F>(
&self,
handle: WaitableHandle,
callback: F,
env: Option<&CallbackEnviron<'_>>,
) -> io::Result<WaitMember<'_>>
where
F: Fn(&WaitActivation<'_>) + Send + Sync + 'static,
{
let mut member_env = self.member_environment(env);
let wait = ThreadpoolWait::new(handle, callback, Some(&mut member_env))?;
let (raw, context, target) = wait.into_parts();
self.adopt(OwnedResource {
ptr: context,
prepare_shutdown: ThreadpoolWait::prepare_shutdown,
free: ThreadpoolWait::drop_context,
});
let target = Box::into_raw(Box::new(target));
self.adopt(OwnedResource {
ptr: target.cast(),
prepare_shutdown: prepare_shutdown_noop,
free: free_boxed::<WaitTarget>,
});
Ok(WaitMember {
handle: raw,
watched: target,
_group: PhantomData,
})
}
pub fn close_members(&mut self, cancel_pending: bool) {
self.release_members(cancel_pending);
}
#[must_use]
pub fn owned_resources(&self) -> usize {
self.resources
.lock()
.unwrap_or_else(|poison| poison.into_inner())
.len()
}
fn release_members(&mut self, cancel_pending: bool) {
{
let resources = self
.resources
.lock()
.unwrap_or_else(|poison| poison.into_inner());
for resource in resources.iter() {
unsafe { (resource.prepare_shutdown)(resource.ptr) };
}
}
unsafe {
CloseThreadpoolCleanupGroupMembers(
self.group,
if cancel_pending { TRUE } else { FALSE },
ptr::null_mut(),
);
}
let resources = std::mem::take(
&mut *self
.resources
.lock()
.unwrap_or_else(|poison| poison.into_inner()),
);
for resource in resources {
unsafe { (resource.free)(resource.ptr) };
}
}
}
impl Drop for CleanupGroup {
fn drop(&mut self) {
self.release_members(false);
unsafe { CloseThreadpoolCleanupGroup(self.group) };
}
}
impl std::fmt::Debug for CleanupGroup {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("CleanupGroup")
.field("owned_resources", &self.owned_resources())
.finish_non_exhaustive()
}
}
#[derive(Debug)]
pub struct WorkMember<'group> {
handle: PTP_WORK,
_group: PhantomData<&'group CleanupGroup>,
}
impl WorkMember<'_> {
pub fn submit(&self) {
unsafe { SubmitThreadpoolWork(self.handle) };
}
pub fn wait(&self) {
unsafe { WaitForThreadpoolWorkCallbacks(self.handle, FALSE) };
}
pub fn cancel_pending(&self) {
unsafe { WaitForThreadpoolWorkCallbacks(self.handle, TRUE) };
}
}
#[derive(Debug)]
pub struct TimerMember<'group> {
handle: PTP_TIMER,
_group: PhantomData<&'group CleanupGroup>,
}
impl TimerMember<'_> {
pub fn set_after(&self, delay: Duration) {
unsafe { arm_raw(self.handle, relative_filetime(delay), 0, 0) };
}
pub fn set_at(&self, when: SystemTime) {
unsafe { arm_raw(self.handle, absolute_filetime(when), 0, 0) };
}
pub fn disarm(&self) {
unsafe { disarm_raw(self.handle) };
}
#[must_use]
pub fn is_set(&self) -> bool {
unsafe { IsThreadpoolTimerSet(self.handle) != 0 }
}
pub fn wait(&self) {
unsafe { WaitForThreadpoolTimerCallbacks(self.handle, FALSE) };
}
pub fn cancel_pending(&self) {
unsafe { WaitForThreadpoolTimerCallbacks(self.handle, TRUE) };
}
}
#[derive(Debug)]
pub struct PeriodicTimerMember<'group> {
handle: PTP_TIMER,
period: Duration,
_group: PhantomData<&'group CleanupGroup>,
}
impl PeriodicTimerMember<'_> {
#[must_use]
pub fn period(&self) -> Duration {
self.period
}
pub fn start(&self) {
self.start_after(self.period);
}
pub fn start_after(&self, first_delay: Duration) {
unsafe {
arm_raw(
self.handle,
relative_filetime(first_delay),
millis_u32(self.period),
0,
);
}
}
pub fn stop(&self) {
unsafe { disarm_raw(self.handle) };
}
#[must_use]
pub fn is_running(&self) -> bool {
unsafe { IsThreadpoolTimerSet(self.handle) != 0 }
}
pub fn wait(&self) {
unsafe { WaitForThreadpoolTimerCallbacks(self.handle, FALSE) };
}
pub fn stop_and_drain(&self) {
self.stop();
unsafe { WaitForThreadpoolTimerCallbacks(self.handle, TRUE) };
}
}
#[derive(Debug)]
pub struct WaitMember<'group> {
handle: PTP_WAIT,
watched: *mut WaitTarget,
_group: PhantomData<&'group CleanupGroup>,
}
unsafe impl Send for WaitMember<'_> {}
unsafe impl Sync for WaitMember<'_> {}
impl WaitMember<'_> {
#[must_use]
pub fn handle(&self) -> BorrowedHandle<'_> {
unsafe { (*self.watched).borrow() }
}
pub fn arm(&self, timeout: Option<Duration>) {
unsafe { crate::wait::arm_member(self.handle, &*self.watched, timeout) };
}
pub fn disarm(&self) {
unsafe { crate::wait::disarm_raw(self.handle) };
}
pub fn wait(&self) {
unsafe { WaitForThreadpoolWaitCallbacks(self.handle, FALSE) };
}
pub fn cancel_pending(&self) {
unsafe { WaitForThreadpoolWaitCallbacks(self.handle, TRUE) };
}
}
#[cfg(test)]
mod tests;