use std::sync::atomic::Ordering;
use crate::error::{Result, RingfireError};
use crate::header::ReaderSlot;
pub const DEFAULT_MAX_READERS: usize = 32;
#[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,
}
pub struct ReaderRegistry {
slots: *mut ReaderSlot,
count: usize,
}
unsafe impl Send for ReaderRegistry {}
unsafe impl Sync for ReaderRegistry {}
impl ReaderRegistry {
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 }
}
pub unsafe fn from_ptr(ptr: *mut u8, count: usize) -> Self {
Self {
slots: ptr as *mut ReaderSlot,
count,
}
}
pub fn capacity(&self) -> usize {
self.count
}
#[inline]
fn slot(&self, i: usize) -> &ReaderSlot {
unsafe { &*self.slots.add(i) }
}
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()
{
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)
}
#[inline]
fn reclaim(&self, i: usize, pid: u32) -> bool {
self.slot(i)
.pid
.compare_exchange(pid, 0, Ordering::AcqRel, Ordering::Relaxed)
.is_ok()
}
#[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
}
pub fn min_reader_seq(&self) -> Option<u64> {
self.prune_dead_readers();
self.min_cursor()
}
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,
}
}
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,
}
}
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
}
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
}
}
pub struct ReaderRegistration {
slot_index: usize,
slot: *mut ReaderSlot,
pid: u32,
}
unsafe impl Send for ReaderRegistration {}
unsafe impl Sync for ReaderRegistration {}
impl ReaderRegistration {
#[inline]
pub fn update_cursor(&self, seq: u64) {
let slot = unsafe { &*self.slot.add(self.slot_index) };
slot.cursor_seq.store(seq, Ordering::Release);
}
#[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);
}
}
#[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) };
ret == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM)
}
#[cfg(test)]
#[path = "../tests/unit/registry.rs"]
mod tests;