use std::{
collections::{BTreeMap, HashSet},
sync::Arc,
time::Duration,
};
use crate::{Error, Hop, runtime::Instant, track};
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(super) struct Candidate {
pub route: u64,
pub first: Option<Hop>,
pub local: bool,
}
#[derive(Clone, Debug)]
pub(super) struct Refusal {
pub err: Error,
pub standing: bool,
}
#[derive(Clone, Debug)]
pub(super) enum Event {
Selected {
best: Option<Candidate>,
serving_closing: bool,
},
Resolved { route: u64, result: Result<u64, Refusal> },
SourceClosed { source: u64 },
TrackAssigned { track: Arc<str> },
TrackInfo {
track: Arc<str>,
source: u64,
closing: bool,
result: Result<track::Info, Error>,
},
TrackEnded {
track: Arc<str>,
source: u64,
closing: bool,
result: Result<(), Error>,
delivered: bool,
},
Used { track: Arc<str> },
Unused { track: Arc<str>, now: Instant },
Deadline { now: Instant },
Closed,
}
#[derive(Clone, Debug)]
pub(super) enum Action {
Reselect,
Request { route: u64 },
Detach { source: u64 },
Resolve,
Query { track: Arc<str>, source: u64 },
Splice { track: Arc<str>, source: u64 },
Park { track: Arc<str> },
Release { track: Arc<str> },
Finish { track: Arc<str> },
Abort { track: Arc<str>, err: Error },
Arm { at: Option<Instant> },
End { err: Error },
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(super) enum Identity {
Undetermined,
Local,
Anonymous { route: u64 },
Publisher(Hop),
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(super) enum Pin {
Any,
Local,
Publisher(Hop),
Route(u64),
}
impl Identity {
fn pin(self) -> Pin {
match self {
Self::Undetermined => Pin::Any,
Self::Local => Pin::Local,
Self::Anonymous { route } => Pin::Route(route),
Self::Publisher(hop) => Pin::Publisher(hop),
}
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
enum TrackState {
Idle,
Querying { source: u64 },
Spliced { source: u64 },
Parked { since: Instant },
}
#[derive(Clone, Debug)]
struct Track {
state: TrackState,
refused: HashSet<u64>,
refusal: Option<Error>,
used: bool,
}
#[derive(Clone, Debug)]
pub(super) struct Front {
identity: Identity,
serving: Option<(u64, u64)>,
serving_closing: bool,
upstream: Option<u64>,
refused: HashSet<u64>,
last_err: Option<Error>,
resolved: bool,
tracks: BTreeMap<Arc<str>, Track>,
info: BTreeMap<Arc<str>, track::Info>,
linger: Duration,
armed: Option<Instant>,
ended: bool,
}
impl Front {
pub(super) fn new(linger: Duration) -> Self {
Self {
identity: Identity::Undetermined,
serving: None,
serving_closing: false,
upstream: None,
refused: HashSet::new(),
last_err: None,
resolved: false,
tracks: BTreeMap::new(),
info: BTreeMap::new(),
linger,
armed: None,
ended: false,
}
}
pub(super) fn pin(&self) -> Pin {
self.identity.pin()
}
pub(super) fn refused_routes(&self) -> &HashSet<u64> {
&self.refused
}
pub(super) fn retain_routes(&mut self, standing: impl Fn(u64) -> bool) {
self.refused.retain(|route| standing(*route));
}
pub(super) fn serving(&self) -> Option<u64> {
self.serving.map(|(source, _)| source)
}
#[cfg(test)]
fn ended(&self) -> bool {
self.ended
}
pub(super) fn step(&mut self, event: Event) -> Vec<Action> {
let mut actions = Vec::new();
if self.ended {
return actions;
}
match event {
Event::Selected { best, serving_closing } => self.selected(best, serving_closing, &mut actions),
Event::Resolved { route, result } => self.resolved(route, result, &mut actions),
Event::SourceClosed { source } => self.source_closed(source, &mut actions),
Event::TrackAssigned { track } => {
self.tracks.insert(
track,
Track {
state: TrackState::Idle,
refused: HashSet::new(),
refusal: None,
used: false,
},
);
}
Event::TrackInfo {
track,
source,
closing,
result,
} => self.track_info(track, source, closing, result, &mut actions),
Event::TrackEnded {
track,
source,
closing,
result,
delivered,
} => self.track_ended(track, source, closing, result, delivered, &mut actions),
Event::Used { track } => self.used(track, &mut actions),
Event::Unused { track, now } => self.unused(track, now, &mut actions),
Event::Deadline { now } => self.deadline(now, &mut actions),
Event::Closed => self.end(Error::Dropped, &mut actions),
}
if !self.ended {
let at = self.next_deadline();
if at != self.armed {
self.armed = at;
actions.push(Action::Arm { at });
}
}
actions
}
fn selected(&mut self, best: Option<Candidate>, serving_closing: bool, actions: &mut Vec<Action>) {
if self.serving.is_some() {
self.serving_closing = serving_closing;
}
let serving_route = self.serving.map(|(_, route)| route);
match best {
Some(candidate) if Some(candidate.route) == serving_route => {
self.upstream = None;
}
Some(candidate) if self.upstream == Some(candidate.route) => {}
Some(candidate) if candidate.local && self.identity == Identity::Local && serving_closing => {
self.end(Error::Dropped, actions);
}
Some(candidate) => {
self.upstream = Some(candidate.route);
actions.push(Action::Request { route: candidate.route });
}
None => {
self.upstream = None;
let err = match self.identity {
Identity::Undetermined => self.last_err.take().unwrap_or(Error::Unroutable),
_ => self.last_err.take().unwrap_or(Error::Dropped),
};
self.end(err, actions);
}
}
}
fn resolved(&mut self, route: u64, result: Result<u64, Refusal>, actions: &mut Vec<Action>) {
match result {
Ok(source) if self.upstream != Some(route) => actions.push(Action::Detach { source }),
Ok(source) => {
self.upstream = None;
self.attach(source, route, actions);
}
Err(_) if self.upstream != Some(route) => {}
Err(Refusal { standing: false, .. }) => {
self.upstream = None;
self.last_err = Some(Error::Unroutable);
actions.push(Action::Reselect);
}
Err(Refusal { err, standing: true }) => {
self.upstream = None;
match self.serving {
Some(_) => {
self.refused.insert(route);
self.last_err = Some(err);
actions.push(Action::Reselect);
}
None => self.end(err, actions),
}
}
}
}
fn attach(&mut self, source: u64, route: u64, actions: &mut Vec<Action>) {
if let Some((old, _)) = self.serving.take() {
actions.push(Action::Detach { source: old });
for track in self.tracks.values_mut() {
if matches!(track.state, TrackState::Querying { source: s } | TrackState::Spliced { source: s } if s == old)
{
track.state = TrackState::Idle;
}
}
}
self.serving = Some((source, route));
self.serving_closing = false;
if !self.resolved {
self.resolved = true;
actions.push(Action::Resolve);
}
for (name, track) in &mut self.tracks {
if track.used && track.state == TrackState::Idle {
track.state = TrackState::Querying { source };
actions.push(Action::Query {
track: name.clone(),
source,
});
}
}
}
pub(super) fn identify(&mut self, candidate: Candidate) {
if self.identity != Identity::Undetermined {
return;
}
self.identity = match (candidate.local, candidate.first) {
(true, _) => Identity::Local,
(false, Some(hop)) if hop != Hop::UNKNOWN => Identity::Publisher(hop),
(false, _) => Identity::Anonymous { route: candidate.route },
};
}
fn source_closed(&mut self, source: u64, actions: &mut Vec<Action>) {
let Some((serving, route)) = self.serving else {
return;
};
if serving != source {
return;
}
self.refused.insert(route);
self.serving = None;
self.serving_closing = false;
actions.push(Action::Detach { source });
for track in self.tracks.values_mut() {
if matches!(track.state, TrackState::Querying { source: s } | TrackState::Spliced { source: s } if s == source)
{
track.state = TrackState::Idle;
}
}
self.last_err = Some(Error::Dropped);
match self.identity {
Identity::Local | Identity::Anonymous { .. } | Identity::Undetermined => self.end(Error::Dropped, actions),
Identity::Publisher(_) => actions.push(Action::Reselect),
}
}
fn track_info(
&mut self,
name: Arc<str>,
source: u64,
closing: bool,
result: Result<track::Info, Error>,
actions: &mut Vec<Action>,
) {
let Some(track) = self.tracks.get_mut(&name) else {
return;
};
if track.state != (TrackState::Querying { source }) {
return;
}
let verdict = match result {
Ok(info) => match self.info.get(&name) {
Some(expected)
if expected.timescale != info.timescale
|| expected.max_age != info.max_age
|| expected.priority != info.priority =>
{
Err(Error::Unsupported)
}
Some(_) => Ok(()),
None => {
self.info.insert(name.clone(), info);
Ok(())
}
},
Err(err) => Err(err),
};
match verdict {
Ok(()) => {
track.state = TrackState::Spliced { source };
actions.push(Action::Splice {
track: name.clone(),
source,
});
}
Err(err) => self.refuse(name, source, closing, err, actions),
}
}
fn track_ended(
&mut self,
name: Arc<str>,
source: u64,
closing: bool,
result: Result<(), Error>,
delivered: bool,
actions: &mut Vec<Action>,
) {
let Some(track) = self.tracks.get_mut(&name) else {
return;
};
if track.state != (TrackState::Spliced { source }) {
return;
}
match result {
Ok(()) => {
self.tracks.remove(&name);
actions.push(Action::Finish { track: name });
}
Err(_) if delivered => {
track.state = TrackState::Idle;
self.redispatch(name, actions);
}
Err(err) => self.refuse(name, source, closing, err, actions),
}
}
fn refuse(&mut self, name: Arc<str>, source: u64, closing: bool, err: Error, actions: &mut Vec<Action>) {
let track = self.tracks.get_mut(&name).expect("refusing a known track");
track.state = TrackState::Idle;
if closing {
return;
}
track.refused.insert(source);
track.refusal = Some(err);
self.redispatch(name, actions);
}
fn redispatch(&mut self, name: Arc<str>, actions: &mut Vec<Action>) {
let Some((source, _)) = self.serving else {
return;
};
let track = self.tracks.get_mut(&name).expect("dispatching a known track");
if !track.used || track.state != TrackState::Idle {
return;
}
if track.refused.contains(&source) {
let err = track.refusal.clone().unwrap_or(Error::NotFound);
self.tracks.remove(&name);
actions.push(Action::Abort { track: name, err });
return;
}
if self.serving_closing {
return;
}
track.state = TrackState::Querying { source };
actions.push(Action::Query { track: name, source });
}
fn used(&mut self, name: Arc<str>, actions: &mut Vec<Action>) {
let Some(track) = self.tracks.get_mut(&name) else {
return;
};
track.used = true;
if let TrackState::Parked { .. } = track.state {
track.state = TrackState::Idle;
}
self.redispatch(name, actions);
}
fn unused(&mut self, name: Arc<str>, now: Instant, actions: &mut Vec<Action>) {
let Some(track) = self.tracks.get_mut(&name) else {
return;
};
track.used = false;
match track.state {
TrackState::Spliced { .. } if self.identity == Identity::Local => {
track.state = TrackState::Idle;
actions.push(Action::Release { track: name });
}
TrackState::Spliced { .. } => {
track.state = TrackState::Parked { since: now };
actions.push(Action::Park { track: name });
}
TrackState::Querying { .. } => track.state = TrackState::Idle,
_ => {}
}
}
fn deadline(&mut self, now: Instant, actions: &mut Vec<Action>) {
self.armed = None;
for (name, track) in &mut self.tracks {
if let TrackState::Parked { since } = track.state
&& since + self.linger <= now
{
track.state = TrackState::Idle;
actions.push(Action::Release { track: name.clone() });
}
}
}
fn next_deadline(&self) -> Option<Instant> {
self.tracks
.values()
.filter_map(|track| match track.state {
TrackState::Parked { since } => Some(since + self.linger),
_ => None,
})
.min()
}
fn end(&mut self, err: Error, actions: &mut Vec<Action>) {
if self.ended {
return;
}
self.ended = true;
self.serving = None;
actions.push(Action::End { err });
}
}
#[cfg(test)]
mod tests {
use super::*;
const LINGER: Duration = Duration::from_secs(30);
fn hop(id: u64) -> Hop {
Hop::new(id).unwrap()
}
fn remote(route: u64, first: u64) -> Candidate {
Candidate {
route,
first: Some(hop(first)),
local: false,
}
}
fn local(route: u64) -> Candidate {
Candidate {
route,
first: None,
local: true,
}
}
fn info() -> track::Info {
track::Info::default()
}
fn name(s: &str) -> Arc<str> {
Arc::from(s)
}
#[track_caller]
fn assert_actions(actual: Vec<Action>, expected: &[Action]) {
assert_eq!(format!("{actual:?}"), format!("{expected:?}"));
}
fn serving(candidate: Candidate, source: u64) -> Front {
let mut front = Front::new(LINGER);
assert_actions(
front.step(Event::Selected {
best: Some(candidate),
serving_closing: false,
}),
&[Action::Request { route: candidate.route }],
);
front.identify(candidate);
assert_actions(
front.step(Event::Resolved {
route: candidate.route,
result: Ok(source),
}),
&[Action::Resolve],
);
assert_actions(front.step(Event::TrackAssigned { track: name("video") }), &[]);
assert_actions(
front.step(Event::Used { track: name("video") }),
&[Action::Query {
track: name("video"),
source,
}],
);
assert_actions(
front.step(Event::TrackInfo {
track: name("video"),
source,
closing: false,
result: Ok(info()),
}),
&[Action::Splice {
track: name("video"),
source,
}],
);
front
}
#[test]
fn first_source_resolves_and_serves_read_tracks() {
let front = serving(remote(1, 10), 100);
assert_eq!(front.identity, Identity::Publisher(hop(10)));
assert_eq!(front.pin(), Pin::Publisher(hop(10)));
}
#[test]
fn nothing_routable_ends_an_unresolved_front() {
let mut front = Front::new(LINGER);
assert_actions(
front.step(Event::Selected {
best: None,
serving_closing: false,
}),
&[Action::End { err: Error::Unroutable }],
);
assert!(front.ended());
assert_actions(front.step(Event::Closed), &[]);
}
#[test]
fn unchanged_selection_is_a_no_op() {
let mut front = serving(remote(1, 10), 100);
assert_actions(
front.step(Event::Selected {
best: Some(remote(1, 10)),
serving_closing: false,
}),
&[],
);
}
#[test]
fn better_route_takes_over_and_resplices() {
let mut front = serving(remote(1, 10), 100);
assert_actions(
front.step(Event::Selected {
best: Some(remote(2, 10)),
serving_closing: false,
}),
&[Action::Request { route: 2 }],
);
assert_actions(
front.step(Event::Resolved {
route: 2,
result: Ok(200),
}),
&[
Action::Detach { source: 100 },
Action::Query {
track: name("video"),
source: 200,
},
],
);
assert_eq!(front.serving, Some((200, 2)));
}
#[test]
fn a_stale_resolution_is_let_go() {
let mut front = serving(remote(1, 10), 100);
front.step(Event::Selected {
best: Some(remote(2, 10)),
serving_closing: false,
});
assert_actions(
front.step(Event::Selected {
best: Some(remote(1, 10)),
serving_closing: false,
}),
&[],
);
assert_actions(
front.step(Event::Resolved {
route: 2,
result: Ok(200),
}),
&[Action::Detach { source: 200 }],
);
}
#[test]
fn dead_source_reselects_through_the_same_publisher() {
let mut front = serving(remote(1, 10), 100);
assert_actions(
front.step(Event::SourceClosed { source: 100 }),
&[Action::Detach { source: 100 }, Action::Reselect],
);
assert!(front.refused_routes().contains(&1));
assert_actions(
front.step(Event::Selected {
best: Some(remote(3, 10)),
serving_closing: false,
}),
&[Action::Request { route: 3 }],
);
assert_actions(
front.step(Event::Resolved {
route: 3,
result: Ok(300),
}),
&[Action::Query {
track: name("video"),
source: 300,
}],
);
}
#[test]
fn retracted_route_with_no_replacement_ends_a_live_front() {
let mut front = serving(local(1), 100);
assert_actions(
front.step(Event::Selected {
best: None,
serving_closing: false,
}),
&[Action::End { err: Error::Dropped }],
);
}
#[test]
fn anonymous_front_keeps_its_standing_route() {
let candidate = Candidate {
route: 1,
first: Some(Hop::UNKNOWN),
local: false,
};
let mut front = serving(candidate, 100);
assert_actions(
front.step(Event::Selected {
best: Some(candidate),
serving_closing: false,
}),
&[],
);
assert_actions(
front.step(Event::Selected {
best: None,
serving_closing: false,
}),
&[Action::End { err: Error::Dropped }],
);
}
#[test]
fn dead_source_with_no_replacement_ends() {
let mut front = serving(remote(1, 10), 100);
front.step(Event::SourceClosed { source: 100 });
assert_actions(
front.step(Event::Selected {
best: None,
serving_closing: false,
}),
&[Action::End { err: Error::Dropped }],
);
}
#[test]
fn anonymous_source_never_resumes() {
let candidate = Candidate {
route: 1,
first: Some(Hop::UNKNOWN),
local: false,
};
let mut front = serving(candidate, 100);
assert_eq!(front.pin(), Pin::Route(1));
assert_actions(
front.step(Event::SourceClosed { source: 100 }),
&[Action::Detach { source: 100 }, Action::End { err: Error::Dropped }],
);
}
#[test]
fn standing_refusal_ends_an_unresolved_front() {
let mut front = Front::new(LINGER);
front.step(Event::Selected {
best: Some(remote(1, 10)),
serving_closing: false,
});
assert_actions(
front.step(Event::Resolved {
route: 1,
result: Err(Refusal {
err: Error::NotFound,
standing: true,
}),
}),
&[Action::End { err: Error::NotFound }],
);
}
#[test]
fn standing_refusal_while_serving_skips_the_refuser() {
let mut front = serving(remote(1, 10), 100);
front.step(Event::Selected {
best: Some(remote(2, 10)),
serving_closing: false,
});
assert_actions(
front.step(Event::Resolved {
route: 2,
result: Err(Refusal {
err: Error::NotFound,
standing: true,
}),
}),
&[Action::Reselect],
);
assert!(front.refused_routes().contains(&2));
assert_eq!(front.serving, Some((100, 1)));
}
#[test]
fn retracted_route_is_not_a_refusal() {
let mut front = Front::new(LINGER);
front.step(Event::Selected {
best: Some(remote(1, 10)),
serving_closing: false,
});
assert_actions(
front.step(Event::Resolved {
route: 1,
result: Err(Refusal {
err: Error::Unroutable,
standing: false,
}),
}),
&[Action::Reselect],
);
assert!(front.refused_routes().is_empty());
assert!(!front.ended());
}
#[test]
fn a_source_refusing_a_track_aborts_it_and_nothing_else() {
let mut front = serving(remote(1, 10), 100);
front.step(Event::TrackAssigned { track: name("audio") });
front.step(Event::Used { track: name("audio") });
assert_actions(
front.step(Event::TrackInfo {
track: name("audio"),
source: 100,
closing: false,
result: Err(Error::NotFound),
}),
&[Action::Abort {
track: name("audio"),
err: Error::NotFound,
}],
);
assert_eq!(front.tracks[&name("video")].state, TrackState::Spliced { source: 100 });
}
#[test]
fn a_closing_source_refusal_is_not_a_verdict() {
let mut front = serving(remote(1, 10), 100);
front.step(Event::TrackAssigned { track: name("audio") });
front.step(Event::Used { track: name("audio") });
assert_actions(
front.step(Event::TrackInfo {
track: name("audio"),
source: 100,
closing: true,
result: Err(Error::NotFound),
}),
&[],
);
assert!(front.tracks[&name("audio")].refused.is_empty());
front.step(Event::SourceClosed { source: 100 });
front.step(Event::Selected {
best: Some(remote(2, 10)),
serving_closing: false,
});
let actions = front.step(Event::Resolved {
route: 2,
result: Ok(200),
});
assert!(
actions
.iter()
.any(|action| matches!(action, Action::Query { track, source: 200 } if track.as_ref() == "audio"))
);
}
#[test]
fn a_copy_dying_after_delivering_resplices() {
let mut front = serving(remote(1, 10), 100);
assert_actions(
front.step(Event::TrackEnded {
track: name("video"),
source: 100,
closing: false,
result: Err(Error::Dropped),
delivered: true,
}),
&[Action::Query {
track: name("video"),
source: 100,
}],
);
}
#[test]
fn a_copy_dying_before_delivering_is_a_refusal() {
let mut front = serving(remote(1, 10), 100);
assert_actions(
front.step(Event::TrackEnded {
track: name("video"),
source: 100,
closing: false,
result: Err(Error::Dropped),
delivered: false,
}),
&[Action::Abort {
track: name("video"),
err: Error::Dropped,
}],
);
}
#[test]
fn a_completed_copy_finishes_the_track() {
let mut front = serving(remote(1, 10), 100);
assert_actions(
front.step(Event::TrackEnded {
track: name("video"),
source: 100,
closing: false,
result: Ok(()),
delivered: true,
}),
&[Action::Finish { track: name("video") }],
);
}
#[test]
fn incompatible_copy_is_refused() {
let mut front = serving(remote(1, 10), 100);
front.step(Event::Selected {
best: Some(remote(2, 10)),
serving_closing: false,
});
front.step(Event::Resolved {
route: 2,
result: Ok(200),
});
let other = track::Info {
max_age: Duration::from_secs(1),
..track::Info::default()
};
assert_actions(
front.step(Event::TrackInfo {
track: name("video"),
source: 200,
closing: false,
result: Ok(other),
}),
&[Action::Abort {
track: name("video"),
err: Error::Unsupported,
}],
);
}
#[test]
fn unread_track_parks_then_releases_after_the_linger() {
let mut front = serving(remote(1, 10), 100);
let t0 = Instant::now();
assert_actions(
front.step(Event::Unused {
track: name("video"),
now: t0,
}),
&[
Action::Park { track: name("video") },
Action::Arm { at: Some(t0 + LINGER) },
],
);
assert_actions(
front.step(Event::Deadline { now: t0 + LINGER / 2 }),
&[Action::Arm { at: Some(t0 + LINGER) }],
);
assert_actions(
front.step(Event::Deadline { now: t0 + LINGER }),
&[Action::Release { track: name("video") }],
);
assert_actions(
front.step(Event::Used { track: name("video") }),
&[Action::Query {
track: name("video"),
source: 100,
}],
);
}
#[test]
fn a_returning_reader_cancels_the_linger() {
let mut front = serving(remote(1, 10), 100);
let t0 = Instant::now();
front.step(Event::Unused {
track: name("video"),
now: t0,
});
assert_actions(
front.step(Event::Used { track: name("video") }),
&[
Action::Query {
track: name("video"),
source: 100,
},
Action::Arm { at: None },
],
);
}
#[test]
fn unread_track_is_never_spliced() {
let mut front = serving(remote(1, 10), 100);
assert_actions(front.step(Event::TrackAssigned { track: name("audio") }), &[]);
assert_eq!(front.tracks[&name("audio")].state, TrackState::Idle);
}
#[test]
fn local_newcomer_takes_over_a_live_incumbent() {
let mut front = serving(local(1), 100);
assert_eq!(front.identity, Identity::Local);
assert_eq!(front.pin(), Pin::Local);
assert_actions(
front.step(Event::Selected {
best: Some(local(2)),
serving_closing: false,
}),
&[Action::Request { route: 2 }],
);
}
#[test]
fn local_newcomer_never_splices_into_a_closing_incumbent() {
let mut front = serving(local(1), 100);
assert_actions(
front.step(Event::Selected {
best: Some(local(2)),
serving_closing: true,
}),
&[Action::End { err: Error::Dropped }],
);
}
#[test]
fn local_incumbent_ending_ends_the_front() {
let mut front = serving(local(1), 100);
assert_actions(
front.step(Event::SourceClosed { source: 100 }),
&[Action::Detach { source: 100 }, Action::End { err: Error::Dropped }],
);
}
#[test]
fn teardown_ends_everything() {
let mut front = serving(remote(1, 10), 100);
assert_actions(front.step(Event::Closed), &[Action::End { err: Error::Dropped }]);
}
#[test]
fn bounded_sequences_hold_the_invariants() {
let t0 = Instant::now();
let alphabet: Vec<Event> = vec![
Event::Selected {
best: Some(remote(1, 10)),
serving_closing: false,
},
Event::Selected {
best: Some(remote(2, 10)),
serving_closing: false,
},
Event::Selected {
best: None,
serving_closing: false,
},
Event::Resolved {
route: 1,
result: Ok(100),
},
Event::Resolved {
route: 2,
result: Ok(200),
},
Event::Resolved {
route: 2,
result: Err(Refusal {
err: Error::NotFound,
standing: true,
}),
},
Event::SourceClosed { source: 100 },
Event::TrackAssigned { track: name("v") },
Event::Used { track: name("v") },
Event::Unused {
track: name("v"),
now: t0,
},
Event::TrackInfo {
track: name("v"),
source: 100,
closing: false,
result: Ok(info()),
},
Event::TrackInfo {
track: name("v"),
source: 100,
closing: false,
result: Err(Error::NotFound),
},
Event::TrackEnded {
track: name("v"),
source: 100,
closing: false,
result: Err(Error::Dropped),
delivered: false,
},
Event::Deadline { now: t0 + LINGER },
];
fn walk(front: &Front, alphabet: &[Event], depth: usize, sequences: &mut usize) {
for event in alphabet {
if let Event::Resolved { result: Ok(source), .. } = event
&& front.serving.map(|(serving, _)| serving) == Some(*source)
{
continue;
}
let mut next = front.clone();
let before = format!("{next:?}");
let actions = next.step(event.clone());
*sequences += 1;
if format!("{next:?}") == before {
assert!(
actions.iter().all(|action| matches!(action, Action::Detach { .. })),
"no-op event {event:?} produced {actions:?}"
);
}
if front.ended() {
assert!(actions.is_empty());
}
for action in &actions {
if let Action::Query { track, source } = action {
assert!(!next.tracks[track].refused.contains(source));
}
}
let splices = actions
.iter()
.filter(|action| matches!(action, Action::Splice { .. }))
.count();
assert!(splices <= 1);
for action in &actions {
if let Action::Detach { source } = action {
assert_ne!(next.serving.map(|(s, _)| s), Some(*source));
}
}
if depth > 1 {
walk(&next, alphabet, depth - 1, sequences);
}
}
}
let mut sequences = 0;
walk(&Front::new(LINGER), &alphabet, 5, &mut sequences);
assert!(sequences > 100_000);
}
}