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::os::fd::{AsFd, AsRawFd, BorrowedFd, RawFd};
13use std::sync::atomic::{AtomicBool, Ordering};
14use std::sync::{Arc, Condvar, Mutex};
15use std::thread::JoinHandle;
16use std::{fmt, io};
17
18use super::shm::ShmRing;
19use crate::readiness::{Readiness, Signal};
20
21/// Process-local fd readiness bridge for one shared Orbit ring.
22///
23/// The established name is retained on macOS, where a nonblocking pipe backs
24/// the fd. Native macOS readiness requires 14.4 or later; older systems return
25/// `ErrorKind::Unsupported` from construction so callers can use polling.
26///
27/// Every subscribing process creates its own instance. Ring publishers bump a
28/// generation stored in SHM and wake all platform waiters; the local driver
29/// then marks this fd readable. Multiple publishes may coalesce into one wake,
30/// so a consumer must drain the fd and poll the ring through its own cursor.
31pub struct RingEventFd {
32    fd: Readiness,
33    ring: Arc<ShmRing>,
34    stop: Arc<AtomicBool>,
35    driver: Option<JoinHandle<()>>
36}
37
38impl RingEventFd {
39    pub(crate) fn new(ring: Arc<ShmRing>) -> io::Result<Self> {
40        let (fd, driver_fd) = local_notification_pair()?;
41        let stop = Arc::new(AtomicBool::new(false));
42        let driver_stop = stop.clone();
43        let driver_ring = ring.clone();
44        let mut observed = ring.notification_generation().load(Ordering::Acquire);
45
46        // A pipe has two distinct ends; an eventfd is one descriptor cloned.
47        #[cfg(target_os = "macos")]
48        debug_assert_ne!(fd.as_raw_fd(), driver_fd.as_raw_fd());
49
50        let driver = std::thread::Builder::new()
51            .name(format!("orbit-ring-{}-eventfd", ring.kind()))
52            .spawn(move || {
53                while !driver_stop.load(Ordering::Acquire) {
54                    let current = driver_ring.notification_generation().load(Ordering::Acquire);
55                    if current != observed {
56                        observed = current;
57                        if driver_fd.signal().is_err() {
58                            break;
59                        }
60                        continue;
61                    }
62                    if crate::sync::wait_word(driver_ring.notification_generation(), observed)
63                        .is_err()
64                    {
65                        break;
66                    }
67                }
68            })?;
69
70        Ok(Self { fd, ring, stop, driver: Some(driver) })
71    }
72
73    pub(crate) fn notify(ring: &ShmRing) -> io::Result<()> {
74        ring.notification_generation().fetch_add(1, Ordering::Release);
75        crate::sync::wake_word(ring.notification_generation())
76    }
77
78    /// Drain coalesced readiness tokens from the nonblocking local fd.
79    ///
80    /// Ring events themselves remain in SHM; the returned number is only the
81    /// local wake count and must not be interpreted as an event count.
82    pub fn drain(&self) -> io::Result<u64> {
83        self.fd.drain()
84    }
85}
86
87impl AsRawFd for RingEventFd {
88    fn as_raw_fd(&self) -> RawFd {
89        self.fd.as_raw_fd()
90    }
91}
92
93impl AsFd for RingEventFd {
94    fn as_fd(&self) -> BorrowedFd<'_> {
95        self.fd.as_fd()
96    }
97}
98
99impl fmt::Debug for RingEventFd {
100    fn fmt(
101        &self,
102        f: &mut fmt::Formatter<'_>
103    ) -> fmt::Result {
104        f.debug_struct("RingEventFd")
105            .field("fd", &self.fd.as_raw_fd())
106            .field("ring_kind", &self.ring.kind())
107            .finish_non_exhaustive()
108    }
109}
110
111impl Drop for RingEventFd {
112    fn drop(&mut self) {
113        self.stop.store(true, Ordering::Release);
114        // Change the generation before waking. If the driver passed its stop
115        // check but has not entered the platform wait yet, the atomic compare
116        // prevents it from parking after our wake and deadlocking join.
117        self.ring.notification_generation().fetch_add(1, Ordering::Release);
118        let _ = crate::sync::wake_word(self.ring.notification_generation());
119        if let Some(driver) = self.driver.take() {
120            let _ = driver.join();
121        }
122    }
123}
124
125/// Readiness for a ring whose publishers wake only parked readers.
126///
127/// [`RingEventFd`]'s driver parks on the generation again as soon as it has
128/// marked its fd readable, so a publisher cannot tell a busy reader from an
129/// idle one and must wake on every publish: one syscall per frame, and one
130/// driver wake and fd signal on the reading side. This driver instead waits
131/// for its consumer to [`drain`](Self::drain) before it parks again, and
132/// counts itself in the ring header while it is parked. A publisher using
133/// [`Fleet::publish_notified_parked`](crate::Fleet::publish_notified_parked)
134/// reads that count and wakes only when it is non-zero — under load the
135/// reader is busy draining and publishing costs no syscall at all.
136///
137/// Every reader of a ring published with the parked variant must use this
138/// type: a [`RingEventFd`] driver does not count itself and would not be
139/// woken. Rings published with [`Fleet::publish_notified`](crate::Fleet)
140/// are unaffected by this type's existence.
141///
142/// The contract is [`RingEventFd`]'s plus one step: drain the fd, then poll
143/// the ring through your own cursor. Draining is what lets the driver park.
144///
145/// **Clear the reactor's readiness before draining, not after.** The driver
146/// signals once per drain: a signal that arrives just after [`drain`](Self::drain)
147/// is the only one until the next drain, and a reactor that clears readiness
148/// after draining (tokio's `clear_ready`) erases it — the consumer then waits
149/// for an edge that never comes, and the driver for a drain that never
150/// happens. [`RingEventFd`] tolerates either order because its driver
151/// signals on every publish.
152pub struct ParkedRingEventFd {
153    fd: Readiness,
154    ring: Arc<ShmRing>,
155    stop: Arc<AtomicBool>,
156    drained: Arc<Drained>,
157    driver: Option<JoinHandle<()>>
158}
159
160/// Whether the consumer has taken the last signal. The driver waits here,
161/// in this process, until it has — not on the shared word.
162#[derive(Default)]
163struct Drained {
164    pending: Mutex<bool>,
165    taken: Condvar
166}
167
168impl ParkedRingEventFd {
169    pub(crate) fn new(ring: Arc<ShmRing>) -> io::Result<Self> {
170        let (fd, driver_fd) = local_notification_pair()?;
171        let stop = Arc::new(AtomicBool::new(false));
172        let drained = Arc::new(Drained::default());
173        let (driver_stop, driver_drained, driver_ring) =
174            (stop.clone(), drained.clone(), ring.clone());
175        let mut observed = ring.notification_generation().load(Ordering::Acquire);
176
177        let driver = std::thread::Builder::new()
178            .name(format!("orbit-ring-{}-parked", ring.kind()))
179            .spawn(move || {
180                let generation = driver_ring.notification_generation();
181                let waiters = driver_ring.notification_waiters();
182                while !driver_stop.load(Ordering::Acquire) {
183                    // Count in, then look once more. A publisher bumps the
184                    // generation before it reads the count; one of the two
185                    // sees the other, so a publish cannot fall between them.
186                    waiters.fetch_add(1, Ordering::SeqCst);
187                    let current = generation.load(Ordering::SeqCst);
188                    if current == observed && !driver_stop.load(Ordering::Acquire) {
189                        let parked = crate::sync::wait_word(generation, observed);
190                        waiters.fetch_sub(1, Ordering::SeqCst);
191                        if parked.is_err() {
192                            break;
193                        }
194                        continue;
195                    }
196                    waiters.fetch_sub(1, Ordering::SeqCst);
197                    if driver_stop.load(Ordering::Acquire) {
198                        break;
199                    }
200                    observed = current;
201                    *driver_drained.pending.lock().unwrap() = true;
202                    if driver_fd.signal().is_err() {
203                        break;
204                    }
205                    let mut pending = driver_drained.pending.lock().unwrap();
206                    while *pending && !driver_stop.load(Ordering::Acquire) {
207                        pending = driver_drained.taken.wait(pending).unwrap();
208                    }
209                }
210            })?;
211
212        Ok(Self { fd, ring, stop, drained, driver: Some(driver) })
213    }
214
215    /// Wake parked readers of `ring`, and only them.
216    pub(crate) fn notify(ring: &ShmRing) -> io::Result<()> {
217        ring.notification_generation().fetch_add(1, Ordering::SeqCst);
218        if ring.notification_waiters().load(Ordering::SeqCst) == 0 {
219            return Ok(());
220        }
221        crate::sync::wake_word(ring.notification_generation())
222    }
223
224    /// Drain the fd and let the driver park again. Poll the ring after this.
225    pub fn drain(&self) -> io::Result<u64> {
226        let tokens = self.fd.drain()?;
227        *self.drained.pending.lock().unwrap() = false;
228        self.drained.taken.notify_one();
229        Ok(tokens)
230    }
231}
232
233impl AsRawFd for ParkedRingEventFd {
234    fn as_raw_fd(&self) -> RawFd {
235        self.fd.as_raw_fd()
236    }
237}
238
239impl AsFd for ParkedRingEventFd {
240    fn as_fd(&self) -> BorrowedFd<'_> {
241        self.fd.as_fd()
242    }
243}
244
245impl fmt::Debug for ParkedRingEventFd {
246    fn fmt(
247        &self,
248        f: &mut fmt::Formatter<'_>
249    ) -> fmt::Result {
250        f.debug_struct("ParkedRingEventFd")
251            .field("fd", &self.fd.as_raw_fd())
252            .field("ring_kind", &self.ring.kind())
253            .finish_non_exhaustive()
254    }
255}
256
257impl Drop for ParkedRingEventFd {
258    fn drop(&mut self) {
259        self.stop.store(true, Ordering::Release);
260        // Release a driver waiting for a drain, then one parked on the word:
261        // the generation changes first so a driver between its check and its
262        // wait does not park after the wake.
263        {
264            let _pending = self.drained.pending.lock().unwrap();
265            self.drained.taken.notify_one();
266        }
267        self.ring.notification_generation().fetch_add(1, Ordering::SeqCst);
268        let _ = crate::sync::wake_word(self.ring.notification_generation());
269        if let Some(driver) = self.driver.take() {
270            let _ = driver.join();
271        }
272    }
273}
274
275#[cfg(any(target_os = "linux", target_os = "freebsd"))]
276fn local_notification_pair() -> io::Result<(Readiness, Signal)> {
277    crate::readiness::pair()
278}
279
280/// macOS needs 14.4 for the shared address wait the driver parks on, so a
281/// pair is refused there rather than handed out with nothing to feed it.
282#[cfg(target_os = "macos")]
283fn local_notification_pair() -> io::Result<(Readiness, Signal)> {
284    if crate::sync::macos::api().is_none() {
285        return Err(io::Error::new(
286            io::ErrorKind::Unsupported,
287            "Orbit native readiness requires macOS 14.4 or later"
288        ));
289    }
290    crate::readiness::pair()
291}