use anyhow::Context;
use hang::moq_net;
pub(super) async fn subscribe(origin: moq_net::origin::Consumer, broadcast: &str) -> anyhow::Result<moq_mux::Source> {
origin
.routed(broadcast)
.await
.with_context(|| format!("origin closed before broadcast `{broadcast}` was announced"))?;
Ok(moq_mux::Source::new(origin, broadcast))
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
#[tokio::test]
async fn subscribe_waits_for_the_announcement() {
tokio::time::pause();
let origin = moq_tokio::origin::spawn();
let consumer = origin.consume();
let unannounced = moq_mux::Source::new(consumer.clone(), "room.hang").broadcast().await;
assert!(unannounced.is_err(), "expected an unroutable broadcast");
let mut waiting = std::pin::pin!(subscribe(consumer, "room.hang"));
let parked = tokio::time::timeout(Duration::from_secs(60), &mut waiting).await;
assert!(parked.is_err(), "expected to still be waiting on the announcement");
let _broadcast = origin.create_broadcast("room.hang").unwrap();
_broadcast.announce(Default::default()).unwrap();
waiting.await.unwrap();
}
}