Skip to main content

orbit_core/ring/
readiness.rs

1//! Native readiness bridge for an SHM ring.
2//!
3//! The shared signal is a generation in the ring header, waited through Linux
4//! futex, FreeBSD umtx, or macOS shared address waits. Each process owns a
5//! private readiness fd (`eventfd`, or a pipe on macOS) and a small
6//! blocking driver thread that converts generation changes into fd readiness.
7//! Async runtimes can therefore wait on their normal reactor without sharing
8//! one drainable eventfd across readers.
9
10#![cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
11
12use std::fmt;
13use std::io;
14use std::os::fd::{AsFd, AsRawFd, BorrowedFd, RawFd};
15use std::sync::Arc;
16use std::sync::atomic::{AtomicBool, Ordering};
17use std::thread::JoinHandle;
18
19use super::shm::ShmRing;
20use crate::readiness::{Readiness, Signal};
21
22/// Process-local fd readiness bridge for one shared Orbit ring.
23///
24/// The established name is retained on macOS, where a nonblocking pipe backs
25/// the fd. Native macOS readiness requires 14.4 or later; older systems return
26/// `ErrorKind::Unsupported` from construction so callers can use polling.
27///
28/// Every subscribing process creates its own instance. Ring publishers bump a
29/// generation stored in SHM and wake all platform waiters; the local driver
30/// then marks this fd readable. Multiple publishes may coalesce into one wake,
31/// so a consumer must drain the fd and poll the ring through its own cursor.
32pub struct RingEventFd {
33    fd: Readiness,
34    ring: Arc<ShmRing>,
35    stop: Arc<AtomicBool>,
36    driver: Option<JoinHandle<()>>,
37}
38
39impl RingEventFd {
40    pub(crate) fn new(ring: Arc<ShmRing>) -> io::Result<Self> {
41        let (fd, driver_fd) = local_notification_pair()?;
42        let stop = Arc::new(AtomicBool::new(false));
43        let driver_stop = stop.clone();
44        let driver_ring = ring.clone();
45        let mut observed = ring.notification_generation().load(Ordering::Acquire);
46
47        // A pipe has two distinct ends; an eventfd is one descriptor cloned.
48        #[cfg(target_os = "macos")]
49        debug_assert_ne!(fd.as_raw_fd(), driver_fd.as_raw_fd());
50
51        let driver = std::thread::Builder::new()
52            .name(format!("orbit-ring-{}-eventfd", ring.kind()))
53            .spawn(move || {
54                while !driver_stop.load(Ordering::Acquire) {
55                    let current = driver_ring
56                        .notification_generation()
57                        .load(Ordering::Acquire);
58                    if current != observed {
59                        observed = current;
60                        if driver_fd.signal().is_err() {
61                            break;
62                        }
63                        continue;
64                    }
65                    if crate::sync::wait_word(driver_ring.notification_generation(), observed)
66                        .is_err()
67                    {
68                        break;
69                    }
70                }
71            })?;
72
73        Ok(Self {
74            fd,
75            ring,
76            stop,
77            driver: Some(driver),
78        })
79    }
80
81    pub(crate) fn notify(ring: &ShmRing) -> io::Result<()> {
82        ring.notification_generation()
83            .fetch_add(1, Ordering::Release);
84        crate::sync::wake_word(ring.notification_generation())
85    }
86
87    /// Drain coalesced readiness tokens from the nonblocking local fd.
88    ///
89    /// Ring events themselves remain in SHM; the returned number is only the
90    /// local wake count and must not be interpreted as an event count.
91    pub fn drain(&self) -> io::Result<u64> {
92        self.fd.drain()
93    }
94}
95
96impl AsRawFd for RingEventFd {
97    fn as_raw_fd(&self) -> RawFd {
98        self.fd.as_raw_fd()
99    }
100}
101
102impl AsFd for RingEventFd {
103    fn as_fd(&self) -> BorrowedFd<'_> {
104        self.fd.as_fd()
105    }
106}
107
108impl fmt::Debug for RingEventFd {
109    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
110        f.debug_struct("RingEventFd")
111            .field("fd", &self.fd.as_raw_fd())
112            .field("ring_kind", &self.ring.kind())
113            .finish_non_exhaustive()
114    }
115}
116
117impl Drop for RingEventFd {
118    fn drop(&mut self) {
119        self.stop.store(true, Ordering::Release);
120        // Change the generation before waking. If the driver passed its stop
121        // check but has not entered the platform wait yet, the atomic compare
122        // prevents it from parking after our wake and deadlocking join.
123        self.ring
124            .notification_generation()
125            .fetch_add(1, Ordering::Release);
126        let _ = crate::sync::wake_word(self.ring.notification_generation());
127        if let Some(driver) = self.driver.take() {
128            let _ = driver.join();
129        }
130    }
131}
132
133#[cfg(any(target_os = "linux", target_os = "freebsd"))]
134fn local_notification_pair() -> io::Result<(Readiness, Signal)> {
135    crate::readiness::pair()
136}
137
138/// macOS needs 14.4 for the shared address wait the driver parks on, so a
139/// pair is refused there rather than handed out with nothing to feed it.
140#[cfg(target_os = "macos")]
141fn local_notification_pair() -> io::Result<(Readiness, Signal)> {
142    if crate::sync::macos::api().is_none() {
143        return Err(io::Error::new(
144            io::ErrorKind::Unsupported,
145            "Orbit native readiness requires macOS 14.4 or later",
146        ));
147    }
148    crate::readiness::pair()
149}