1#![cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
17
18use std::io;
19#[cfg(target_os = "macos")]
20use std::mem::size_of;
21use std::sync::atomic::AtomicU32;
22use std::time::Duration;
23
24pub fn supported() -> bool {
30 #[cfg(any(target_os = "linux", target_os = "freebsd"))]
31 {
32 true
33 }
34 #[cfg(target_os = "macos")]
35 {
36 macos::api().is_some()
37 }
38}
39
40#[cfg(target_os = "linux")]
41pub fn wait_word(word: &AtomicU32, expected: u32) -> io::Result<()> {
42 let result = unsafe {
43 libc::syscall(
44 libc::SYS_futex,
45 word.as_ptr(),
46 libc::FUTEX_WAIT,
47 expected,
48 std::ptr::null::<libc::timespec>(),
49 std::ptr::null::<u32>(),
50 0,
51 )
52 };
53 if result == 0 {
54 return Ok(());
55 }
56
57 let error = io::Error::last_os_error();
58 match error.raw_os_error() {
59 Some(libc::EAGAIN) | Some(libc::EINTR) => Ok(()),
62 _ => Err(error),
63 }
64}
65
66#[cfg(target_os = "linux")]
75pub fn wait_word_timeout(word: &AtomicU32, expected: u32, timeout: Duration) -> io::Result<bool> {
76 let left = libc::timespec {
78 tv_sec: timeout.as_secs().min(i64::MAX as u64) as libc::time_t,
79 tv_nsec: timeout.subsec_nanos() as libc::c_long,
80 };
81 let result = unsafe {
82 libc::syscall(
83 libc::SYS_futex,
84 word.as_ptr(),
85 libc::FUTEX_WAIT,
86 expected,
87 &left as *const libc::timespec,
88 std::ptr::null::<u32>(),
89 0,
90 )
91 };
92 if result == 0 {
93 return Ok(true);
94 }
95
96 let error = io::Error::last_os_error();
97 match error.raw_os_error() {
98 Some(libc::EAGAIN) | Some(libc::EINTR) => Ok(true),
99 Some(libc::ETIMEDOUT) => Ok(false),
100 _ => Err(error),
101 }
102}
103
104#[cfg(target_os = "linux")]
105pub fn wake_word(word: &AtomicU32) -> io::Result<()> {
106 let result = unsafe {
107 libc::syscall(
108 libc::SYS_futex,
109 word.as_ptr(),
110 libc::FUTEX_WAKE,
111 i32::MAX,
112 std::ptr::null::<libc::timespec>(),
113 std::ptr::null::<u32>(),
114 0,
115 )
116 };
117 if result >= 0 {
118 Ok(())
119 } else {
120 Err(io::Error::last_os_error())
121 }
122}
123
124#[cfg(target_os = "freebsd")]
125pub fn wait_word(word: &AtomicU32, expected: u32) -> io::Result<()> {
126 let result = unsafe {
127 libc::_umtx_op(
128 word.as_ptr().cast(),
129 libc::UMTX_OP_WAIT_UINT,
130 expected as libc::c_ulong,
131 std::ptr::null_mut(),
132 std::ptr::null_mut(),
133 )
134 };
135 if result == 0 {
136 return Ok(());
137 }
138
139 let error = io::Error::last_os_error();
140 match error.raw_os_error() {
141 Some(libc::EINTR) => Ok(()),
144 _ => Err(error),
145 }
146}
147
148#[cfg(target_os = "freebsd")]
157pub fn wait_word_timeout(word: &AtomicU32, expected: u32, timeout: Duration) -> io::Result<bool> {
158 let left = libc::timespec {
162 tv_sec: timeout.as_secs().min(i64::MAX as u64) as libc::time_t,
163 tv_nsec: timeout.subsec_nanos() as libc::c_long,
164 };
165 let result = unsafe {
166 libc::_umtx_op(
167 word.as_ptr().cast(),
168 libc::UMTX_OP_WAIT_UINT,
169 expected as libc::c_ulong,
170 size_of::<libc::timespec>() as *mut libc::c_void,
171 &left as *const libc::timespec as *mut libc::c_void,
172 )
173 };
174 if result == 0 {
175 return Ok(true);
176 }
177
178 let error = io::Error::last_os_error();
179 match error.raw_os_error() {
180 Some(libc::EINTR) => Ok(true),
181 Some(libc::ETIMEDOUT) => Ok(false),
182 _ => Err(error),
183 }
184}
185
186#[cfg(target_os = "freebsd")]
187pub fn wake_word(word: &AtomicU32) -> io::Result<()> {
188 let result = unsafe {
189 libc::_umtx_op(
190 word.as_ptr().cast(),
191 libc::UMTX_OP_WAKE,
192 i32::MAX as libc::c_ulong,
193 std::ptr::null_mut(),
194 std::ptr::null_mut(),
195 )
196 };
197 if result == 0 {
198 Ok(())
199 } else {
200 Err(io::Error::last_os_error())
201 }
202}
203
204#[cfg(target_os = "macos")]
205pub fn wait_word(word: &AtomicU32, expected: u32) -> io::Result<()> {
206 let api = macos::api().ok_or_else(|| {
207 io::Error::new(
208 io::ErrorKind::Unsupported,
209 "macOS shared address waits unavailable",
210 )
211 })?;
212 debug_assert_eq!(
213 (word.as_ptr() as usize) % size_of::<u32>(),
214 0,
215 "shared wait word must be naturally aligned"
216 );
217 let result = unsafe {
220 (api.wait)(
221 word.as_ptr().cast(),
222 u64::from(expected),
223 size_of::<u32>(),
224 macos::SHARED,
225 )
226 };
227 if result >= 0 {
228 return Ok(());
229 }
230 let error = io::Error::last_os_error();
231 match error.raw_os_error() {
232 Some(libc::EINTR) => Ok(()),
233 _ => Err(error),
234 }
235}
236
237#[cfg(target_os = "macos")]
246pub fn wait_word_timeout(word: &AtomicU32, expected: u32, timeout: Duration) -> io::Result<bool> {
247 let api = macos::api().ok_or_else(|| {
248 io::Error::new(
249 io::ErrorKind::Unsupported,
250 "macOS shared address waits unavailable",
251 )
252 })?;
253 let nanos = timeout.as_nanos().min(u128::from(u64::MAX)) as u64;
254 let result = unsafe {
256 (api.wait_timeout)(
257 word.as_ptr().cast(),
258 u64::from(expected),
259 size_of::<u32>(),
260 macos::SHARED,
261 macos::MACH_ABSOLUTE_TIME,
262 nanos,
263 )
264 };
265 if result >= 0 {
266 return Ok(true);
267 }
268 let error = io::Error::last_os_error();
269 match error.raw_os_error() {
270 Some(libc::EINTR) => Ok(true),
271 Some(libc::ETIMEDOUT) => Ok(false),
272 _ => Err(error),
273 }
274}
275
276#[cfg(target_os = "macos")]
277pub fn wake_word(word: &AtomicU32) -> io::Result<()> {
278 debug_assert_eq!(
279 (word.as_ptr() as usize) % size_of::<u32>(),
280 0,
281 "shared wake word must be naturally aligned"
282 );
283 let Some(api) = macos::api() else {
286 return Ok(());
287 };
288 loop {
289 let result =
290 unsafe { (api.wake_all)(word.as_ptr().cast(), size_of::<u32>(), macos::SHARED) };
291 if result >= 0 {
292 return Ok(());
293 }
294 let error = io::Error::last_os_error();
295 match error.raw_os_error() {
296 Some(libc::ENOENT) => return Ok(()),
299 Some(libc::EINTR) => continue,
300 _ => return Err(error),
301 }
302 }
303}
304
305#[cfg(target_os = "macos")]
306pub(crate) mod macos {
307 use std::sync::atomic::{AtomicU8, AtomicUsize, Ordering};
308
309 pub(super) const SHARED: u32 = 1;
312
313 pub(super) const MACH_ABSOLUTE_TIME: u32 = 32;
316
317 type Wait = unsafe extern "C" fn(*mut libc::c_void, u64, usize, u32) -> libc::c_int;
318 type WaitTimeout =
319 unsafe extern "C" fn(*mut libc::c_void, u64, usize, u32, u32, u64) -> libc::c_int;
320 type Wake = unsafe extern "C" fn(*mut libc::c_void, usize, u32) -> libc::c_int;
321
322 #[derive(Clone, Copy)]
323 pub(crate) struct Api {
324 pub(super) wait: Wait,
325 pub(super) wait_timeout: WaitTimeout,
326 pub(super) wake_all: Wake,
327 }
328
329 const UNRESOLVED: u8 = 0;
330 const UNAVAILABLE: u8 = 1;
331 const READY: u8 = 2;
332
333 static STATE: AtomicU8 = AtomicU8::new(UNRESOLVED);
334 static WAIT: AtomicUsize = AtomicUsize::new(0);
335 static WAIT_TIMEOUT: AtomicUsize = AtomicUsize::new(0);
336 static WAKE: AtomicUsize = AtomicUsize::new(0);
337
338 pub(crate) fn api() -> Option<Api> {
354 match STATE.load(Ordering::Acquire) {
355 READY => Some(load()),
356 UNAVAILABLE => None,
357 _ => resolve(),
358 }
359 }
360
361 fn load() -> Api {
362 unsafe {
363 Api {
364 wait: std::mem::transmute::<usize, Wait>(WAIT.load(Ordering::Acquire)),
365 wait_timeout: std::mem::transmute::<usize, WaitTimeout>(
366 WAIT_TIMEOUT.load(Ordering::Acquire),
367 ),
368 wake_all: std::mem::transmute::<usize, Wake>(WAKE.load(Ordering::Acquire)),
369 }
370 }
371 }
372
373 fn resolve() -> Option<Api> {
374 let wait = unsafe { libc::dlsym(libc::RTLD_DEFAULT, c"os_sync_wait_on_address".as_ptr()) };
375 let wait_timeout = unsafe {
377 libc::dlsym(
378 libc::RTLD_DEFAULT,
379 c"os_sync_wait_on_address_with_timeout".as_ptr(),
380 )
381 };
382 let wake =
383 unsafe { libc::dlsym(libc::RTLD_DEFAULT, c"os_sync_wake_by_address_all".as_ptr()) };
384 if wait.is_null() || wait_timeout.is_null() || wake.is_null() {
385 STATE.store(UNAVAILABLE, Ordering::Release);
386 return None;
387 }
388 WAIT.store(wait as usize, Ordering::Release);
390 WAIT_TIMEOUT.store(wait_timeout as usize, Ordering::Release);
391 WAKE.store(wake as usize, Ordering::Release);
392 STATE.store(READY, Ordering::Release);
393 Some(load())
394 }
395}
396
397#[cfg(test)]
398mod timeout_tests {
399 use std::sync::Arc;
400 use std::sync::atomic::Ordering;
401 use std::time::Instant;
402
403 use super::*;
404
405 #[test]
409 fn a_timeout_is_a_timeout() {
410 if !supported() {
411 return;
412 }
413 let word = AtomicU32::new(7);
414 let started = Instant::now();
415 assert!(!wait_word_timeout(&word, 7, Duration::from_millis(200)).expect("wait"));
416 let waited = started.elapsed();
417 assert!(waited >= Duration::from_millis(150), "returned after {waited:?}");
418 assert!(waited < Duration::from_secs(2), "returned after {waited:?}");
419 }
420
421 #[test]
422 fn a_wake_beats_the_timeout() {
423 if !supported() {
424 return;
425 }
426 let word = Arc::new(AtomicU32::new(0));
427 let waker = Arc::clone(&word);
428 std::thread::spawn(move || {
429 std::thread::sleep(Duration::from_millis(50));
430 waker.store(1, Ordering::SeqCst);
431 let _ = wake_word(&waker);
432 });
433 let started = Instant::now();
434 assert!(wait_word_timeout(&word, 0, Duration::from_secs(10)).expect("wait"));
435 assert!(started.elapsed() < Duration::from_secs(5));
436 }
437
438 #[test]
439 fn a_word_that_already_moved_does_not_park_at_all() {
440 if !supported() {
441 return;
442 }
443 let word = AtomicU32::new(3);
444 let started = Instant::now();
445 assert!(wait_word_timeout(&word, 9, Duration::from_secs(30)).expect("wait"));
446 assert!(started.elapsed() < Duration::from_secs(1));
447 }
448}