use std::collections::VecDeque;
#[derive(Debug, Clone)]
pub struct RequestRing<T> {
rows: VecDeque<(u64, T)>,
capacity: usize,
next_seq: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RingPage<'a, T> {
pub rows: Vec<&'a T>,
pub cursor: u64,
pub missed: u64,
}
impl<T> RequestRing<T> {
pub fn new(capacity: usize) -> Self {
RequestRing {
rows: VecDeque::new(),
capacity: capacity.max(1),
next_seq: 0,
}
}
pub fn push(&mut self, row: T) -> u64 {
let seq = self.next_seq;
self.next_seq += 1;
if self.rows.len() == self.capacity {
self.rows.pop_front();
}
self.rows.push_back((seq, row));
seq
}
pub fn since(&self, since: u64, limit: usize) -> RingPage<'_, T> {
let oldest = self
.rows
.front()
.map(|(seq, _)| *seq)
.unwrap_or(self.next_seq);
let missed = oldest.saturating_sub(since);
let matched: Vec<&(u64, T)> = self.rows.iter().filter(|(seq, _)| *seq >= since).collect();
let returned = &matched[..limit.min(matched.len())];
let cursor = if returned.len() < matched.len() {
returned.last().map(|(seq, _)| seq + 1).unwrap_or(since)
} else {
self.next_seq
};
RingPage {
rows: returned.iter().map(|(_, row)| row).collect(),
cursor,
missed,
}
}
pub fn rows(&self) -> impl Iterator<Item = &T> {
self.rows.iter().map(|(_, row)| row)
}
pub fn recorded_total(&self) -> u64 {
self.next_seq
}
pub fn len(&self) -> usize {
self.rows.len()
}
pub fn is_empty(&self) -> bool {
self.rows.is_empty()
}
pub fn capacity(&self) -> usize {
self.capacity
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_ring_keeps_the_newest_rows_and_numbers_them_for_all_time() {
let mut ring = RequestRing::new(3);
for i in 0..5u32 {
assert_eq!(ring.push(i), i as u64);
}
assert_eq!(ring.len(), 3);
assert_eq!(ring.rows().copied().collect::<Vec<_>>(), vec![2, 3, 4]);
assert_eq!(ring.recorded_total(), 5);
}
#[test]
fn a_poller_that_keeps_up_reads_every_row_exactly_once() {
let mut ring = RequestRing::new(10);
let mut cursor = 0;
let mut seen = Vec::new();
for round in 0..3u32 {
for i in 0..2 {
ring.push(round * 2 + i);
}
let page = ring.since(cursor, 100);
assert_eq!(page.missed, 0);
seen.extend(page.rows.iter().copied().copied());
cursor = page.cursor;
}
assert_eq!(seen, vec![0, 1, 2, 3, 4, 5]);
}
#[test]
fn a_truncated_page_resumes_at_the_row_after_the_last_one_delivered() {
let mut ring = RequestRing::new(10);
for i in 0..6u32 {
ring.push(i);
}
let page = ring.since(0, 2);
assert_eq!(page.rows, [&0, &1]);
assert_eq!(page.cursor, 2, "not 6");
let page = ring.since(page.cursor, 2);
assert_eq!(page.rows, [&2, &3]);
assert_eq!(page.cursor, 4);
let page = ring.since(page.cursor, 100);
assert_eq!(page.rows, [&4, &5]);
assert_eq!(page.cursor, 6);
}
#[test]
fn a_poller_that_falls_behind_is_told_how_much_it_lost() {
let mut ring = RequestRing::new(3);
for i in 0..10u32 {
ring.push(i);
}
let page = ring.since(0, 100);
assert_eq!(page.rows.len(), 3);
assert_eq!(page.missed, 7);
assert_eq!(page.cursor, 10);
let page = ring.since(page.cursor, 100);
assert!(page.rows.is_empty());
assert_eq!(page.missed, 0);
assert_eq!(page.cursor, 10);
}
}