use crate::{cache, stats, track};
use std::{
collections::{HashMap, VecDeque},
sync::Arc,
task::{Poll, ready},
};
use crate::Error;
use crate::origin::Route;
use super::origin_impl::Announcer;
use super::{Requests, WeakCache};
#[derive(Clone, Debug)]
#[non_exhaustive]
pub struct Info {
pub pool: cache::Pool,
pub cache_duration: std::time::Duration,
pub path: crate::PathOwned,
}
impl Default for Info {
fn default() -> Self {
Self {
pool: cache::Pool::new(cache::Config::default().with_expiry(cache::DEFAULT_EXPIRY)),
cache_duration: std::time::Duration::MAX,
path: crate::PathOwned::default(),
}
}
}
impl Info {
pub fn new() -> Self {
Self::default()
}
pub fn produce(self) -> Producer {
Producer::new(self)
}
}
#[derive(Default)]
struct BroadcastState {
tracks: WeakCache<Arc<str>, track::TrackWeak>,
unique: u64,
requests: Requests<Arc<str>, track::Request>,
spliced: Option<SplicedState>,
closing: bool,
finished: bool,
abort: Option<Error>,
}
#[derive(Default)]
struct SplicedState {
tracks: HashMap<Arc<str>, super::resume::Producer>,
pending: VecDeque<Arc<str>>,
}
impl BroadcastState {
fn insert_track(&mut self, weak: track::TrackWeak) -> Result<(), Error> {
match self.tracks.insert(weak.name().clone(), weak) {
Some(_) => Err(Error::Duplicate),
None => Ok(()),
}
}
fn reject_unserved(&mut self, err: Error) {
for request in self.requests.drain_queued() {
request.reject(err.clone());
}
for track in self.tracks.iter() {
track.reject(err.clone());
}
}
fn is_used(&self) -> bool {
if let Some(spliced) = &self.spliced {
return spliced.tracks.values().any(|track| track.is_used());
}
!self.requests.is_empty() || self.tracks.iter().any(|track| track.is_used())
}
fn register_demand(&self, waiter: &kio::Waiter, want: bool) {
if let Some(spliced) = &self.spliced {
for track in spliced.tracks.values() {
let _ = match want {
true => track.poll_used(waiter),
false => track.poll_unused(waiter),
};
}
return;
}
for track in self.tracks.iter() {
match want {
true => track.poll_used(waiter),
false => track.poll_unused(waiter),
}
}
}
}
#[derive(Clone)]
pub struct Producer {
info: Arc<Info>,
alive: Arc<Alive>,
state: kio::Shared<BroadcastState>,
stats: stats::Scope,
}
impl Producer {
pub fn new(info: Info) -> Self {
let state = kio::Shared::<BroadcastState>::default();
Self {
info: Arc::new(info),
alive: Alive::new(state.clone()),
state,
stats: stats::Scope::default(),
}
}
pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
self.stats = scope;
self
}
pub(crate) fn with_announcer(self, announcer: Announcer) -> Self {
*self.alive.announcer.lock() = Some(announcer);
self
}
pub fn announce(&self, route: Route) -> Result<(), Error> {
let mut announcer = self.alive.announcer.lock();
let announcer = announcer.as_mut().ok_or(Error::Closed)?;
announcer.announce(route)
}
pub fn unannounce(&self) {
self.alive.unannounce();
}
pub(crate) fn new_spliced(info: Info) -> Self {
let state = kio::Shared::new(BroadcastState {
spliced: Some(SplicedState::default()),
..Default::default()
});
Self {
info: Arc::new(info),
alive: Alive::new(state.clone()),
state,
stats: stats::Scope::default(),
}
}
pub fn info(&self) -> &Info {
&self.info
}
pub fn demand(&self) -> Demand {
Demand {
alive: self.alive.token.consume().weak(),
state: self.state.clone(),
}
}
pub fn create_track(
&self,
name: impl Into<Arc<str>>,
info: impl Into<Option<track::Info>>,
) -> Result<track::Producer, Error> {
let name = name.into();
let info = info.into().unwrap_or_default();
let mut state = self.state.lock();
if let Some(request) = state.requests.take(name.as_ref()) {
let track = request.with_stats(self.stats.clone()).accept(info);
let _ = state.tracks.insert(name, track.weak());
return Ok(track);
}
let track = track::Producer::new(self.info.clone(), name, info).with_stats(self.stats.clone());
state.insert_track(track.weak())?;
Ok(track)
}
pub fn reserve_track(&self, name: impl Into<Arc<str>>) -> Result<track::Request, Error> {
let request = track::Request::new(self.info.clone(), name).with_stats(self.stats.clone());
self.state.lock().insert_track(request.weak())?;
Ok(request)
}
pub fn unique_track(&self, suffix: &str, info: impl Into<Option<track::Info>>) -> Result<track::Producer, Error> {
let name = self.unique_name(suffix);
self.create_track(name, info)
}
pub fn unique_name(&self, suffix: &str) -> String {
let mut state = self.state.lock();
let separator = if suffix.starts_with(|c: char| c.is_ascii_digit()) {
"-"
} else {
""
};
loop {
let id = state.unique;
state.unique = id.checked_add(1).expect("unique track IDs exhausted");
let name = format!("{id}{separator}{suffix}");
if !state.tracks.contains_key(name.as_str()) {
return name;
}
}
}
pub fn dynamic(&self) -> Dynamic {
Dynamic::new(
self.info.clone(),
self.alive.clone(),
self.state.clone(),
self.stats.clone(),
)
}
pub(crate) fn poll_spliced_assigned(&self, waiter: &kio::Waiter) -> Poll<(Arc<str>, super::resume::Producer)> {
let mut state = ready!(self.state.poll(waiter, |state| {
match &state.spliced {
Some(spliced) if !spliced.pending.is_empty() => Poll::Ready(()),
_ => Poll::Pending,
}
}));
let spliced = state.spliced.as_mut().expect("predicate guaranteed spliced");
let name = spliced.pending.pop_front().expect("predicate guaranteed a request");
let producer = spliced.tracks.get(&name).expect("pending name without a track").clone();
Poll::Ready((name, producer))
}
pub(crate) fn release_spliced(&self, err: Error) {
let mut state = self.state.lock();
if let Some(spliced) = state.spliced.as_mut() {
for name in std::mem::take(&mut spliced.pending) {
if let Some(producer) = spliced.tracks.get_mut(&name) {
let _ = producer.abort(err.clone());
}
}
spliced.tracks.clear();
}
}
pub fn consume(&self) -> Consumer {
Consumer {
info: self.info.clone(),
alive: self.alive.token.consume(),
state: self.state.clone(),
stats: stats::Scope::default(),
}
}
pub fn close(&self) {
self.alive.close();
}
#[doc(hidden)]
#[deprecated(note = "use close(); a broadcast end carries no cause")]
pub fn finish(&self) {
self.alive.end(true);
}
#[doc(hidden)]
#[deprecated(note = "use close(); a broadcast end carries no cause")]
pub fn abort(self, err: Error) -> Result<(), Error> {
{
let mut state = self.state.lock();
if state.closing {
return Err(Error::Closed);
}
state.closing = true;
state.abort = Some(err.clone());
state.reject_unserved(err);
}
let _ = self.alive.token.close();
self.alive.retire();
Ok(())
}
pub fn is_clone(&self, other: &Self) -> bool {
self.state.same_channel(&other.state)
}
}
struct Alive {
token: kio::Producer<()>,
state: kio::Shared<BroadcastState>,
announcer: kio::Lock<Option<Announcer>>,
}
impl Alive {
fn new(state: kio::Shared<BroadcastState>) -> Arc<Self> {
Arc::new(Self {
token: kio::Producer::default(),
state,
announcer: kio::Lock::new(None),
})
}
fn unannounce(&self) {
if let Some(announcer) = self.announcer.lock().as_mut() {
announcer.withdraw();
}
}
fn close(&self) {
self.end(false);
}
fn end(&self, finished: bool) {
{
let mut state = self.state.lock();
if std::mem::replace(&mut state.closing, true) {
return;
}
state.finished = finished;
state.reject_unserved(Error::Unroutable);
}
let _ = self.token.close();
self.retire();
}
fn retire(&self) {
let announcer = self.announcer.lock().take();
drop(announcer);
}
}
impl Drop for Alive {
fn drop(&mut self) {
self.close();
}
}
#[cfg(test)]
#[allow(missing_docs)] impl Producer {
pub fn assert_create_track(
&mut self,
name: impl Into<Arc<str>>,
info: impl Into<Option<track::Info>>,
) -> track::Producer {
self.create_track(name, info).expect("should not have errored")
}
}
pub(crate) struct SourceGuard(Producer);
impl SourceGuard {
pub fn new(producer: Producer) -> Self {
Self(producer)
}
}
impl Drop for SourceGuard {
fn drop(&mut self) {
self.0.close();
}
}
#[derive(Clone)]
pub struct Dynamic {
info: Arc<Info>,
alive: Arc<Alive>,
state: kio::Shared<BroadcastState>,
stats: stats::Scope,
_handler: Handler,
}
struct Handler(kio::Shared<BroadcastState>);
impl Handler {
fn new(state: kio::Shared<BroadcastState>) -> Self {
state.lock().requests.add_handler();
Self(state)
}
}
impl Clone for Handler {
fn clone(&self) -> Self {
Self::new(self.0.clone())
}
}
impl Drop for Handler {
fn drop(&mut self) {
let mut state = self.0.lock();
if state.requests.remove_handler() {
for request in state.requests.drain_queued() {
request.reject(Error::Dropped);
}
}
}
}
impl Dynamic {
fn new(info: Arc<Info>, alive: Arc<Alive>, state: kio::Shared<BroadcastState>, stats: stats::Scope) -> Self {
Self {
info,
alive,
_handler: Handler::new(state.clone()),
state,
stats,
}
}
pub fn info(&self) -> &Info {
&self.info
}
pub fn poll_requested_track(&mut self, waiter: &kio::Waiter) -> Poll<Result<track::Request, Error>> {
let mut state = ready!(self.state.poll(waiter, |state| {
if state.requests.has_queued() || state.closing {
Poll::Ready(())
} else {
Poll::Pending
}
}));
if state.closing && !state.requests.has_queued() {
return Poll::Ready(Err(Error::Closed));
}
let name = state.requests.pop().expect("predicate guaranteed a request");
let pending = state.requests.remove(&name).expect("popped key must be pending");
let _ = state.tracks.insert(name, pending.weak());
Poll::Ready(Ok(pending.claim().with_stats(self.stats.clone())))
}
pub async fn requested_track(&mut self) -> Result<track::Request, Error> {
kio::wait(|waiter| self.poll_requested_track(waiter)).await
}
pub fn consume(&self) -> Consumer {
Consumer {
info: self.info.clone(),
alive: self.alive.token.consume(),
state: self.state.clone(),
stats: stats::Scope::default(),
}
}
pub async fn closed(&self) -> Error {
kio::wait(|waiter| self.poll_closed(waiter)).await
}
pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
ready!(self.alive.token.poll_closed(waiter));
Poll::Ready(self.state.read().abort.clone().unwrap_or(Error::Dropped))
}
pub fn is_clone(&self, other: &Self) -> bool {
self.state.same_channel(&other.state)
}
}
#[cfg(test)]
use futures::FutureExt;
#[cfg(test)]
#[allow(missing_docs)] impl Dynamic {
pub fn assert_request(&mut self) -> track::Request {
self.requested_track()
.now_or_never()
.expect("should not have blocked")
.expect("should not have errored")
}
pub fn assert_no_request(&mut self) {
assert!(self.requested_track().now_or_never().is_none(), "should have blocked");
}
}
pub struct Consumer {
info: Arc<Info>,
alive: kio::Consumer<()>,
state: kio::Shared<BroadcastState>,
stats: stats::Scope,
}
impl Clone for Consumer {
fn clone(&self) -> Self {
Self {
info: self.info.clone(),
alive: self.alive.clone(),
state: self.state.clone(),
stats: self.stats.clone(),
}
}
}
impl Consumer {
pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
self.stats = scope;
self
}
pub(crate) fn with_path(mut self, path: crate::PathOwned) -> Self {
if self.info.path != path {
let mut info = (*self.info).clone();
info.path = path;
self.info = Arc::new(info);
}
self
}
pub fn info(&self) -> &Info {
&self.info
}
pub fn track(&self, name: &str) -> Result<track::Consumer, Error> {
self.track_inner(name)
.map(|track| track.with_broadcast(self.info.clone()).with_stats(self.stats.clone()))
}
fn track_inner(&self, name: &str) -> Result<track::Consumer, Error> {
let mut state = self.state.lock();
if state.closing {
return Err(Error::Unroutable);
}
if let Some(spliced) = state.spliced.as_mut() {
if spliced.tracks.get(name).is_some_and(|track| track.is_aborted()) {
spliced.tracks.remove(name);
}
if let Some(producer) = spliced.tracks.get(name) {
return Ok(track::Consumer::spliced(
name.into(),
self.info.clone(),
producer.consume(),
));
}
let name: Arc<str> = name.into();
let producer = super::resume::Producer::new();
let consumer = producer.consume();
spliced.tracks.insert(name.clone(), producer);
spliced.pending.push_back(name.clone());
return Ok(track::Consumer::spliced(name, self.info.clone(), consumer));
}
if let Some(weak) = state.tracks.get(name) {
match weak.try_consume() {
Some(consumer) => return Ok(consumer),
None => {
state.tracks.remove(name);
}
}
}
if let Some(pending) = state.requests.join(name) {
return Ok(pending.consume());
}
let name: Arc<str> = name.into();
let request = track::Request::new(self.info.clone(), name.clone());
let consumer = request.consume();
if state.requests.insert(name, request).is_err() {
return Err(Error::NotFound);
}
Ok(consumer)
}
pub fn demand(&self) -> Demand {
Demand {
alive: self.alive.weak(),
state: self.state.clone(),
}
}
pub async fn closed(&self) -> Error {
self.alive.closed().await;
self.state.read().abort.clone().unwrap_or(Error::Dropped)
}
pub fn is_closed(&self) -> bool {
self.alive.is_closed()
}
pub(crate) fn is_closing(&self) -> bool {
self.state.read().closing
}
#[doc(hidden)]
#[deprecated(note = "a broadcast end carries no cause")]
pub fn is_finished(&self) -> bool {
self.state.read().finished
}
pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<()> {
self.alive.poll_closed(waiter)
}
pub fn is_clone(&self, other: &Self) -> bool {
self.state.same_channel(&other.state)
}
pub(crate) fn weak(&self) -> WeakConsumer {
WeakConsumer {
info: self.info.clone(),
alive: self.alive.weak(),
state: self.state.clone(),
}
}
}
#[derive(Clone)]
pub(crate) struct WeakConsumer {
info: Arc<Info>,
alive: kio::ConsumerWeak<()>,
state: kio::Shared<BroadcastState>,
}
impl WeakConsumer {
pub fn consume(&self) -> Consumer {
Consumer {
info: self.info.clone(),
alive: self.alive.consume(),
state: self.state.clone(),
stats: stats::Scope::default(),
}
}
}
impl super::WeakEntry for WeakConsumer {
fn is_closed(&self) -> bool {
self.alive.is_closed()
}
fn same_channel(&self, other: &Self) -> bool {
self.state.same_channel(&other.state)
}
}
#[derive(Clone)]
pub struct Demand {
alive: kio::ConsumerWeak<()>,
state: kio::Shared<BroadcastState>,
}
impl Demand {
pub fn is_used(&self) -> bool {
self.state.read().is_used()
}
pub async fn used(&self) -> Result<(), Error> {
kio::wait(|waiter| self.poll_used(waiter)).await
}
pub async fn unused(&self) -> Result<(), Error> {
kio::wait(|waiter| self.poll_unused(waiter)).await
}
pub fn poll_used(&self, waiter: &kio::Waiter) -> Poll<Result<(), Error>> {
self.poll_demand(waiter, true)
}
pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<Result<(), Error>> {
self.poll_demand(waiter, false)
}
fn poll_demand(&self, waiter: &kio::Waiter, want: bool) -> Poll<Result<(), Error>> {
if self.alive.poll_closed(waiter).is_ready() {
return Poll::Ready(Err(Error::Dropped));
}
let ready = self.state.poll(waiter, |state| {
state.register_demand(waiter, want);
match state.is_used() == want {
true => Poll::Ready(()),
false => Poll::Pending,
}
});
match ready {
Poll::Ready(_) => Poll::Ready(Ok(())),
Poll::Pending => Poll::Pending,
}
}
}
#[cfg(test)]
#[allow(missing_docs)] impl Consumer {
pub fn assert_not_closed(&self) {
assert!(self.closed().now_or_never().is_none(), "should not be closed");
}
pub fn assert_closed(&self) {
assert!(self.closed().now_or_never().is_some(), "should be closed");
}
}
#[cfg(test)]
mod test {
use super::*;
use std::time::Duration;
#[test]
fn unique_names_are_never_reused() {
let producer = Info::new().produce();
let name = producer.unique_name(".opus");
assert_eq!(name, "0.opus");
let track = producer.create_track(name.clone(), None).unwrap();
assert_eq!(producer.unique_name(".opus"), "1.opus");
drop(track);
}
#[test]
fn unique_names_survive_closed_track_pruning() {
let producer = Info::new().produce();
let consumer = producer.consume();
let track = producer.unique_track(".opus", None).unwrap();
assert_eq!(track.name(), "0.opus");
drop(track);
assert!(matches!(consumer.track_inner("0.opus"), Err(Error::NotFound)));
assert_eq!(producer.unique_name(".opus"), "1.opus");
}
#[test]
fn unique_name_skips_a_live_collision() {
let producer = Info::new().produce();
let track = producer.create_track("0.opus", None).unwrap();
assert_eq!(producer.unique_name(".opus"), "1.opus");
drop(track);
assert_eq!(producer.unique_name(".opus"), "2.opus");
}
#[test]
fn unique_names_share_a_counter() {
let producer = Info::new().produce();
assert_eq!(producer.unique_name("-video"), "0-video");
assert_eq!(producer.clone().unique_name("-audio"), "1-audio");
assert_eq!(producer.unique_name("-video"), "2-video");
}
#[test]
fn unique_names_separate_numeric_suffixes() {
let producer = Info::new().produce();
assert_eq!(producer.unique_name(""), "0");
let name = producer.unique_name("2");
assert_eq!(name, "1-2");
for _ in 2..12 {
producer.unique_name("");
}
assert_eq!(producer.unique_name(""), "12");
}
async fn expect<T>(fut: impl Future<Output = T>) -> T {
tokio::time::timeout(Duration::from_secs(1), fut)
.await
.expect("timed out waiting for a demand edge")
}
#[tokio::test]
async fn demand_ordinary() {
tokio::time::pause();
let producer = Info::new().produce();
let consumer = producer.consume();
let demand = producer.demand();
assert!(!demand.is_used());
demand.unused().await.unwrap();
let _track = producer.create_track("a", None).unwrap();
assert!(!demand.is_used());
let (used, handle) = tokio::join!(expect(demand.used()), async { consumer.track("a").unwrap() });
used.unwrap();
assert!(demand.is_used());
let (unused, ()) = tokio::join!(expect(demand.unused()), async { drop(handle) });
unused.unwrap();
assert!(!demand.is_used());
producer.close();
assert!(matches!(demand.used().await, Err(Error::Dropped)));
assert!(matches!(demand.unused().await, Err(Error::Dropped)));
}
#[tokio::test]
async fn demand_spliced() {
tokio::time::pause();
let producer = Producer::new_spliced(Info::new());
let consumer = producer.consume();
let demand = producer.demand();
let watched = consumer.demand();
assert!(!demand.is_used());
assert!(!watched.is_used());
let track = consumer.track("video").unwrap();
assert!(demand.is_used());
assert!(watched.is_used());
let (unused, ()) = tokio::join!(expect(watched.unused()), async { drop(track) });
unused.unwrap();
assert!(!demand.is_used());
assert!(!watched.is_used());
let _track = consumer.track("video").unwrap();
assert!(demand.is_used());
}
#[tokio::test]
async fn consumer_demand_reports_dropped_producer() {
let producer = Producer::new_spliced(Info::new());
let consumer = producer.consume();
let watched = consumer.demand();
let track = consumer.track("video").unwrap();
assert!(watched.is_used());
let (unused, ()) = tokio::join!(expect(watched.unused()), async { drop(track) });
unused.unwrap();
drop(producer);
assert!(matches!(watched.used().await, Err(Error::Dropped)));
assert!(matches!(watched.unused().await, Err(Error::Dropped)));
}
macro_rules! subscribe_pending {
($consumer:expr, $name:expr) => {{
let pending = $consumer.track($name).unwrap().subscribe(None);
assert!(
pending.poll_ok(&kio::Waiter::noop()).is_pending(),
"subscribe should stay pending until the request is accepted"
);
pending
}};
}
#[tokio::test]
async fn insert() {
let mut producer = Info::new().produce();
let track1 = producer.assert_create_track("track1", None);
track1.append_group().unwrap();
let consumer = producer.consume();
let mut track1_sub = consumer.track("track1").unwrap().subscribe(None).await.unwrap();
track1_sub.assert_group();
let track2 = producer.assert_create_track("track2", None);
let consumer2 = producer.consume();
let mut track2_consumer = consumer2.track("track2").unwrap().subscribe(None).await.unwrap();
track2_consumer.assert_no_group();
track2.append_group().unwrap();
track2_consumer.assert_group();
}
#[tokio::test]
async fn closed() {
let mut producer = Info::new().produce();
let dynamic = producer.dynamic();
let consumer = producer.consume();
consumer.assert_not_closed();
let track1 = producer.assert_create_track("track1", None);
let mut track1c = consumer.track("track1").unwrap().subscribe(None).await.unwrap();
let track2_fut = subscribe_pending!(consumer, "track2");
drop(dynamic);
assert!(track2_fut.await.is_err());
assert!(!track1.is_closed());
track1c.assert_not_closed();
}
#[tokio::test]
async fn close_ends_every_clone() {
let producer = Info::new().produce();
let clone = producer.clone();
let consumer = producer.consume();
producer.close();
assert!(matches!(consumer.closed().await, Error::Dropped));
assert!(matches!(consumer.track("video"), Err(Error::Unroutable)));
assert!(matches!(clone.consume().track("video"), Err(Error::Unroutable)));
clone.close();
producer.close();
}
#[tokio::test]
async fn drop_ends_like_close() {
let producer = Info::new().produce();
let consumer = producer.consume();
drop(producer);
assert!(matches!(consumer.closed().await, Error::Dropped));
assert!(matches!(consumer.track("video"), Err(Error::Unroutable)));
}
#[tokio::test]
#[allow(deprecated)]
async fn deprecated_end_causes() {
let producer = Info::new().produce();
let consumer = producer.consume();
producer.abort(Error::Timeout).unwrap();
assert!(matches!(consumer.closed().await, Error::Timeout));
assert!(!consumer.is_finished());
let producer = Info::new().produce();
let consumer = producer.consume();
producer.finish();
assert!(matches!(consumer.closed().await, Error::Dropped));
assert!(consumer.is_finished());
}
#[tokio::test]
async fn requests() {
let mut producer = Info::new().produce().dynamic();
let consumer = producer.consume();
let consumer2 = consumer.clone();
let track1_fut = subscribe_pending!(consumer, "track1");
let track2_fut = subscribe_pending!(consumer2, "track1");
let request = producer.assert_request();
producer.assert_no_request();
assert_eq!(request.name(), "track1");
let track3 = request.accept(None);
let mut track1 = track1_fut.await.unwrap();
let mut track2 = track2_fut.await.unwrap();
track1.assert_not_closed();
track1.assert_is_clone(&track2);
track3.subscribe(None).assert_is_clone(&track1);
track3.append_group().unwrap();
track1.assert_group();
track2.assert_group();
let track4_fut = subscribe_pending!(consumer, "track2");
drop(producer);
assert!(track4_fut.await.is_err());
let track5 = consumer2.track("track3");
assert!(track5.is_err(), "should have errored");
}
#[tokio::test]
async fn stale_producer() {
let mut broadcast = Info::new().produce().dynamic();
let consumer = broadcast.consume();
let track1_fut = subscribe_pending!(consumer, "track1");
let producer1 = broadcast.assert_request().accept(None);
let mut track1 = track1_fut.await.unwrap();
producer1.append_group().unwrap();
producer1.finish().unwrap();
drop(producer1);
track1.assert_closed();
let track2_fut = subscribe_pending!(consumer, "track1");
let producer2 = broadcast.assert_request().accept(None);
let mut track2 = track2_fut.await.unwrap();
track2.assert_not_closed();
track2.assert_not_clone(&track1);
producer2.append_group().unwrap();
track2.assert_group();
}
#[tokio::test(start_paused = true)]
async fn requested_unused() {
let mut broadcast = Info::new().produce().dynamic();
let bc = broadcast.consume();
let c1_fut = subscribe_pending!(bc, "unknown_track");
let producer1 = broadcast.assert_request().accept(None);
let consumer1 = c1_fut.await.unwrap();
assert!(
producer1.unused().now_or_never().is_none(),
"track producer should be used"
);
let consumer2 = bc.track("unknown_track").unwrap().subscribe(None).await.unwrap();
consumer2.assert_is_clone(&consumer1);
drop(consumer1);
assert!(
producer1.unused().now_or_never().is_none(),
"track producer should be used"
);
drop(consumer2);
assert!(
producer1.unused().now_or_never().is_some(),
"track producer should be unused after all consumers are dropped"
);
let consumer3 = bc.track("unknown_track").unwrap().subscribe(None).await.unwrap();
consumer3.assert_is_clone(&producer1.subscribe(None));
broadcast.assert_no_request();
drop(consumer3);
producer1.abort(Error::Cancel).unwrap();
let c4_fut = subscribe_pending!(bc, "unknown_track");
let producer2 = broadcast.assert_request().accept(None);
let consumer4 = c4_fut.await.unwrap();
drop(consumer4);
assert!(
producer2.unused().now_or_never().is_some(),
"new track producer should be unused after its consumer is dropped"
);
}
#[tokio::test]
async fn create_track_fulfills_queued_request() {
let producer = Info::new().produce();
let mut dynamic = producer.dynamic();
let bc = dynamic.consume();
let subscribing = subscribe_pending!(bc, "video");
let track = producer.create_track("video", None).unwrap();
let mut sub = subscribing.await.expect("fulfilled by create_track");
track.append_group().unwrap();
sub.recv_group().await.expect("recv").expect("group");
dynamic.assert_no_request();
let again = bc.track("video").unwrap().subscribe(None).await.unwrap();
again.assert_is_clone(&track.subscribe(None));
}
#[tokio::test]
async fn dynamic_clone_keeps_alive() {
let broadcast = Info::new().produce().dynamic();
let consumer = broadcast.consume();
let clone = broadcast.clone();
drop(clone);
let _fut = subscribe_pending!(consumer, "track1");
}
#[tokio::test]
async fn close_resolves_a_reserved_name() {
let producer = Info::new().produce();
let consumer = producer.consume();
let _request = producer.reserve_track("track1").unwrap();
let pending = subscribe_pending!(consumer, "track1");
producer.close();
assert!(matches!(pending.await, Err(Error::Unroutable)));
}
#[tokio::test]
#[allow(deprecated)]
async fn abort_resolves_a_reserved_name_with_its_reason() {
let producer = Info::new().produce();
let consumer = producer.consume();
let request = producer.reserve_track("track1").unwrap();
let pending = subscribe_pending!(consumer, "track1");
producer.abort(Error::Cancel).unwrap();
assert!(matches!(pending.await, Err(Error::Cancel)));
let track = request.accept(None);
let mut subscriber = track.subscribe(None);
assert!(matches!(subscriber.recv_group().await, Err(Error::Cancel)));
}
#[tokio::test]
async fn close_resolves_a_queued_request() {
let producer = Info::new().produce();
let dynamic = producer.dynamic();
let consumer = dynamic.consume();
let pending = subscribe_pending!(consumer, "track1");
producer.close();
assert!(matches!(pending.await, Err(Error::Unroutable)));
drop(dynamic);
}
#[tokio::test]
async fn dropping_the_last_handle_resolves_a_queued_request() {
let dynamic = Info::new().produce().dynamic();
let consumer = dynamic.consume();
let pending = subscribe_pending!(consumer, "track1");
drop(dynamic);
assert!(matches!(pending.await, Err(Error::Unroutable)));
}
#[tokio::test]
async fn dropping_the_last_handler_resolves_a_queued_request_dropped() {
let producer = Info::new().produce();
let dynamic = producer.dynamic();
let consumer = dynamic.consume();
let pending = subscribe_pending!(consumer, "track1");
drop(dynamic);
assert!(matches!(pending.await, Err(Error::Dropped)));
producer.close();
}
#[tokio::test]
async fn close_leaves_a_claimed_request_to_its_handler() {
let producer = Info::new().produce();
let mut dynamic = producer.dynamic();
let consumer = dynamic.consume();
let accepted = subscribe_pending!(consumer, "track1");
let request = dynamic.requested_track().await.unwrap();
let dropped = subscribe_pending!(consumer, "track2");
let abandoned = dynamic.requested_track().await.unwrap();
producer.close();
assert!(
accepted.poll_ok(&kio::Waiter::noop()).is_pending(),
"close rejected a claimed request"
);
assert!(
dropped.poll_ok(&kio::Waiter::noop()).is_pending(),
"close rejected a claimed request"
);
let _track = request.accept(None);
assert!(accepted.await.is_ok(), "the handler's accept reaches the consumer");
drop(abandoned);
assert!(dropped.await.is_err(), "the handler dropping it rejects the consumer");
drop(dynamic);
}
#[tokio::test]
async fn close_resolves_an_unaccepted_track_with_fetched_info() {
let producer = Info::new().produce();
let consumer = producer.consume();
let request = producer.reserve_track("track1").unwrap();
let dynamic = request.dynamic();
let track = consumer.track("track1").unwrap();
let pending_fetch = track.fetch_group(0, None);
let fetch = dynamic.requested_group().await.unwrap();
let group = fetch.accept(None).unwrap();
group.finish().unwrap();
pending_fetch.await.unwrap();
let mut subscriber = track.subscribe(None).await.unwrap();
producer.close();
assert!(matches!(subscriber.recv_group().await, Err(Error::Unroutable)));
let stale = request.accept(None);
assert!(stale.append_group().is_err());
}
#[tokio::test]
async fn close_spares_a_served_track() {
let producer = Info::new().produce();
let consumer = producer.consume();
let track = producer.create_track("track1", None).unwrap();
let mut subscriber = consumer.track("track1").unwrap().subscribe(None).await.unwrap();
producer.close();
assert!(matches!(consumer.track("track1"), Err(Error::Unroutable)));
track.append_group().unwrap();
subscriber.assert_group();
track.finish().unwrap();
}
#[tokio::test]
async fn close_leaves_a_stale_reservation_inert() {
let producer = Info::new().produce();
let consumer = producer.consume();
let request = producer.reserve_track("track1").unwrap();
let pending = subscribe_pending!(consumer, "track1");
producer.close();
assert!(matches!(pending.await, Err(Error::Unroutable)));
let track = request.accept(None);
assert!(track.append_group().is_err());
let mut subscriber = track.subscribe(None);
assert!(matches!(subscriber.recv_group().await, Err(Error::Unroutable)));
assert!(consumer.track("track1").is_err());
}
#[tokio::test]
async fn dropping_a_reserved_request_resolves_dropped() {
let producer = Info::new().produce();
let consumer = producer.consume();
let request = producer.reserve_track("track1").unwrap();
let pending = subscribe_pending!(consumer, "track1");
drop(request);
assert!(matches!(pending.await, Err(Error::Dropped)));
producer.close();
}
#[tokio::test]
async fn rejecting_a_reserved_request_carries_the_reason() {
let producer = Info::new().produce();
let consumer = producer.consume();
let request = producer.reserve_track("track1").unwrap();
let pending = subscribe_pending!(consumer, "track1");
request.reject(Error::NotFound);
assert!(matches!(pending.await, Err(Error::NotFound)));
producer.close();
}
#[tokio::test]
async fn an_idle_teardown_yields_to_a_returning_viewer() {
let producer = Info::new().produce();
let consumer = producer.consume();
let track = producer.create_track("video", None).unwrap();
assert!(track.poll_unused(&kio::Waiter::noop()).is_ready());
let viewer = consumer.track("video").unwrap();
let track = track
.abort_unused(Error::Cancel)
.expect_err("viewer keeps the track alive");
assert!(!track.is_closed());
let mut subscriber = viewer.subscribe(None).await.unwrap();
subscriber.assert_no_group();
track.append_group().unwrap();
assert!(subscriber.recv_group().await.unwrap().is_some());
drop(subscriber);
drop(viewer);
assert!(track.abort_unused(Error::Cancel).is_ok());
assert!(matches!(consumer.track("video"), Err(Error::NotFound)));
producer.close();
}
#[test]
fn abort_unused_accepts_an_already_closed_track_with_consumers() {
let producer = Info::new().produce();
let consumer = producer.consume();
let track = producer.create_track("video", None).unwrap();
let _viewer = consumer.track("video").unwrap();
assert!(track.is_used());
track.clone().abort(Error::Cancel).unwrap();
assert!(!track.is_used());
assert!(track.abort_unused(Error::Cancel).is_ok());
producer.close();
}
}