use std::collections::{BTreeMap, HashSet};
use std::ops::Bound;
use std::task::{Poll, ready};
use crate::{Datagram, Error, Result, Timestamp, frame, group, track};
use track::{Anchor, LiveEdge};
use super::subscription::{Cap, Position, Subscription, max_some, min_some};
#[derive(Clone)]
struct Segment {
id: u64,
start: Option<Position>,
end: Option<Position>,
track: track::Consumer,
ask: Option<Position>,
warm: bool,
}
impl Segment {
fn covers(&self, position: Position) -> bool {
self.start.is_none_or(|start| position >= start) && self.end.is_none_or(|end| position < end)
}
fn produced(&self) -> Option<Position> {
let produced = self.track.resume_position()?;
if produced <= self.start.unwrap_or_default() {
return None;
}
Some(match self.end {
Some(end) => produced.min(end),
None => produced,
})
}
fn warm_edge(&self) -> Option<Position> {
let group = self.track.peek_latest()?;
Some(Position {
group: group.sequence,
frame: group.frame_count() as u64,
})
}
fn peek_below(&self, before: Option<u64>) -> Option<group::Consumer> {
let limit = min_some(before, last_group(self.end));
let group = match limit {
Some(limit) => self.track.peek_before(limit)?,
None => self.track.peek_latest()?,
};
if self.start.is_some_and(|start| group.sequence < start.group) {
return None;
}
Some(group)
}
}
fn slice(prefs: &Subscription, start: Option<Position>, end: Option<Position>) -> Subscription {
Subscription {
start: max_some(prefs.start, start),
end: min_some(prefs.end, end),
..prefs.clone()
}
}
fn poll_first_start<'a>(
segments: impl IntoIterator<Item = &'a Segment>,
waiter: &kio::Waiter,
from: u64,
cap: Option<u64>,
) -> Option<Option<Timestamp>> {
segments.into_iter().find_map(|segment| {
let start = segment.start.map_or(0, |start| start.group).max(from);
segment
.track
.poll_first_start(waiter, start, min_some(cap, last_group(segment.end)))
})
}
#[derive(Clone)]
pub(crate) struct Successor {
state: kio::ConsumerWeak<ResumeState>,
after: u64,
boundary: u64,
cap: Option<u64>,
}
impl Successor {
pub(crate) fn poll_start(&self, waiter: &kio::Waiter) -> Option<Timestamp> {
let mut start = None;
let _ = self.state.poll(waiter, |state| {
let later = state.segments.iter().filter(|segment| segment.id > self.after);
start = poll_first_start(later, waiter, self.boundary, self.cap).flatten();
Poll::<()>::Pending
});
start
}
}
impl PartialEq for Successor {
fn eq(&self, other: &Self) -> bool {
self.after == other.after
&& self.boundary == other.boundary
&& self.cap == other.cap
&& self.state.same_channel(&other.state)
}
}
const MAX_SEGMENTS: usize = 3;
struct ResumeState {
segments: Vec<Segment>,
pruned: Option<Position>,
epoch: u64,
finished: bool,
abort: Option<Error>,
}
impl Default for ResumeState {
fn default() -> Self {
Self {
segments: Vec::new(),
pruned: None,
epoch: 1,
finished: false,
abort: None,
}
}
}
struct Snapshot {
epoch: u64,
finished: bool,
abort: Option<Error>,
segments: Vec<Segment>,
}
impl ResumeState {
fn snapshot(&self) -> Snapshot {
Snapshot {
epoch: self.epoch,
finished: self.finished,
abort: self.abort.clone(),
segments: self.segments.clone(),
}
}
fn resume_position(&self) -> Option<Position> {
self.segments
.iter()
.filter_map(Segment::produced)
.chain(self.pruned)
.max()
}
fn latest(&self) -> Option<u64> {
let position = self.resume_position()?;
match position.frame {
0 => position.group.checked_sub(1),
_ => Some(position.group),
}
}
fn live_edge(&self, cap: Option<u64>) -> Option<LiveEdge> {
self.segments
.iter()
.filter_map(|segment| {
let edge = segment.track.live_edge(min_some(cap, last_group(segment.end)))?;
let start = segment.start.map_or(0, |start| start.group);
(edge.sequence >= start).then_some(edge)
})
.max_by_key(|edge| edge.sequence)
}
fn switch(&mut self, track: track::Consumer, start: Option<Position>) -> Result<()> {
if !self.segments.is_empty() {
let Some(start) = start else {
return Err(crate::coding::BoundsExceeded.into());
};
while let Some(prev) = self.segments.last() {
let prev_start = prev.start.unwrap_or_default();
if start > prev_start {
break;
}
if prev.produced().is_some() {
return Err(crate::coding::BoundsExceeded.into());
}
self.segments.pop();
}
if let Some(prev) = self.segments.last_mut() {
prev.end = Some(start);
}
}
let id = self.epoch;
self.segments.push(Segment {
id,
start,
end: None,
track,
ask: start,
warm: false,
});
self.epoch += 1;
self.prune();
Ok(())
}
fn prune(&mut self) {
while self.segments.len() > MAX_SEGMENTS {
let front = &self.segments[0];
let Some(end) = front.end else { break };
let owes_more =
front.produced() < Some(end) && front.track.poll_complete(&kio::Waiter::noop()).is_pending();
if owes_more {
break;
}
self.pruned = self.pruned.max(Some(end));
self.segments.remove(0);
}
}
}
#[derive(Clone, Default)]
pub struct Producer {
state: kio::Producer<ResumeState>,
}
impl Producer {
pub fn new() -> Self {
Self::default()
}
#[cfg_attr(not(test), expect(dead_code))]
pub fn switch(
&mut self,
track: impl super::origin_impl::Consume<track::Consumer>,
start: impl Into<Option<Position>>,
) -> Result<()> {
let track = track.consume();
let start = start.into();
let mut state = self.state.write().map_err(|_| Error::Dropped)?;
if state.finished || state.abort.is_some() {
return Err(Error::Closed);
}
state.switch(track, start)
}
pub fn takeover(&mut self, track: impl super::origin_impl::Consume<track::Consumer>) -> Result<()> {
let track = track.consume();
let mut state = self.state.write().map_err(|_| Error::Dropped)?;
if state.finished || state.abort.is_some() {
return Err(Error::Closed);
}
let ask = state
.segments
.last()
.filter(|last| last.warm)
.and_then(Segment::warm_edge);
let start = state.resume_position();
if start.is_none() {
state.segments.clear();
}
state.switch(track, start)?;
if let Some(ask) = ask
&& start.is_some()
&& let Some(last) = state.segments.last_mut()
{
last.ask = Some(ask);
}
Ok(())
}
pub(crate) fn release(&mut self) -> Result<()> {
let mut state = self.state.write().map_err(|_| Error::Dropped)?;
if state.finished || state.abort.is_some() {
return Err(Error::Closed);
}
if state.segments.is_empty() {
return Ok(());
}
state.segments.clear();
state.pruned = None;
state.epoch += 1;
Ok(())
}
pub(crate) fn park(&mut self, warm: impl super::origin_impl::Consume<track::Consumer>) -> Result<()> {
let track = warm.consume();
let mut state = self.state.write().map_err(|_| Error::Dropped)?;
if state.finished || state.abort.is_some() {
return Err(Error::Closed);
}
state.segments.clear();
state.pruned = None;
state.switch(track, None)?;
if let Some(segment) = state.segments.last_mut() {
segment.warm = true;
}
Ok(())
}
pub(crate) fn is_spliced(&self) -> bool {
!self.state.read().segments.is_empty()
}
pub(crate) fn resume_position(&self) -> Option<Position> {
self.state.read().resume_position()
}
pub fn finish(&mut self) -> Result<()> {
let mut state = self.state.write().map_err(|_| Error::Dropped)?;
if state.finished || state.abort.is_some() {
return Err(Error::Closed);
}
state.finished = true;
state.epoch += 1;
Ok(())
}
pub fn abort(&mut self, err: Error) -> Result<()> {
let mut state = self.state.write().map_err(|_| Error::Dropped)?;
if state.finished || state.abort.is_some() {
return Err(Error::Closed);
}
state.abort = Some(err);
state.epoch += 1;
Ok(())
}
pub(crate) fn is_aborted(&self) -> bool {
self.state.read().abort.is_some()
}
pub(crate) fn is_clone(&self, other: &Self) -> bool {
self.state.same_channel(&other.state)
}
pub fn consume(&self) -> Consumer {
Consumer {
state: self.state.consume(),
}
}
pub fn is_used(&self) -> bool {
self.state.is_used()
}
pub(crate) fn poll_used(&self, waiter: &kio::Waiter) -> Poll<()> {
self.state.poll_used(waiter).map(|_| ())
}
pub(crate) fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
self.state.poll_unused(waiter).map(|_| ())
}
}
#[derive(Clone)]
pub struct Consumer {
state: kio::Consumer<ResumeState>,
}
impl Consumer {
pub(crate) fn cached_groups(&self) -> Vec<(group::Producer, bool)> {
let segments = self.state.read().segments.clone();
let mut seen = HashSet::new();
let mut groups = Vec::new();
for segment in segments {
let start = segment
.start
.map(|start| start.group.saturating_add(u64::from(start.frame != 0)));
let end = segment.end.map(|end| end.group);
for (group, visible) in segment.track.cached_groups() {
if start.is_some_and(|start| group.sequence < start)
|| end.is_some_and(|end| group.sequence >= end)
|| !seen.insert(group.sequence)
{
continue;
}
groups.push((group, visible));
}
}
groups
}
pub(crate) fn cached_info(&self) -> Option<track::Info> {
self.state.read().segments.first()?.track.cached_info()
}
pub(crate) fn cached_group(&self, sequence: u64, frame_start: u64) -> Option<group::Consumer> {
let segments = self.state.read().segments.clone();
segments.iter().rev().find_map(|segment| {
let (start, end) = frames(segment.start, segment.end, sequence)?;
(start <= frame_start && end.is_none())
.then(|| segment.track.cached_group(sequence, frame_start))
.flatten()
})
}
#[cfg(test)]
pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> Subscriber {
let prefs = kio::Producer::new(subscription.into().unwrap_or_default());
self.subscribe_shared(prefs)
}
pub(crate) fn subscribe_shared(&self, prefs: kio::Producer<Subscription>) -> Subscriber {
let last_prefs = prefs.read().clone();
Subscriber {
state: self.state.clone(),
prefs,
last_prefs,
epoch: 0,
finished: false,
abort: None,
closed: false,
segments: Vec::new(),
next_sequence: 0,
min_sequence: 0,
end_sequence: None,
outer: Anchor::default(),
drift_anchor: kio::Producer::new(Anchor::default()),
}
}
pub fn poll_info(&self, waiter: &kio::Waiter) -> Poll<Result<track::Info>> {
let track = match ready!(self.state.poll(waiter, |state| {
if state.abort.is_some() || !state.segments.is_empty() {
Poll::Ready(
state
.abort
.clone()
.map_or_else(|| Ok(state.segments[0].track.clone()), Err),
)
} else {
Poll::Pending
}
})) {
Ok(res) => res?,
Err(state) => match (&state.abort, state.segments.first()) {
(Some(err), _) => return Poll::Ready(Err(err.clone())),
(None, Some(segment)) => segment.track.clone(),
(None, None) => return Poll::Ready(Err(Error::Dropped)),
},
};
track.query().poll_ok(waiter)
}
#[cfg(test)]
pub async fn info(&self) -> Result<track::Info> {
kio::wait(|waiter| self.poll_info(waiter)).await
}
pub fn fetch_group(&self, sequence: u64, options: impl Into<Option<group::Fetch>>) -> kio::Pending<Fetching> {
kio::Pending::new(Fetching {
state: self.state.clone(),
sequence,
options: options.into().unwrap_or_default(),
inner: kio::Lock::new(None),
})
}
pub fn latest(&self) -> Option<u64> {
self.state.read().latest()
}
pub(crate) fn resume_position(&self) -> Option<Position> {
self.state.read().resume_position()
}
pub(crate) fn live_edge(&self, cap: Option<u64>) -> Option<LiveEdge> {
self.state.read().live_edge(cap)
}
pub(crate) fn poll_first_start(
&self,
waiter: &kio::Waiter,
from: u64,
cap: Option<u64>,
) -> Option<Option<Timestamp>> {
let mut start = None;
let _ = self.state.poll(waiter, |state| {
start = poll_first_start(&state.segments, waiter, from, cap);
Poll::<()>::Pending
});
start
}
pub(crate) fn peek_latest(&self) -> Option<group::Consumer> {
self.peek_below(None)
}
pub(crate) fn peek_group(&self, sequence: u64) -> Option<group::Consumer> {
let track = self.state.read().segments.last().map(|segment| segment.track.clone())?;
track.peek_group(sequence)
}
pub(crate) fn peek_before(&self, sequence: u64) -> Option<group::Consumer> {
self.peek_below(Some(sequence))
}
fn peek_below(&self, before: Option<u64>) -> Option<group::Consumer> {
let segments = self.state.read().segments.clone();
segments.iter().rev().find_map(|segment| segment.peek_below(before))
}
}
pub struct Fetching {
state: kio::Consumer<ResumeState>,
sequence: u64,
options: group::Fetch,
#[allow(clippy::type_complexity)]
inner: kio::Lock<Option<(u64, track::Consumer, kio::Pending<track::Fetching>)>>,
}
impl Fetching {
fn poll_latch(&self, waiter: &kio::Waiter, latched: Option<u64>) -> Poll<Result<Option<(u64, track::Consumer)>>> {
let newest = move |segments: &[Segment]| {
segments
.last()
.filter(|segment| latched.is_none_or(|latched| segment.id > latched))
.map(|segment| (segment.id, segment.track.clone()))
};
match self.state.poll(waiter, |s| match (&s.abort, newest(&s.segments)) {
(Some(err), _) => Poll::Ready(Err(err.clone())),
(None, Some(next)) => Poll::Ready(Ok(next)),
(None, None) => Poll::Pending,
}) {
Poll::Ready(Ok(res)) => Poll::Ready(res.map(Some)),
Poll::Ready(Err(state)) => Poll::Ready(match &state.abort {
Some(err) => Err(err.clone()),
None => Ok(newest(&state.segments)),
}),
Poll::Pending => Poll::Pending,
}
}
}
impl kio::Pollable for Fetching {
type Output = Result<group::Consumer>;
fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
if let Some(group) = (Consumer {
state: self.state.clone(),
})
.cached_group(self.sequence, self.options.frame_start)
{
return Poll::Ready(Ok(group));
}
let mut inner = self.inner.lock();
loop {
if inner.is_none() {
let (id, track) = match ready!(self.poll_latch(waiter, None))? {
Some(next) => next,
None => return Poll::Ready(Err(Error::NotFound)),
};
let fetch = track.fetch_group(self.sequence, self.options.clone());
*inner = Some((id, track, fetch));
}
let (latched, track, fetch) = inner.as_ref().expect("latched above");
let err = match kio::Pollable::poll(&**fetch, waiter) {
Poll::Ready(Err(err)) => err,
Poll::Ready(Ok(group)) => return Poll::Ready(Ok(group)),
Poll::Pending => {
let latched = *latched;
match self.poll_latch(waiter, Some(latched)) {
Poll::Ready(Ok(Some((id, track)))) => {
let fetch = track.fetch_group(self.sequence, self.options.clone());
*inner = Some((id, track, fetch));
continue;
}
Poll::Ready(Err(err)) => return Poll::Ready(Err(err)),
Poll::Ready(Ok(None)) | Poll::Pending => {}
}
return match self.state.poll(waiter, |s| match &s.abort {
Some(err) => Poll::Ready(err.clone()),
None => Poll::Pending,
}) {
Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
_ => Poll::Pending,
};
}
};
let next = match self.poll_latch(waiter, Some(*latched)) {
Poll::Ready(res) => res?,
Poll::Pending => match track.poll_complete(&kio::Waiter::noop()) {
Poll::Ready(Err(_)) => return Poll::Pending,
_ => return Poll::Ready(Err(err)),
},
};
let Some((id, track)) = next else {
return Poll::Ready(Err(err));
};
let fetch = track.fetch_group(self.sequence, self.options.clone());
*inner = Some((id, track, fetch));
}
}
}
fn frames(start: Option<Position>, end: Option<Position>, sequence: u64) -> Option<(u64, Option<u64>)> {
let start = match start {
Some(start) if start.group > sequence => return None,
Some(start) if start.group == sequence => start.frame,
_ => 0,
};
let end = match end {
Some(end) if end.group > sequence => None,
Some(end) if end.group == sequence => Some(end.frame),
Some(_) => return None,
None => None,
};
Some((start, end))
}
fn last_group(end: Option<Position>) -> Option<u64> {
end.and_then(|end| Cap::from(end.group_end()).exclusive())
}
pub(crate) struct Group {
state: kio::Consumer<ResumeState>,
subscription: kio::Consumer<Subscription>,
anchor: kio::Consumer<Anchor>,
sequence: u64,
index: u64,
end: Option<u64>,
current: Option<Current>,
dead: Option<(u64, Error)>,
stale_stats: crate::stats::Meter,
}
struct Current {
segment: u64,
cap: Option<u64>,
bound: Option<u64>,
group: group::Consumer,
}
type Covering = (u64, track::Consumer, Option<Option<u64>>, Option<u64>);
impl Clone for Group {
fn clone(&self) -> Self {
let current = self.current.as_ref().and_then(|current| {
let mut group = current.group.clone();
group.start_at(self.index);
(group.index() == self.index).then_some(Current {
segment: current.segment,
cap: current.cap,
bound: current.bound,
group,
})
});
Self {
state: self.state.clone(),
subscription: self.subscription.clone(),
anchor: self.anchor.clone(),
sequence: self.sequence,
index: self.index,
end: self.end,
current,
dead: self.dead.clone(),
stale_stats: self.stale_stats.clone(),
}
}
}
impl Group {
fn new(
state: kio::Consumer<ResumeState>,
subscription: kio::Consumer<Subscription>,
anchor: kio::Consumer<Anchor>,
sequence: u64,
index: u64,
) -> Self {
Self {
state,
subscription,
anchor,
sequence,
index,
end: None,
current: None,
dead: None,
stale_stats: Default::default(),
}
}
pub(crate) fn frame_count(&self) -> usize {
let count = self
.current
.as_ref()
.map_or(self.index, |current| current.group.frame_count() as u64)
.max(self.index);
usize::try_from(count).unwrap_or(usize::MAX)
}
pub(crate) fn set_stale_meter(&mut self, meter: crate::stats::Meter) {
if let Some(current) = &mut self.current {
current.group.set_stale_meter(meter.clone());
}
self.stale_stats = meter;
}
fn latched(mut self, segment: u64, cap: Option<u64>, bound: Option<u64>, mut group: group::Consumer) -> Self {
group.end_at(cap.map_or(Bound::Unbounded, Bound::Excluded));
group.start_at(self.index);
group.set_stale_meter(self.stale_stats.clone());
if group.index() == self.index {
self.current = Some(Current {
segment,
cap,
bound,
group,
});
}
self
}
pub fn index(&self) -> u64 {
self.index
}
pub fn start_at(&mut self, index: u64) {
if index <= self.index {
return;
}
self.index = index;
if let Some(current) = &mut self.current {
current.group.start_at(index);
if current.group.index() != index || current.cap.is_some_and(|cap| index >= cap) {
self.current = None;
}
}
}
pub fn end_at(&mut self, index: Option<u64>) {
self.end = index;
}
fn poll_covering(&self, position: Position, dead: Option<u64>, waiter: &kio::Waiter) -> Poll<Option<Covering>> {
let sequence = self.sequence;
let located = self.state.poll(waiter, |state| {
let stranded = dead.is_some() && state.resume_position().is_some_and(|resume| resume > position);
let lost = state.pruned.is_some_and(|floor| position < floor);
let terminal = state.finished || state.abort.is_some() || stranded || lost;
match state.segments.iter().find(|segment| segment.covers(position)) {
Some(segment) if dead == Some(segment.id) => {
match segment.track.poll_serving_group(sequence, position.frame, waiter) {
Poll::Ready(()) => Poll::Ready(Some((
segment.id,
segment.track.clone(),
frames(segment.start, segment.end, sequence).map(|(_, end)| end),
last_group(segment.end),
))),
Poll::Pending => match terminal {
true => Poll::Ready(None),
false => Poll::Pending,
},
}
}
Some(segment) => Poll::Ready(Some((
segment.id,
segment.track.clone(),
frames(segment.start, segment.end, sequence).map(|(_, end)| end),
last_group(segment.end),
))),
None if terminal => Poll::Ready(None),
None => Poll::Pending,
}
});
match located {
Poll::Ready(Ok(found)) => Poll::Ready(found),
Poll::Ready(Err(_)) => Poll::Ready(None),
Poll::Pending => Poll::Pending,
}
}
fn poll_current(&mut self, waiter: &kio::Waiter) -> Poll<Result<bool>> {
loop {
let position = Position {
group: self.sequence,
frame: self.index,
};
let dead = self.dead.as_ref().map(|(segment, _)| *segment);
let sequence = self.sequence;
let found = ready!(self.poll_covering(position, dead, waiter));
let Some((segment, track, Some(cap), bound)) = found else {
if self.current.is_some() {
return Poll::Ready(Ok(true));
}
return Poll::Ready(self.give_up());
};
if self
.current
.as_ref()
.is_some_and(|current| current.segment == segment && current.cap == cap && current.bound == bound)
{
return Poll::Ready(Ok(true));
}
let Some(group) = ready!(track.poll_peek_group(sequence, waiter)) else {
self.dead = Some((segment, Error::NotFound));
continue;
};
let mut group = track.guard_group(group, self.subscription.clone(), self.anchor.clone(), bound);
group.set_stale_meter(self.stale_stats.clone());
group.start_at(self.index);
if group.index() != self.index {
self.dead = Some((segment, Error::Lagged));
continue;
}
group.end_at(cap.map_or(Bound::Unbounded, Bound::Excluded));
self.current = Some(Current {
segment,
cap,
bound,
group,
});
self.dead = None;
return Poll::Ready(Ok(true));
}
}
fn give_up(&self) -> Result<bool> {
match &self.dead {
Some((_, err)) => {
if !matches!(crate::StreamError::from(err), crate::StreamError::Old) {
tracing::warn!(
group = self.sequence,
frame = self.index,
%err,
"no route can serve the rest of this group"
);
}
Err(err.clone())
}
None => Ok(false),
}
}
fn roll(&mut self) -> bool {
let Some(current) = &self.current else {
return false;
};
match current.cap {
Some(cap) => {
debug_assert_eq!(self.index, cap, "bounded copy ended below its boundary");
self.index = self.index.max(cap);
self.current = None;
true
}
None => false,
}
}
fn bury(&mut self, err: Error) {
self.dead = self.current.take().map(|current| (current.segment, err));
}
pub fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
loop {
if self.end.is_some_and(|end| self.index >= end) {
return Poll::Ready(Ok(None));
}
if !ready!(self.poll_current(waiter))? {
return Poll::Ready(Ok(None));
}
let result = {
let current = self.current.as_mut().expect("resolved above");
ready!(current.group.poll_read_frame(waiter))
};
let latency_expired = self
.current
.as_ref()
.is_some_and(|current| current.group.latency_expired());
match result {
Ok(Some(frame)) => {
self.index += 1;
return Poll::Ready(Ok(Some(frame)));
}
Ok(None) if self.roll() => continue,
Ok(None) => return Poll::Ready(Ok(None)),
Err(err) if latency_expired => return Poll::Ready(Err(err)),
Err(err) => self.bury(err),
}
}
}
pub fn poll_next_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Consumer>>> {
loop {
if self.end.is_some_and(|end| self.index >= end) {
return Poll::Ready(Ok(None));
}
if !ready!(self.poll_current(waiter))? {
return Poll::Ready(Ok(None));
}
let result = {
let current = self.current.as_mut().expect("resolved above");
ready!(current.group.poll_next_frame(waiter))
};
let latency_expired = self
.current
.as_ref()
.is_some_and(|current| current.group.latency_expired());
match result {
Ok(Some(frame)) => {
self.index += 1;
return Poll::Ready(Ok(Some(frame)));
}
Ok(None) if self.roll() => continue,
Ok(None) => return Poll::Ready(Ok(None)),
Err(err) if latency_expired => return Poll::Ready(Err(err)),
Err(err) => self.bury(err),
}
}
}
pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
if !ready!(self.poll_current(waiter))? {
return Poll::Ready(Ok(self.index));
}
let current = self.current.as_mut().expect("resolved above");
let Some(cap) = current.cap else {
return current.group.poll_finished(waiter);
};
let seam = Position {
group: self.sequence,
frame: cap,
};
loop {
let dead = self.dead.as_ref().map(|(segment, _)| *segment);
let Some((segment, track, _, bound)) = ready!(self.poll_covering(seam, dead, waiter)) else {
return Poll::Ready(Ok(cap));
};
match ready!(track.poll_peek_group(self.sequence, waiter)) {
Some(continuation) => {
let mut continuation =
track.guard_group(continuation, self.subscription.clone(), self.anchor.clone(), bound);
continuation.set_stale_meter(self.stale_stats.clone());
return continuation.poll_finished(waiter);
}
None => self.dead = Some((segment, Error::NotFound)),
}
}
}
}
struct SegmentSub {
id: u64,
start: Option<Position>,
end: Option<Position>,
ask: Option<Position>,
sub: SubState,
terminal: Option<track::Subscriber>,
anchor: Anchor,
pruned: bool,
parked: BTreeMap<u64, group::Consumer>,
warm: Option<Warm>,
}
struct Warm {
edge: Option<u64>,
next: Option<track::Consumer>,
}
impl SegmentSub {
fn stale_sub_mut(&mut self) -> Option<&mut track::Subscriber> {
match &mut self.sub {
SubState::Active(sub) => Some(sub.as_mut()),
_ => self.terminal.as_mut(),
}
}
fn complete(&mut self, end: Result<u64>) {
let previous = std::mem::replace(&mut self.sub, SubState::Done(end));
if let SubState::Active(sub) = previous {
self.terminal = Some(*sub);
}
}
fn first_group(&self) -> u64 {
self.start.map_or(0, |start| start.group)
}
fn last_group(&self) -> Option<u64> {
debug_assert!(
self.end != Some(Position::default()),
"a segment cannot end at the start of the track"
);
last_group(self.end)
}
fn retired(&self) -> bool {
self.pruned && (self.end.is_none() || (matches!(self.sub, SubState::Done(_)) && self.parked.is_empty()))
}
}
enum SubState {
Pending(kio::Pending<track::Subscribing>),
Active(Box<track::Subscriber>),
Done(Result<u64>),
}
pub struct Subscriber {
state: kio::Consumer<ResumeState>,
prefs: kio::Producer<Subscription>,
last_prefs: Subscription,
epoch: u64,
finished: bool,
abort: Option<Error>,
closed: bool,
segments: Vec<SegmentSub>,
next_sequence: u64,
min_sequence: u64,
end_sequence: Option<u64>,
outer: Anchor,
drift_anchor: kio::Producer<Anchor>,
}
impl Subscriber {
fn poll_sync(&mut self, waiter: &kio::Waiter) {
self.sync(waiter);
self.reap();
self.refresh_anchor();
}
fn reap(&mut self) {
self.segments.retain(|s| !s.retired());
let mut cut = self
.segments
.iter()
.filter(|s| s.pruned)
.count()
.saturating_sub(MAX_SEGMENTS);
if cut > 0 {
for seg in &mut self.segments {
if cut == 0 {
break;
}
if seg.pruned {
seg.sub = SubState::Done(Err(Error::Dropped));
seg.terminal = None;
seg.parked.clear();
cut -= 1;
}
}
self.segments.retain(|s| !s.retired());
}
}
fn sync(&mut self, waiter: &kio::Waiter) {
loop {
let prefs = {
let last = &self.last_prefs;
match self
.prefs
.poll(waiter, |p| if **p != *last { Poll::Ready(()) } else { Poll::Pending })
{
Poll::Ready(Ok(guard)) => (*guard).clone(),
Poll::Ready(Err(_)) | Poll::Pending => break,
}
};
self.last_prefs = prefs;
for seg in &mut self.segments {
let prefs = slice(&self.last_prefs, seg.ask, seg.end);
if let Some(sub) = seg.stale_sub_mut() {
let _ = sub.update(prefs);
}
}
}
loop {
let epoch = self.epoch;
let (snapshot, closed) = match self.state.poll(waiter, |state| {
if state.epoch != epoch {
Poll::Ready(state.snapshot())
} else {
Poll::Pending
}
}) {
Poll::Ready(Ok(snapshot)) => (Some(snapshot), false),
Poll::Ready(Err(state)) => {
let snapshot = (state.epoch != epoch).then(|| state.snapshot());
(snapshot, true)
}
Poll::Pending => return,
};
if let Some(snapshot) = snapshot {
self.apply(snapshot);
}
if closed {
self.closed = true;
return;
}
}
}
fn apply(&mut self, snapshot: Snapshot) {
let Snapshot {
epoch,
finished,
abort,
segments,
} = snapshot;
self.epoch = epoch;
self.finished = finished;
self.abort = abort;
for s in &mut self.segments {
s.pruned = !segments.iter().any(|n| n.id == s.id);
}
self.segments.retain(|s| !s.retired());
let nexts: Vec<_> = segments.iter().skip(1).map(|next| Some(next.track.clone())).collect();
let nexts = nexts.into_iter().chain(std::iter::once(None));
for (segment, next) in segments.into_iter().zip(nexts) {
match self.segments.iter_mut().find(|s| s.id == segment.id) {
Some(existing) => {
if let Some(warm) = &mut existing.warm {
warm.next = next;
}
if existing.end != segment.end {
existing.end = segment.end;
if let Some(sub) = existing.stale_sub_mut() {
let _ = sub.update(slice(&self.last_prefs, segment.ask, segment.end));
}
}
}
None => {
let sub = segment
.track
.subscribe(slice(&self.last_prefs, segment.ask, segment.end));
self.segments.push(SegmentSub {
id: segment.id,
start: segment.start,
end: segment.end,
ask: segment.ask,
sub: SubState::Pending(sub),
terminal: None,
anchor: Anchor::default(),
pruned: false,
parked: BTreeMap::new(),
warm: segment.warm.then(|| Warm {
edge: segment.track.latest(),
next,
}),
});
}
}
}
}
fn hand_out(&self, segment: usize, group: group::Consumer) -> Option<group::Consumer> {
let seg = &self.segments[segment];
let sequence = group.sequence;
let (start, end) = frames(seg.start, seg.end, sequence)?;
if start != 0 {
return None;
}
let spliced = Group::new(
self.state.clone(),
self.prefs.consume(),
self.drift_anchor.consume(),
sequence,
0,
)
.latched(seg.id, end, seg.last_group(), group.clone());
Some(group.into_spliced(spliced))
}
fn segment_anchor(seg: &SegmentSub, anchor: Anchor, state: &kio::Consumer<ResumeState>) -> Anchor {
let Some(boundary) = seg.last_group() else {
return anchor;
};
let cap = anchor.cap;
let mut capped = anchor.capped(Some(boundary));
if capped.cap != cap {
capped.successor = Some(Successor {
state: state.weak(),
after: seg.id,
boundary,
cap,
});
}
capped
}
fn anchor(&self) -> Anchor {
self.drift_anchor.read().clone()
}
fn refresh_anchor(&mut self) {
let outer = self.outer.clone().capped(self.end_sequence);
let edge = self
.state
.read()
.live_edge(outer.cap)
.into_iter()
.chain(outer.edge.clone())
.max_by_key(|edge| edge.sequence);
let anchor = Anchor { edge, ..outer };
if self.anchor() != anchor
&& let Ok(mut current) = self.drift_anchor.write()
{
*current = anchor.clone();
}
for seg in &mut self.segments {
let anchor = Self::segment_anchor(seg, anchor.clone(), &self.state);
seg.anchor = anchor.clone();
if let Some(sub) = seg.stale_sub_mut() {
sub.set_anchor(anchor);
}
}
}
pub(crate) fn set_anchor(&mut self, anchor: Anchor) {
self.outer = anchor;
self.refresh_anchor();
}
pub(crate) fn commit_seek_stale(&mut self, committed: u64) {
for seg in &mut self.segments {
if let Some(sub) = seg.stale_sub_mut() {
sub.commit_seek_stale(committed);
}
}
}
pub(crate) fn discard_seek_conviction(&mut self, sequence: u64) {
for seg in &mut self.segments {
if let Some(sub) = seg.stale_sub_mut() {
sub.discard_seek_conviction(sequence);
}
}
}
pub(crate) fn poll_stale(&mut self, group: &group::Consumer, waiter: &kio::Waiter) -> Poll<Result<bool>> {
for seg in &mut self.segments {
if seg.first_group() <= group.sequence
&& super::subscription::before_end(group.sequence, seg.last_group())
&& let Some(sub) = seg.stale_sub_mut()
{
return sub.poll_stale(group, waiter);
}
}
Poll::Ready(Ok(false))
}
fn poll_activate(seg: &mut SegmentSub, prefs: &Subscription, min_sequence: u64, waiter: &kio::Waiter) -> Poll<()> {
if matches!(seg.sub, SubState::Pending(_))
&& let Some(warm) = &seg.warm
{
let Some(next) = &warm.next else {
return Poll::Pending;
};
let start = ready!(next.poll_start(waiter));
let edge = warm.edge;
seg.warm = None;
if start.is_some_and(|start| edge.is_some_and(|edge| start > edge)) {
seg.complete(Err(Error::Dropped));
return Poll::Ready(());
}
}
if let SubState::Pending(pending) = &mut seg.sub {
match ready!(pending.poll_ok(waiter)) {
Ok(mut sub) => {
sub.raise_start_to(seg.first_group().max(min_sequence));
sub.set_anchor(seg.anchor.clone());
let _ = sub.update(slice(prefs, seg.ask, seg.end));
seg.sub = SubState::Active(Box::new(sub));
}
Err(err) => seg.sub = SubState::Done(Err(err)),
}
}
Poll::Ready(())
}
fn poll_segment(
seg: &mut SegmentSub,
prefs: &Subscription,
min_sequence: u64,
waiter: &kio::Waiter,
) -> Poll<Option<group::Consumer>> {
loop {
match &mut seg.sub {
SubState::Pending(_) => {
ready!(Self::poll_activate(seg, prefs, min_sequence, waiter));
}
SubState::Active(sub) => match ready!(sub.poll_recv_group(waiter)) {
Ok(Some(group)) => {
if !super::subscription::before_end(group.sequence, seg.last_group()) {
continue;
}
return Poll::Ready(Some(group));
}
Ok(None) => {
let end = match sub.poll_finished(waiter) {
Poll::Ready(end) => end,
Poll::Pending => Err(Error::Dropped),
};
seg.complete(end);
return Poll::Ready(None);
}
Err(err) => {
seg.complete(Err(err));
return Poll::Ready(None);
}
},
SubState::Done(_) => return Poll::Ready(None),
}
}
}
fn poll_seek(
&mut self,
floor: u64,
end: Option<u64>,
deliver: bool,
waiter: &kio::Waiter,
) -> Poll<Result<Option<group::Consumer>>> {
self.poll_sync(waiter);
if deliver {
let committed = self.next_sequence;
self.commit_seek_stale(committed);
}
let mut floor = floor;
'retry: loop {
let mut all_done = true;
let mut best: Option<(usize, group::Consumer)> = None;
for index in 0..self.segments.len() {
if matches!(self.segments[index].sub, SubState::Pending(_))
&& Self::poll_activate(&mut self.segments[index], &self.last_prefs, self.min_sequence, waiter)
.is_pending()
{
all_done = false;
continue;
}
let seg = &mut self.segments[index];
if seg.last_group().is_some_and(|last| last <= floor) {
continue;
}
let cap = min_some(end, seg.last_group());
match &mut seg.sub {
SubState::Active(sub) => match sub.poll_seek_group(floor, cap, waiter) {
Poll::Ready(Ok(Some(group))) => {
all_done = false;
if best.as_ref().is_none_or(|(_, best)| group.sequence < best.sequence) {
best = Some((index, group));
}
}
Poll::Ready(Ok(None)) => {
let end = match sub.poll_finished(waiter) {
Poll::Ready(end) => end,
Poll::Pending => Err(Error::Dropped),
};
seg.complete(end);
}
Poll::Ready(Err(err)) => seg.complete(Err(err)),
Poll::Pending => all_done = false,
},
SubState::Pending(_) => all_done = false,
SubState::Done(_) => {}
}
}
if let Some((index, group)) = best {
let sequence = group.sequence;
if deliver {
group.cache_refresh();
}
match self.hand_out(index, group) {
Some(group) => return Poll::Ready(Ok(Some(group))),
None => {
floor = sequence.saturating_add(1);
continue 'retry;
}
}
}
if let Some(err) = &self.abort {
return Poll::Ready(Err(err.clone()));
}
if all_done && (self.finished || self.closed) {
return Poll::Ready(self.final_end().map(|()| None));
}
return Poll::Pending;
}
}
pub fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
self.poll_sync(waiter);
let end_sequence = self.end_sequence;
let min_sequence = self.min_sequence;
let beyond_cap = |sequence: u64| !super::subscription::before_end(sequence, end_sequence);
let mut all_done = true;
let watch = |group: &group::Consumer| match group.poll_closed(waiter) {
Poll::Pending => true,
Poll::Ready(()) => !group.is_aborted(),
};
for index in 0..self.segments.len() {
self.segments[index]
.parked
.retain(|sequence, group| *sequence >= min_sequence && watch(group));
while let Some(&sequence) = self.segments[index].parked.keys().next() {
if beyond_cap(sequence) {
break;
}
let group = self.segments[index]
.parked
.remove(&sequence)
.expect("parked key just observed");
if let Some(sub) = self.segments[index].stale_sub_mut()
&& matches!(sub.poll_stale(&group, waiter), Poll::Ready(Ok(true)))
{
continue;
}
group.cache_refresh();
if let Some(group) = self.hand_out(index, group) {
self.next_sequence = self.next_sequence.max(sequence.saturating_add(1));
return Poll::Ready(Ok(Some(group)));
}
}
loop {
let polled = Self::poll_segment(&mut self.segments[index], &self.last_prefs, min_sequence, waiter);
match polled {
Poll::Ready(Some(group)) => {
if beyond_cap(group.sequence) {
if watch(&group) {
self.segments[index].parked.insert(group.sequence, group);
}
continue;
}
if group.sequence < min_sequence {
continue;
}
let sequence = group.sequence;
let Some(group) = self.hand_out(index, group) else {
continue;
};
self.next_sequence = self.next_sequence.max(sequence.saturating_add(1));
return Poll::Ready(Ok(Some(group)));
}
Poll::Ready(None) => break,
Poll::Pending => break,
}
}
let seg = &self.segments[index];
if !seg.parked.is_empty() || !matches!(seg.sub, SubState::Done(_)) {
all_done = false;
}
}
if let Some(err) = &self.abort {
return Poll::Ready(Err(err.clone()));
}
if all_done && (self.finished || self.closed) {
return Poll::Ready(self.final_end().map(|()| None));
}
Poll::Pending
}
#[cfg(test)]
pub async fn recv_group(&mut self) -> Result<Option<group::Consumer>> {
kio::wait(|waiter| self.poll_recv_group(waiter)).await
}
pub fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
let floor = self.next_sequence.max(self.min_sequence);
let Some(group) = ready!(self.poll_seek(floor, self.end_sequence, true, waiter))? else {
return Poll::Ready(Ok(None));
};
self.next_sequence = self.next_sequence.max(group.sequence.saturating_add(1));
self.discard_seek_conviction(group.sequence);
self.commit_seek_stale(self.next_sequence);
Poll::Ready(Ok(Some(group)))
}
pub(crate) fn poll_seek_group(
&mut self,
floor: u64,
end: Option<u64>,
waiter: &kio::Waiter,
) -> Poll<Result<Option<group::Consumer>>> {
let floor = floor.max(self.min_sequence);
let end = min_some(end, self.end_sequence);
self.poll_seek(floor, end, false, waiter)
}
#[cfg(test)]
pub async fn next_group(&mut self) -> Result<Option<group::Consumer>> {
kio::wait(|waiter| self.poll_next_group(waiter)).await
}
pub fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
self.poll_sync(waiter);
let mut pending_activation = false;
if let Some(seg) = self.segments.last_mut() {
if Self::poll_activate(seg, &self.last_prefs, self.min_sequence, waiter).is_pending() {
pending_activation = true;
} else if let SubState::Active(sub) = &mut seg.sub
&& let Ok(Some(datagram)) = ready!(sub.poll_recv_datagram(waiter))
{
return Poll::Ready(Ok(Some(datagram)));
}
}
if let Some(err) = &self.abort {
return Poll::Ready(Err(err.clone()));
}
if self.finished {
if pending_activation {
return Poll::Ready(Ok(None));
}
return match ready!(self.poll_final(waiter)) {
Some(end) => Poll::Ready(end.map(|_| None)),
None => Poll::Ready(Ok(None)),
};
}
if self.closed && !pending_activation {
return match ready!(self.poll_final(waiter)) {
Some(end) => Poll::Ready(end.map(|_| None)),
None => Poll::Ready(Err(Error::Dropped)),
};
}
Poll::Pending
}
#[cfg(test)]
pub async fn closed(&mut self) -> Result<()> {
kio::wait(|waiter| self.poll_finished(waiter)).await.map(|_| ())
}
pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
self.poll_sync(waiter);
if let Some(err) = &self.abort {
return Poll::Ready(Err(err.clone()));
}
if !self.finished && !self.closed {
return Poll::Pending;
}
match ready!(self.poll_final(waiter)) {
Some(end) => Poll::Ready(end),
None if self.finished => Poll::Ready(Ok(0)),
None => Poll::Ready(Err(Error::Dropped)),
}
}
pub(crate) fn poll_start(&mut self, waiter: &kio::Waiter) -> Poll<Option<u64>> {
self.poll_sync(waiter);
let floor = self.min_sequence;
for seg in &mut self.segments {
if seg.last_group().is_some_and(|last| last <= floor) {
continue;
}
ready!(Self::poll_activate(seg, &self.last_prefs, floor, waiter));
if let SubState::Active(sub) = &mut seg.sub {
return Poll::Ready(ready!(sub.poll_start(waiter)).map(|start| start.max(floor)));
}
}
Poll::Ready(None)
}
fn poll_final(&mut self, waiter: &kio::Waiter) -> Poll<Option<Result<u64>>> {
let Some(seg) = self.segments.last_mut() else {
return Poll::Ready(None);
};
ready!(Self::poll_activate(seg, &self.last_prefs, self.min_sequence, waiter));
match &mut seg.sub {
SubState::Done(end) => Poll::Ready(Some(end.clone())),
SubState::Active(sub) => Poll::Ready(Some(ready!(sub.poll_finished(waiter)))),
SubState::Pending(_) => unreachable!("poll_activate resolved above"),
}
}
fn final_end(&self) -> Result<()> {
match self.segments.last().map(|seg| &seg.sub) {
Some(SubState::Done(end)) => end.clone().map(|_| ()),
None if self.finished => Ok(()),
_ => Err(Error::Dropped),
}
}
#[cfg(test)]
pub async fn finished(&mut self) -> Result<u64> {
kio::wait(|waiter| self.poll_finished(waiter)).await
}
pub fn start_at(&mut self, sequence: u64) {
self.min_sequence = sequence;
for seg in &mut self.segments {
let floor = seg.first_group().max(sequence);
if let SubState::Active(sub) = &mut seg.sub {
sub.start_at(floor);
}
}
}
pub(crate) fn raise_start_to(&mut self, sequence: u64) {
self.min_sequence = self.min_sequence.max(sequence);
for seg in &mut self.segments {
let floor = seg.first_group().max(sequence);
if let SubState::Active(sub) = &mut seg.sub {
sub.raise_start_to(floor);
}
}
}
pub fn end_at(&mut self, end: impl Into<Cap>) {
self.end_sequence = end.into().exclusive();
self.refresh_anchor();
}
pub(crate) fn prefs(&self) -> kio::Producer<Subscription> {
self.prefs.clone()
}
pub(crate) fn take_stale(&mut self) -> crate::stats::Content {
let mut stale = crate::stats::Content::default();
for seg in &mut self.segments {
if let Some(sub) = seg.stale_sub_mut() {
stale.add(sub.take_stale());
}
}
stale
}
pub fn update(&mut self, subscription: Subscription) {
if let Ok(mut prefs) = self.prefs.write() {
*prefs = subscription;
}
}
pub fn latest(&self) -> Option<u64> {
self.state.read().latest()
}
pub fn is_clone(&self, other: &Self) -> bool {
self.state.same_channel(&other.state)
}
}
#[cfg(test)]
mod test {
use super::*;
use crate::broadcast;
use futures::FutureExt;
use std::sync::Arc;
use std::time::Duration;
fn track_pair(name: &str) -> (track::Producer, track::Consumer) {
let producer = track::Producer::new(Arc::new(broadcast::Info::default()), name, None);
let consumer = producer.consume();
(producer, consumer)
}
fn track_pair_with(name: &str, info: track::Info) -> (track::Producer, track::Consumer) {
let producer = track::Producer::new(Arc::new(broadcast::Info::default()), name, info);
let consumer = producer.consume();
(producer, consumer)
}
fn replay() -> Subscription {
Subscription::default().with_max_age(std::time::Duration::from_secs(30))
}
fn write_group(producer: &mut track::Producer, sequence: u64, payload: &str) {
let mut group = producer.create_group(group::Info { sequence }).unwrap();
group.write_frame(Timestamp::ZERO, payload.as_bytes().to_vec()).unwrap();
group.finish().unwrap();
}
fn write_group_at(producer: &mut track::Producer, sequence: u64, payload: &str, at: Duration) {
let mut group = producer.create_group(group::Info { sequence }).unwrap();
group
.write_frame(at.try_into().unwrap(), payload.as_bytes().to_vec())
.unwrap();
group.finish().unwrap();
}
fn recv(sub: &mut Subscriber) -> u64 {
sub.recv_group()
.now_or_never()
.expect("should not block")
.expect("should not error")
.expect("should not be finished")
.sequence
}
fn read(group: &mut group::Consumer) -> Vec<u8> {
group
.read_frame()
.now_or_never()
.expect("should not block")
.expect("should not error")
.expect("should not be finished")
.payload
.to_vec()
}
fn recv_pending(sub: &mut Subscriber) {
assert!(sub.recv_group().now_or_never().is_none(), "should have blocked");
}
struct CountWaker(std::sync::atomic::AtomicUsize);
impl std::task::Wake for CountWaker {
fn wake(self: Arc<Self>) {
self.0.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
}
}
impl CountWaker {
fn new() -> (Arc<Self>, std::task::Waker) {
let counter = Arc::new(Self(std::sync::atomic::AtomicUsize::new(0)));
(counter.clone(), std::task::Waker::from(counter))
}
fn count(&self) -> usize {
self.0.load(std::sync::atomic::Ordering::SeqCst)
}
}
#[tokio::test]
async fn switch_splices_groups() {
let (mut track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(replay());
write_group(&mut track_a, 0, "a0");
write_group(&mut track_a, 1, "a1");
assert_eq!(recv(&mut sub), 0);
assert_eq!(recv(&mut sub), 1);
producer.switch(&consumer_b, Position::group(2)).unwrap();
write_group(&mut track_a, 2, "a2-over-cap");
write_group(&mut track_b, 2, "b2");
write_group(&mut track_b, 3, "b3");
assert_eq!(recv(&mut sub), 2);
assert_eq!(recv(&mut sub), 3);
recv_pending(&mut sub);
}
#[tokio::test]
async fn poll_start_skips_a_segment_below_the_floor() {
let (mut track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
track_a.request_start(Some(0)).unwrap();
track_b.request_start(Some(5)).unwrap();
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
producer.switch(&consumer_b, Position::group(5)).unwrap();
let mut sub = producer.consume().subscribe(replay());
sub.start_at(10);
assert!(
kio::wait(|waiter| sub.poll_start(waiter)).now_or_never().is_none(),
"B's source has not resolved yet"
);
track_b.start_at(12).unwrap();
let start = kio::wait(|waiter| sub.poll_start(waiter)).now_or_never();
assert_eq!(start, Some(Some(12)), "waited on a segment below the floor");
}
#[tokio::test]
async fn poll_start_skips_an_ended_segment() {
let (track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
track_b.request_start(Some(5)).unwrap();
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
producer.switch(&consumer_b, Position::group(5)).unwrap();
track_a.abort(Error::Cancel).unwrap();
let mut sub = producer.consume().subscribe(replay());
recv_pending(&mut sub);
assert!(
kio::wait(|waiter| sub.poll_start(waiter)).now_or_never().is_none(),
"B's source has not resolved yet"
);
track_b.start_at(5).unwrap();
let start = kio::wait(|waiter| sub.poll_start(waiter)).now_or_never();
assert_eq!(start, Some(Some(5)));
}
#[tokio::test]
async fn demand_reflects_boundaries() {
let (track_a, consumer_a) = track_pair("a");
let (track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer
.consume()
.subscribe(Subscription::default().with_start(Position::group(0)));
recv_pending(&mut sub);
assert_eq!(track_a.subscription().unwrap().end, None);
producer.switch(&consumer_b, Position::group(5)).unwrap();
recv_pending(&mut sub);
assert_eq!(track_a.subscription().unwrap().end, Some(Position::group(5)));
assert_eq!(track_b.subscription().unwrap().start, Some(Position::group(5)));
}
#[tokio::test]
async fn update_reslices_demand() {
let (track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
recv_pending(&mut sub);
assert_eq!(track_a.subscription().unwrap().priority, 0);
sub.update(Subscription::default().with_priority(7));
recv_pending(&mut sub);
assert_eq!(track_a.subscription().unwrap().priority, 7);
}
#[tokio::test]
async fn dead_segment_stalls_until_switch() {
let (mut track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
write_group(&mut track_a, 0, "a0");
assert_eq!(recv(&mut sub), 0);
track_a.abort(Error::Dropped).unwrap();
recv_pending(&mut sub);
producer.switch(&consumer_b, Position::group(1)).unwrap();
write_group(&mut track_b, 1, "b1");
assert_eq!(recv(&mut sub), 1);
}
#[tokio::test]
async fn takeover_computes_boundary() {
let (mut track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(replay());
write_group(&mut track_a, 0, "a0");
write_group(&mut track_a, 1, "a1");
assert_eq!(recv(&mut sub), 0);
assert_eq!(recv(&mut sub), 1);
track_a.abort(Error::Dropped).unwrap();
producer.takeover(&consumer_b).unwrap();
write_group(&mut track_b, 2, "b2");
assert_eq!(recv(&mut sub), 2);
}
#[tokio::test]
async fn a_takeover_backlog_is_bounded_by_the_budget() {
let retain = track::Info::default().with_max_age(std::time::Duration::from_secs(60));
let (mut track_a, consumer_a) = track_pair_with("a", retain.clone());
let (mut track_b, consumer_b) = track_pair_with("b", retain);
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut live = producer.consume().subscribe(None);
let mut patient = producer
.consume()
.subscribe(Subscription::default().with_max_age(std::time::Duration::from_secs(60)));
write_group_at(&mut track_a, 0, "a0", std::time::Duration::ZERO);
assert_eq!(recv(&mut live), 0);
assert_eq!(recv(&mut patient), 0);
track_a.abort(Error::Dropped).unwrap();
producer.takeover(&consumer_b).unwrap();
recv_pending(&mut live);
for second in 1..=30 {
write_group_at(&mut track_b, second, "b", std::time::Duration::from_secs(second));
}
assert_eq!(recv(&mut live), 30);
recv_pending(&mut live);
let mut backfill = Vec::new();
while let Some(Ok(Some(group))) = patient.recv_group().now_or_never() {
backfill.push(group.sequence);
}
assert_eq!(backfill, (1..=30).collect::<Vec<_>>());
}
#[tokio::test]
async fn a_capped_spliced_subscriber_measures_drift_against_its_cap() {
let retain = track::Info::default().with_max_age(std::time::Duration::from_secs(60));
let (mut track_a, consumer_a) = track_pair_with("a", retain);
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
sub.end_at(..1);
write_group_at(&mut track_a, 0, "a0", std::time::Duration::ZERO);
write_group_at(&mut track_a, 1, "a1", std::time::Duration::from_secs(30));
assert_eq!(recv(&mut sub), 0);
}
#[tokio::test]
async fn a_raised_cap_rechecks_parked_groups() {
let retain = track::Info::default().with_max_age(std::time::Duration::from_secs(60));
let (mut track_a, consumer_a) = track_pair_with("a", retain);
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
sub.end_at(..1);
write_group_at(&mut track_a, 0, "a0", std::time::Duration::ZERO);
assert_eq!(recv(&mut sub), 0);
write_group_at(&mut track_a, 1, "a1", std::time::Duration::from_secs(1));
recv_pending(&mut sub);
write_group_at(&mut track_a, 2, "a2", std::time::Duration::from_secs(30));
sub.end_at(..);
assert_eq!(recv(&mut sub), 2);
recv_pending(&mut sub);
}
#[tokio::test]
async fn a_raised_cap_rechecks_parked_groups_after_segment_finish() {
let retain = track::Info::default().with_max_age(std::time::Duration::from_secs(60));
let (mut track_a, consumer_a) = track_pair_with("a", retain);
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
sub.end_at(..1);
write_group_at(&mut track_a, 0, "a0", std::time::Duration::ZERO);
assert_eq!(recv(&mut sub), 0);
write_group_at(&mut track_a, 1, "a1", std::time::Duration::from_secs(1));
write_group_at(&mut track_a, 2, "a2", std::time::Duration::from_secs(30));
track_a.finish().unwrap();
recv_pending(&mut sub);
sub.end_at(..);
assert_eq!(recv(&mut sub), 2, "the stale parked group is skipped after finish");
recv_pending(&mut sub);
}
#[tokio::test]
async fn takeover_replaces_empty_segment() {
let (track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
recv_pending(&mut sub);
drop(track_a);
producer.takeover(&consumer_b).unwrap();
write_group(&mut track_b, 0, "b0");
assert_eq!(recv(&mut sub), 0);
}
#[tokio::test]
async fn finish_ends_after_final_segment() {
let (mut track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
write_group(&mut track_a, 0, "a0");
assert_eq!(recv(&mut sub), 0);
producer.finish().unwrap();
recv_pending(&mut sub);
track_a.finish().unwrap();
assert!(
sub.recv_group()
.now_or_never()
.expect("should not block")
.expect("should not error")
.is_none(),
"should be finished"
);
assert_eq!(sub.finished().now_or_never().unwrap().unwrap(), 1);
assert!(sub.closed().now_or_never().unwrap().is_ok());
}
#[tokio::test]
async fn next_group_across_segments() {
let (mut track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(replay());
write_group(&mut track_a, 0, "a0");
producer.switch(&consumer_b, Position::group(1)).unwrap();
write_group(&mut track_b, 1, "b1");
let mut read = |expected: &[u8]| {
let mut group = sub.next_group().now_or_never().unwrap().unwrap().unwrap();
let frame = group.read_frame().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(&frame.payload[..], expected);
};
read(b"a0");
read(b"b1");
}
#[tokio::test]
async fn info_from_first_segment() {
let (_track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
let consumer = producer.consume();
assert!(consumer.info().now_or_never().is_none());
producer.switch(&consumer_a, None).unwrap();
let info = consumer.info().now_or_never().unwrap().unwrap();
assert_eq!(info.timescale, crate::Timescale::default());
}
#[tokio::test]
async fn fetch_uses_an_older_complete_cache_before_the_newest_segment() {
let (mut track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
producer.switch(&consumer_b, Position::group(10)).unwrap();
write_group(&mut track_a, 3, "a3");
let consumer = producer.consume();
let mut group = consumer
.fetch_group(3, None)
.now_or_never()
.expect("cached fetch should resolve")
.unwrap();
assert_eq!(group.sequence, 3);
assert_eq!(read(&mut group), b"a3");
write_group(&mut track_b, 4, "b4");
let mut group = consumer
.fetch_group(4, None)
.now_or_never()
.expect("newest cached fetch should resolve")
.unwrap();
assert_eq!(read(&mut group), b"b4");
}
#[tokio::test]
async fn fetch_waits_for_first_segment() {
let (mut track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
let consumer = producer.consume();
let fetch = consumer.fetch_group(0, None);
let mut fetch = std::pin::pin!(fetch);
assert!(futures::poll!(fetch.as_mut()).is_pending(), "fetch should wait");
write_group(&mut track_a, 0, "a0");
producer.switch(&consumer_a, None).unwrap();
let group = fetch.await.expect("fetch should resolve");
assert_eq!(group.sequence, 0);
}
#[tokio::test]
async fn pending_fetch_follows_a_newer_segment() {
let (track_a, consumer_a) = track_pair("a");
let dynamic_a = track_a.dynamic();
let (track_b, consumer_b) = track_pair("b");
let dynamic_b = track_b.dynamic();
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let consumer = producer.consume();
let fetch = tokio::spawn(async move { consumer.fetch_group(0, None).await });
let _request_a = tokio::time::timeout(Duration::from_secs(1), dynamic_a.requested_group())
.await
.expect("first segment should receive the fetch")
.unwrap();
producer.takeover(&consumer_b).unwrap();
let request_b = tokio::time::timeout(Duration::from_secs(1), dynamic_b.requested_group())
.await
.expect("pending fetch should follow the takeover")
.unwrap();
let mut served = request_b.accept(None).unwrap();
served.write_frame(Timestamp::ZERO, b"b0".to_vec()).unwrap();
served.finish().unwrap();
let mut group = tokio::time::timeout(Duration::from_secs(1), fetch)
.await
.expect("replacement should answer the fetch")
.unwrap()
.unwrap();
assert_eq!(read(&mut group), b"b0");
}
#[test]
fn warm_snapshot_skips_only_the_frame_seam() {
let (mut track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
write_group(&mut track_a, 0, "a0");
let mut seam = track_a.create_group(group::Info { sequence: 1 }).unwrap();
seam.write_frame(Timestamp::ZERO, b"a1".to_vec()).unwrap();
producer.takeover(&consumer_b).unwrap();
let mut seam = track_b.create_group(group::Info { sequence: 1 }).unwrap();
seam.start_at(1).unwrap();
seam.write_frame(Timestamp::ZERO, b"b1".to_vec()).unwrap();
seam.finish().unwrap();
write_group(&mut track_b, 2, "b2");
let cached = producer.consume().cached_groups();
assert_eq!(
cached.iter().map(|(group, _)| group.sequence).collect::<Vec<_>>(),
vec![0, 2],
"only the group split across both segments is unsafe to snapshot"
);
}
#[tokio::test]
async fn takeover_survives_dead_empty_segment() {
let (mut track_a, consumer_a) = track_pair("a");
let (track_b, consumer_b) = track_pair("b");
let (mut track_c, consumer_c) = track_pair("c");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
write_group(&mut track_a, 0, "a0");
assert_eq!(recv(&mut sub), 0);
track_a.abort(Error::Dropped).unwrap();
producer.takeover(&consumer_b).unwrap();
drop(track_b);
producer.takeover(&consumer_c).unwrap();
write_group(&mut track_c, 1, "c1");
assert_eq!(recv(&mut sub), 1);
}
#[tokio::test]
async fn finished_does_not_consume_groups() {
let (mut track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
write_group(&mut track_a, 0, "a0");
producer.finish().unwrap();
assert!(sub.finished().now_or_never().is_none(), "final segment still open");
assert_eq!(recv(&mut sub), 0);
track_a.finish().unwrap();
assert_eq!(sub.finished().now_or_never().unwrap().unwrap(), 1);
}
#[tokio::test]
async fn datagram_only_subscriber_activates() {
let (mut track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
assert!(
kio::wait(|waiter| sub.poll_recv_datagram(waiter))
.now_or_never()
.is_none(),
"no datagram yet"
);
track_a.append_datagram(Timestamp::ZERO, b"d0".as_ref()).unwrap();
let datagram = kio::wait(|waiter| sub.poll_recv_datagram(waiter))
.now_or_never()
.expect("datagram should be ready")
.expect("should not error")
.expect("track should not be finished");
assert_eq!(&datagram.payload[..], b"d0");
}
#[tokio::test]
async fn end_at_parks_at_cap() {
let (mut track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
write_group(&mut track_a, 0, "a0");
write_group(&mut track_a, 1, "a1");
sub.end_at(..1);
assert_eq!(recv(&mut sub), 0);
recv_pending(&mut sub);
sub.end_at(..2);
assert_eq!(recv(&mut sub), 1);
}
#[tokio::test]
async fn end_at_reoffers_reordered_arrivals() {
let (mut track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(replay());
sub.end_at(..2);
write_group(&mut track_a, 2, "a2");
write_group(&mut track_a, 0, "a0");
write_group(&mut track_a, 1, "a1");
assert_eq!(recv(&mut sub), 0);
assert_eq!(recv(&mut sub), 1);
recv_pending(&mut sub);
sub.end_at(..3);
assert_eq!(recv(&mut sub), 2);
}
#[tokio::test]
async fn evicted_parked_groups_are_dropped() {
let (mut track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
sub.end_at(..1);
write_group(&mut track_a, 0, "a0");
assert_eq!(recv(&mut sub), 0);
let straggler = track_a.create_group(group::Info { sequence: 1 }).unwrap();
recv_pending(&mut sub);
straggler.abort(Error::Old).unwrap();
sub.end_at(..);
write_group(&mut track_a, 2, "a2");
assert_eq!(recv(&mut sub), 2, "the evicted parked group is dropped, not re-offered");
}
#[tokio::test]
async fn next_group_skips_boundary_duplicate() {
let (mut track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(replay());
let next = |sub: &mut Subscriber| {
kio::wait(|waiter| sub.poll_next_group(waiter))
.now_or_never()
.expect("should not block")
.expect("should not error")
.expect("should not be finished")
.sequence
};
write_group(&mut track_a, 0, "a0");
write_group(&mut track_a, 1, "a1");
assert_eq!(next(&mut sub), 0);
assert_eq!(next(&mut sub), 1);
producer.switch(&consumer_b, Position::group(1)).unwrap();
write_group(&mut track_b, 1, "b1");
write_group(&mut track_b, 2, "b2");
assert_eq!(next(&mut sub), 2);
}
#[tokio::test]
async fn next_group_cap_holds_a_reordered_group_without_losing_late_arrivals() {
let (mut track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(replay());
let next = |sub: &mut Subscriber| {
kio::wait(|waiter| sub.poll_next_group(waiter))
.now_or_never()
.expect("should not block")
.expect("should not error")
.expect("should not be finished")
.sequence
};
let next_pending = |sub: &mut Subscriber| {
assert!(
kio::wait(|waiter| sub.poll_next_group(waiter)).now_or_never().is_none(),
"should have blocked"
);
};
sub.end_at(..2);
write_group(&mut track_a, 2, "a2");
next_pending(&mut sub);
write_group(&mut track_a, 0, "a0");
write_group(&mut track_a, 1, "a1");
assert_eq!(next(&mut sub), 0);
assert_eq!(next(&mut sub), 1);
next_pending(&mut sub);
sub.end_at(..3);
assert_eq!(next(&mut sub), 2);
}
#[tokio::test]
async fn next_group_cap_lowering_after_a_lookahead_keeps_late_arrivals() {
let (mut track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
producer.switch(&consumer_b, Position::group(1)).unwrap();
let mut sub = producer.consume().subscribe(replay());
let next = |sub: &mut Subscriber| {
kio::wait(|waiter| sub.poll_next_group(waiter))
.now_or_never()
.expect("should not block")
.expect("should not error")
.expect("should not be finished")
.sequence
};
let next_pending = |sub: &mut Subscriber| {
assert!(
kio::wait(|waiter| sub.poll_next_group(waiter)).now_or_never().is_none(),
"should have blocked"
);
};
write_group(&mut track_b, 2, "b2");
write_group(&mut track_a, 0, "a0");
assert_eq!(next(&mut sub), 0);
sub.end_at(..2);
next_pending(&mut sub);
write_group(&mut track_b, 1, "b1");
assert_eq!(next(&mut sub), 1);
next_pending(&mut sub);
sub.end_at(..);
assert_eq!(next(&mut sub), 2);
}
#[tokio::test]
async fn next_group_sheds_a_stale_backlog_like_a_plain_track() {
let stamp = |sequence: u64| Duration::from_secs(10 * sequence);
let (mut plain, _plain_consumer) = track_pair("plain");
for sequence in 0..4 {
write_group_at(&mut plain, sequence, "data", stamp(sequence));
}
let mut plain_ordered = plain.subscribe(None).ordered();
let baseline: Vec<u64> = std::iter::from_fn(|| {
plain_ordered
.next_group()
.now_or_never()?
.expect("should not error")
.map(|group| group.sequence)
})
.collect();
assert_eq!(baseline, vec![3], "a plain sequence cursor sheds the backlog");
let mut arrival = plain.subscribe(None);
assert_eq!(
arrival.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3,
"the arrival path skips straight to the live edge"
);
let retain = track::Info::default().with_max_age(Duration::from_secs(60));
let (mut track_a, consumer_a) = track_pair_with("a", retain.clone());
let (mut track_b, consumer_b) = track_pair_with("b", retain);
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
producer.switch(&consumer_b, Position::group(2)).unwrap();
let mut sub = producer.consume().subscribe(None);
write_group_at(&mut track_a, 0, "a0", stamp(0));
write_group_at(&mut track_a, 1, "a1", stamp(1));
write_group_at(&mut track_b, 2, "b2", stamp(2));
write_group_at(&mut track_b, 3, "b3", stamp(3));
let spliced: Vec<u64> = std::iter::from_fn(|| {
kio::wait(|waiter| sub.poll_next_group(waiter))
.now_or_never()?
.expect("should not error")
.map(|group| group.sequence)
})
.collect();
assert_eq!(spliced, baseline, "a splice sheds the backlog like one track");
let mut arrival = producer.consume().subscribe(None);
let arrival: Vec<u64> = std::iter::from_fn(|| {
kio::wait(|waiter| arrival.poll_recv_group(waiter))
.now_or_never()?
.expect("should not error")
.map(|group| group.sequence)
})
.collect();
assert_eq!(arrival, spliced, "both spliced cursors shed the same backlog");
let mut replay = producer
.consume()
.subscribe(Subscription::default().with_max_age(Duration::from_secs(60)));
let replayed: Vec<u64> = std::iter::from_fn(|| {
kio::wait(|waiter| replay.poll_next_group(waiter))
.now_or_never()?
.expect("should not error")
.map(|group| group.sequence)
})
.collect();
assert_eq!(replayed, vec![0, 1, 2, 3], "a backlog inside the budget crosses whole");
}
#[tokio::test]
async fn unstamped_successor_segment_keeps_the_previous_group_unbounded() {
let (mut a, a_read) = track_pair("a");
let (b, b_read) = track_pair("b");
let (mut c, c_read) = track_pair("c");
write_group_at(&mut a, 0, "a0", Duration::ZERO);
let _unstamped = b.create_group(1u64.into()).unwrap();
write_group_at(&mut c, 2, "c2", Duration::from_secs(10));
write_group_at(&mut c, 3, "c3", Duration::from_secs(20));
let mut producer = Producer::new();
producer.switch(a_read, None).unwrap();
producer.switch(b_read, Position::group(1)).unwrap();
producer.switch(c_read, Position::group(2)).unwrap();
let mut sub = producer.consume().subscribe(None);
assert_eq!(
recv(&mut sub),
0,
"the unstamped successor must not borrow segment C's start"
);
}
#[tokio::test]
async fn unstamped_successor_first_frame_wakes_a_parked_read() {
use std::task::Context;
let (a, a_read) = track_pair("a");
let (mut b, b_read) = track_pair("b");
let mut head = a.create_group(0u64.into()).unwrap();
head.write_frame(Timestamp::ZERO, b"a0".to_vec()).unwrap();
let mut successor = b.create_group(1u64.into()).unwrap();
write_group_at(&mut b, 2, "edge", Duration::from_secs(20));
let mut producer = Producer::new();
producer.switch(a_read, None).unwrap();
producer.switch(b_read, Position::group(1)).unwrap();
let mut sub = producer.consume().subscribe(None);
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(reading.sequence, 0, "an unbounded reach is not stale");
assert_eq!(read(&mut reading), b"a0");
let (counter, waker) = CountWaker::new();
let mut cx = Context::from_waker(&waker);
let mut next = std::pin::pin!(reading.read_frame());
assert!(next.as_mut().poll(&mut cx).is_pending());
let before = counter.count();
successor
.write_frame(Duration::from_secs(1).try_into().unwrap(), b"b1".to_vec())
.unwrap();
assert!(counter.count() > before, "the successor's first frame lost its wakeup");
let result = next.as_mut().poll(&mut cx);
assert!(matches!(result, Poll::Ready(Ok(None))), "the head is stale: {result:?}");
}
#[tokio::test]
async fn aborted_unstamped_successor_re_resolves_a_parked_read() {
use std::task::Context;
let (a, a_read) = track_pair("a");
let (mut b, b_read) = track_pair("b");
let mut head = a.create_group(0u64.into()).unwrap();
head.write_frame(Timestamp::ZERO, b"a0".to_vec()).unwrap();
let successor = b.create_group(1u64.into()).unwrap();
write_group_at(&mut b, 3, "edge", Duration::from_secs(20));
let mut producer = Producer::new();
producer.switch(a_read, None).unwrap();
producer.switch(b_read, Position::group(1)).unwrap();
let budget = Subscription::default().with_max_age(Duration::from_secs(10));
let mut sub = producer.consume().subscribe(budget);
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(reading.sequence, 0, "an unbounded reach is not stale");
assert_eq!(read(&mut reading), b"a0");
let (counter, waker) = CountWaker::new();
let mut cx = Context::from_waker(&waker);
let mut next = std::pin::pin!(reading.read_frame());
assert!(next.as_mut().poll(&mut cx).is_pending());
successor.abort(Error::Cancel).unwrap();
assert!(next.as_mut().poll(&mut cx).is_pending());
let before = counter.count();
write_group_at(&mut b, 2, "b2", Duration::from_secs(1));
assert!(counter.count() > before, "the replacement successor lost its wakeup");
let result = next.as_mut().poll(&mut cx);
assert!(matches!(result, Poll::Ready(Ok(None))), "the head is stale: {result:?}");
}
#[tokio::test]
async fn aborted_unstamped_successor_re_resolves_across_segments() {
use std::task::Context;
let (a, a_read) = track_pair("a");
let (b, b_read) = track_pair("b");
let (mut c, c_read) = track_pair("c");
let mut head = a.create_group(0u64.into()).unwrap();
head.write_frame(Timestamp::ZERO, b"a0".to_vec()).unwrap();
let successor = b.create_group(1u64.into()).unwrap();
write_group_at(&mut c, 2, "c2", Duration::from_secs(1));
write_group_at(&mut c, 3, "c3", Duration::from_secs(20));
let mut producer = Producer::new();
producer.switch(a_read, None).unwrap();
producer.switch(b_read, Position::group(1)).unwrap();
producer.switch(c_read, Position::group(2)).unwrap();
let mut sub = producer.consume().subscribe(None);
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(reading.sequence, 0, "an unbounded reach is not stale");
assert_eq!(read(&mut reading), b"a0");
let (counter, waker) = CountWaker::new();
let mut cx = Context::from_waker(&waker);
let mut next = std::pin::pin!(reading.read_frame());
assert!(next.as_mut().poll(&mut cx).is_pending());
let before = counter.count();
successor.abort(Error::Cancel).unwrap();
assert!(counter.count() > before, "the successor's abort lost its wakeup");
let result = next.as_mut().poll(&mut cx);
assert!(matches!(result, Poll::Ready(Ok(None))), "the head is stale: {result:?}");
}
#[tokio::test]
async fn pruned_segment_boundary_is_judged_against_later_segments() {
let (mut a, a_read) = track_pair("a");
let mut producer = Producer::new();
producer.switch(a_read, None).unwrap();
let mut sub = producer.consume().subscribe(None);
write_group_at(&mut a, 0, "a0", Duration::ZERO);
assert_eq!(recv(&mut sub), 0);
write_group_at(&mut a, 3, "past-boundary", Duration::from_secs(3));
for sequence in 3..=(2 + MAX_SEGMENTS as u64) {
let (mut track, consumer) = track_pair("later");
producer.switch(consumer, Position::group(sequence)).unwrap();
write_group_at(&mut track, sequence, "later", Duration::from_secs(sequence * 10));
assert_eq!(recv(&mut sub), sequence);
}
assert!(producer.state.read().pruned.is_some());
write_group_at(&mut a, 2, "a2", Duration::from_secs(2));
recv_pending(&mut sub);
}
#[tokio::test]
async fn a_capped_segment_is_judged_like_one_track() {
let budget = Subscription::default().with_max_age(Duration::from_millis(100));
let old = [(0, 0), (1, 20), (2, 40), (3, 60)];
let new = [(20, 400), (21, 420), (22, 440), (23, 460), (24, 480), (25, 500)];
let (mut plain, _plain_consumer) = track_pair("plain");
for (sequence, millis) in old.into_iter().chain(new) {
write_group_at(&mut plain, sequence, "p", Duration::from_millis(millis));
}
let mut baseline = plain.subscribe(budget.clone());
let baseline: Vec<u64> = std::iter::from_fn(|| {
baseline
.recv_group()
.now_or_never()?
.expect("should not error")
.map(|group| group.sequence)
})
.collect();
assert_eq!(baseline, vec![20, 21, 22, 23, 24, 25]);
let (mut track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
producer.switch(&consumer_b, Position::group(4)).unwrap();
for (sequence, millis) in old {
write_group_at(&mut track_a, sequence, "a", Duration::from_millis(millis));
}
for (sequence, millis) in new {
write_group_at(&mut track_b, sequence, "b", Duration::from_millis(millis));
}
let mut arrival = producer.consume().subscribe(budget.clone());
let arrival: Vec<u64> = std::iter::from_fn(|| {
kio::wait(|waiter| arrival.poll_recv_group(waiter))
.now_or_never()?
.expect("should not error")
.map(|group| group.sequence)
})
.collect();
assert_eq!(arrival, baseline);
let mut ordered = producer.consume().subscribe(budget);
let ordered: Vec<u64> = std::iter::from_fn(|| {
kio::wait(|waiter| ordered.poll_next_group(waiter))
.now_or_never()?
.expect("should not error")
.map(|group| group.sequence)
})
.collect();
assert_eq!(ordered, baseline);
}
#[tokio::test]
async fn nested_splice_judges_within_the_outer_boundary() {
let stamp = |sequence: u64| Duration::from_secs(10 * sequence);
let retain = track::Info::default().with_max_age(Duration::from_secs(60));
let budget = Subscription::default().with_max_age(Duration::from_secs(15));
let (mut track_a, consumer_a) = track_pair_with("a", retain.clone());
let mut inner = Producer::new();
inner.switch(&consumer_a, None).unwrap();
write_group_at(&mut track_a, 0, "a0", stamp(0));
write_group_at(&mut track_a, 1, "a1", stamp(1));
write_group_at(&mut track_a, 2, "a2", stamp(0));
let inner_track =
track::Consumer::spliced("inner".into(), Arc::new(broadcast::Info::default()), inner.consume());
let (mut track_b, consumer_b) = track_pair_with("b", retain);
let mut outer = Producer::new();
outer.switch(&inner_track, None).unwrap();
outer.switch(&consumer_b, Position::group(2)).unwrap();
write_group_at(&mut track_b, 2, "b2", stamp(2));
write_group_at(&mut track_b, 3, "b3", stamp(3));
let mut sub = outer.consume().subscribe(budget.clone());
let sequences: Vec<u64> = std::iter::from_fn(|| {
kio::wait(|waiter| sub.poll_next_group(waiter))
.now_or_never()?
.expect("should not error")
.map(|group| group.sequence)
})
.collect();
assert_eq!(
sequences,
vec![1, 2, 3],
"the nested segment is judged within the outer boundary"
);
let mut arrival = outer.consume().subscribe(budget);
let arrival: Vec<u64> = std::iter::from_fn(|| {
kio::wait(|waiter| arrival.poll_recv_group(waiter))
.now_or_never()?
.expect("should not error")
.map(|group| group.sequence)
})
.collect();
assert_eq!(arrival, sequences, "the arrival path sheds the same backlog");
}
#[tokio::test]
async fn seek_conviction_counts_only_once_committed() {
let stamp = |sequence: u64| Duration::from_secs([0, 10, 40, 30][sequence as usize]);
let (mut track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
producer.switch(&consumer_b, Position::group(2)).unwrap();
write_group_at(&mut track_a, 0, "a0", stamp(0));
write_group_at(&mut track_a, 1, "a1", stamp(1));
write_group_at(&mut track_b, 2, "b2", stamp(2));
write_group_at(&mut track_b, 3, "b3", stamp(3));
let mut sub = producer.consume().subscribe(None);
let next = |sub: &mut Subscriber| {
kio::wait(|waiter| sub.poll_next_group(waiter))
.now_or_never()
.expect("should not block")
.expect("should not error")
.expect("should not be finished")
.sequence
};
assert_eq!(next(&mut sub), 1);
assert_eq!(
sub.take_stale().groups,
1,
"only the conviction the delivery passed counts"
);
sub.update(Subscription::default().with_max_age(Duration::from_secs(60)));
assert_eq!(next(&mut sub), 2);
assert_eq!(next(&mut sub), 3);
assert_eq!(sub.take_stale().groups, 0, "delivered content never counts as stale");
}
#[tokio::test]
async fn a_reversible_floor_does_not_commit_a_conviction() {
let stamp = |sequence: u64| Duration::from_secs([0, 10, 40, 30][sequence as usize]);
let (mut track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
producer.switch(&consumer_b, Position::group(2)).unwrap();
write_group_at(&mut track_a, 0, "a0", stamp(0));
write_group_at(&mut track_a, 1, "a1", stamp(1));
write_group_at(&mut track_b, 2, "b2", stamp(2));
write_group_at(&mut track_b, 3, "b3", stamp(3));
let mut sub = producer.consume().subscribe(None);
let next = |sub: &mut Subscriber| {
kio::wait(|waiter| sub.poll_next_group(waiter))
.now_or_never()
.expect("should not block")
.expect("should not error")
.expect("should not be finished")
.sequence
};
assert_eq!(next(&mut sub), 1);
assert_eq!(sub.take_stale().groups, 1, "the delivery commits past group 0 only");
sub.start_at(10);
assert!(
kio::wait(|waiter| sub.poll_next_group(waiter)).now_or_never().is_none(),
"nothing at or above the raised floor"
);
sub.start_at(0);
sub.update(Subscription::default().with_max_age(Duration::from_secs(60)));
assert_eq!(next(&mut sub), 2);
assert_eq!(next(&mut sub), 3);
assert_eq!(sub.take_stale().groups, 0, "a floor jump commits nothing");
}
#[tokio::test]
async fn a_delivered_continuation_is_not_counted_stale() {
let stamp = |sequence: u64| Duration::from_secs([0, 10, 40, 30][sequence as usize]);
let (mut track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
producer.switch(&consumer_b, Position { group: 1, frame: 1 }).unwrap();
write_group_at(&mut track_a, 0, "a0", stamp(0));
write_group_at(&mut track_a, 1, "a1", stamp(1));
write_group_at(&mut track_b, 1, "b1", stamp(1));
write_group_at(&mut track_b, 2, "b2", stamp(2));
write_group_at(&mut track_b, 3, "b3", stamp(3));
let mut sub = producer.consume().subscribe(None);
let next = |sub: &mut Subscriber| {
kio::wait(|waiter| sub.poll_next_group(waiter))
.now_or_never()
.expect("should not block")
.expect("should not error")
.expect("should not be finished")
.sequence
};
assert_eq!(next(&mut sub), 1);
assert_eq!(next(&mut sub), 3);
assert_eq!(
sub.take_stale().groups,
2,
"groups 0 and 2 are stale; the delivered continuation of 1 is not"
);
}
#[tokio::test]
async fn a_finalized_segment_flushes_nothing_without_a_delivery() {
let stamp = |sequence: u64| Duration::from_secs([0, 40, 20, 30][sequence as usize]);
let retain = track::Info::default().with_max_age(Duration::from_secs(60));
let (mut track_c, consumer_c) = track_pair_with("c", retain.clone());
let (mut track_a, consumer_a) = track_pair_with("a", retain.clone());
let (mut track_b, consumer_b) = track_pair_with("b", retain);
let mut producer = Producer::new();
producer.switch(&consumer_c, None).unwrap();
producer.switch(&consumer_a, Position::group(1)).unwrap();
producer.switch(&consumer_b, Position { group: 1, frame: 1 }).unwrap();
write_group_at(&mut track_c, 0, "c0", stamp(0));
write_group_at(&mut track_a, 1, "a1", stamp(1));
write_group_at(&mut track_b, 1, "b1", stamp(1));
write_group_at(&mut track_b, 2, "b2", stamp(2));
write_group_at(&mut track_b, 3, "b3", stamp(3));
track_b.finish().unwrap();
let mut sub = producer.consume().subscribe(None);
let next = |sub: &mut Subscriber| {
kio::wait(|waiter| sub.poll_next_group(waiter))
.now_or_never()
.expect("should not block")
.expect("should not error")
.expect("should not be finished")
.sequence
};
assert_eq!(next(&mut sub), 0);
sub.start_at(10);
assert!(
kio::wait(|waiter| sub.poll_next_group(waiter)).now_or_never().is_none(),
"nothing at or above the raised floor"
);
sub.start_at(0);
sub.update(Subscription::default().with_max_age(Duration::from_secs(60)));
assert_eq!(next(&mut sub), 1);
assert_eq!(
sub.take_stale().groups,
0,
"no delivery ever committed past a conviction"
);
}
#[tokio::test]
async fn nested_parked_group_is_rechecked_when_the_cap_rises() {
let stamp = |sequence: u64| Duration::from_secs(10 * sequence);
let (mut track_a, consumer_a) = track_pair("a");
let mut inner = Producer::new();
inner.switch(&consumer_a, None).unwrap();
write_group_at(&mut track_a, 0, "a0", stamp(0));
write_group_at(&mut track_a, 1, "a1", stamp(1));
write_group_at(&mut track_a, 2, "a2", stamp(2));
let inner_track =
track::Consumer::spliced("inner".into(), Arc::new(broadcast::Info::default()), inner.consume());
let mut outer = Producer::new();
outer.switch(&inner_track, None).unwrap();
let mut sub = outer.consume().subscribe(None);
sub.end_at(..1);
let recv = |sub: &mut Subscriber| {
kio::wait(|waiter| sub.poll_recv_group(waiter))
.now_or_never()
.map(|group| {
group
.expect("should not error")
.expect("should not be finished")
.sequence
})
};
assert_eq!(recv(&mut sub), Some(0));
assert_eq!(recv(&mut sub), None, "groups beyond the cap park");
sub.end_at(..);
assert_eq!(recv(&mut sub), Some(2), "the re-offered backlog is re-checked");
assert_eq!(sub.take_stale().groups, 1, "the skipped park is counted once");
}
#[tokio::test]
async fn consecutive_updates_wake() {
use std::task::Context;
let (track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
let prefs = sub.prefs();
let (counter, waker) = CountWaker::new();
let mut cx = Context::from_waker(&waker);
let mut fut = std::pin::pin!(sub.recv_group());
assert!(fut.as_mut().poll(&mut cx).is_pending());
*prefs.write().ok().unwrap() = Subscription::default().with_priority(1);
assert_eq!(counter.count(), 1);
assert!(fut.as_mut().poll(&mut cx).is_pending());
assert_eq!(track_a.subscription().unwrap().priority, 1);
let before = counter.count();
*prefs.write().ok().unwrap() = Subscription::default().with_priority(2);
assert!(counter.count() > before, "second update lost its wakeup");
assert!(fut.as_mut().poll(&mut cx).is_pending());
assert_eq!(track_a.subscription().unwrap().priority, 2);
}
#[tokio::test]
async fn takeover_splices_mid_group() {
let (track_a, consumer_a) = track_pair("a");
let (track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
let mut group = track_a.create_group(group::Info { sequence: 0 }).unwrap();
group.write_frame(Timestamp::ZERO, b"a0".to_vec()).unwrap();
group.write_frame(Timestamp::ZERO, b"a1".to_vec()).unwrap();
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(reading.sequence, 0);
assert_eq!(read(&mut reading), b"a0", "first");
assert_eq!(read(&mut reading), b"a1", "second");
assert!(reading.read_frame().now_or_never().is_none(), "group is still open");
producer.takeover(&consumer_b).unwrap();
track_a.abort(Error::Dropped).unwrap();
recv_pending(&mut sub);
let demand = track_b.subscription().unwrap();
assert_eq!(
demand.start,
Some(Position { group: 0, frame: 2 }),
"resumes in the same group, at the frame the old route stopped on"
);
let mut group = track_b.create_group(group::Info { sequence: 0 }).unwrap();
group.start_at(2).unwrap();
group.write_frame(Timestamp::ZERO, b"b2".to_vec()).unwrap();
group.finish().unwrap();
assert_eq!(read(&mut reading), b"b2");
assert!(reading.read_frame().now_or_never().unwrap().unwrap().is_none());
recv_pending(&mut sub);
}
#[tokio::test]
async fn a_replacement_copy_keeps_the_handed_out_group_latency_budget() {
let (track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
let mut head = track_a.create_group(group::Info { sequence: 0 }).unwrap();
head.write_frame(Timestamp::ZERO, b"a0".to_vec()).unwrap();
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().expect("head group");
assert_eq!(read(&mut reading), b"a0");
producer.takeover(&consumer_b).unwrap();
track_a.abort(Error::Dropped).unwrap();
recv_pending(&mut sub);
let mut continuation = track_b.create_group(group::Info { sequence: 0 }).unwrap();
continuation.start_at(1).unwrap();
continuation.write_frame(Timestamp::ZERO, b"b1".to_vec()).unwrap();
assert_eq!(read(&mut reading), b"b1");
assert!(
reading.read_frame().now_or_never().is_none(),
"the replacement group is still the live edge"
);
crate::model::clock::advance(std::time::Duration::from_secs(1));
write_group_at(&mut track_b, 1, "edge", std::time::Duration::from_secs(1));
let result = reading
.read_frame()
.now_or_never()
.expect("the newer group changes the verdict");
assert!(matches!(result, Ok(None)), "the replacement ends: {result:?}");
}
#[tokio::test]
async fn takeover_splices_a_replacement_that_resends_the_head() {
let (track_a, consumer_a) = track_pair("a");
let (track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
let mut group = track_a.create_group(group::Info { sequence: 0 }).unwrap();
group.write_frame(Timestamp::ZERO, b"a0".to_vec()).unwrap();
group.write_frame(Timestamp::ZERO, b"a1".to_vec()).unwrap();
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(read(&mut reading), b"a0");
assert_eq!(read(&mut reading), b"a1");
producer.takeover(&consumer_b).unwrap();
track_a.abort(Error::Dropped).unwrap();
recv_pending(&mut sub);
let mut group = track_b.create_group(group::Info { sequence: 0 }).unwrap();
group.write_frame(Timestamp::ZERO, b"dup0".to_vec()).unwrap();
group.write_frame(Timestamp::ZERO, b"dup1".to_vec()).unwrap();
group.write_frame(Timestamp::ZERO, b"b2".to_vec()).unwrap();
group.finish().unwrap();
assert_eq!(read(&mut reading), b"b2");
assert!(reading.read_frame().now_or_never().unwrap().unwrap().is_none());
recv_pending(&mut sub);
}
#[tokio::test]
async fn takeover_redelivers_an_incomplete_frame() {
let (track_a, consumer_a) = track_pair("a");
let (track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
let mut group = track_a.create_group(group::Info { sequence: 0 }).unwrap();
group.write_frame(Timestamp::ZERO, b"a0".to_vec()).unwrap();
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(read(&mut reading), b"a0");
{
let mut frame = group
.create_frame(frame::Info {
size: 6,
timestamp: Timestamp::ZERO,
})
.unwrap();
frame.write(b"foo".to_vec()).unwrap();
}
track_a.abort(Error::Dropped).unwrap();
producer.takeover(&consumer_b).unwrap();
recv_pending(&mut sub);
let demand = track_b.subscription().unwrap();
assert_eq!(
demand.start,
Some(Position { group: 0, frame: 1 }),
"the half-written frame must be redelivered, not skipped"
);
let mut group = track_b.create_group(group::Info { sequence: 0 }).unwrap();
group.start_at(1).unwrap();
group.write_frame(Timestamp::ZERO, b"b1".to_vec()).unwrap();
assert_eq!(read(&mut reading), b"b1");
}
#[tokio::test]
async fn mid_group_boundary_caps_the_old_route() {
let (track_a, consumer_a) = track_pair("a");
let (track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
let mut group_a = track_a.create_group(group::Info { sequence: 0 }).unwrap();
group_a.write_frame(Timestamp::ZERO, b"a0".to_vec()).unwrap();
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(read(&mut reading), b"a0");
producer.takeover(&consumer_b).unwrap();
recv_pending(&mut sub);
let demand = track_a.subscription().unwrap();
assert_eq!(demand.end, Some(Position { group: 0, frame: 1 }));
group_a.write_frame(Timestamp::ZERO, b"a1-over-cap".to_vec()).unwrap();
let mut group_b = track_b.create_group(group::Info { sequence: 0 }).unwrap();
group_b.start_at(1).unwrap();
group_b.write_frame(Timestamp::ZERO, b"b1".to_vec()).unwrap();
assert_eq!(read(&mut reading), b"b1");
}
#[tokio::test]
async fn takeover_rolls_past_a_finished_group() {
let (mut track_a, consumer_a) = track_pair("a");
let (track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
write_group(&mut track_a, 0, "a0");
assert_eq!(recv(&mut sub), 0);
producer.takeover(&consumer_b).unwrap();
recv_pending(&mut sub);
assert_eq!(producer.resume_position(), Some(Position::group(1)));
assert_eq!(track_b.subscription().unwrap().start, Some(Position::group(1)));
}
#[tokio::test]
async fn dead_copy_stalls_until_the_continuation() {
let (track_a, consumer_a) = track_pair("a");
let (track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
let mut group = track_a.create_group(group::Info { sequence: 0 }).unwrap();
group.write_frame(Timestamp::ZERO, b"a0".to_vec()).unwrap();
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(read(&mut reading), b"a0");
group.abort(Error::Dropped).unwrap();
assert!(reading.read_frame().now_or_never().is_none(), "should stall, not error");
producer.takeover(&consumer_b).unwrap();
let mut group = track_b.create_group(group::Info { sequence: 0 }).unwrap();
group.start_at(1).unwrap();
group.write_frame(Timestamp::ZERO, b"b1".to_vec()).unwrap();
assert_eq!(read(&mut reading), b"b1");
}
#[tokio::test]
async fn dead_copy_ends_once_the_track_aborts() {
let (track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
let mut group = track_a.create_group(group::Info { sequence: 0 }).unwrap();
group.write_frame(Timestamp::ZERO, b"a0".to_vec()).unwrap();
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(read(&mut reading), b"a0");
group.abort(Error::Dropped).unwrap();
assert!(
reading.read_frame().now_or_never().is_none(),
"a switch could still come"
);
producer.abort(Error::Cancel).unwrap();
assert!(
matches!(reading.read_frame().now_or_never(), Some(Err(_))),
"an aborted track must not leave the reader parked"
);
}
#[tokio::test]
async fn lost_group_stays_lost() {
let (mut track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(replay());
let mut group = track_a.create_group(group::Info { sequence: 0 }).unwrap();
group.write_frame(Timestamp::ZERO, b"a0".to_vec()).unwrap();
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(read(&mut reading), b"a0");
group.abort(Error::Stream(crate::StreamError::Old)).unwrap();
write_group(&mut track_a, 1, "a1");
assert!(matches!(
reading.read_frame().now_or_never(),
Some(Err(Error::Stream(crate::StreamError::Old)))
));
assert!(
matches!(
reading.finished().now_or_never(),
Some(Err(Error::Stream(crate::StreamError::Old)))
),
"a lost group must not park or change its answer"
);
assert!(matches!(
reading.read_frame().now_or_never(),
Some(Err(Error::Stream(crate::StreamError::Old)))
));
}
#[tokio::test]
async fn copy_missing_the_head_is_lost() {
let (track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
let mut group = track_a.create_group(group::Info { sequence: 0 }).unwrap();
group.start_at(2).unwrap();
group.write_frame(Timestamp::ZERO, b"a2".to_vec()).unwrap();
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(reading.sequence, 0);
assert_eq!(reading.index(), 0, "the reader still wants frame 0");
assert!(
matches!(reading.read_frame().now_or_never(), Some(Err(Error::Lagged))),
"must report the loss rather than serve the tail as the head"
);
}
#[tokio::test]
async fn live_route_ahead_of_the_seam_still_parks() {
let (track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(replay());
let mut group = track_a.create_group(group::Info { sequence: 0 }).unwrap();
group.write_frame(Timestamp::ZERO, b"a0".to_vec()).unwrap();
group.write_frame(Timestamp::ZERO, b"a1".to_vec()).unwrap();
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(read(&mut reading), b"a0");
assert_eq!(read(&mut reading), b"a1");
producer.takeover(&consumer_b).unwrap();
write_group(&mut track_b, 1, "b1");
assert_eq!(recv(&mut sub), 1);
assert!(
reading.read_frame().now_or_never().is_none(),
"a live route may still fill the seam out of order"
);
let mut group = track_b.create_group(group::Info { sequence: 0 }).unwrap();
group.start_at(2).unwrap();
group.write_frame(Timestamp::ZERO, b"b2".to_vec()).unwrap();
assert_eq!(read(&mut reading), b"b2");
}
#[tokio::test]
async fn takeover_after_empty_segment_keeps_live_edge() {
let (track_a, consumer_a) = track_pair("a");
let (track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
recv_pending(&mut sub);
assert_eq!(track_a.subscription().unwrap().start, None);
drop(track_a);
producer.takeover(&consumer_b).unwrap();
recv_pending(&mut sub);
assert_eq!(track_b.subscription().unwrap().start, None);
}
#[tokio::test]
async fn fetch_fails_over_to_a_newer_segment() {
let (track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let consumer = producer.consume();
track_a.abort(Error::Dropped).unwrap();
let fetch = consumer.fetch_group(0, None);
let mut fetch = std::pin::pin!(fetch);
assert!(futures::poll!(fetch.as_mut()).is_pending(), "a dead copy should park");
write_group(&mut track_b, 0, "b0");
producer.takeover(&consumer_b).unwrap();
let group = fetch.await.expect("fetch should fail over");
assert_eq!(group.sequence, 0);
}
#[tokio::test]
async fn fetch_aborts_with_the_track() {
let (track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let consumer = producer.consume();
track_a.abort(Error::Dropped).unwrap();
let fetch = consumer.fetch_group(0, None);
let mut fetch = std::pin::pin!(fetch);
assert!(futures::poll!(fetch.as_mut()).is_pending(), "a dead copy should park");
producer.abort(Error::Cancel).unwrap();
assert!(matches!(fetch.await, Err(Error::Cancel)));
}
#[tokio::test]
async fn fetch_pending_ends_when_the_track_aborts() {
let (track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let consumer = producer.consume();
let handler = track_a.dynamic();
let fetch = consumer.fetch_group(0, None);
let mut fetch = std::pin::pin!(fetch);
assert!(
futures::poll!(fetch.as_mut()).is_pending(),
"unanswered fetch should park"
);
producer.abort(Error::Cancel).unwrap();
assert!(matches!(fetch.await, Err(Error::Cancel)));
drop(handler);
}
#[tokio::test]
async fn fetch_error_from_a_live_copy_is_authoritative() {
let (_track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let consumer = producer.consume();
let result = consumer
.fetch_group(0, None)
.now_or_never()
.expect("a live copy's answer must resolve immediately");
assert!(matches!(result, Err(Error::NotFound)));
}
#[tokio::test]
async fn takeover_after_produced_segment_resumes_at_the_boundary() {
let (mut track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(replay());
recv_pending(&mut sub);
assert_eq!(track_a.subscription().unwrap().start, None);
write_group(&mut track_a, 0, "a0");
write_group(&mut track_a, 1, "a1");
assert_eq!(recv(&mut sub), 0);
assert_eq!(recv(&mut sub), 1);
drop(track_a);
producer.takeover(&consumer_b).unwrap();
recv_pending(&mut sub);
assert_eq!(track_b.subscription().unwrap().start, Some(Position::group(2)));
write_group(&mut track_b, 1, "b1");
recv_pending(&mut sub);
write_group(&mut track_b, 5, "b5");
assert_eq!(recv(&mut sub), 5);
}
#[tokio::test]
async fn takeover_keeps_an_explicit_start() {
let (mut track_a, consumer_a) = track_pair("a");
let (track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(
Subscription::default()
.with_start(Position::group(0))
.with_max_age(std::time::Duration::from_secs(60)),
);
recv_pending(&mut sub);
assert_eq!(track_a.subscription().unwrap().start, Some(Position::group(0)));
write_group(&mut track_a, 0, "a0");
write_group(&mut track_a, 1, "a1");
assert_eq!(recv(&mut sub), 0);
assert_eq!(recv(&mut sub), 1);
drop(track_a);
producer.takeover(&consumer_b).unwrap();
recv_pending(&mut sub);
assert_eq!(track_b.subscription().unwrap().start, Some(Position::group(2)));
}
#[tokio::test]
async fn switch_validates_boundaries() {
let (mut track_a, consumer_a) = track_pair("a");
let (_track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
assert!(producer.switch(&consumer_b, None).is_err());
write_group(&mut track_a, 0, "a0");
assert!(producer.switch(&consumer_b, Position::group(0)).is_err());
producer.switch(&consumer_b, Position::group(1)).unwrap();
}
#[tokio::test]
async fn prune_bounds_segments_and_keeps_the_boundary() {
let mut producer = Producer::new();
let mut sub = producer.consume().subscribe(None);
let rounds = 2 * MAX_SEGMENTS as u64;
for sequence in 0..rounds {
let (mut track, consumer) = track_pair("t");
producer.takeover(&consumer).unwrap();
write_group(&mut track, sequence, "payload");
assert_eq!(recv(&mut sub), sequence);
track.abort(Error::Dropped).unwrap();
}
assert_eq!(
producer.state.read().segments.len(),
MAX_SEGMENTS,
"dead predecessors should have been pruned"
);
let (mut track, consumer) = track_pair("final");
producer.takeover(&consumer).unwrap();
write_group(&mut track, 0, "below-the-floor");
recv_pending(&mut sub);
write_group(&mut track, rounds, "resumed");
assert_eq!(recv(&mut sub), rounds);
}
#[tokio::test]
async fn prune_retires_live_predecessors() {
let mut producer = Producer::new();
let mut sub = producer.consume().subscribe(None);
let mut tracks = Vec::new();
let rounds = 2 * MAX_SEGMENTS as u64;
for sequence in 0..rounds {
let (mut track, consumer) = track_pair("t");
producer.takeover(&consumer).unwrap();
write_group(&mut track, sequence, "payload");
assert_eq!(recv(&mut sub), sequence);
tracks.push(track);
}
assert_eq!(
producer.state.read().segments.len(),
MAX_SEGMENTS,
"live predecessors should still be pruned"
);
}
#[tokio::test]
async fn declared_start_fails_over_a_skipped_group() {
let (track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
let mut group = track_a.create_group(group::Info { sequence: 0 }).unwrap();
group.write_frame(Timestamp::ZERO, b"a0".to_vec()).unwrap();
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(read(&mut reading), b"a0");
assert!(
reading.read_frame().now_or_never().is_none(),
"the seam parks for a continuation"
);
drop(group);
track_a.abort(Error::Dropped).unwrap();
producer.takeover(&consumer_b).unwrap();
track_b.start_at(1).unwrap();
write_group(&mut track_b, 1, "b1");
assert_eq!(recv(&mut sub), 1);
assert!(
reading
.read_frame()
.now_or_never()
.expect("the skipped seam must resolve")
.is_err(),
"the skipped frames are a loss, not a clean end"
);
}
#[tokio::test]
async fn latched_reader_follows_a_moved_boundary() {
let (track_a, consumer_a) = track_pair("a");
let (track_b, consumer_b) = track_pair("b");
let (track_c, consumer_c) = track_pair("c");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
let mut group = track_a.create_group(group::Info { sequence: 0 }).unwrap();
group.write_frame(Timestamp::ZERO, b"a0".to_vec()).unwrap();
group.write_frame(Timestamp::ZERO, b"a1".to_vec()).unwrap();
producer.switch(&consumer_b, Position { group: 0, frame: 2 }).unwrap();
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(read(&mut reading), b"a0");
drop(track_b);
producer.switch(&consumer_c, Position { group: 0, frame: 1 }).unwrap();
let mut group = track_c.create_group(group::Info { sequence: 0 }).unwrap();
group.start_at(1).unwrap();
group.write_frame(Timestamp::ZERO, b"c1".to_vec()).unwrap();
assert_eq!(read(&mut reading), b"c1");
}
#[tokio::test]
async fn pruned_cursors_stay_bounded() {
let mut producer = Producer::new();
let mut sub = producer.consume().subscribe(None);
let mut tracks = Vec::new();
let rounds = 3 * MAX_SEGMENTS as u64;
for sequence in 0..rounds {
let (mut track, consumer) = track_pair("t");
producer.takeover(&consumer).unwrap();
write_group(&mut track, sequence, "payload");
assert_eq!(recv(&mut sub), sequence);
tracks.push(track);
}
recv_pending(&mut sub);
recv_pending(&mut sub);
assert_eq!(
sub.segments.len(),
2 * MAX_SEGMENTS,
"pruned cursors beyond the bound must be cut"
);
assert!(tracks[0].subscription().is_none(), "a cut cursor releases its demand");
assert!(
tracks[rounds as usize - MAX_SEGMENTS - 1].subscription().is_some(),
"a pruned cursor within the bound keeps draining"
);
}
#[tokio::test]
async fn capped_subscriber_bounds_parked_segments() {
let mut producer = Producer::new();
let mut sub = producer.consume().subscribe(replay());
sub.end_at(..1);
let mut tracks = Vec::new();
let rounds = 3 * MAX_SEGMENTS as u64;
for round in 0..rounds {
let (mut track, consumer) = track_pair("t");
producer.takeover(&consumer).unwrap();
write_group(&mut track, round + 1, "beyond-the-cap");
recv_pending(&mut sub);
tracks.push(track);
}
recv_pending(&mut sub);
assert_eq!(
sub.segments.len(),
2 * MAX_SEGMENTS,
"parked entries must stay bounded, not accumulate"
);
assert!(tracks[0].subscription().is_none(), "a cut entry releases its demand");
sub.end_at(..);
for sequence in (rounds - 2 * MAX_SEGMENTS as u64 + 1)..=rounds {
assert_eq!(recv(&mut sub), sequence);
}
recv_pending(&mut sub);
}
#[tokio::test]
async fn buried_route_revives_when_the_copy_lands() {
let (track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
let mut group = track_a.create_group(group::Info { sequence: 0 }).unwrap();
group.write_frame(Timestamp::ZERO, b"a0".to_vec()).unwrap();
producer.takeover(&consumer_b).unwrap();
track_b.start_at(1).unwrap();
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(read(&mut reading), b"a0");
assert!(reading.read_frame().now_or_never().is_none(), "the seam parks");
track_b.start_at(None).unwrap();
let mut cont = track_b.create_group(group::Info { sequence: 0 }).unwrap();
cont.start_at(1).unwrap();
cont.write_frame(Timestamp::ZERO, b"b1".to_vec()).unwrap();
assert_eq!(read(&mut reading), b"b1");
}
#[tokio::test]
async fn misaligned_copy_is_lost_without_spinning() {
let (track, consumer) = track_pair("t");
let mut producer = Producer::new();
producer.takeover(&consumer).unwrap();
let mut sub = producer.consume().subscribe(None);
let mut group = track.create_group(group::Info { sequence: 0 }).unwrap();
group.start_at(5).unwrap();
group.write_frame(Timestamp::ZERO, b"f5".to_vec()).unwrap();
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(reading.index(), 0);
let err = reading
.read_frame()
.now_or_never()
.expect("the loss must resolve rather than park")
.expect_err("a misaligned copy is a loss, never misnumbered frames");
assert!(matches!(err, Error::Lagged), "expected a lagged copy, got {err:?}");
}
#[tokio::test]
async fn peeks_walk_past_an_empty_segment() {
let (mut track_a, consumer_a) = track_pair("a");
let (_track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
write_group(&mut track_a, 0, "a0");
write_group(&mut track_a, 1, "a1");
producer.switch(&consumer_b, Position { group: 2, frame: 0 }).unwrap();
let consumer = producer.consume();
assert_eq!(consumer.peek_latest().map(|group| group.sequence), Some(1));
assert_eq!(consumer.peek_before(1).map(|group| group.sequence), Some(0));
assert_eq!(consumer.peek_before(0).map(|group| group.sequence), None);
}
#[tokio::test]
async fn peeks_respect_segment_bounds() {
let (mut track_a, consumer_a) = track_pair("a");
let (_track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
write_group(&mut track_a, 0, "a0");
producer.switch(&consumer_b, Position { group: 1, frame: 0 }).unwrap();
write_group(&mut track_a, 5, "a5");
let consumer = producer.consume();
assert_eq!(
consumer.peek_latest().map(|group| group.sequence),
Some(0),
"group 5 is above A's cap and belongs to B's range"
);
}
#[tokio::test]
async fn seek_keeps_a_pruned_latch() {
let (track, consumer) = track_pair("t");
let mut producer = Producer::new();
producer.takeover(&consumer).unwrap();
let mut sub = producer.consume().subscribe(None);
let mut group = track.create_group(group::Info { sequence: 0 }).unwrap();
for payload in [b"f0", b"f1", b"f2"] {
group.write_frame(Timestamp::ZERO, payload.to_vec()).unwrap();
}
group.finish().unwrap();
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(read(&mut reading), b"f0");
track.abort(Error::Dropped).unwrap();
for sequence in 1..=MAX_SEGMENTS as u64 {
let (mut track, consumer) = track_pair("t");
producer.takeover(&consumer).unwrap();
write_group(&mut track, sequence, "payload");
assert_eq!(recv(&mut sub), sequence);
track.abort(Error::Dropped).unwrap();
}
assert!(
producer.state.read().pruned.is_some(),
"the first segment should be pruned"
);
reading.start_at(2);
assert_eq!(read(&mut reading), b"f2");
}
#[tokio::test]
async fn datagram_poller_bounds_pruned_cursors() {
let mut producer = Producer::new();
let mut sub = producer.consume().subscribe(None);
for sequence in 0..(3 * MAX_SEGMENTS as u64) {
let (mut track, consumer) = track_pair("t");
producer.takeover(&consumer).unwrap();
write_group(&mut track, sequence, "payload");
assert!(
kio::wait(|waiter| sub.poll_recv_datagram(waiter))
.now_or_never()
.is_none(),
"no datagram expected"
);
}
assert_eq!(
sub.segments.len(),
2 * MAX_SEGMENTS,
"the reap must run on the datagram path too"
);
}
#[tokio::test]
async fn late_group_drains_from_a_pruned_cursor() {
let mut producer = Producer::new();
let mut sub = producer.consume().subscribe(replay());
let (mut track_a, consumer_a) = track_pair("a");
producer.takeover(&consumer_a).unwrap();
write_group(&mut track_a, 1, "a1");
assert_eq!(recv(&mut sub), 1);
for sequence in 2..=(1 + MAX_SEGMENTS as u64) {
let (mut track, consumer) = track_pair("t");
producer.takeover(&consumer).unwrap();
write_group(&mut track, sequence, "payload");
assert_eq!(recv(&mut sub), sequence);
}
assert!(
producer.state.read().pruned.is_some(),
"the first segment should be pruned"
);
write_group(&mut track_a, 0, "a0");
assert_eq!(recv(&mut sub), 0);
}
#[tokio::test]
async fn finished_resolves_for_a_pruned_bounded_group() {
let (track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
let mut group = track_a.create_group(group::Info { sequence: 0 }).unwrap();
group.write_frame(Timestamp::ZERO, b"f0".to_vec()).unwrap();
group.write_frame(Timestamp::ZERO, b"f1".to_vec()).unwrap();
let (mut track_b, consumer_b) = track_pair("b");
producer.takeover(&consumer_b).unwrap();
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
drop(group);
track_a.abort(Error::Dropped).unwrap();
write_group(&mut track_b, 1, "b1");
assert_eq!(recv(&mut sub), 1);
assert!(
reading.finished().now_or_never().is_none(),
"the seam is still coverable"
);
for sequence in 2..=(1 + MAX_SEGMENTS as u64) {
let (mut track, consumer) = track_pair("t");
producer.takeover(&consumer).unwrap();
write_group(&mut track, sequence, "payload");
assert_eq!(recv(&mut sub), sequence);
}
assert_eq!(
reading
.finished()
.now_or_never()
.expect("the lost seam must resolve the count")
.unwrap(),
2
);
}
#[tokio::test]
async fn group_reader_gives_up_below_the_pruned_floor() {
let mut producer = Producer::new();
let mut sub = producer.consume().subscribe(None);
let (mut track, consumer) = track_pair("t0");
producer.takeover(&consumer).unwrap();
write_group(&mut track, 0, "kept");
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(reading.sequence, 0);
let mut cloned = reading.clone();
track.abort(Error::Dropped).unwrap();
for sequence in 1..=MAX_SEGMENTS as u64 {
let (mut track, consumer) = track_pair("t");
producer.takeover(&consumer).unwrap();
write_group(&mut track, sequence, "payload");
assert_eq!(recv(&mut sub), sequence);
track.abort(Error::Dropped).unwrap();
}
assert!(
producer.state.read().pruned.is_some(),
"the first segment should be pruned"
);
for reader in [&mut reading, &mut cloned] {
assert_eq!(read(reader), b"kept");
assert!(
reader.read_frame().now_or_never().unwrap().unwrap().is_none(),
"the group ends at what the route produced"
);
}
}
#[tokio::test]
async fn finished_resolves_when_the_successor_skips_the_seam() {
let (track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
let mut sub = producer.consume().subscribe(None);
let mut group = track_a.create_group(group::Info { sequence: 0 }).unwrap();
group.write_frame(Timestamp::ZERO, b"f0".to_vec()).unwrap();
group.write_frame(Timestamp::ZERO, b"f1".to_vec()).unwrap();
producer.takeover(&consumer_b).unwrap();
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
drop(group);
track_a.abort(Error::Dropped).unwrap();
assert!(reading.finished().now_or_never().is_none(), "the seam is coverable");
track_b.start_at(1).unwrap();
write_group(&mut track_b, 1, "b1");
assert_eq!(recv(&mut sub), 1);
assert_eq!(
reading
.finished()
.now_or_never()
.expect("a skip-declared seam must resolve the count")
.unwrap(),
2
);
}
#[tokio::test]
async fn reader_drains_a_pruned_segments_copy() {
let mut producer = Producer::new();
let mut sub = producer.consume().subscribe(None);
let (track, consumer) = track_pair("t0");
producer.takeover(&consumer).unwrap();
let mut group = track.create_group(group::Info { sequence: 0 }).unwrap();
group.write_frame(Timestamp::ZERO, b"f0".to_vec()).unwrap();
group.write_frame(Timestamp::ZERO, b"f1".to_vec()).unwrap();
group.finish().unwrap();
let mut reading = sub.recv_group().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(read(&mut reading), b"f0");
track.abort(Error::Dropped).unwrap();
for sequence in 1..=MAX_SEGMENTS as u64 {
let (mut track, consumer) = track_pair("t");
producer.takeover(&consumer).unwrap();
write_group(&mut track, sequence, "payload");
assert_eq!(recv(&mut sub), sequence);
track.abort(Error::Dropped).unwrap();
}
assert!(
producer.state.read().pruned.is_some(),
"the first segment should be pruned"
);
assert_eq!(read(&mut reading), b"f1");
}
#[tokio::test]
async fn switch_replaces_a_run_of_empty_segments() {
let (mut track_a, consumer_a) = track_pair("a");
let (_track_b, consumer_b) = track_pair("b");
let (_track_c, consumer_c) = track_pair("c");
let (mut track_d, consumer_d) = track_pair("d");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
write_group(&mut track_a, 0, "a0");
assert_eq!(recv(&mut sub), 0);
producer.switch(&consumer_b, Position::group(1)).unwrap();
producer.switch(&consumer_c, Position::group(2)).unwrap();
producer.switch(&consumer_d, Position::group(1)).unwrap();
write_group(&mut track_d, 1, "d1");
assert_eq!(recv(&mut sub), 1);
assert_eq!(producer.state.read().segments.len(), 2);
}
#[tokio::test]
async fn abort_drains_before_erroring() {
let (mut track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
write_group(&mut track_a, 0, "a0");
producer.abort(Error::Cancel).unwrap();
assert!(matches!(sub.finished().now_or_never().unwrap(), Err(Error::Cancel)));
assert_eq!(recv(&mut sub), 0);
assert!(matches!(sub.recv_group().now_or_never().unwrap(), Err(Error::Cancel)));
}
#[tokio::test]
async fn terminal_states_are_exclusive() {
let (_track_a, consumer_a) = track_pair("a");
let (_track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
producer.finish().unwrap();
assert!(matches!(producer.abort(Error::Cancel), Err(Error::Closed)));
assert!(matches!(producer.finish(), Err(Error::Closed)));
assert!(matches!(
producer.switch(&consumer_b, Position::group(1)),
Err(Error::Closed)
));
assert!(matches!(producer.takeover(&consumer_b), Err(Error::Closed)));
assert!(matches!(producer.release(), Err(Error::Closed)));
let mut producer = Producer::new();
producer.abort(Error::Cancel).unwrap();
assert!(matches!(producer.finish(), Err(Error::Closed)));
assert!(matches!(producer.abort(Error::Cancel), Err(Error::Closed)));
assert!(matches!(producer.takeover(&consumer_b), Err(Error::Closed)));
}
#[tokio::test]
async fn dropped_producer_errors_once_drained() {
let (mut track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
write_group(&mut track_a, 0, "a0");
assert_eq!(recv(&mut sub), 0);
track_a.abort(Error::Timeout).unwrap();
drop(producer);
let result = sub.recv_group().now_or_never().expect("must not stall forever");
assert!(matches!(result, Err(Error::Timeout)));
}
#[tokio::test]
async fn finished_producer_ends_with_a_dead_final_segment() {
let (mut track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
write_group(&mut track_a, 0, "a0");
assert_eq!(recv(&mut sub), 0);
producer.finish().unwrap();
track_a.abort(Error::Timeout).unwrap();
let result = sub.recv_group().now_or_never().expect("must not stall forever");
assert!(matches!(result, Err(Error::Timeout)));
}
#[tokio::test]
async fn finished_producer_ends_datagrams_with_a_dead_final_segment() {
let (track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
assert!(
kio::wait(|waiter| sub.poll_recv_datagram(waiter))
.now_or_never()
.is_none(),
"no datagram yet"
);
producer.finish().unwrap();
track_a.abort(Error::Timeout).unwrap();
let result = kio::wait(|waiter| sub.poll_recv_datagram(waiter))
.now_or_never()
.expect("must not stall forever");
assert!(matches!(result, Err(Error::Timeout)));
}
#[tokio::test]
async fn dropped_producer_keeps_a_live_segment_serving() {
let (mut track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
drop(producer);
write_group(&mut track_a, 0, "a0");
assert_eq!(recv(&mut sub), 0);
recv_pending(&mut sub);
track_a.finish().unwrap();
let result = sub.recv_group().now_or_never().expect("must not stall forever");
assert!(matches!(result, Ok(None)));
}
#[tokio::test]
async fn dropped_producer_ends_finished_waiters_with_the_segment() {
for clean in [true, false] {
let (track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
recv_pending(&mut sub);
drop(producer);
assert!(
sub.finished().now_or_never().is_none(),
"ended while the segment is live"
);
match clean {
true => track_a.finish().unwrap(),
false => track_a.abort(Error::Timeout).unwrap(),
}
let result = sub.finished().now_or_never().expect("must not stall forever");
match clean {
true => assert!(matches!(result, Ok(0))),
false => assert!(matches!(result, Err(Error::Timeout))),
}
}
}
#[tokio::test]
async fn end_waiters_leave_queued_groups_readable() {
for dropped in [false, true] {
let (mut track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(replay());
write_group(&mut track_a, 0, "a0");
track_a.finish().unwrap();
match dropped {
true => drop(producer),
false => producer.finish().unwrap(),
}
let count = sub.finished().now_or_never().expect("the end is known");
assert!(matches!(count, Ok(1)), "{count:?}");
let datagram = kio::wait(|waiter| sub.poll_recv_datagram(waiter)).now_or_never();
assert!(matches!(datagram, Some(Ok(None))), "{datagram:?}");
assert_eq!(recv(&mut sub), 0, "dropped={dropped}");
let end = sub.recv_group().now_or_never().expect("must not stall forever");
assert!(matches!(end, Ok(None)), "dropped={dropped}");
}
}
#[tokio::test]
async fn dropped_producer_ends_datagrams() {
let (track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
drop(track_a);
drop(producer);
let result = kio::wait(|waiter| sub.poll_recv_datagram(waiter))
.now_or_never()
.expect("must not stall forever");
assert!(matches!(result, Err(Error::Dropped)));
}
#[tokio::test]
async fn evicted_parked_group_wakes_the_clean_end() {
use std::task::Context;
let (mut track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
sub.end_at(..1);
write_group(&mut track_a, 0, "a0");
let straggler = track_a.create_group(group::Info { sequence: 1 }).unwrap();
assert_eq!(recv(&mut sub), 0);
recv_pending(&mut sub);
track_a.finish().unwrap();
producer.finish().unwrap();
let (counter, waker) = CountWaker::new();
let mut cx = Context::from_waker(&waker);
let mut fut = std::pin::pin!(sub.recv_group());
assert!(
fut.as_mut().poll(&mut cx).is_pending(),
"the parked group holds the end open"
);
straggler.abort(Error::Old).unwrap();
assert!(counter.count() > 0, "the eviction wakeup was lost");
let result = fut.as_mut().poll(&mut cx);
assert!(matches!(result, Poll::Ready(Ok(None))));
}
#[tokio::test]
async fn straggler_parked_after_finish_still_wakes() {
use std::task::Context;
let (mut track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
sub.end_at(..1);
write_group(&mut track_a, 0, "a0");
assert_eq!(recv(&mut sub), 0);
let straggler = track_a.create_group(group::Info { sequence: 1 }).unwrap();
track_a.finish().unwrap();
producer.finish().unwrap();
let (counter, waker) = CountWaker::new();
let mut cx = Context::from_waker(&waker);
let mut fut = std::pin::pin!(sub.recv_group());
assert!(
fut.as_mut().poll(&mut cx).is_pending(),
"the parked group holds the end open"
);
straggler.abort(Error::Old).unwrap();
assert!(counter.count() > 0, "the abort wakeup was lost");
assert!(matches!(fut.as_mut().poll(&mut cx), Poll::Ready(Ok(None))));
}
#[tokio::test]
async fn release_restarts_unbounded() {
let (mut track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.takeover(&consumer_a).unwrap();
write_group(&mut track_a, 7, "a7");
assert!(producer.is_spliced());
producer.release().unwrap();
assert!(!producer.is_spliced());
assert_eq!(producer.resume_position(), None);
producer.takeover(&consumer_b).unwrap();
let mut sub = producer.consume().subscribe(None);
recv_pending(&mut sub);
assert_eq!(track_b.subscription().unwrap().start, None);
write_group(&mut track_b, 2, "b2");
assert_eq!(recv(&mut sub), 2);
}
#[tokio::test]
async fn release_resets_the_pruned_floor() {
let mut producer = Producer::new();
let count = 2 * MAX_SEGMENTS as u64;
let mut tracks = Vec::new();
for sequence in 0..count {
let (mut track, consumer) = track_pair("t");
producer.takeover(&consumer).unwrap();
write_group(&mut track, sequence, "g");
tracks.push(track);
}
producer.release().unwrap();
let (mut track, consumer) = track_pair("fresh");
producer.takeover(&consumer).unwrap();
let mut sub = producer.consume().subscribe(None);
write_group(&mut track, 0, "g0");
assert_eq!(recv(&mut sub), 0);
}
#[tokio::test]
async fn fetch_fails_when_the_producer_dies_segmentless() {
let producer = Producer::new();
let consumer = producer.consume();
let fetch = consumer.fetch_group(0, None);
let mut fetch = std::pin::pin!(fetch);
assert!(futures::poll!(fetch.as_mut()).is_pending(), "fetch should wait");
drop(producer);
assert!(matches!(fetch.await, Err(Error::NotFound)));
}
#[tokio::test]
async fn info_fails_when_the_producer_dies_segmentless() {
let producer = Producer::new();
let consumer = producer.consume();
let info = consumer.info();
let mut info = std::pin::pin!(info);
assert!(futures::poll!(info.as_mut()).is_pending(), "info should wait");
drop(producer);
assert!(matches!(info.await, Err(Error::Dropped)));
}
#[tokio::test]
async fn start_at_drops_parked_groups_below_the_floor() {
let (mut track_a, consumer_a) = track_pair("a");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
sub.end_at(..1);
write_group(&mut track_a, 0, "a0");
write_group(&mut track_a, 1, "a1");
write_group(&mut track_a, 2, "a2");
assert_eq!(recv(&mut sub), 0);
recv_pending(&mut sub);
sub.start_at(2);
sub.end_at(..);
assert_eq!(recv(&mut sub), 2, "group 1 was overtaken by start_at");
recv_pending(&mut sub);
}
#[tokio::test]
async fn demand_intersects_subscriber_end_with_boundary() {
let (track_a, consumer_a) = track_pair("a");
let (track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(
Subscription::default()
.with_start(Position::group(0))
.with_end(Position::after_group(3)),
);
recv_pending(&mut sub);
assert_eq!(track_a.subscription().unwrap().end, Position::after_group(3));
producer.switch(&consumer_b, Position::group(2)).unwrap();
recv_pending(&mut sub);
assert_eq!(track_a.subscription().unwrap().end, Some(Position::group(2)));
assert_eq!(track_b.subscription().unwrap().end, Position::after_group(3));
}
#[tokio::test]
async fn subscribers_read_independently() {
let (mut track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let consumer = producer.consume();
let mut sub1 = consumer.subscribe(replay());
let mut sub2 = consumer.subscribe(replay());
recv_pending(&mut sub1);
recv_pending(&mut sub2);
write_group(&mut track_a, 0, "a0");
producer.switch(&consumer_b, Position::group(1)).unwrap();
write_group(&mut track_b, 1, "b1");
assert_eq!(recv(&mut sub1), 0);
assert_eq!(recv(&mut sub1), 1);
assert_eq!(recv(&mut sub2), 0);
assert_eq!(recv(&mut sub2), 1);
assert!(sub1.is_clone(&sub2));
}
#[tokio::test]
async fn latest_clamps_to_segment_bounds() {
let (mut track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let consumer = producer.consume();
assert_eq!(consumer.latest(), None);
write_group(&mut track_a, 0, "a0");
assert_eq!(consumer.latest(), Some(0));
producer.switch(&consumer_b, Position::group(2)).unwrap();
write_group(&mut track_a, 5, "a5");
assert_eq!(consumer.latest(), Some(1));
write_group(&mut track_b, 0, "b0");
assert_eq!(consumer.latest(), Some(1));
write_group(&mut track_b, 3, "b3");
assert_eq!(consumer.latest(), Some(3));
}
#[tokio::test]
async fn datagrams_come_from_the_newest_segment() {
let (mut track_a, consumer_a) = track_pair("a");
let (mut track_b, consumer_b) = track_pair("b");
let mut producer = Producer::new();
producer.switch(&consumer_a, None).unwrap();
let mut sub = producer.consume().subscribe(None);
assert!(
kio::wait(|waiter| sub.poll_recv_datagram(waiter))
.now_or_never()
.is_none(),
"no datagram yet"
);
producer.switch(&consumer_b, Position::group(1)).unwrap();
track_a.append_datagram(Timestamp::ZERO, b"old".as_ref()).unwrap();
assert!(
kio::wait(|waiter| sub.poll_recv_datagram(waiter))
.now_or_never()
.is_none(),
"stale datagram must not surface"
);
track_b.append_datagram(Timestamp::ZERO, b"new".as_ref()).unwrap();
let datagram = kio::wait(|waiter| sub.poll_recv_datagram(waiter))
.now_or_never()
.expect("datagram should be ready")
.expect("should not error")
.expect("track should not be finished");
assert_eq!(&datagram.payload[..], b"new");
}
}