use std::fs::{File, OpenOptions};
use std::path::Path;
use std::sync::atomic::{AtomicU32, AtomicU64, Ordering};
use std::time::{Duration, Instant};
use memmap2::{MmapMut, MmapOptions};
pub const WAKER_MAGIC: u64 = 0xE7E7_5742_4B45_5201;
pub const MAX_WAITERS_DEFAULT: usize = 32;
const STATE_FREE: u32 = 0;
const STATE_RESERVED: u32 = 1;
const STATE_PARKED: u32 = 2;
const STATE_WOKEN: u32 = 3;
#[repr(C, align(64))]
struct WakerHeader {
magic: u64,
capacity: u32,
_pad0: [u8; 4],
parked_mask: AtomicU64,
_pad: [u8; 64 - 24],
}
#[repr(C, align(64))]
struct WakerSlot {
state: AtomicU32,
_pad1: [u8; 4],
target_seq: AtomicU64,
_pad2: [u8; 64 - 16],
}
const _: () = {
assert!(std::mem::size_of::<WakerHeader>() == 64);
assert!(std::mem::size_of::<WakerSlot>() == 64);
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum WakerError {
Full,
Timeout,
LayoutMismatch,
IoError(std::io::ErrorKind),
}
impl From<std::io::Error> for WakerError {
fn from(e: std::io::Error) -> Self { Self::IoError(e.kind()) }
}
impl std::fmt::Display for WakerError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Full => write!(f, "no free waker slot"),
Self::Timeout => write!(f, "wait timed out"),
Self::LayoutMismatch => write!(f, "waker layout mismatch on open"),
Self::IoError(k) => write!(f, "waker mmap io error: {k:?}"),
}
}
}
impl std::error::Error for WakerError {}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct WakerToken {
slot: u32,
}
impl WakerToken {
pub fn slot_index(&self) -> u32 { self.slot }
}
pub const fn waker_region_size(capacity: usize) -> usize {
std::mem::size_of::<WakerHeader>() + capacity * std::mem::size_of::<WakerSlot>()
}
pub struct CrossProcessWaker {
_backing: WakerBacking,
raw_ptr: *mut u8,
capacity: usize,
}
unsafe impl Send for CrossProcessWaker {}
unsafe impl Sync for CrossProcessWaker {}
#[allow(dead_code)]
enum WakerBacking {
Anon(MmapMut),
File(File, MmapMut),
Shm(crate::shm_file::ShmFile),
}
unsafe fn init_waker_layout_raw(ptr: *mut u8, capacity: usize) {
let hdr_ptr = ptr as *mut WakerHeader;
unsafe {
std::ptr::write_bytes(hdr_ptr as *mut u8, 0, std::mem::size_of::<WakerHeader>());
(*hdr_ptr).magic = WAKER_MAGIC;
(*hdr_ptr).capacity = capacity as u32;
}
let slots_base = unsafe { ptr.add(std::mem::size_of::<WakerHeader>()) };
for i in 0..capacity {
let slot_ptr = unsafe {
slots_base.add(i * std::mem::size_of::<WakerSlot>())
} as *mut WakerSlot;
unsafe {
std::ptr::write(slot_ptr, WakerSlot {
state: AtomicU32::new(STATE_FREE),
_pad1: [0; 4],
target_seq: AtomicU64::new(0),
_pad2: [0; 64 - 16],
});
}
}
}
impl CrossProcessWaker {
pub fn create_anon(capacity: usize) -> Result<Self, WakerError> {
assert!(capacity >= 1, "capacity must be >= 1");
let total = waker_region_size(capacity);
let mut mmap = MmapOptions::new().len(total).map_anon()?;
let raw_ptr = mmap.as_mut_ptr();
unsafe { init_waker_layout_raw(raw_ptr, capacity); }
Ok(Self {
_backing: WakerBacking::Anon(mmap),
raw_ptr,
capacity,
})
}
pub fn create(path: impl AsRef<Path>, capacity: usize) -> Result<Self, WakerError> {
assert!(capacity >= 1, "capacity must be >= 1");
let total = waker_region_size(capacity);
let (file, mut mmap) = crate::mmf_attach::create_or_attach(
path.as_ref(),
total,
|ptr| unsafe { init_waker_layout_raw(ptr, capacity) },
|ptr| unsafe { (*(ptr as *const WakerHeader)).magic == WAKER_MAGIC },
)?;
let raw_ptr = mmap.as_mut_ptr();
let hdr = unsafe { &*(raw_ptr as *const WakerHeader) };
if hdr.capacity as usize != capacity {
return Err(WakerError::LayoutMismatch);
}
Ok(Self {
_backing: WakerBacking::File(file, mmap),
raw_ptr,
capacity,
})
}
pub fn reset(path: impl AsRef<Path>, capacity: usize) -> Result<Self, WakerError> {
assert!(capacity >= 1, "capacity must be >= 1");
let total = waker_region_size(capacity);
let (file, mut mmap) = crate::mmf_attach::reset(path.as_ref(), total, |ptr| unsafe {
init_waker_layout_raw(ptr, capacity)
})?;
let raw_ptr = mmap.as_mut_ptr();
Ok(Self {
_backing: WakerBacking::File(file, mmap),
raw_ptr,
capacity,
})
}
pub fn open(
path: impl AsRef<Path>,
expected_capacity: usize,
) -> Result<Self, WakerError> {
let total = waker_region_size(expected_capacity);
let file = OpenOptions::new().read(true).write(true).open(path.as_ref())?;
if (file.metadata()?.len() as usize) < total {
return Err(WakerError::LayoutMismatch);
}
let mut mmap = unsafe { MmapOptions::new().len(total).map_mut(&file)? };
let raw_ptr = mmap.as_mut_ptr();
let hdr = unsafe { &*(raw_ptr as *const WakerHeader) };
if hdr.magic != WAKER_MAGIC || hdr.capacity as usize != expected_capacity {
return Err(WakerError::LayoutMismatch);
}
Ok(Self {
_backing: WakerBacking::File(file, mmap),
raw_ptr,
capacity: expected_capacity,
})
}
pub fn create_from_shm(
mut shm: crate::shm_file::ShmFile,
capacity: usize,
) -> Result<Self, WakerError> {
assert!(capacity >= 1, "capacity must be >= 1");
let total = waker_region_size(capacity);
if shm.len() < total {
return Err(WakerError::LayoutMismatch);
}
let raw_ptr = shm.as_mut_slice().as_mut_ptr();
unsafe { init_waker_layout_raw(raw_ptr, capacity); }
Ok(Self {
_backing: WakerBacking::Shm(shm),
raw_ptr,
capacity,
})
}
pub fn open_from_shm(
mut shm: crate::shm_file::ShmFile,
expected_capacity: usize,
) -> Result<Self, WakerError> {
let total = waker_region_size(expected_capacity);
if shm.len() < total {
return Err(WakerError::LayoutMismatch);
}
let raw_ptr = shm.as_mut_slice().as_mut_ptr();
let hdr = unsafe { &*(raw_ptr as *const WakerHeader) };
if hdr.magic != WAKER_MAGIC || hdr.capacity as usize != expected_capacity {
return Err(WakerError::LayoutMismatch);
}
Ok(Self {
_backing: WakerBacking::Shm(shm),
raw_ptr,
capacity: expected_capacity,
})
}
pub fn capacity(&self) -> usize { self.capacity }
#[inline]
fn slot(&self, idx: usize) -> &WakerSlot {
let base = unsafe { self.raw_ptr.add(std::mem::size_of::<WakerHeader>()) };
unsafe { &*(base.add(idx * std::mem::size_of::<WakerSlot>()) as *const WakerSlot) }
}
pub fn try_park(&self, target_seq: u64) -> Result<WakerToken, WakerError> {
for idx in 0..self.capacity {
let slot = self.slot(idx);
if slot
.state
.compare_exchange(
STATE_FREE,
STATE_RESERVED,
Ordering::Acquire,
Ordering::Relaxed,
)
.is_ok()
{
slot.target_seq.store(target_seq, Ordering::Relaxed);
slot.state.store(STATE_PARKED, Ordering::Release);
self.mask_set(idx);
return Ok(WakerToken { slot: idx as u32 });
}
}
Err(WakerError::Full)
}
fn is_cross_process(&self) -> bool {
!matches!(self._backing, WakerBacking::Anon(_))
}
fn header(&self) -> &WakerHeader {
unsafe { &*(self.raw_ptr as *const WakerHeader) }
}
#[inline]
fn mask_set(&self, idx: usize) {
if idx < 64 {
self.header()
.parked_mask
.fetch_or(1u64 << idx, Ordering::Release);
}
}
#[inline]
fn mask_clear(&self, idx: usize) {
if idx < 64 {
self.header()
.parked_mask
.fetch_and(!(1u64 << idx), Ordering::Release);
}
}
#[inline]
fn wake_candidates(&self) -> WakeCandidates {
if self.capacity <= 64 {
WakeCandidates::Mask(
self.header().parked_mask.load(Ordering::Acquire),
)
} else {
WakeCandidates::Range(0, self.capacity)
}
}
pub fn wait(
&self,
token: WakerToken,
timeout: Option<Duration>,
) -> Result<(), WakerError> {
let slot = self.slot(token.slot as usize);
let cross_process = self.is_cross_process();
let result = match timeout {
None => {
loop {
let cur = slot.state.load(Ordering::Acquire);
if cur != STATE_PARKED {
break Ok(());
}
platform_wait::wait_forever(
&slot.state, STATE_PARKED, cross_process,
);
}
}
Some(d) => {
let deadline = Instant::now() + d;
loop {
if slot.state.load(Ordering::Acquire) != STATE_PARKED {
break Ok(());
}
match deadline.checked_duration_since(Instant::now()) {
Some(remaining) if !remaining.is_zero() => {
platform_wait::wait_with_timeout(
&slot.state, STATE_PARKED, remaining, cross_process,
);
}
_ => {
break if slot.state.load(Ordering::Acquire) != STATE_PARKED {
Ok(())
} else {
Err(WakerError::Timeout)
};
}
}
}
}
};
slot.state.store(STATE_FREE, Ordering::Release);
self.mask_clear(token.slot as usize);
result
}
pub fn release(&self, token: WakerToken) {
let slot = self.slot(token.slot as usize);
slot.state.store(STATE_FREE, Ordering::Release);
self.mask_clear(token.slot as usize);
}
pub fn wake_up_to(&self, seq: u64) -> usize {
let mut count = 0usize;
for idx in self.wake_candidates() {
let slot = self.slot(idx);
let cur_state = slot.state.load(Ordering::Acquire);
if cur_state != STATE_PARKED {
continue;
}
let tgt = slot.target_seq.load(Ordering::Relaxed);
if seq < tgt {
continue;
}
if slot
.state
.compare_exchange(
STATE_PARKED,
STATE_WOKEN,
Ordering::AcqRel,
Ordering::Relaxed,
)
.is_ok()
{
platform_wait::wake_one(&slot.state, self.is_cross_process());
count += 1;
}
}
count
}
pub fn wake_one_up_to(&self, seq: u64) -> usize {
for idx in self.wake_candidates() {
let slot = self.slot(idx);
let cur_state = slot.state.load(Ordering::Acquire);
if cur_state != STATE_PARKED {
continue;
}
let tgt = slot.target_seq.load(Ordering::Relaxed);
if seq < tgt {
continue;
}
if slot
.state
.compare_exchange(
STATE_PARKED,
STATE_WOKEN,
Ordering::AcqRel,
Ordering::Relaxed,
)
.is_ok()
{
platform_wait::wake_one(&slot.state, self.is_cross_process());
return 1;
}
}
0
}
pub fn wake_all(&self) -> usize {
let mut count = 0usize;
for idx in self.wake_candidates() {
let slot = self.slot(idx);
if slot
.state
.compare_exchange(
STATE_PARKED,
STATE_WOKEN,
Ordering::AcqRel,
Ordering::Relaxed,
)
.is_ok()
{
platform_wait::wake_one(&slot.state, self.is_cross_process());
count += 1;
}
}
count
}
}
enum WakeCandidates {
Mask(u64),
Range(usize, usize),
}
impl Iterator for WakeCandidates {
type Item = usize;
#[inline]
fn next(&mut self) -> Option<usize> {
match self {
WakeCandidates::Mask(m) => {
if *m == 0 {
return None;
}
let idx = m.trailing_zeros() as usize;
*m &= *m - 1;
Some(idx)
}
WakeCandidates::Range(next, end) => {
if next < end {
let idx = *next;
*next += 1;
Some(idx)
} else {
None
}
}
}
}
}
mod platform_wait {
use std::sync::atomic::AtomicU32;
use std::time::Duration;
#[cfg(target_os = "macos")]
mod os_sync_dyn {
use std::ffi::c_void;
use std::sync::atomic::{AtomicUsize, Ordering};
pub type WaitFn = unsafe extern "C" fn(*mut c_void, u64, usize, u32) -> i32;
pub type WaitTimeoutFn =
unsafe extern "C" fn(*mut c_void, u64, usize, u32, u32, u64) -> i32;
pub type WakeFn = unsafe extern "C" fn(*mut c_void, usize, u32) -> i32;
static WAIT: AtomicUsize = AtomicUsize::new(1);
static WAIT_TO: AtomicUsize = AtomicUsize::new(1);
static WAKE: AtomicUsize = AtomicUsize::new(1);
fn cached(slot: &AtomicUsize, name: &[u8]) -> usize {
let v = slot.load(Ordering::Relaxed);
if v != 1 {
return v;
}
let r = unsafe { libc::dlsym(libc::RTLD_DEFAULT, name.as_ptr() as *const _) } as usize;
slot.store(r, Ordering::Relaxed);
r
}
pub fn wait() -> Option<WaitFn> {
match cached(&WAIT, b"os_sync_wait_on_address\0") {
0 => None,
p => Some(unsafe { std::mem::transmute::<usize, WaitFn>(p) }),
}
}
pub fn wait_timeout() -> Option<WaitTimeoutFn> {
match cached(&WAIT_TO, b"os_sync_wait_on_address_with_timeout\0") {
0 => None,
p => Some(unsafe { std::mem::transmute::<usize, WaitTimeoutFn>(p) }),
}
}
pub fn wake() -> Option<WakeFn> {
match cached(&WAKE, b"os_sync_wake_by_address_any\0") {
0 => None,
p => Some(unsafe { std::mem::transmute::<usize, WakeFn>(p) }),
}
}
}
#[inline]
fn monitor_tier(atomic: &AtomicU32, expected: u32) -> bool {
crate::monitor_wait::monitor_wait_u32(
atomic,
expected,
crate::monitor_wait::monitor_wait_budget_cycles(),
)
}
pub fn wait_forever(atomic: &AtomicU32, expected: u32, _cross_process: bool) {
if monitor_tier(atomic, expected) {
return;
}
#[cfg(windows)]
if _cross_process
&& crate::monitor_wait::monitor_wait_kind().is_some()
{
let budget = crate::monitor_wait::monitor_wait_budget_cycles();
while atomic.load(std::sync::atomic::Ordering::Acquire) == expected {
crate::monitor_wait::monitor_wait_u32(atomic, expected, budget);
}
return;
}
#[cfg(any(target_os = "linux", target_os = "android"))]
{
unsafe {
libc::syscall(
libc::SYS_futex,
atomic.as_ptr(),
libc::FUTEX_WAIT,
expected as libc::c_int,
std::ptr::null::<libc::timespec>(),
);
}
}
#[cfg(target_os = "freebsd")]
{
unsafe {
libc::_umtx_op(
atomic.as_ptr() as *mut libc::c_void,
libc::UMTX_OP_WAIT_UINT,
expected as libc::c_ulong,
std::ptr::null_mut(),
std::ptr::null_mut(),
);
}
}
#[cfg(windows)]
{
use windows_sys::Win32::System::Threading::{WaitOnAddress, INFINITE};
let expected_local = expected;
unsafe {
WaitOnAddress(
atomic.as_ptr() as *const std::ffi::c_void,
&expected_local as *const u32 as *const std::ffi::c_void,
std::mem::size_of::<u32>(),
INFINITE,
);
}
}
#[cfg(target_os = "macos")]
{
if let Some(wait_fn) = os_sync_dyn::wait() {
let flags = if _cross_process {
libc::OS_SYNC_WAIT_ON_ADDRESS_SHARED
} else {
libc::OS_SYNC_WAIT_ON_ADDRESS_NONE
};
unsafe {
wait_fn(
atomic.as_ptr() as *mut libc::c_void,
expected as u64,
std::mem::size_of::<u32>(),
flags,
);
}
} else {
std::thread::yield_now();
while atomic.load(std::sync::atomic::Ordering::Acquire) == expected {
std::thread::sleep(Duration::from_millis(1));
}
}
}
#[cfg(not(any(target_os = "linux", target_os = "android",
target_os = "freebsd", target_os = "macos", windows)))]
{
std::thread::yield_now();
while atomic.load(std::sync::atomic::Ordering::Acquire) == expected {
std::thread::sleep(Duration::from_millis(1));
}
}
}
pub fn wait_with_timeout(
atomic: &AtomicU32,
expected: u32,
timeout: Duration,
_cross_process: bool,
) -> bool {
let monitor_start = std::time::Instant::now();
if monitor_tier(atomic, expected) {
return true;
}
let timeout = match timeout.checked_sub(monitor_start.elapsed()) {
Some(rest) if !rest.is_zero() => rest,
_ => {
return atomic.load(std::sync::atomic::Ordering::Acquire)
!= expected;
}
};
#[cfg(windows)]
if _cross_process
&& crate::monitor_wait::monitor_wait_kind().is_some()
{
let deadline = std::time::Instant::now() + timeout;
let budget = crate::monitor_wait::monitor_wait_budget_cycles();
loop {
if crate::monitor_wait::monitor_wait_u32(atomic, expected, budget) {
return true;
}
if std::time::Instant::now() >= deadline {
return atomic.load(std::sync::atomic::Ordering::Acquire)
!= expected;
}
}
}
#[cfg(any(target_os = "linux", target_os = "android"))]
{
let ts = libc::timespec {
tv_sec: timeout.as_secs() as libc::time_t,
tv_nsec: timeout.subsec_nanos() as libc::c_long,
};
unsafe {
let rc = libc::syscall(
libc::SYS_futex,
atomic.as_ptr(),
libc::FUTEX_WAIT,
expected as libc::c_int,
&ts as *const libc::timespec,
std::ptr::null::<()>(),
0u32,
);
if rc == -1 {
let err = *libc::__errno_location();
return err != libc::ETIMEDOUT;
}
}
true
}
#[cfg(target_os = "freebsd")]
{
let mut ts = libc::timespec {
tv_sec: timeout.as_secs() as libc::time_t,
tv_nsec: timeout.subsec_nanos() as libc::c_long,
};
let rc = unsafe {
libc::_umtx_op(
atomic.as_ptr() as *mut libc::c_void,
libc::UMTX_OP_WAIT_UINT,
expected as libc::c_ulong,
std::mem::size_of::<libc::timespec>() as *mut libc::c_void,
&mut ts as *mut libc::timespec as *mut libc::c_void,
)
};
if rc == -1 {
let err = unsafe { *libc::__error() };
return err != libc::ETIMEDOUT;
}
true
}
#[cfg(windows)]
{
use windows_sys::Win32::System::Threading::WaitOnAddress;
let expected_local = expected;
let ms = timeout.as_millis().min(u32::MAX as u128) as u32;
let rc = unsafe {
WaitOnAddress(
atomic.as_ptr() as *const std::ffi::c_void,
&expected_local as *const u32 as *const std::ffi::c_void,
std::mem::size_of::<u32>(),
ms,
)
};
rc != 0
}
#[cfg(target_os = "macos")]
{
if let Some(wait_fn) = os_sync_dyn::wait_timeout() {
let flags = if _cross_process {
libc::OS_SYNC_WAIT_ON_ADDRESS_SHARED
} else {
libc::OS_SYNC_WAIT_ON_ADDRESS_NONE
};
let rc = unsafe {
wait_fn(
atomic.as_ptr() as *mut libc::c_void,
expected as u64,
std::mem::size_of::<u32>(),
flags,
libc::OS_CLOCK_MACH_ABSOLUTE_TIME,
timeout.as_nanos().min(u64::MAX as u128) as u64,
)
};
if rc < 0 {
return std::io::Error::last_os_error().raw_os_error()
!= Some(libc::ETIMEDOUT);
}
true
} else {
let deadline = std::time::Instant::now() + timeout;
loop {
if atomic.load(std::sync::atomic::Ordering::Acquire) != expected {
return true;
}
if std::time::Instant::now() >= deadline {
return atomic.load(std::sync::atomic::Ordering::Acquire) != expected;
}
std::thread::sleep(Duration::from_millis(2));
}
}
}
#[cfg(not(any(target_os = "linux", target_os = "android",
target_os = "freebsd", target_os = "macos", windows)))]
{
let deadline = std::time::Instant::now() + timeout;
let step = Duration::from_millis(2);
loop {
if atomic.load(std::sync::atomic::Ordering::Acquire) != expected {
return true;
}
let now = std::time::Instant::now();
if now >= deadline {
return false;
}
let remaining = deadline - now;
std::thread::sleep(remaining.min(step));
}
}
}
pub fn wake_one(atomic: &AtomicU32, _cross_process: bool) {
#[cfg(any(target_os = "linux", target_os = "android"))]
{
unsafe {
libc::syscall(
libc::SYS_futex,
atomic.as_ptr(),
libc::FUTEX_WAKE,
1i32,
);
}
}
#[cfg(target_os = "freebsd")]
{
unsafe {
libc::_umtx_op(
atomic.as_ptr() as *mut libc::c_void,
libc::UMTX_OP_WAKE,
1 as libc::c_ulong,
std::ptr::null_mut(),
std::ptr::null_mut(),
);
}
}
#[cfg(windows)]
{
use windows_sys::Win32::System::Threading::WakeByAddressSingle;
unsafe {
WakeByAddressSingle(atomic.as_ptr() as *const std::ffi::c_void);
}
}
#[cfg(target_os = "macos")]
{
if let Some(wake_fn) = os_sync_dyn::wake() {
let flags = if _cross_process {
libc::OS_SYNC_WAKE_BY_ADDRESS_SHARED
} else {
libc::OS_SYNC_WAKE_BY_ADDRESS_NONE
};
unsafe {
wake_fn(
atomic.as_ptr() as *mut libc::c_void,
std::mem::size_of::<u32>(),
flags,
);
}
}
}
#[cfg(not(any(target_os = "linux", target_os = "android",
target_os = "freebsd", target_os = "macos", windows)))]
{
drop(atomic);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
use std::thread;
use std::time::Instant;
#[test]
fn anon_round_trip() {
let waker = CrossProcessWaker::create_anon(4).expect("create");
let token = waker.try_park(10).expect("park");
let waker_c = Arc::new(waker);
let waker_p = Arc::clone(&waker_c);
let h = thread::spawn(move || {
thread::sleep(Duration::from_millis(20));
let n = waker_p.wake_up_to(15);
assert!(n >= 1);
});
waker_c.wait(token, Some(Duration::from_secs(2))).expect("wake");
h.join().unwrap();
}
#[test]
fn wait_returns_immediately_if_already_woken() {
let waker = CrossProcessWaker::create_anon(2).expect("create");
let token = waker.try_park(5).expect("park");
assert_eq!(waker.wake_up_to(10), 1);
let t0 = Instant::now();
waker.wait(token, Some(Duration::from_secs(1))).expect("ok");
assert!(t0.elapsed() < Duration::from_millis(50),
"wait should return fast since wake fired before wait entered");
}
#[test]
fn timeout_works() {
let waker = CrossProcessWaker::create_anon(2).expect("create");
let token = waker.try_park(100).expect("park");
let t0 = Instant::now();
let err = waker.wait(token, Some(Duration::from_millis(80)));
assert_eq!(err, Err(WakerError::Timeout));
assert!(t0.elapsed() >= Duration::from_millis(70));
}
#[test]
fn full_when_all_slots_taken() {
let waker = CrossProcessWaker::create_anon(2).expect("create");
let _a = waker.try_park(1).expect("park 0");
let _b = waker.try_park(2).expect("park 1");
assert_eq!(waker.try_park(3), Err(WakerError::Full));
}
#[test]
fn release_lets_others_park() {
let waker = CrossProcessWaker::create_anon(2).expect("create");
let a = waker.try_park(1).expect("park 0");
let _b = waker.try_park(2).expect("park 1");
waker.release(a);
let _c = waker.try_park(3).expect("re-park 0");
}
#[test]
fn wake_all_drains_blocked_consumers() {
let waker = Arc::new(CrossProcessWaker::create_anon(4).expect("create"));
let mut handles = Vec::new();
for target in 0..4u64 {
let w = Arc::clone(&waker);
let token = w.try_park(target + 1000).expect("park");
handles.push(thread::spawn(move || {
w.wait(token, Some(Duration::from_secs(2))).expect("woken");
}));
}
thread::sleep(Duration::from_millis(20));
assert_eq!(waker.wake_all(), 4);
for h in handles { h.join().unwrap(); }
}
}