use crate::ring::Frame;
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum RingRead {
Ready(Frame),
Pending,
Unavailable
}
pub trait RingFrameSource {
fn kind(&self) -> u8;
fn head(&self) -> u64;
fn capacity(&self) -> usize;
fn read_at(
&self,
counter: u64
) -> Option<Frame>;
fn read_state_at(
&self,
counter: u64
) -> RingRead {
self.read_at(counter).map_or(RingRead::Unavailable, RingRead::Ready)
}
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct RingCursor {
next_counter: u64
}
impl RingCursor {
pub const fn from_start() -> Self {
Self { next_counter: 0 }
}
pub const fn from_counter(next_counter: u64) -> Self {
Self { next_counter }
}
pub const fn next_counter(self) -> u64 {
self.next_counter
}
pub(crate) fn set_next_counter(
&mut self,
next_counter: u64
) {
self.next_counter = next_counter;
}
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct RingLoss {
pub overwritten: u64,
pub unavailable: u64
}
impl RingLoss {
pub const fn total(self) -> u64 {
self.overwritten.saturating_add(self.unavailable)
}
pub const fn is_empty(self) -> bool {
self.overwritten == 0 && self.unavailable == 0
}
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct RingPoll {
pub frames: Vec<Frame>,
pub loss: RingLoss,
pub from_counter: u64,
pub to_counter: u64
}
impl RingPoll {
pub fn is_empty(&self) -> bool {
self.frames.is_empty() && self.loss.is_empty()
}
}
pub fn poll_ring<S: RingFrameSource>(
source: &S,
cursor: &mut RingCursor
) -> RingPoll {
let head = source.head();
let from_counter = cursor.next_counter();
if from_counter >= head {
cursor.set_next_counter(head);
return RingPoll { from_counter, to_counter: head, ..RingPoll::default() };
}
let capacity = source.capacity() as u64;
let oldest_available = head.saturating_sub(capacity);
let mut next = from_counter;
let mut loss = RingLoss::default();
if next < oldest_available {
loss.overwritten = oldest_available - next;
next = oldest_available;
}
let kind = source.kind();
let mut frames = Vec::new();
while next < head {
match source.read_state_at(next) {
RingRead::Ready(frame) => {
if frame.id.kind() != kind || frame.id.counter() != next {
loss.unavailable = loss.unavailable.saturating_add(1);
} else {
frames.push(frame);
}
next = next.saturating_add(1);
}
RingRead::Pending => break,
RingRead::Unavailable => {
loss.unavailable = loss.unavailable.saturating_add(1);
next = next.saturating_add(1);
}
}
}
cursor.set_next_counter(next);
RingPoll { frames, loss, from_counter, to_counter: next }
}