use std::collections::{BTreeMap, HashSet};
use std::ops::Bound;
use std::task::{Poll, ready};
use crate::{Datagram, Error, Result, frame, group, track};
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,
}
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 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()
}
}
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 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,
});
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 start = state.resume_position();
if start.is_none() {
state.segments.clear();
}
state.switch(track, start)
}
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 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 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,
stale_cap: None,
drift_cap: kio::Producer::new(None),
}
}
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 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>,
cap: kio::Consumer<Option<u64>>,
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(),
cap: self.cap.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>,
cap: kio::Consumer<Option<u64>>,
sequence: u64,
index: u64,
) -> Self {
Self {
state,
subscription,
cap,
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.cap.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(&mut self) -> Result<bool> {
match self.dead.take() {
Some((_, err)) => {
tracing::warn!(
group = self.sequence,
frame = self.index,
%err,
"no route can serve the rest of this group"
);
Err(err)
}
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.cap.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>,
sub: SubState,
terminal: Option<track::Subscriber>,
pruned: bool,
parked: BTreeMap<u64, group::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, count: Option<u64>) {
let previous = std::mem::replace(&mut self.sub, SubState::Done(count));
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(Option<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>,
stale_cap: Option<u64>,
drift_cap: kio::Producer<Option<u64>>,
}
impl Subscriber {
fn poll_sync(&mut self, waiter: &kio::Waiter) {
self.sync(waiter);
self.reap();
}
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(None);
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.start, 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 anchor = self.anchor_end();
for segment in segments {
match self.segments.iter_mut().find(|s| s.id == segment.id) {
Some(existing) => {
if existing.end != segment.end {
existing.end = segment.end;
let cap = Self::stale_cap(existing, anchor);
if let Some(sub) = existing.stale_sub_mut() {
sub.set_stale_cap(cap);
let _ = sub.update(slice(&self.last_prefs, segment.start, segment.end));
}
}
}
None => {
let sub = segment
.track
.subscribe(slice(&self.last_prefs, segment.start, segment.end));
self.segments.push(SegmentSub {
id: segment.id,
start: segment.start,
end: segment.end,
sub: SubState::Pending(sub),
terminal: None,
pruned: false,
parked: BTreeMap::new(),
});
}
}
}
}
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_cap.consume(),
sequence,
0,
)
.latched(seg.id, end, seg.last_group(), group.clone());
Some(group.into_spliced(spliced))
}
fn stale_cap(seg: &SegmentSub, anchor_end: Option<u64>) -> Option<u64> {
min_some(anchor_end, seg.last_group())
}
fn anchor_end(&self) -> Option<u64> {
min_some(self.stale_cap, self.end_sequence)
}
pub(crate) fn set_stale_cap(&mut self, cap: Option<u64>) {
self.stale_cap = cap;
self.update_drift_cap();
let anchor = self.anchor_end();
for seg in &mut self.segments {
let cap = Self::stale_cap(seg, anchor);
if let Some(sub) = seg.stale_sub_mut() {
sub.set_stale_cap(cap);
}
}
}
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 update_drift_cap(&mut self) {
if let Ok(mut cap) = self.drift_cap.write() {
*cap = self.anchor_end();
}
}
fn poll_activate(
seg: &mut SegmentSub,
prefs: &Subscription,
min_sequence: u64,
anchor_end: Option<u64>,
waiter: &kio::Waiter,
) -> Poll<()> {
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_stale_cap(Self::stale_cap(seg, anchor_end));
let _ = sub.update(slice(prefs, seg.start, seg.end));
seg.sub = SubState::Active(Box::new(sub));
}
Err(_) => seg.sub = SubState::Done(None),
}
}
Poll::Ready(())
}
fn poll_segment(
seg: &mut SegmentSub,
prefs: &Subscription,
min_sequence: u64,
anchor_end: Option<u64>,
waiter: &kio::Waiter,
) -> Poll<Option<group::Consumer>> {
loop {
match &mut seg.sub {
SubState::Pending(_) => {
ready!(Self::poll_activate(seg, prefs, min_sequence, anchor_end, 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 count = sub.poll_finished(waiter).map(|res| res.ok());
let count = match count {
Poll::Ready(count) => count,
Poll::Pending => None,
};
seg.complete(count);
return Poll::Ready(None);
}
Err(_) => {
seg.complete(None);
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 anchor = self.anchor_end();
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,
anchor,
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 count = match sub.poll_finished(waiter) {
Poll::Ready(count) => count.ok(),
Poll::Pending => None,
};
seg.complete(count);
}
Poll::Ready(Err(_)) => seg.complete(None),
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 {
if self.finished {
return Poll::Ready(Ok(None));
}
if self.closed {
return Poll::Ready(self.orphan_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 anchor = self.anchor_end();
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,
anchor,
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 {
if self.finished {
return Poll::Ready(Ok(None));
}
if self.closed {
return Poll::Ready(self.orphan_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;
let anchor = self.anchor_end();
if let Some(seg) = self.segments.last_mut() {
if Self::poll_activate(seg, &self.last_prefs, self.min_sequence, anchor, 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 {
return Poll::Ready(Ok(None));
}
if self.closed && !pending_activation {
return match ready!(self.poll_final(waiter)) {
Some(_) => Poll::Ready(Ok(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;
}
let end = ready!(self.poll_final(waiter));
match self.finished {
true => Poll::Ready(Ok(end.unwrap_or(0))),
false => Poll::Ready(end.ok_or(Error::Dropped)),
}
}
fn poll_final(&mut self, waiter: &kio::Waiter) -> Poll<Option<u64>> {
let anchor = min_some(self.stale_cap, self.end_sequence);
let Some(seg) = self.segments.last_mut() else {
return Poll::Ready(None);
};
ready!(Self::poll_activate(
seg,
&self.last_prefs,
self.min_sequence,
anchor,
waiter
));
match &mut seg.sub {
SubState::Done(count) => Poll::Ready(*count),
SubState::Active(sub) => Poll::Ready(ready!(sub.poll_finished(waiter)).ok()),
SubState::Pending(_) => unreachable!("poll_activate resolved above"),
}
}
fn orphan_end(&self) -> Result<()> {
match self.segments.last().map(|seg| &seg.sub) {
Some(SubState::Done(Some(_))) => 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.update_drift_cap();
let anchor = self.anchor_end();
for seg in &mut self.segments {
let cap = Self::stale_cap(seg, anchor);
if let Some(sub) = seg.stale_sub_mut() {
sub.set_stale_cap(cap);
}
}
}
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::{Timestamp, 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 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(None);
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 (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();
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, vec![1, 3], "each segment is judged within its own boundary");
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 nested_splice_judges_within_the_outer_boundary() {
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 track_b, consumer_b) = track_pair("b");
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(None);
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, 3],
"the nested segment is judged within the outer boundary"
);
let mut arrival = outer.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, 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(10 * sequence);
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(10 * sequence);
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(10 * sequence);
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(10 * sequence);
let (mut track_c, consumer_c) = track_pair("c");
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_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 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(None);
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::Cancel).unwrap();
drop(producer);
let result = sub.recv_group().now_or_never().expect("must not stall forever");
assert!(matches!(result, Err(Error::Dropped)));
}
#[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::Cancel).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::Dropped))),
}
}
}
#[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(None);
let mut sub2 = consumer.subscribe(None);
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");
}
}