orbit_core/ring/
readiness.rs1#![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
21pub 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 #[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 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 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
125pub struct ParkedRingEventFd {
153 fd: Readiness,
154 ring: Arc<ShmRing>,
155 stop: Arc<AtomicBool>,
156 drained: Arc<Drained>,
157 driver: Option<JoinHandle<()>>
158}
159
160#[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 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 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 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 {
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#[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}