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`] wakes every waiter on it. The
6//! word may live in shared memory: Linux futex, FreeBSD umtx and macOS
7//! `os_sync_wait_on_address` all key waiters by the physical location, so a
8//! wake in one process reaches a waiter in another with nothing carried
9//! between them. No descriptor, no channel: the memory is the signal.
10//!
11//! macOS needs 14.4 for the shared form; below that [`wait_word`] answers
12//! `ErrorKind::Unsupported` and a caller polls instead.
13
14#![cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
15
16use std::io;
17#[cfg(target_os = "macos")]
18use std::mem::size_of;
19use std::sync::atomic::AtomicU32;
20
21#[cfg(target_os = "linux")]
22pub fn wait_word(word: &AtomicU32, expected: u32) -> io::Result<()> {
23    let result = unsafe {
24        libc::syscall(
25            libc::SYS_futex,
26            word.as_ptr(),
27            libc::FUTEX_WAIT,
28            expected,
29            std::ptr::null::<libc::timespec>(),
30            std::ptr::null::<u32>(),
31            0,
32        )
33    };
34    if result == 0 {
35        return Ok(());
36    }
37
38    let error = io::Error::last_os_error();
39    match error.raw_os_error() {
40        // The generation changed before the kernel parked us, or the driver
41        // was interrupted. The outer loop re-checks both generation and stop.
42        Some(libc::EAGAIN) | Some(libc::EINTR) => Ok(()),
43        _ => Err(error),
44    }
45}
46
47#[cfg(target_os = "linux")]
48pub fn wake_word(word: &AtomicU32) -> io::Result<()> {
49    let result = unsafe {
50        libc::syscall(
51            libc::SYS_futex,
52            word.as_ptr(),
53            libc::FUTEX_WAKE,
54            i32::MAX,
55            std::ptr::null::<libc::timespec>(),
56            std::ptr::null::<u32>(),
57            0,
58        )
59    };
60    if result >= 0 {
61        Ok(())
62    } else {
63        Err(io::Error::last_os_error())
64    }
65}
66
67#[cfg(target_os = "freebsd")]
68pub fn wait_word(word: &AtomicU32, expected: u32) -> io::Result<()> {
69    let result = unsafe {
70        libc::_umtx_op(
71            word.as_ptr().cast(),
72            libc::UMTX_OP_WAIT_UINT,
73            expected as libc::c_ulong,
74            std::ptr::null_mut(),
75            std::ptr::null_mut(),
76        )
77    };
78    if result == 0 {
79        return Ok(());
80    }
81
82    let error = io::Error::last_os_error();
83    match error.raw_os_error() {
84        // The generation changed before the kernel parked us, or the driver
85        // was interrupted. The outer loop re-checks generation and stop.
86        Some(libc::EINTR) => Ok(()),
87        _ => Err(error),
88    }
89}
90
91#[cfg(target_os = "freebsd")]
92pub fn wake_word(word: &AtomicU32) -> io::Result<()> {
93    let result = unsafe {
94        libc::_umtx_op(
95            word.as_ptr().cast(),
96            libc::UMTX_OP_WAKE,
97            i32::MAX as libc::c_ulong,
98            std::ptr::null_mut(),
99            std::ptr::null_mut(),
100        )
101    };
102    if result == 0 {
103        Ok(())
104    } else {
105        Err(io::Error::last_os_error())
106    }
107}
108
109#[cfg(target_os = "macos")]
110pub fn wait_word(word: &AtomicU32, expected: u32) -> io::Result<()> {
111    let api = macos::api().ok_or_else(|| {
112        io::Error::new(
113            io::ErrorKind::Unsupported,
114            "macOS shared address waits unavailable",
115        )
116    })?;
117    debug_assert_eq!(
118        (word.as_ptr() as usize) % size_of::<u32>(),
119        0,
120        "shared wait word must be naturally aligned"
121    );
122    // The SHM word is AtomicU32, not u64. Wait and wake must agree on size
123    // and shared mode. Apple returns a nonnegative waiter count on success.
124    let result = unsafe {
125        (api.wait)(
126            word.as_ptr().cast(),
127            u64::from(expected),
128            size_of::<u32>(),
129            macos::SHARED,
130        )
131    };
132    if result >= 0 {
133        return Ok(());
134    }
135    let error = io::Error::last_os_error();
136    match error.raw_os_error() {
137        Some(libc::EINTR) => Ok(()),
138        _ => Err(error),
139    }
140}
141
142#[cfg(target_os = "macos")]
143pub fn wake_word(word: &AtomicU32) -> io::Result<()> {
144    debug_assert_eq!(
145        (word.as_ptr() as usize) % size_of::<u32>(),
146        0,
147        "shared wake word must be naturally aligned"
148    );
149    // Older macOS has no native subscribers; publication still succeeds and
150    // polling readers see the committed frames.
151    let Some(api) = macos::api() else {
152        return Ok(());
153    };
154    loop {
155        let result =
156            unsafe { (api.wake_all)(word.as_ptr().cast(), size_of::<u32>(), macos::SHARED) };
157        if result >= 0 {
158            return Ok(());
159        }
160        let error = io::Error::last_os_error();
161        match error.raw_os_error() {
162            // No waiter is normal: publication can precede subscription or
163            // race the driver's compare-and-wait. The generation persists.
164            Some(libc::ENOENT) => return Ok(()),
165            Some(libc::EINTR) => continue,
166            _ => return Err(error),
167        }
168    }
169}
170
171#[cfg(target_os = "macos")]
172pub(crate) mod macos {
173    use std::sync::atomic::{AtomicU8, AtomicUsize, Ordering};
174
175    // OS_SYNC_WAIT_ON_ADDRESS_SHARED and OS_SYNC_WAKE_BY_ADDRESS_SHARED
176    // have the same ABI value in <os/os_sync_wait_on_address.h>.
177    pub(super) const SHARED: u32 = 1;
178
179    type Wait = unsafe extern "C" fn(*mut libc::c_void, u64, usize, u32) -> libc::c_int;
180    type Wake = unsafe extern "C" fn(*mut libc::c_void, usize, u32) -> libc::c_int;
181
182    #[derive(Clone, Copy)]
183    pub(crate) struct Api {
184        pub(super) wait: Wait,
185        pub(super) wake_all: Wake,
186    }
187
188    const UNRESOLVED: u8 = 0;
189    const UNAVAILABLE: u8 = 1;
190    const READY: u8 = 2;
191
192    static STATE: AtomicU8 = AtomicU8::new(UNRESOLVED);
193    static WAIT: AtomicUsize = AtomicUsize::new(0);
194    static WAKE: AtomicUsize = AtomicUsize::new(0);
195
196    /// Resolve the 14.4 entry points, once per process but never by waiting.
197    ///
198    /// Deliberately not a `OnceLock`. A publisher reaches this from
199    /// `wake_all_generation_waiters`, and a publisher may be a process that
200    /// was just forked while another thread of its parent was inside the
201    /// resolution -- the test harness runs tests on parallel threads, and a
202    /// supervisor forks workers with readiness threads alive. The child
203    /// inherits the "initializing" state and none of the thread that would
204    /// finish it, so a blocking once-cell parks forever. Resolution here is
205    /// idempotent: racing callers each `dlsym` the same two symbols and store
206    /// the same values, and nobody waits for anybody.
207    ///
208    /// Resolving lazily rather than linking avoids hard references to
209    /// 14.4-only symbols on older deployment targets. libSystem stays loaded
210    /// for the life of the process, so the pointers never dangle.
211    pub(crate) fn api() -> Option<Api> {
212        match STATE.load(Ordering::Acquire) {
213            READY => Some(load()),
214            UNAVAILABLE => None,
215            _ => resolve(),
216        }
217    }
218
219    fn load() -> Api {
220        unsafe {
221            Api {
222                wait: std::mem::transmute::<usize, Wait>(WAIT.load(Ordering::Acquire)),
223                wake_all: std::mem::transmute::<usize, Wake>(WAKE.load(Ordering::Acquire)),
224            }
225        }
226    }
227
228    fn resolve() -> Option<Api> {
229        let wait = unsafe { libc::dlsym(libc::RTLD_DEFAULT, c"os_sync_wait_on_address".as_ptr()) };
230        let wake =
231            unsafe { libc::dlsym(libc::RTLD_DEFAULT, c"os_sync_wake_by_address_all".as_ptr()) };
232        if wait.is_null() || wake.is_null() {
233            STATE.store(UNAVAILABLE, Ordering::Release);
234            return None;
235        }
236        // Pointers first, then the state that publishes them.
237        WAIT.store(wait as usize, Ordering::Release);
238        WAKE.store(wake as usize, Ordering::Release);
239        STATE.store(READY, Ordering::Release);
240        Some(load())
241    }
242}