#![cfg(windows)]
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Condvar, Mutex};
use std::time::Duration;
use windows_sys::Win32::Foundation::{CloseHandle, HANDLE};
use windows_sys::Win32::System::Threading::{CreateEventW, SetEvent};
use windows_threadpool_sys::cleanup_group::CleanupGroup;
use windows_threadpool_sys::wait::{ThreadpoolWait, WaitCloseFn, WaitableHandle};
const WAITS: usize = 256;
const CALLBACK_DWELL: Duration = Duration::from_millis(50);
const START_TIMEOUT: Duration = Duration::from_secs(60);
struct Probe {
closed: &'static Mutex<Vec<usize>>,
violations: &'static AtomicUsize,
started: &'static (Mutex<usize>, Condvar),
}
impl Probe {
fn closed_handles(&self) -> Vec<usize> {
self.closed
.lock()
.unwrap_or_else(|poison| poison.into_inner())
.clone()
}
fn wait_for_started(&self, target: usize) {
let (count, condvar) = self.started;
let count = count.lock().unwrap_or_else(|poison| poison.into_inner());
let (count, timeout) = condvar
.wait_timeout_while(count, START_TIMEOUT, |count| *count < target)
.unwrap_or_else(|poison| poison.into_inner());
assert!(
!timeout.timed_out(),
"timed out waiting for {target} callback(s) to start; saw {count}"
);
}
fn assert_closed_each_exactly_once(&self, expected: usize, path: &str) {
let mut closed = self.closed_handles();
assert_eq!(
closed.len(),
expected,
"{path}: every custom target must be closed exactly once"
);
closed.sort_unstable();
closed.dedup();
assert_eq!(
closed.len(),
expected,
"{path}: a handle was closed more than once"
);
assert_eq!(
self.violations.load(Ordering::SeqCst),
0,
"{path}: a handle was closed while its own callback was executing"
);
}
}
macro_rules! probe {
() => {{
static CLOSED: Mutex<Vec<usize>> = Mutex::new(Vec::new());
static VIOLATIONS: AtomicUsize = AtomicUsize::new(0);
static STARTED: (Mutex<usize>, Condvar) = (Mutex::new(0), Condvar::new());
unsafe extern "system" fn close(handle: HANDLE) -> i32 {
CLOSED
.lock()
.unwrap_or_else(|poison| poison.into_inner())
.push(handle as usize);
unsafe { CloseHandle(handle) }
}
let probe = Probe {
closed: &CLOSED,
violations: &VIOLATIONS,
started: &STARTED,
};
let close: WaitCloseFn = close;
(probe, close, &CLOSED, &VIOLATIONS, &STARTED)
}};
}
fn raw_event() -> HANDLE {
let raw = unsafe { CreateEventW(std::ptr::null(), 1, 0, std::ptr::null()) };
assert!(!raw.is_null(), "CreateEventW failed");
raw
}
unsafe fn signal(raw: HANDLE) {
let ok = unsafe { SetEvent(raw) };
assert_ne!(ok, 0, "SetEvent failed");
}
fn observe(
key: usize,
closed: &'static Mutex<Vec<usize>>,
violations: &'static AtomicUsize,
started: &'static (Mutex<usize>, Condvar),
) {
{
let (count, condvar) = started;
let mut count = count.lock().unwrap_or_else(|poison| poison.into_inner());
*count += 1;
condvar.notify_all();
}
std::thread::sleep(CALLBACK_DWELL);
let already_closed = closed
.lock()
.unwrap_or_else(|poison| poison.into_inner())
.contains(&key);
if already_closed {
violations.fetch_add(1, Ordering::SeqCst);
}
}
#[test]
fn dropping_owned_waits_closes_every_custom_target_exactly_once() {
let (probe, close, closed, violations, started) = probe!();
let mut waits = Vec::with_capacity(WAITS);
for _ in 0..WAITS {
let raw = raw_event();
let key = raw as usize;
let handle = unsafe { WaitableHandle::assume_waitable_with(raw, close) };
let wait = ThreadpoolWait::new(
handle,
move |_| observe(key, closed, violations, started),
None,
)
.expect("create wait");
wait.arm(None);
unsafe { signal(raw) };
waits.push(wait);
}
probe.wait_for_started(WAITS / 2);
assert!(
probe.closed_handles().is_empty(),
"nothing may be closed while the waits are alive"
);
let entered_teardown = std::time::Instant::now();
drop(waits);
let blocked_for = entered_teardown.elapsed();
assert!(
blocked_for >= CALLBACK_DWELL / 2,
"teardown returned in {blocked_for:?}, so it did not drain running callbacks"
);
probe.assert_closed_each_exactly_once(WAITS, "direct drop");
}
#[test]
fn releasing_a_group_closes_every_custom_target_exactly_once() {
let (probe, close, closed, violations, started) = probe!();
let mut group = CleanupGroup::new().expect("create group");
let mut raws = Vec::with_capacity(WAITS);
for _ in 0..WAITS {
let raw = raw_event();
let key = raw as usize;
let handle = unsafe { WaitableHandle::assume_waitable_with(raw, close) };
let member = group
.create_wait(
handle,
move |_| observe(key, closed, violations, started),
None,
)
.expect("create wait");
member.arm(None);
unsafe { signal(raw) };
raws.push(raw);
}
probe.wait_for_started(WAITS / 2);
assert!(
probe.closed_handles().is_empty(),
"nothing may be closed while the members are live"
);
assert_eq!(
group.owned_resources(),
WAITS * 2,
"each member parks a context and a target on the group"
);
let entered_teardown = std::time::Instant::now();
group.close_members(false);
let blocked_for = entered_teardown.elapsed();
assert!(
blocked_for >= CALLBACK_DWELL / 2,
"the release returned in {blocked_for:?}, so it did not drain running callbacks"
);
probe.assert_closed_each_exactly_once(WAITS, "group release");
assert_eq!(group.owned_resources(), 0, "the group holds nothing after");
group.close_members(false);
drop(group);
probe.assert_closed_each_exactly_once(WAITS, "group release, repeated");
}
#[test]
fn releasing_a_group_with_cancel_pending_closes_every_custom_target_exactly_once() {
let (probe, close, closed, violations, started) = probe!();
let mut group = CleanupGroup::new().expect("create group");
for _ in 0..WAITS {
let raw = raw_event();
let key = raw as usize;
let handle = unsafe { WaitableHandle::assume_waitable_with(raw, close) };
let member = group
.create_wait(
handle,
move |_| observe(key, closed, violations, started),
None,
)
.expect("create wait");
member.arm(None);
unsafe { signal(raw) };
}
probe.wait_for_started(1);
assert!(
probe.closed_handles().is_empty(),
"nothing may be closed while the members are live"
);
group.close_members(true);
probe.assert_closed_each_exactly_once(WAITS, "group release, cancel_pending");
}
#[test]
fn a_group_releases_custom_and_default_targets_together() {
let (probe, close, closed, violations, started) = probe!();
let mut group = CleanupGroup::new().expect("create group");
for index in 0..WAITS {
if index % 2 == 0 {
let raw = raw_event();
let key = raw as usize;
let handle = unsafe { WaitableHandle::assume_waitable_with(raw, close) };
let member = group
.create_wait(
handle,
move |_| observe(key, closed, violations, started),
None,
)
.expect("create custom wait");
member.arm(None);
unsafe { signal(raw) };
} else {
let handle = WaitableHandle::event(true, false).expect("create an event");
let member = group
.create_wait(handle, |_| std::thread::sleep(CALLBACK_DWELL), None)
.expect("create owned wait");
member.arm(None);
}
}
group.close_members(false);
probe.assert_closed_each_exactly_once(WAITS / 2, "mixed group release");
assert_eq!(group.owned_resources(), 0, "the group holds nothing after");
}