1#![cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
18
19use std::io;
20#[cfg(target_os = "macos")]
21use std::mem::size_of;
22use std::sync::atomic::AtomicU32;
23use std::time::Duration;
24
25pub fn supported() -> bool {
31 #[cfg(any(target_os = "linux", target_os = "freebsd"))]
32 {
33 true
34 }
35 #[cfg(target_os = "macos")]
36 {
37 macos::api().is_some()
38 }
39}
40
41#[cfg(target_os = "linux")]
42pub fn wait_word(word: &AtomicU32, expected: u32) -> io::Result<()> {
43 let result = unsafe {
44 libc::syscall(
45 libc::SYS_futex,
46 word.as_ptr(),
47 libc::FUTEX_WAIT,
48 expected,
49 std::ptr::null::<libc::timespec>(),
50 std::ptr::null::<u32>(),
51 0,
52 )
53 };
54 if result == 0 {
55 return Ok(());
56 }
57
58 let error = io::Error::last_os_error();
59 match error.raw_os_error() {
60 Some(libc::EAGAIN) | Some(libc::EINTR) => Ok(()),
63 _ => Err(error),
64 }
65}
66
67#[cfg(target_os = "linux")]
76pub fn wait_word_timeout(word: &AtomicU32, expected: u32, timeout: Duration) -> io::Result<bool> {
77 let left = libc::timespec {
79 tv_sec: timeout.as_secs().min(i64::MAX as u64) as libc::time_t,
80 tv_nsec: timeout.subsec_nanos() as libc::c_long,
81 };
82 let result = unsafe {
83 libc::syscall(
84 libc::SYS_futex,
85 word.as_ptr(),
86 libc::FUTEX_WAIT,
87 expected,
88 &left as *const libc::timespec,
89 std::ptr::null::<u32>(),
90 0,
91 )
92 };
93 if result == 0 {
94 return Ok(true);
95 }
96
97 let error = io::Error::last_os_error();
98 match error.raw_os_error() {
99 Some(libc::EAGAIN) | Some(libc::EINTR) => Ok(true),
100 Some(libc::ETIMEDOUT) => Ok(false),
101 _ => Err(error),
102 }
103}
104
105#[cfg(target_os = "linux")]
106pub fn wake_word_one(word: &AtomicU32) -> io::Result<()> {
107 let result = unsafe {
108 libc::syscall(
109 libc::SYS_futex,
110 word.as_ptr(),
111 libc::FUTEX_WAKE,
112 1,
113 std::ptr::null::<libc::timespec>(),
114 std::ptr::null::<u32>(),
115 0,
116 )
117 };
118 if result >= 0 {
119 Ok(())
120 } else {
121 Err(io::Error::last_os_error())
122 }
123}
124
125#[cfg(target_os = "linux")]
126pub fn wake_word(word: &AtomicU32) -> io::Result<()> {
127 let result = unsafe {
128 libc::syscall(
129 libc::SYS_futex,
130 word.as_ptr(),
131 libc::FUTEX_WAKE,
132 i32::MAX,
133 std::ptr::null::<libc::timespec>(),
134 std::ptr::null::<u32>(),
135 0,
136 )
137 };
138 if result >= 0 {
139 Ok(())
140 } else {
141 Err(io::Error::last_os_error())
142 }
143}
144
145#[cfg(target_os = "freebsd")]
146pub fn wait_word(word: &AtomicU32, expected: u32) -> io::Result<()> {
147 let result = unsafe {
148 libc::_umtx_op(
149 word.as_ptr().cast(),
150 libc::UMTX_OP_WAIT_UINT,
151 expected as libc::c_ulong,
152 std::ptr::null_mut(),
153 std::ptr::null_mut(),
154 )
155 };
156 if result == 0 {
157 return Ok(());
158 }
159
160 let error = io::Error::last_os_error();
161 match error.raw_os_error() {
162 Some(libc::EINTR) => Ok(()),
165 _ => Err(error),
166 }
167}
168
169#[cfg(target_os = "freebsd")]
178pub fn wait_word_timeout(word: &AtomicU32, expected: u32, timeout: Duration) -> io::Result<bool> {
179 let left = libc::timespec {
183 tv_sec: timeout.as_secs().min(i64::MAX as u64) as libc::time_t,
184 tv_nsec: timeout.subsec_nanos() as libc::c_long,
185 };
186 let result = unsafe {
187 libc::_umtx_op(
188 word.as_ptr().cast(),
189 libc::UMTX_OP_WAIT_UINT,
190 expected as libc::c_ulong,
191 size_of::<libc::timespec>() as *mut libc::c_void,
192 &left as *const libc::timespec as *mut libc::c_void,
193 )
194 };
195 if result == 0 {
196 return Ok(true);
197 }
198
199 let error = io::Error::last_os_error();
200 match error.raw_os_error() {
201 Some(libc::EINTR) => Ok(true),
202 Some(libc::ETIMEDOUT) => Ok(false),
203 _ => Err(error),
204 }
205}
206
207#[cfg(target_os = "freebsd")]
208pub fn wake_word_one(word: &AtomicU32) -> io::Result<()> {
209 let result = unsafe {
210 libc::_umtx_op(
211 word.as_ptr().cast(),
212 libc::UMTX_OP_WAKE,
213 1,
214 std::ptr::null_mut(),
215 std::ptr::null_mut(),
216 )
217 };
218 if result == 0 {
219 Ok(())
220 } else {
221 Err(io::Error::last_os_error())
222 }
223}
224
225#[cfg(target_os = "freebsd")]
226pub fn wake_word(word: &AtomicU32) -> io::Result<()> {
227 let result = unsafe {
228 libc::_umtx_op(
229 word.as_ptr().cast(),
230 libc::UMTX_OP_WAKE,
231 i32::MAX as libc::c_ulong,
232 std::ptr::null_mut(),
233 std::ptr::null_mut(),
234 )
235 };
236 if result == 0 {
237 Ok(())
238 } else {
239 Err(io::Error::last_os_error())
240 }
241}
242
243#[cfg(target_os = "macos")]
244pub fn wait_word(word: &AtomicU32, expected: u32) -> io::Result<()> {
245 let api = macos::api().ok_or_else(|| {
246 io::Error::new(
247 io::ErrorKind::Unsupported,
248 "macOS shared address waits unavailable",
249 )
250 })?;
251 debug_assert_eq!(
252 (word.as_ptr() as usize) % size_of::<u32>(),
253 0,
254 "shared wait word must be naturally aligned"
255 );
256 let result = unsafe {
259 (api.wait)(
260 word.as_ptr().cast(),
261 u64::from(expected),
262 size_of::<u32>(),
263 macos::SHARED,
264 )
265 };
266 if result >= 0 {
267 return Ok(());
268 }
269 let error = io::Error::last_os_error();
270 match error.raw_os_error() {
271 Some(libc::EINTR) => Ok(()),
272 _ => Err(error),
273 }
274}
275
276#[cfg(target_os = "macos")]
285pub fn wait_word_timeout(word: &AtomicU32, expected: u32, timeout: Duration) -> io::Result<bool> {
286 let api = macos::api().ok_or_else(|| {
287 io::Error::new(
288 io::ErrorKind::Unsupported,
289 "macOS shared address waits unavailable",
290 )
291 })?;
292 let nanos = timeout.as_nanos().min(u128::from(u64::MAX)) as u64;
293 let result = unsafe {
295 (api.wait_timeout)(
296 word.as_ptr().cast(),
297 u64::from(expected),
298 size_of::<u32>(),
299 macos::SHARED,
300 macos::MACH_ABSOLUTE_TIME,
301 nanos,
302 )
303 };
304 if result >= 0 {
305 return Ok(true);
306 }
307 let error = io::Error::last_os_error();
308 match error.raw_os_error() {
309 Some(libc::EINTR) => Ok(true),
310 Some(libc::ETIMEDOUT) => Ok(false),
311 _ => Err(error),
312 }
313}
314
315#[cfg(target_os = "macos")]
316pub fn wake_word_one(word: &AtomicU32) -> io::Result<()> {
317 wake_macos(word, false)
318}
319
320#[cfg(target_os = "macos")]
321pub fn wake_word(word: &AtomicU32) -> io::Result<()> {
322 wake_macos(word, true)
323}
324
325#[cfg(target_os = "macos")]
326fn wake_macos(word: &AtomicU32, all: bool) -> io::Result<()> {
327 debug_assert_eq!(
328 (word.as_ptr() as usize) % size_of::<u32>(),
329 0,
330 "shared wake word must be naturally aligned"
331 );
332 let Some(api) = macos::api() else {
335 return Ok(());
336 };
337 loop {
338 let wake = if all { api.wake_all } else { api.wake_one };
339 let result = unsafe { wake(word.as_ptr().cast(), size_of::<u32>(), macos::SHARED) };
340 if result >= 0 {
341 return Ok(());
342 }
343 let error = io::Error::last_os_error();
344 match error.raw_os_error() {
345 Some(libc::ENOENT) => return Ok(()),
348 Some(libc::EINTR) => continue,
349 _ => return Err(error),
350 }
351 }
352}
353
354#[cfg(target_os = "macos")]
355pub(crate) mod macos {
356 use std::sync::atomic::{AtomicU8, AtomicUsize, Ordering};
357
358 pub(super) const SHARED: u32 = 1;
361
362 pub(super) const MACH_ABSOLUTE_TIME: u32 = 32;
365
366 type Wait = unsafe extern "C" fn(*mut libc::c_void, u64, usize, u32) -> libc::c_int;
367 type WaitTimeout =
368 unsafe extern "C" fn(*mut libc::c_void, u64, usize, u32, u32, u64) -> libc::c_int;
369 type Wake = unsafe extern "C" fn(*mut libc::c_void, usize, u32) -> libc::c_int;
370
371 #[derive(Clone, Copy)]
372 pub(crate) struct Api {
373 pub(super) wait: Wait,
374 pub(super) wait_timeout: WaitTimeout,
375 pub(super) wake_one: Wake,
376 pub(super) wake_all: Wake,
377 }
378
379 const UNRESOLVED: u8 = 0;
380 const UNAVAILABLE: u8 = 1;
381 const READY: u8 = 2;
382
383 static STATE: AtomicU8 = AtomicU8::new(UNRESOLVED);
384 static WAIT: AtomicUsize = AtomicUsize::new(0);
385 static WAIT_TIMEOUT: AtomicUsize = AtomicUsize::new(0);
386 static WAKE_ONE: AtomicUsize = AtomicUsize::new(0);
387 static WAKE_ALL: AtomicUsize = AtomicUsize::new(0);
388
389 pub(crate) fn api() -> Option<Api> {
405 match STATE.load(Ordering::Acquire) {
406 READY => Some(load()),
407 UNAVAILABLE => None,
408 _ => resolve(),
409 }
410 }
411
412 fn load() -> Api {
413 unsafe {
414 Api {
415 wait: std::mem::transmute::<usize, Wait>(WAIT.load(Ordering::Acquire)),
416 wait_timeout: std::mem::transmute::<usize, WaitTimeout>(
417 WAIT_TIMEOUT.load(Ordering::Acquire),
418 ),
419 wake_one: std::mem::transmute::<usize, Wake>(WAKE_ONE.load(Ordering::Acquire)),
420 wake_all: std::mem::transmute::<usize, Wake>(WAKE_ALL.load(Ordering::Acquire)),
421 }
422 }
423 }
424
425 fn resolve() -> Option<Api> {
426 let wait = unsafe { libc::dlsym(libc::RTLD_DEFAULT, c"os_sync_wait_on_address".as_ptr()) };
427 let wait_timeout = unsafe {
429 libc::dlsym(
430 libc::RTLD_DEFAULT,
431 c"os_sync_wait_on_address_with_timeout".as_ptr(),
432 )
433 };
434 let wake_one =
435 unsafe { libc::dlsym(libc::RTLD_DEFAULT, c"os_sync_wake_by_address_any".as_ptr()) };
436 let wake_all =
437 unsafe { libc::dlsym(libc::RTLD_DEFAULT, c"os_sync_wake_by_address_all".as_ptr()) };
438 if wait.is_null() || wait_timeout.is_null() || wake_one.is_null() || wake_all.is_null() {
439 STATE.store(UNAVAILABLE, Ordering::Release);
440 return None;
441 }
442 WAIT.store(wait as usize, Ordering::Release);
444 WAIT_TIMEOUT.store(wait_timeout as usize, Ordering::Release);
445 WAKE_ONE.store(wake_one as usize, Ordering::Release);
446 WAKE_ALL.store(wake_all as usize, Ordering::Release);
447 STATE.store(READY, Ordering::Release);
448 Some(load())
449 }
450}
451
452#[cfg(test)]
453mod timeout_tests {
454 use std::sync::Arc;
455 use std::sync::atomic::Ordering;
456 use std::time::Instant;
457
458 use super::*;
459
460 #[test]
464 fn a_timeout_is_a_timeout() {
465 if !supported() {
466 return;
467 }
468 let word = AtomicU32::new(7);
469 let started = Instant::now();
470 assert!(!wait_word_timeout(&word, 7, Duration::from_millis(200)).expect("wait"));
471 let waited = started.elapsed();
472 assert!(waited >= Duration::from_millis(150), "returned after {waited:?}");
473 assert!(waited < Duration::from_secs(2), "returned after {waited:?}");
474 }
475
476 #[test]
477 fn a_wake_beats_the_timeout() {
478 if !supported() {
479 return;
480 }
481 let word = Arc::new(AtomicU32::new(0));
482 let waker = Arc::clone(&word);
483 std::thread::spawn(move || {
484 std::thread::sleep(Duration::from_millis(50));
485 waker.store(1, Ordering::SeqCst);
486 let _ = wake_word(&waker);
487 });
488 let started = Instant::now();
489 assert!(wait_word_timeout(&word, 0, Duration::from_secs(10)).expect("wait"));
490 assert!(started.elapsed() < Duration::from_secs(5));
491 }
492
493 #[test]
494 fn wake_one_releases_only_one_parked_waiter() {
495 if !supported() {
496 return;
497 }
498 let word = Arc::new(AtomicU32::new(0));
499 let (sent, received) = std::sync::mpsc::channel();
500 let waiters = (0..2)
501 .map(|_| {
502 let word = Arc::clone(&word);
503 let sent = sent.clone();
504 std::thread::spawn(move || {
505 let outcome =
506 wait_word_timeout(&word, 0, Duration::from_millis(500)).expect("wait");
507 sent.send(outcome).unwrap();
508 })
509 })
510 .collect::<Vec<_>>();
511
512 std::thread::sleep(Duration::from_millis(100));
513 word.store(1, Ordering::SeqCst);
514 wake_word_one(&word).expect("wake one");
515
516 assert!(received.recv_timeout(Duration::from_millis(200)).unwrap());
517 assert!(received.recv_timeout(Duration::from_millis(100)).is_err());
518 assert!(!received.recv_timeout(Duration::from_millis(400)).unwrap());
519 for waiter in waiters {
520 waiter.join().unwrap();
521 }
522 }
523
524 #[test]
525 fn a_word_that_already_moved_does_not_park_at_all() {
526 if !supported() {
527 return;
528 }
529 let word = AtomicU32::new(3);
530 let started = Instant::now();
531 assert!(wait_word_timeout(&word, 9, Duration::from_secs(30)).expect("wait"));
532 assert!(started.elapsed() < Duration::from_secs(1));
533 }
534}