ringfire 0.5.2

Zero-copy lock-free inter-process communication (IPC) ring buffer and shared memory bus in Rust, with ring mirroring across hosts
Documentation
//! # ReaderRegistry & Reader Tracking
//!
//! Shared-memory registry for monitoring active consumers, tracking reader lag,
//! calculating safe write headroom, and automatically reclaiming dead processes.
//!
//! A slot is owned while `pid != 0` and counts as a live reader while additionally
//! `active == 1`. Ownership changes only through compare-and-swap on `pid`, so a scanner
//! reclaiming a dead reader can never wipe out a registration that replaced it.
//!
//! Liveness is judged by PID: readers in another PID namespace (e.g. a different
//! container) look dead to the producer and are reclaimed. Share the PID namespace
//! (`--pid=host` / a shared pod) when using lossless flow control across containers.

use std::sync::atomic::Ordering;
use crate::error::{Result, RingfireError};
use crate::header::ReaderSlot;

pub const DEFAULT_MAX_READERS: usize = 32;

/// Snapshot of an active reader's status for monitoring.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ReaderInfo {
    pub slot_index: usize,
    pub pid: u32,
    pub name: String,
    pub cursor_seq: u64,
    pub lag: u64,
}

/// Shared memory registry holding reader registration slots.
pub struct ReaderRegistry {
    slots: *mut ReaderSlot,
    count: usize,
}

unsafe impl Send for ReaderRegistry {}
unsafe impl Sync for ReaderRegistry {}

impl ReaderRegistry {
    /// Initialize a new `ReaderRegistry` over pre-mapped shared memory.
    ///
    /// # Safety
    /// `ptr` must point to writable memory of at least `count * size_of::<ReaderSlot>()` bytes.
    pub unsafe fn init(ptr: *mut u8, count: usize) -> Self {
        let slots = ptr as *mut ReaderSlot;
        for i in 0..count {
            let slot = unsafe { &mut *slots.add(i) };
            slot.pid.store(0, Ordering::Relaxed);
            slot.active.store(0, Ordering::Relaxed);
            slot.cursor_seq.store(0, Ordering::Relaxed);
            slot.heartbeat_tsc.store(0, Ordering::Relaxed);
            slot.name = [0u8; 32];
            slot._pad = [0u8; 8];
        }

        Self { slots, count }
    }

    /// Open an existing `ReaderRegistry` from pre-mapped shared memory.
    ///
    /// # Safety
    /// `ptr` must point to an initialized `ReaderRegistry` region of `count` slots.
    pub unsafe fn from_ptr(ptr: *mut u8, count: usize) -> Self {
        Self {
            slots: ptr as *mut ReaderSlot,
            count,
        }
    }

    /// Total number of reader slots.
    pub fn capacity(&self) -> usize {
        self.count
    }

    #[inline]
    fn slot(&self, i: usize) -> &ReaderSlot {
        unsafe { &*self.slots.add(i) }
    }

    /// Register a new reader in the shared memory registry.
    pub fn register(&self, name: &str, initial_cursor: u64) -> Result<ReaderRegistration> {
        let current_pid = std::process::id();

        for i in 0..self.count {
            let slot = self.slot(i);
            let pid = slot.pid.load(Ordering::Acquire);

            let is_free = pid == 0;
            let is_dead = pid != 0 && !is_process_alive(pid);

            if (is_free || is_dead)
                && slot
                    .pid
                    .compare_exchange(pid, current_pid, Ordering::AcqRel, Ordering::Relaxed)
                    .is_ok()
            {
                // The slot may still carry the previous owner's active flag and cursor:
                // hide it from scanners until the new cursor is in place.
                slot.active.store(0, Ordering::Relaxed);
                slot.cursor_seq.store(initial_cursor, Ordering::Relaxed);
                slot.heartbeat_tsc.store(0, Ordering::Relaxed);

                let slot_mut = unsafe { &mut *self.slots.add(i) };
                let mut name_buf = [0u8; 32];
                let bytes = name.as_bytes();
                let len = bytes.len().min(31);
                name_buf[..len].copy_from_slice(&bytes[..len]);
                slot_mut.name = name_buf;

                slot.active.store(1, Ordering::Release);

                return Ok(ReaderRegistration {
                    slot_index: i,
                    slot: self.slots,
                    pid: current_pid,
                });
            }
        }

        Err(RingfireError::NoAvailableReaderSlots)
    }

    /// Releases slot `i` if it is still owned by the dead process `pid`.
    #[inline]
    fn reclaim(&self, i: usize, pid: u32) -> bool {
        self.slot(i)
            .pid
            .compare_exchange(pid, 0, Ordering::AcqRel, Ordering::Relaxed)
            .is_ok()
    }

    /// Minimum cursor across registered readers, without checking process liveness.
    ///
    /// Cheap enough for the producer's flow-control path: one atomic load per slot, no
    /// syscalls. Dead readers are reclaimed separately by [`Self::prune_dead_readers`].
    #[inline]
    pub fn min_cursor(&self) -> Option<u64> {
        let mut min_seq: Option<u64> = None;
        for i in 0..self.count {
            let slot = self.slot(i);
            if slot.active.load(Ordering::Acquire) == 1 && slot.pid.load(Ordering::Acquire) != 0 {
                let seq = slot.cursor_seq.load(Ordering::Acquire);
                min_seq = Some(min_seq.map_or(seq, |curr| curr.min(seq)));
            }
        }
        min_seq
    }

    /// Find the minimum sequence cursor across all active, alive readers,
    /// reclaiming slots of dead processes on the way.
    pub fn min_reader_seq(&self) -> Option<u64> {
        self.prune_dead_readers();
        self.min_cursor()
    }

    /// Calculate the lag of the slowest active reader relative to `write_seq`.
    pub fn reader_lag(&self, write_seq: u64) -> u64 {
        match self.min_reader_seq() {
            Some(min_seq) => write_seq.saturating_sub(min_seq.saturating_sub(1)),
            None => 0,
        }
    }

    /// Calculate how many messages the writer can write before the slowest active reader is lapped.
    pub fn headroom(&self, write_seq: u64, ring_capacity: u64) -> u64 {
        match self.min_reader_seq() {
            Some(min_seq) => {
                let unread = write_seq.saturating_sub(min_seq.saturating_sub(1));
                ring_capacity.saturating_sub(unread)
            }
            None => ring_capacity,
        }
    }

    /// List all currently active readers and their lag.
    pub fn active_readers(&self, write_seq: u64) -> Vec<ReaderInfo> {
        self.prune_dead_readers();
        let mut result = Vec::new();

        for i in 0..self.count {
            let slot = self.slot(i);
            let pid = slot.pid.load(Ordering::Acquire);
            if pid != 0 && slot.active.load(Ordering::Acquire) == 1 {
                let cursor_seq = slot.cursor_seq.load(Ordering::Acquire);
                let name_len = slot.name.iter().position(|&b| b == 0).unwrap_or(32);
                let name = String::from_utf8_lossy(&slot.name[..name_len]).into_owned();

                result.push(ReaderInfo {
                    slot_index: i,
                    pid,
                    name,
                    cursor_seq,
                    lag: write_seq.saturating_sub(cursor_seq.saturating_sub(1)),
                });
            }
        }

        result
    }

    /// Scan and clear any abandoned slots from dead processes.
    pub fn prune_dead_readers(&self) -> usize {
        let mut pruned = 0;
        for i in 0..self.count {
            let pid = self.slot(i).pid.load(Ordering::Acquire);
            if pid != 0 && !is_process_alive(pid) && self.reclaim(i, pid) {
                pruned += 1;
            }
        }
        pruned
    }
}

/// RAII registration handle for a consumer in `ReaderRegistry`.
pub struct ReaderRegistration {
    slot_index: usize,
    slot: *mut ReaderSlot,
    pid: u32,
}

unsafe impl Send for ReaderRegistration {}
unsafe impl Sync for ReaderRegistration {}

impl ReaderRegistration {
    /// Update the reader's cursor: the sequence it will read next.
    ///
    /// Release ordering: the producer may overwrite everything below `seq` as soon as it
    /// observes this store, so all reads of those slots must have completed.
    #[inline]
    pub fn update_cursor(&self, seq: u64) {
        let slot = unsafe { &*self.slot.add(self.slot_index) };
        slot.cursor_seq.store(seq, Ordering::Release);
    }

    /// Index of this reader slot.
    #[inline]
    pub fn slot_index(&self) -> usize {
        self.slot_index
    }
}

impl Drop for ReaderRegistration {
    fn drop(&mut self) {
        let slot = unsafe { &*self.slot.add(self.slot_index) };
        slot.active.store(0, Ordering::Release);
        let _ = slot
            .pid
            .compare_exchange(self.pid, 0, Ordering::AcqRel, Ordering::Relaxed);
    }
}

/// Check if a process with `pid` is currently alive (in this PID namespace).
#[inline]
pub fn is_process_alive(pid: u32) -> bool {
    if pid == 0 || pid > i32::MAX as u32 {
        return false;
    }
    let ret = unsafe { libc::kill(pid as libc::pid_t, 0) };
    // EPERM: the process exists but belongs to another user.
    ret == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM)
}

#[cfg(test)]
#[path = "../tests/unit/registry.rs"]
mod tests;