use super::ring::{OutputEvent, RecentOutputSnapshot};
use std::ops::Range;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct OutputCursor {
next_sequence: u64,
missed_events: u64,
}
impl OutputCursor {
#[must_use]
pub const fn new(next_sequence: u64) -> Self {
Self {
next_sequence,
missed_events: 0,
}
}
#[must_use]
pub const fn next_sequence(&self) -> u64 {
self.next_sequence
}
#[must_use]
pub const fn missed_events(&self) -> u64 {
self.missed_events
}
pub(super) fn advance_to(&mut self, next_sequence: u64) {
self.next_sequence = next_sequence;
}
pub(super) fn record_gap(&mut self, missed: u64, resume_sequence: u64) {
self.missed_events = self.missed_events.saturating_add(missed);
self.next_sequence = resume_sequence;
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum OutputCursorItem {
Event(OutputEvent),
Gap(OutputGap),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct OutputGap {
expected_sequence: u64,
resume_sequence: u64,
missed_events: u64,
newest_sequence: u64,
recent_snapshot: RecentOutputSnapshot,
}
impl OutputGap {
pub(super) const fn new(
expected_sequence: u64,
resume_sequence: u64,
missed_events: u64,
newest_sequence: u64,
recent_snapshot: RecentOutputSnapshot,
) -> Self {
Self {
expected_sequence,
resume_sequence,
missed_events,
newest_sequence,
recent_snapshot,
}
}
#[must_use]
pub const fn expected_sequence(&self) -> u64 {
self.expected_sequence
}
#[must_use]
pub const fn resume_sequence(&self) -> u64 {
self.resume_sequence
}
#[must_use]
pub const fn missed_events(&self) -> u64 {
self.missed_events
}
#[must_use]
pub fn missed_range(&self) -> Range<u64> {
self.expected_sequence..self.resume_sequence
}
#[must_use]
pub const fn newest_sequence(&self) -> u64 {
self.newest_sequence
}
#[must_use]
pub const fn recent_snapshot(&self) -> &RecentOutputSnapshot {
&self.recent_snapshot
}
}
#[cfg(test)]
mod tests {
use super::{OutputCursor, OutputCursorItem};
use crate::events::OutputRing;
#[test]
fn cursor_advances_independently_through_retained_events() {
let mut ring = OutputRing::new(8, 64);
ring.push(b"one".to_vec());
ring.push(b"two".to_vec());
let mut first = ring.cursor_from_oldest();
let mut second = ring.cursor_from_oldest();
assert_eq!(
ring.poll_cursor(&mut first),
Some(OutputCursorItem::Event(ring.retained_events()[0].clone()))
);
assert_eq!(first.next_sequence(), 1);
assert_eq!(second.next_sequence(), 0);
assert_eq!(
ring.poll_cursor(&mut first),
Some(OutputCursorItem::Event(ring.retained_events()[1].clone()))
);
assert_eq!(ring.poll_cursor(&mut first), None);
assert_eq!(first.next_sequence(), ring.next_sequence());
assert_eq!(
ring.poll_cursor(&mut second),
Some(OutputCursorItem::Event(ring.retained_events()[0].clone()))
);
assert_eq!(second.next_sequence(), 1);
}
#[test]
fn lagged_cursor_reports_explicit_gap_and_resumes_at_oldest_event() {
let mut ring = OutputRing::new(2, 64);
let mut cursor = OutputCursor::new(0);
for bytes in [b"zero".as_slice(), b"one".as_slice(), b"two".as_slice()] {
ring.push(bytes.to_vec());
}
let Some(OutputCursorItem::Gap(gap)) = ring.poll_cursor(&mut cursor) else {
panic!("cursor should report lag");
};
assert_eq!(gap.expected_sequence(), 0);
assert_eq!(gap.resume_sequence(), 1);
assert_eq!(gap.missed_events(), 1);
assert_eq!(gap.missed_range(), 0..1);
assert_eq!(gap.newest_sequence(), 2);
assert_eq!(gap.recent_snapshot().oldest_sequence(), Some(0));
assert_eq!(gap.recent_snapshot().newest_sequence(), Some(2));
assert_eq!(cursor.missed_events(), 1);
assert_eq!(cursor.next_sequence(), 1);
let Some(OutputCursorItem::Event(event)) = ring.poll_cursor(&mut cursor) else {
panic!("cursor should resume with oldest retained event");
};
assert_eq!(event.sequence(), 1);
assert_eq!(event.bytes(), b"one");
}
}