use std::{ops::Range, task::Poll, time::Duration};
use crate::runtime::{Deadline, Instant};
pub(crate) const GRACE: Duration = Duration::from_secs(1);
pub(crate) fn grace(max_age: Duration) -> Duration {
match max_age.is_zero() {
true => GRACE,
false => max_age,
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct Run {
groups: Range<u64>,
since: Instant,
}
#[derive(Debug)]
pub(crate) struct Tail {
runs: Vec<Run>,
grace: Duration,
streams: u64,
active: u64,
}
impl Default for Tail {
fn default() -> Self {
Self {
runs: Vec::new(),
grace: GRACE,
streams: 0,
active: 0,
}
}
}
impl Tail {
pub fn new(grace: Duration) -> Self {
Self {
grace,
..Self::default()
}
}
pub fn set_grace(&mut self, grace: Duration) {
self.grace = grace;
}
pub fn account(&mut self, groups: Range<u64>, now: Instant) {
if groups.is_empty() {
return;
}
let first = self.runs.partition_point(|run| run.groups.end < groups.start);
let last = self.runs.partition_point(|run| run.groups.start <= groups.end);
let merged = match &self.runs[first..last] {
[] => Run {
groups,
since: self.runs.get(first).map_or(now, |above| above.since),
},
[head, .., tail] => Run {
groups: head.groups.start.min(groups.start)..tail.groups.end.max(groups.end),
since: head.since,
},
[only] => Run {
groups: only.groups.start.min(groups.start)..only.groups.end.max(groups.end),
since: only.since,
},
};
self.runs.splice(first..last, [merged]);
self.expire(now);
}
pub fn demand(&mut self, groups: Range<u64>, now: Instant) {
if groups.is_empty() {
return;
}
let mut below = 0;
for run in &mut self.runs {
if below < groups.end && groups.start < run.groups.start {
run.since = now;
}
below = run.groups.end;
}
}
pub fn expire(&mut self, now: Instant) {
let grace = self.grace;
self.runs.dedup_by(|run, below| {
let expired = now.duration_since(run.since) >= grace;
if expired {
below.groups.end = run.groups.end;
}
expired
});
}
pub fn covers(&self, groups: Range<u64>) -> bool {
groups.is_empty()
|| self
.runs
.iter()
.any(|run| run.groups.start <= groups.start && groups.end <= run.groups.end)
}
pub fn streams(&self) -> u64 {
self.streams
}
}
pub(crate) struct Reading {
tail: Option<kio::Producer<Tail>>,
active: bool,
}
impl Reading {
pub fn open(tail: &kio::Producer<Tail>, group: Option<u64>, now: Instant) -> Self {
let Ok(mut state) = tail.write() else {
return Self {
tail: None,
active: false,
};
};
state.streams += 1;
state.active += 1;
if let Some(group) = group {
state.account(group..group.saturating_add(1), now);
}
Self {
tail: Some(tail.clone()),
active: true,
}
}
pub fn park(&mut self) {
self.set_active(false);
}
pub fn resume(&mut self) {
self.set_active(true);
}
fn set_active(&mut self, active: bool) {
if self.active == active {
return;
}
self.active = active;
if let Some(tail) = &self.tail
&& let Ok(mut state) = tail.write()
{
match active {
true => state.active += 1,
false => state.active -= 1,
}
}
}
}
impl Drop for Reading {
fn drop(&mut self) {
self.set_active(false);
}
}
pub(crate) struct Settle {
tail: kio::Consumer<Tail>,
grace: Deadline<crate::time::Clock>,
}
impl Settle {
pub fn new(runtime: &crate::time::Clock, tail: kio::Consumer<Tail>) -> Self {
let grace = tail.read().grace;
Self {
tail,
grace: Deadline::after(runtime, grace),
}
}
pub fn poll(&mut self, waiter: &kio::Waiter, mut complete: impl FnMut(&Tail) -> bool) -> Poll<()> {
let expired = self.grace.poll(waiter).is_ready();
self.tail
.poll(waiter, |tail| match tail.active == 0 && (expired || complete(tail)) {
true => Poll::Ready(()),
false => Poll::Pending,
})
.map(|_| ())
}
}
#[cfg(test)]
mod tests {
use super::*;
fn groups(tail: &Tail) -> Vec<Range<u64>> {
tail.runs.iter().map(|run| run.groups.clone()).collect()
}
#[test]
fn runs_merge_across_gaps() {
let now = Instant::now();
let mut tail = Tail::default();
tail.account(0..1, now);
tail.account(2..3, now);
assert!(!tail.covers(0..3), "group 1 is missing");
assert_eq!(groups(&tail), vec![0..1, 2..3]);
tail.account(1..2, now);
assert!(tail.covers(0..3) && tail.runs.len() == 1, "adjacent runs merge");
tail.account(5..7, now);
tail.account(9..10, now);
tail.account(4..9, now);
assert_eq!(groups(&tail), vec![0..3, 4..10], "one insert swallows several runs");
assert!(tail.covers(4..10));
assert!(!tail.covers(2..5));
assert!(tail.covers(7..7), "an empty range is always covered");
}
#[test]
fn a_gap_past_the_grace_folds_away() {
let start = Instant::now();
let mut tail = Tail::new(GRACE);
tail.account(0..1, start);
tail.account(2..3, start);
tail.account(4..5, start + GRACE / 2);
assert_eq!(groups(&tail), vec![0..1, 2..3, 4..5]);
tail.account(6..7, start + GRACE / 2);
tail.account(8..9, start + GRACE);
assert_eq!(
groups(&tail),
vec![0..3, 4..5, 6..7, 8..9],
"only the oldest gap expired"
);
tail.expire(start + GRACE * 2);
assert_eq!(groups(&tail), vec![0..9]);
assert!(tail.covers(0..9));
}
#[test]
fn a_split_gap_keeps_its_age() {
let start = Instant::now();
let mut tail = Tail::new(GRACE);
tail.account(0..1, start);
tail.account(9..10, start);
tail.account(5..6, start + GRACE / 2);
assert_eq!(groups(&tail), vec![0..1, 5..6, 9..10]);
tail.expire(start + GRACE);
assert_eq!(groups(&tail), vec![0..10], "both halves opened at the start");
}
#[test]
fn a_lowered_floor_restarts_the_gap_age() {
let start = Instant::now();
let mut tail = Tail::new(GRACE);
tail.account(3..4, start);
tail.account(9..10, start);
let lowered = start + GRACE * 2;
tail.demand(1..3, lowered);
tail.account(1..2, lowered);
assert_eq!(
groups(&tail),
vec![1..2, 3..10],
"only the gap the floor reached restarts"
);
assert!(!tail.covers(1..4));
tail.expire(lowered + GRACE);
assert!(tail.covers(1..10));
}
#[tokio::test(start_paused = true)]
async fn the_grace_waits_for_a_stream_being_read() {
use crate::runtime::Timers as _;
let runtime = crate::time::Clock::tokio();
let tail = kio::Producer::new(Tail::default());
let reading = Reading::open(&tail, Some(0), runtime.now());
let mut settle = Settle::new(&runtime, tail.consume());
let mut settled = std::pin::pin!(kio::wait(|waiter| settle.poll(waiter, |_| false)));
tokio::time::sleep(GRACE * 2).await;
assert!(
futures::poll!(settled.as_mut()).is_pending(),
"a stream is still being read"
);
drop(reading);
settled.await;
}
#[tokio::test(start_paused = true)]
async fn a_parked_stream_does_not_hold_the_end() {
use crate::runtime::Timers as _;
let runtime = crate::time::Clock::tokio();
let tail = kio::Producer::new(Tail::default());
let mut reading = Reading::open(&tail, None, runtime.now());
reading.park();
let mut settle = Settle::new(&runtime, tail.consume());
kio::wait(|waiter| settle.poll(waiter, |tail| tail.streams() == 1)).await;
reading.resume();
assert_eq!(tail.read().active, 1);
drop(reading);
assert_eq!(tail.read().active, 0);
}
}