#![cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
use std::os::fd::{AsFd, AsRawFd, BorrowedFd, RawFd};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Condvar, Mutex};
use std::thread::JoinHandle;
use std::{fmt, io};
use super::shm::ShmRing;
use crate::readiness::{Readiness, Signal};
pub struct RingEventFd {
fd: Readiness,
ring: Arc<ShmRing>,
stop: Arc<AtomicBool>,
driver: Option<JoinHandle<()>>
}
impl RingEventFd {
pub(crate) fn new(ring: Arc<ShmRing>) -> io::Result<Self> {
let (fd, driver_fd) = local_notification_pair()?;
let stop = Arc::new(AtomicBool::new(false));
let driver_stop = stop.clone();
let driver_ring = ring.clone();
let mut observed = ring.notification_generation().load(Ordering::Acquire);
#[cfg(target_os = "macos")]
debug_assert_ne!(fd.as_raw_fd(), driver_fd.as_raw_fd());
let driver = std::thread::Builder::new()
.name(format!("orbit-ring-{}-eventfd", ring.kind()))
.spawn(move || {
while !driver_stop.load(Ordering::Acquire) {
let current = driver_ring.notification_generation().load(Ordering::Acquire);
if current != observed {
observed = current;
if driver_fd.signal().is_err() {
break;
}
continue;
}
if crate::sync::wait_word(driver_ring.notification_generation(), observed)
.is_err()
{
break;
}
}
})?;
Ok(Self { fd, ring, stop, driver: Some(driver) })
}
pub(crate) fn notify(ring: &ShmRing) -> io::Result<()> {
ring.notification_generation().fetch_add(1, Ordering::Release);
crate::sync::wake_word(ring.notification_generation())
}
pub fn drain(&self) -> io::Result<u64> {
self.fd.drain()
}
}
impl AsRawFd for RingEventFd {
fn as_raw_fd(&self) -> RawFd {
self.fd.as_raw_fd()
}
}
impl AsFd for RingEventFd {
fn as_fd(&self) -> BorrowedFd<'_> {
self.fd.as_fd()
}
}
impl fmt::Debug for RingEventFd {
fn fmt(
&self,
f: &mut fmt::Formatter<'_>
) -> fmt::Result {
f.debug_struct("RingEventFd")
.field("fd", &self.fd.as_raw_fd())
.field("ring_kind", &self.ring.kind())
.finish_non_exhaustive()
}
}
impl Drop for RingEventFd {
fn drop(&mut self) {
self.stop.store(true, Ordering::Release);
self.ring.notification_generation().fetch_add(1, Ordering::Release);
let _ = crate::sync::wake_word(self.ring.notification_generation());
if let Some(driver) = self.driver.take() {
let _ = driver.join();
}
}
}
pub struct ParkedRingEventFd {
fd: Readiness,
ring: Arc<ShmRing>,
stop: Arc<AtomicBool>,
drained: Arc<Drained>,
driver: Option<JoinHandle<()>>
}
#[derive(Default)]
struct Drained {
pending: Mutex<bool>,
taken: Condvar
}
impl ParkedRingEventFd {
pub(crate) fn new(ring: Arc<ShmRing>) -> io::Result<Self> {
let (fd, driver_fd) = local_notification_pair()?;
let stop = Arc::new(AtomicBool::new(false));
let drained = Arc::new(Drained::default());
let (driver_stop, driver_drained, driver_ring) =
(stop.clone(), drained.clone(), ring.clone());
let mut observed = ring.notification_generation().load(Ordering::Acquire);
let driver = std::thread::Builder::new()
.name(format!("orbit-ring-{}-parked", ring.kind()))
.spawn(move || {
let generation = driver_ring.notification_generation();
let waiters = driver_ring.notification_waiters();
while !driver_stop.load(Ordering::Acquire) {
waiters.fetch_add(1, Ordering::SeqCst);
let current = generation.load(Ordering::SeqCst);
if current == observed && !driver_stop.load(Ordering::Acquire) {
let parked = crate::sync::wait_word(generation, observed);
waiters.fetch_sub(1, Ordering::SeqCst);
if parked.is_err() {
break;
}
continue;
}
waiters.fetch_sub(1, Ordering::SeqCst);
if driver_stop.load(Ordering::Acquire) {
break;
}
observed = current;
*driver_drained.pending.lock().unwrap() = true;
if driver_fd.signal().is_err() {
break;
}
let mut pending = driver_drained.pending.lock().unwrap();
while *pending && !driver_stop.load(Ordering::Acquire) {
pending = driver_drained.taken.wait(pending).unwrap();
}
}
})?;
Ok(Self { fd, ring, stop, drained, driver: Some(driver) })
}
pub(crate) fn notify(ring: &ShmRing) -> io::Result<()> {
ring.notification_generation().fetch_add(1, Ordering::SeqCst);
if ring.notification_waiters().load(Ordering::SeqCst) == 0 {
return Ok(());
}
crate::sync::wake_word(ring.notification_generation())
}
pub fn drain(&self) -> io::Result<u64> {
let tokens = self.fd.drain()?;
*self.drained.pending.lock().unwrap() = false;
self.drained.taken.notify_one();
Ok(tokens)
}
}
impl AsRawFd for ParkedRingEventFd {
fn as_raw_fd(&self) -> RawFd {
self.fd.as_raw_fd()
}
}
impl AsFd for ParkedRingEventFd {
fn as_fd(&self) -> BorrowedFd<'_> {
self.fd.as_fd()
}
}
impl fmt::Debug for ParkedRingEventFd {
fn fmt(
&self,
f: &mut fmt::Formatter<'_>
) -> fmt::Result {
f.debug_struct("ParkedRingEventFd")
.field("fd", &self.fd.as_raw_fd())
.field("ring_kind", &self.ring.kind())
.finish_non_exhaustive()
}
}
impl Drop for ParkedRingEventFd {
fn drop(&mut self) {
self.stop.store(true, Ordering::Release);
{
let _pending = self.drained.pending.lock().unwrap();
self.drained.taken.notify_one();
}
self.ring.notification_generation().fetch_add(1, Ordering::SeqCst);
let _ = crate::sync::wake_word(self.ring.notification_generation());
if let Some(driver) = self.driver.take() {
let _ = driver.join();
}
}
}
#[cfg(any(target_os = "linux", target_os = "freebsd"))]
fn local_notification_pair() -> io::Result<(Readiness, Signal)> {
crate::readiness::pair()
}
#[cfg(target_os = "macos")]
fn local_notification_pair() -> io::Result<(Readiness, Signal)> {
if crate::sync::macos::api().is_none() {
return Err(io::Error::new(
io::ErrorKind::Unsupported,
"Orbit native readiness requires macOS 14.4 or later"
));
}
crate::readiness::pair()
}