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(
43    word: &AtomicU32,
44    expected: u32
45) -> io::Result<()> {
46    let result = unsafe {
47        libc::syscall(
48            libc::SYS_futex,
49            word.as_ptr(),
50            libc::FUTEX_WAIT,
51            expected,
52            std::ptr::null::<libc::timespec>(),
53            std::ptr::null::<u32>(),
54            0
55        )
56    };
57    if result == 0 {
58        return Ok(());
59    }
60
61    let error = io::Error::last_os_error();
62    match error.raw_os_error() {
63        // The generation changed before the kernel parked us, or the driver
64        // was interrupted. The outer loop re-checks both generation and stop.
65        Some(libc::EAGAIN) | Some(libc::EINTR) => Ok(()),
66        _ => Err(error)
67    }
68}
69
70/// Park until the word no longer holds `expected`, or `timeout` passes.
71///
72/// `Ok(true)` means something may have changed — a wake, a value that
73/// moved before the kernel parked us, or a signal — and the caller
74/// re-checks as it does after [`wait_word`]. `Ok(false)` means the
75/// timeout passed and nothing else. The timeout is relative and measured
76/// on a monotonic clock, so a caller holding a deadline recomputes what
77/// is left on each turn of its loop.
78#[cfg(target_os = "linux")]
79pub fn wait_word_timeout(
80    word: &AtomicU32,
81    expected: u32,
82    timeout: Duration
83) -> io::Result<bool> {
84    // FUTEX_WAIT reads this as relative, on CLOCK_MONOTONIC.
85    let left = libc::timespec {
86        tv_sec: timeout.as_secs().min(i64::MAX as u64) as libc::time_t,
87        tv_nsec: timeout.subsec_nanos() as libc::c_long
88    };
89    let result = unsafe {
90        libc::syscall(
91            libc::SYS_futex,
92            word.as_ptr(),
93            libc::FUTEX_WAIT,
94            expected,
95            &left as *const libc::timespec,
96            std::ptr::null::<u32>(),
97            0
98        )
99    };
100    if result == 0 {
101        return Ok(true);
102    }
103
104    let error = io::Error::last_os_error();
105    match error.raw_os_error() {
106        Some(libc::EAGAIN) | Some(libc::EINTR) => Ok(true),
107        Some(libc::ETIMEDOUT) => Ok(false),
108        _ => Err(error)
109    }
110}
111
112#[cfg(target_os = "linux")]
113pub fn wake_word_one(word: &AtomicU32) -> io::Result<()> {
114    let result = unsafe {
115        libc::syscall(
116            libc::SYS_futex,
117            word.as_ptr(),
118            libc::FUTEX_WAKE,
119            1,
120            std::ptr::null::<libc::timespec>(),
121            std::ptr::null::<u32>(),
122            0
123        )
124    };
125    if result >= 0 { Ok(()) } else { Err(io::Error::last_os_error()) }
126}
127
128#[cfg(target_os = "linux")]
129pub fn wake_word(word: &AtomicU32) -> io::Result<()> {
130    let result = unsafe {
131        libc::syscall(
132            libc::SYS_futex,
133            word.as_ptr(),
134            libc::FUTEX_WAKE,
135            i32::MAX,
136            std::ptr::null::<libc::timespec>(),
137            std::ptr::null::<u32>(),
138            0
139        )
140    };
141    if result >= 0 { Ok(()) } else { Err(io::Error::last_os_error()) }
142}
143
144#[cfg(target_os = "freebsd")]
145pub fn wait_word(
146    word: &AtomicU32,
147    expected: u32
148) -> io::Result<()> {
149    let result = unsafe {
150        libc::_umtx_op(
151            word.as_ptr().cast(),
152            libc::UMTX_OP_WAIT_UINT,
153            expected as libc::c_ulong,
154            std::ptr::null_mut(),
155            std::ptr::null_mut()
156        )
157    };
158    if result == 0 {
159        return Ok(());
160    }
161
162    let error = io::Error::last_os_error();
163    match error.raw_os_error() {
164        // The generation changed before the kernel parked us, or the driver
165        // was interrupted. The outer loop re-checks generation and stop.
166        Some(libc::EINTR) => Ok(()),
167        _ => Err(error)
168    }
169}
170
171/// Park until the word no longer holds `expected`, or `timeout` passes.
172///
173/// `Ok(true)` means something may have changed — a wake, a value that
174/// moved before the kernel parked us, or a signal — and the caller
175/// re-checks as it does after [`wait_word`]. `Ok(false)` means the
176/// timeout passed and nothing else. The timeout is relative and measured
177/// on a monotonic clock, so a caller holding a deadline recomputes what
178/// is left on each turn of its loop.
179#[cfg(target_os = "freebsd")]
180pub fn wait_word_timeout(
181    word: &AtomicU32,
182    expected: u32,
183    timeout: Duration
184) -> io::Result<bool> {
185    // For the UMTX_OP_WAIT family the fourth argument is the size of the
186    // timeout structure and the fifth points at it; a bare `timespec` is
187    // read as relative.
188    let left = libc::timespec {
189        tv_sec: timeout.as_secs().min(i64::MAX as u64) as libc::time_t,
190        tv_nsec: timeout.subsec_nanos() as libc::c_long
191    };
192    let result = unsafe {
193        libc::_umtx_op(
194            word.as_ptr().cast(),
195            libc::UMTX_OP_WAIT_UINT,
196            expected as libc::c_ulong,
197            size_of::<libc::timespec>() as *mut libc::c_void,
198            &left as *const libc::timespec as *mut libc::c_void
199        )
200    };
201    if result == 0 {
202        return Ok(true);
203    }
204
205    let error = io::Error::last_os_error();
206    match error.raw_os_error() {
207        Some(libc::EINTR) => Ok(true),
208        Some(libc::ETIMEDOUT) => Ok(false),
209        _ => Err(error)
210    }
211}
212
213#[cfg(target_os = "freebsd")]
214pub fn wake_word_one(word: &AtomicU32) -> io::Result<()> {
215    let result = unsafe {
216        libc::_umtx_op(
217            word.as_ptr().cast(),
218            libc::UMTX_OP_WAKE,
219            1,
220            std::ptr::null_mut(),
221            std::ptr::null_mut()
222        )
223    };
224    if result == 0 { Ok(()) } else { Err(io::Error::last_os_error()) }
225}
226
227#[cfg(target_os = "freebsd")]
228pub fn wake_word(word: &AtomicU32) -> io::Result<()> {
229    let result = unsafe {
230        libc::_umtx_op(
231            word.as_ptr().cast(),
232            libc::UMTX_OP_WAKE,
233            i32::MAX as libc::c_ulong,
234            std::ptr::null_mut(),
235            std::ptr::null_mut()
236        )
237    };
238    if result == 0 { Ok(()) } else { Err(io::Error::last_os_error()) }
239}
240
241#[cfg(target_os = "macos")]
242pub fn wait_word(
243    word: &AtomicU32,
244    expected: u32
245) -> io::Result<()> {
246    let api = macos::api().ok_or_else(|| {
247        io::Error::new(io::ErrorKind::Unsupported, "macOS shared address waits unavailable")
248    })?;
249    debug_assert_eq!(
250        (word.as_ptr() as usize) % size_of::<u32>(),
251        0,
252        "shared wait word must be naturally aligned"
253    );
254    // The SHM word is AtomicU32, not u64. Wait and wake must agree on size
255    // and shared mode. Apple returns a nonnegative waiter count on success.
256    let result = unsafe {
257        (api.wait)(word.as_ptr().cast(), u64::from(expected), size_of::<u32>(), macos::SHARED)
258    };
259    if result >= 0 {
260        return Ok(());
261    }
262    let error = io::Error::last_os_error();
263    match error.raw_os_error() {
264        Some(libc::EINTR) => Ok(()),
265        _ => Err(error)
266    }
267}
268
269/// Park until the word no longer holds `expected`, or `timeout` passes.
270///
271/// `Ok(true)` means something may have changed — a wake, a value that
272/// moved before the kernel parked us, or a signal — and the caller
273/// re-checks as it does after [`wait_word`]. `Ok(false)` means the
274/// timeout passed and nothing else. The timeout is relative and measured
275/// on a monotonic clock, so a caller holding a deadline recomputes what
276/// is left on each turn of its loop.
277#[cfg(target_os = "macos")]
278pub fn wait_word_timeout(
279    word: &AtomicU32,
280    expected: u32,
281    timeout: Duration
282) -> io::Result<bool> {
283    let api = macos::api().ok_or_else(|| {
284        io::Error::new(io::ErrorKind::Unsupported, "macOS shared address waits unavailable")
285    })?;
286    let nanos = timeout.as_nanos().min(u128::from(u64::MAX)) as u64;
287    // SAFETY: the same word, size and shared flag as the untimed wait.
288    let result = unsafe {
289        (api.wait_timeout)(
290            word.as_ptr().cast(),
291            u64::from(expected),
292            size_of::<u32>(),
293            macos::SHARED,
294            macos::MACH_ABSOLUTE_TIME,
295            nanos
296        )
297    };
298    if result >= 0 {
299        return Ok(true);
300    }
301    let error = io::Error::last_os_error();
302    match error.raw_os_error() {
303        Some(libc::EINTR) => Ok(true),
304        Some(libc::ETIMEDOUT) => Ok(false),
305        _ => Err(error)
306    }
307}
308
309#[cfg(target_os = "macos")]
310pub fn wake_word_one(word: &AtomicU32) -> io::Result<()> {
311    wake_macos(word, false)
312}
313
314#[cfg(target_os = "macos")]
315pub fn wake_word(word: &AtomicU32) -> io::Result<()> {
316    wake_macos(word, true)
317}
318
319#[cfg(target_os = "macos")]
320fn wake_macos(
321    word: &AtomicU32,
322    all: bool
323) -> io::Result<()> {
324    debug_assert_eq!(
325        (word.as_ptr() as usize) % size_of::<u32>(),
326        0,
327        "shared wake word must be naturally aligned"
328    );
329    // Older macOS has no native subscribers; publication still succeeds and
330    // polling readers see the committed frames.
331    let Some(api) = macos::api() else {
332        return Ok(());
333    };
334    loop {
335        let wake = if all { api.wake_all } else { api.wake_one };
336        let result = unsafe { wake(word.as_ptr().cast(), size_of::<u32>(), macos::SHARED) };
337        if result >= 0 {
338            return Ok(());
339        }
340        let error = io::Error::last_os_error();
341        match error.raw_os_error() {
342            // No waiter is normal: publication can precede subscription or
343            // race the driver's compare-and-wait. The generation persists.
344            Some(libc::ENOENT) => return Ok(()),
345            Some(libc::EINTR) => continue,
346            _ => return Err(error)
347        }
348    }
349}
350
351#[cfg(target_os = "macos")]
352pub(crate) mod macos {
353    use std::sync::atomic::{AtomicU8, AtomicUsize, Ordering};
354
355    // OS_SYNC_WAIT_ON_ADDRESS_SHARED and OS_SYNC_WAKE_BY_ADDRESS_SHARED
356    // have the same ABI value in <os/os_sync_wait_on_address.h>.
357    pub(super) const SHARED: u32 = 1;
358
359    /// `os_clockid_t` in <os/clock.h>: the only clock the timed wait
360    /// takes, and the one a relative timeout is measured on.
361    pub(super) const MACH_ABSOLUTE_TIME: u32 = 32;
362
363    type Wait = unsafe extern "C" fn(*mut libc::c_void, u64, usize, u32) -> libc::c_int;
364    type WaitTimeout =
365        unsafe extern "C" fn(*mut libc::c_void, u64, usize, u32, u32, u64) -> libc::c_int;
366    type Wake = unsafe extern "C" fn(*mut libc::c_void, usize, u32) -> libc::c_int;
367
368    #[derive(Clone, Copy)]
369    pub(crate) struct Api {
370        pub(super) wait: Wait,
371        pub(super) wait_timeout: WaitTimeout,
372        pub(super) wake_one: Wake,
373        pub(super) wake_all: Wake
374    }
375
376    const UNRESOLVED: u8 = 0;
377    const UNAVAILABLE: u8 = 1;
378    const READY: u8 = 2;
379
380    static STATE: AtomicU8 = AtomicU8::new(UNRESOLVED);
381    static WAIT: AtomicUsize = AtomicUsize::new(0);
382    static WAIT_TIMEOUT: AtomicUsize = AtomicUsize::new(0);
383    static WAKE_ONE: AtomicUsize = AtomicUsize::new(0);
384    static WAKE_ALL: AtomicUsize = AtomicUsize::new(0);
385
386    /// Resolve the 14.4 entry points, once per process but never by waiting.
387    ///
388    /// Deliberately not a `OnceLock`. A publisher reaches this from
389    /// `wake_all_generation_waiters`, and a publisher may be a process that
390    /// was just forked while another thread of its parent was inside the
391    /// resolution -- the test harness runs tests on parallel threads, and a
392    /// supervisor forks workers with readiness threads alive. The child
393    /// inherits the "initializing" state and none of the thread that would
394    /// finish it, so a blocking once-cell parks forever. Resolution here is
395    /// idempotent: racing callers each `dlsym` the same two symbols and store
396    /// the same values, and nobody waits for anybody.
397    ///
398    /// Resolving lazily rather than linking avoids hard references to
399    /// 14.4-only symbols on older deployment targets. libSystem stays loaded
400    /// for the life of the process, so the pointers never dangle.
401    pub(crate) fn api() -> Option<Api> {
402        match STATE.load(Ordering::Acquire) {
403            READY => Some(load()),
404            UNAVAILABLE => None,
405            _ => resolve()
406        }
407    }
408
409    fn load() -> Api {
410        unsafe {
411            Api {
412                wait: std::mem::transmute::<usize, Wait>(WAIT.load(Ordering::Acquire)),
413                wait_timeout: std::mem::transmute::<usize, WaitTimeout>(
414                    WAIT_TIMEOUT.load(Ordering::Acquire)
415                ),
416                wake_one: std::mem::transmute::<usize, Wake>(WAKE_ONE.load(Ordering::Acquire)),
417                wake_all: std::mem::transmute::<usize, Wake>(WAKE_ALL.load(Ordering::Acquire))
418            }
419        }
420    }
421
422    fn resolve() -> Option<Api> {
423        let wait = unsafe { libc::dlsym(libc::RTLD_DEFAULT, c"os_sync_wait_on_address".as_ptr()) };
424        // Shipped in the same release as the other two; all or nothing.
425        let wait_timeout = unsafe {
426            libc::dlsym(libc::RTLD_DEFAULT, c"os_sync_wait_on_address_with_timeout".as_ptr())
427        };
428        let wake_one =
429            unsafe { libc::dlsym(libc::RTLD_DEFAULT, c"os_sync_wake_by_address_any".as_ptr()) };
430        let wake_all =
431            unsafe { libc::dlsym(libc::RTLD_DEFAULT, c"os_sync_wake_by_address_all".as_ptr()) };
432        if wait.is_null() || wait_timeout.is_null() || wake_one.is_null() || wake_all.is_null() {
433            STATE.store(UNAVAILABLE, Ordering::Release);
434            return None;
435        }
436        // Pointers first, then the state that publishes them.
437        WAIT.store(wait as usize, Ordering::Release);
438        WAIT_TIMEOUT.store(wait_timeout as usize, Ordering::Release);
439        WAKE_ONE.store(wake_one as usize, Ordering::Release);
440        WAKE_ALL.store(wake_all as usize, Ordering::Release);
441        STATE.store(READY, Ordering::Release);
442        Some(load())
443    }
444}
445
446#[cfg(test)]
447mod timeout_tests {
448    use std::sync::Arc;
449    use std::sync::atomic::Ordering;
450    use std::time::Instant;
451
452    use super::*;
453
454    /// Also pins the unit the platform reads the timeout in: a wrong one
455    /// shows up here as a wait that is orders of magnitude off, not as a
456    /// wrong answer.
457    #[test]
458    fn a_timeout_is_a_timeout() {
459        if !supported() {
460            return;
461        }
462        let word = AtomicU32::new(7);
463        let started = Instant::now();
464        assert!(!wait_word_timeout(&word, 7, Duration::from_millis(200)).expect("wait"));
465        let waited = started.elapsed();
466        assert!(waited >= Duration::from_millis(150), "returned after {waited:?}");
467        assert!(waited < Duration::from_secs(2), "returned after {waited:?}");
468    }
469
470    #[test]
471    fn a_wake_beats_the_timeout() {
472        if !supported() {
473            return;
474        }
475        let word = Arc::new(AtomicU32::new(0));
476        let waker = Arc::clone(&word);
477        std::thread::spawn(move || {
478            std::thread::sleep(Duration::from_millis(50));
479            waker.store(1, Ordering::SeqCst);
480            let _ = wake_word(&waker);
481        });
482        let started = Instant::now();
483        assert!(wait_word_timeout(&word, 0, Duration::from_secs(10)).expect("wait"));
484        assert!(started.elapsed() < Duration::from_secs(5));
485    }
486
487    #[test]
488    fn wake_one_releases_only_one_parked_waiter() {
489        if !supported() {
490            return;
491        }
492        let word = Arc::new(AtomicU32::new(0));
493        let (sent, received) = std::sync::mpsc::channel();
494        let waiters = (0..2)
495            .map(|_| {
496                let word = Arc::clone(&word);
497                let sent = sent.clone();
498                std::thread::spawn(move || {
499                    let outcome =
500                        wait_word_timeout(&word, 0, Duration::from_millis(500)).expect("wait");
501                    sent.send(outcome).unwrap();
502                })
503            })
504            .collect::<Vec<_>>();
505
506        std::thread::sleep(Duration::from_millis(100));
507        word.store(1, Ordering::SeqCst);
508        wake_word_one(&word).expect("wake one");
509
510        assert!(received.recv_timeout(Duration::from_millis(200)).unwrap());
511        assert!(received.recv_timeout(Duration::from_millis(100)).is_err());
512        assert!(!received.recv_timeout(Duration::from_millis(400)).unwrap());
513        for waiter in waiters {
514            waiter.join().unwrap();
515        }
516    }
517
518    #[test]
519    fn a_word_that_already_moved_does_not_park_at_all() {
520        if !supported() {
521            return;
522        }
523        let word = AtomicU32::new(3);
524        let started = Instant::now();
525        assert!(wait_word_timeout(&word, 9, Duration::from_secs(30)).expect("wait"));
526        assert!(started.elapsed() < Duration::from_secs(1));
527    }
528}