Skip to main content

orbit_core/
sync.rs

1//! Waiting on a shared 32-bit word, the primitive under every "wake me when
2//! this changes" in Orbit: ring readiness, cell changes.
3//!
4//! [`wait_word`] parks the caller until the word no longer holds `expected`
5//! (or spuriously; callers loop). [`wake_word_one`] wakes one waiter and
6//! [`wake_word`] wakes every waiter on it. The
7//! word may live in shared memory: Linux futex, FreeBSD umtx and macOS
8//! `os_sync_wait_on_address` all key waiters by the physical location, so a
9//! wake in one process reaches a waiter in another with nothing carried
10//! between them. No descriptor, no channel: the memory is the signal.
11//!
12//! macOS needs 14.4 for the shared form. Below that there is no wait to
13//! be had, and [`supported`] says so: a crate that parks on words refuses
14//! to open rather than pretending, because a sleep loop wearing the shape
15//! of a wait is worse than a clear no.
16
17#![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
25/// Whether this build can park on a shared word at all.
26///
27/// Linux and FreeBSD always can. macOS can from 14.4; below it, and on
28/// any other target, nothing here works and callers should refuse at
29/// their own front door rather than degrade quietly.
30pub 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        // The generation changed before the kernel parked us, or the driver
61        // was interrupted. The outer loop re-checks both generation and stop.
62        Some(libc::EAGAIN) | Some(libc::EINTR) => Ok(()),
63        _ => Err(error),
64    }
65}
66
67/// Park until the word no longer holds `expected`, or `timeout` passes.
68///
69/// `Ok(true)` means something may have changed — a wake, a value that
70/// moved before the kernel parked us, or a signal — and the caller
71/// re-checks as it does after [`wait_word`]. `Ok(false)` means the
72/// timeout passed and nothing else. The timeout is relative and measured
73/// on a monotonic clock, so a caller holding a deadline recomputes what
74/// is left on each turn of its loop.
75#[cfg(target_os = "linux")]
76pub fn wait_word_timeout(word: &AtomicU32, expected: u32, timeout: Duration) -> io::Result<bool> {
77    // FUTEX_WAIT reads this as relative, on CLOCK_MONOTONIC.
78    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        // The generation changed before the kernel parked us, or the driver
163        // was interrupted. The outer loop re-checks generation and stop.
164        Some(libc::EINTR) => Ok(()),
165        _ => Err(error),
166    }
167}
168
169/// Park until the word no longer holds `expected`, or `timeout` passes.
170///
171/// `Ok(true)` means something may have changed — a wake, a value that
172/// moved before the kernel parked us, or a signal — and the caller
173/// re-checks as it does after [`wait_word`]. `Ok(false)` means the
174/// timeout passed and nothing else. The timeout is relative and measured
175/// on a monotonic clock, so a caller holding a deadline recomputes what
176/// is left on each turn of its loop.
177#[cfg(target_os = "freebsd")]
178pub fn wait_word_timeout(word: &AtomicU32, expected: u32, timeout: Duration) -> io::Result<bool> {
179    // For the UMTX_OP_WAIT family the fourth argument is the size of the
180    // timeout structure and the fifth points at it; a bare `timespec` is
181    // read as relative.
182    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    // The SHM word is AtomicU32, not u64. Wait and wake must agree on size
257    // and shared mode. Apple returns a nonnegative waiter count on success.
258    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/// Park until the word no longer holds `expected`, or `timeout` passes.
277///
278/// `Ok(true)` means something may have changed — a wake, a value that
279/// moved before the kernel parked us, or a signal — and the caller
280/// re-checks as it does after [`wait_word`]. `Ok(false)` means the
281/// timeout passed and nothing else. The timeout is relative and measured
282/// on a monotonic clock, so a caller holding a deadline recomputes what
283/// is left on each turn of its loop.
284#[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    // SAFETY: the same word, size and shared flag as the untimed wait.
294    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    // Older macOS has no native subscribers; publication still succeeds and
333    // polling readers see the committed frames.
334    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            // No waiter is normal: publication can precede subscription or
346            // race the driver's compare-and-wait. The generation persists.
347            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    // OS_SYNC_WAIT_ON_ADDRESS_SHARED and OS_SYNC_WAKE_BY_ADDRESS_SHARED
359    // have the same ABI value in <os/os_sync_wait_on_address.h>.
360    pub(super) const SHARED: u32 = 1;
361
362    /// `os_clockid_t` in <os/clock.h>: the only clock the timed wait
363    /// takes, and the one a relative timeout is measured on.
364    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    /// Resolve the 14.4 entry points, once per process but never by waiting.
390    ///
391    /// Deliberately not a `OnceLock`. A publisher reaches this from
392    /// `wake_all_generation_waiters`, and a publisher may be a process that
393    /// was just forked while another thread of its parent was inside the
394    /// resolution -- the test harness runs tests on parallel threads, and a
395    /// supervisor forks workers with readiness threads alive. The child
396    /// inherits the "initializing" state and none of the thread that would
397    /// finish it, so a blocking once-cell parks forever. Resolution here is
398    /// idempotent: racing callers each `dlsym` the same two symbols and store
399    /// the same values, and nobody waits for anybody.
400    ///
401    /// Resolving lazily rather than linking avoids hard references to
402    /// 14.4-only symbols on older deployment targets. libSystem stays loaded
403    /// for the life of the process, so the pointers never dangle.
404    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        // Shipped in the same release as the other two; all or nothing.
428        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        // Pointers first, then the state that publishes them.
443        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    /// Also pins the unit the platform reads the timeout in: a wrong one
461    /// shows up here as a wait that is orders of magnitude off, not as a
462    /// wrong answer.
463    #[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}