use moq_native::moq_net::{self, Origin};
use std::time::Duration;
const TIMEOUT: Duration = Duration::from_secs(10);
async fn broadcast_test(scheme: &str, client_version: Option<&str>, server_version: Option<&str>) {
let client_version: Option<moq_net::Version> = client_version.map(|v| v.parse().expect("invalid client version"));
let server_version: Option<moq_net::Version> = server_version.map(|v| v.parse().expect("invalid server version"));
let pub_origin = Origin::random().produce();
let mut broadcast = pub_origin
.create_broadcast("test", moq_net::broadcast::Route::new().with_announce(true))
.expect("failed to create broadcast");
let mut track = broadcast.create_track("video", None).expect("failed to create track");
let mut group = track.append_group().expect("failed to append group");
group
.write_frame(moq_native::moq_net::Timestamp::ZERO, b"hello".as_ref())
.expect("failed to write frame");
group.finish().expect("failed to finish group");
let mut server_config = moq_native::ServerConfig::default();
server_config.bind = Some("[::]:0".to_string());
server_config.tls.generate = vec!["localhost".into()];
if let Some(v) = server_version {
server_config.version = vec![v];
}
let mut server = server_config.init().expect("failed to init server");
let addr = server.local_addr().expect("failed to get local addr");
let sub_origin = Origin::random().produce();
let mut announcements = sub_origin.consume().announced();
let mut client_config = moq_native::ClientConfig::default();
client_config.tls.disable_verify = Some(true);
if let Some(v) = client_version {
client_config.version = vec![v];
}
let client = client_config.init().expect("failed to init client");
let url: url::Url = format!("{scheme}://localhost:{}", addr.port()).parse().unwrap();
let server_handle = tokio::spawn(async move {
let request = server.accept().await.expect("no incoming connection");
let session = request.with_publisher(&pub_origin).ok().await?;
let _broadcast = broadcast;
let _track = track;
let _ = session.closed().await;
Ok::<_, anyhow::Error>(())
});
let client = client.with_subscriber(sub_origin);
let session = tokio::time::timeout(TIMEOUT, client.connect(url))
.await
.expect("client connect timed out")
.expect("client connect failed");
let moq_native::moq_net::announce::Update { path, broadcast: bc } =
tokio::time::timeout(TIMEOUT, announcements.next())
.await
.expect("announce timed out")
.expect("origin closed");
assert_eq!(path.as_str(), "test");
let bc = bc.expect("expected announce, got unannounce");
let mut track_sub = bc
.track("video")
.unwrap()
.subscribe(None)
.await
.expect("consume_track failed");
let mut group_sub = tokio::time::timeout(TIMEOUT, track_sub.recv_group())
.await
.expect("recv_group timed out")
.expect("recv_group failed")
.expect("track closed prematurely");
let frame = tokio::time::timeout(TIMEOUT, group_sub.read_frame())
.await
.expect("read_frame timed out")
.expect("read_frame failed")
.expect("group closed prematurely");
assert_eq!(&frame.payload[..], b"hello");
drop(session);
server_handle
.await
.expect("server task panicked")
.expect("server task failed");
}
async fn lite05_timestamp_roundtrip(scheme: &str) {
use moq_native::moq_net::{Timescale, Timestamp};
let pub_origin = Origin::random().produce();
let mut broadcast = pub_origin
.create_broadcast("test", moq_net::broadcast::Route::new().with_announce(true))
.expect("failed to create broadcast");
let mut track = broadcast
.create_track(
"video",
moq_net::track::Info::default().with_timescale(Timescale::MICRO),
)
.expect("failed to create track");
let frames = [10_000u64, 30_000, 20_000];
let mut group = track.append_group().expect("failed to append group");
for &us in &frames {
let payload = format!("frame@{us}").into_bytes();
let frame = moq_native::moq_net::frame::Info {
size: payload.len() as u64,
timestamp: Timestamp::new(us, Timescale::MICRO).unwrap(),
};
let mut writer = group.create_frame(frame).expect("failed to create frame");
writer
.write(bytes::Bytes::from(payload))
.expect("failed to write frame");
writer.finish().expect("failed to finish frame");
}
group.finish().expect("failed to finish group");
let mut server_config = moq_native::ServerConfig::default();
server_config.bind = Some("[::]:0".to_string());
server_config.tls.generate = vec!["localhost".into()];
server_config.version = vec!["moq-lite-05".parse().unwrap()];
let mut server = server_config.init().expect("failed to init server");
let addr = server.local_addr().expect("failed to get local addr");
let sub_origin = Origin::random().produce();
let mut announcements = sub_origin.consume().announced();
let mut client_config = moq_native::ClientConfig::default();
client_config.tls.disable_verify = Some(true);
client_config.version = vec!["moq-lite-05".parse().unwrap()];
let client = client_config.init().expect("failed to init client");
let url: url::Url = format!("{scheme}://localhost:{}", addr.port()).parse().unwrap();
let server_handle = tokio::spawn(async move {
let request = server.accept().await.expect("no incoming connection");
let session = request.with_publisher(&pub_origin).ok().await?;
let _broadcast = broadcast;
let _track = track;
let _ = session.closed().await;
Ok::<_, anyhow::Error>(())
});
let client = client.with_subscriber(sub_origin);
let session = tokio::time::timeout(TIMEOUT, client.connect(url))
.await
.expect("client connect timed out")
.expect("client connect failed");
let moq_native::moq_net::announce::Update { path, broadcast: bc } =
tokio::time::timeout(TIMEOUT, announcements.next())
.await
.expect("announce timed out")
.expect("origin closed");
assert_eq!(path.as_str(), "test");
let bc = bc.expect("expected announce, got unannounce");
let mut track_sub = bc
.track("video")
.unwrap()
.subscribe(None)
.await
.expect("consume_track failed");
let mut group_sub = tokio::time::timeout(TIMEOUT, track_sub.recv_group())
.await
.expect("recv_group timed out")
.expect("recv_group failed")
.expect("track closed prematurely");
for &expected_us in &frames {
let mut frame_sub = tokio::time::timeout(TIMEOUT, group_sub.next_frame())
.await
.expect("next_frame timed out")
.expect("next_frame failed")
.expect("group closed prematurely");
let ts = frame_sub.timestamp;
assert_eq!(ts.scale(), Timescale::MICRO);
assert_eq!(ts.value(), expected_us);
let _ = frame_sub.read_all().await;
}
drop(session);
server_handle
.await
.expect("server task panicked")
.expect("server task failed");
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_moq_lite_05_timestamps_webtransport() {
lite05_timestamp_roundtrip("https").await;
}
async fn lite05_fetch_roundtrip(scheme: &str) {
use moq_native::moq_net::{Timescale, Timestamp};
let pub_origin = Origin::random().produce();
let mut broadcast = pub_origin
.create_broadcast("test", moq_net::broadcast::Route::new().with_announce(true))
.expect("failed to create broadcast");
let mut track = broadcast
.create_track(
"video",
moq_net::track::Info::default().with_timescale(Timescale::MICRO),
)
.expect("failed to create track");
let frames = [10_000u64, 30_000, 20_000];
let mut group = track.append_group().expect("failed to append group"); for &us in &frames {
let payload = format!("frame@{us}").into_bytes();
let frame = moq_native::moq_net::frame::Info {
size: payload.len() as u64,
timestamp: Timestamp::new(us, Timescale::MICRO).unwrap(),
};
let mut writer = group.create_frame(frame).expect("failed to create frame");
writer
.write(bytes::Bytes::from(payload))
.expect("failed to write frame");
writer.finish().expect("failed to finish frame");
}
group.finish().expect("failed to finish group");
let mut server_config = moq_native::ServerConfig::default();
server_config.bind = Some("[::]:0".to_string());
server_config.tls.generate = vec!["localhost".into()];
server_config.version = vec!["moq-lite-05".parse().unwrap()];
let mut server = server_config.init().expect("failed to init server");
let addr = server.local_addr().expect("failed to get local addr");
let sub_origin = Origin::random().produce();
let mut announcements = sub_origin.consume().announced();
let mut client_config = moq_native::ClientConfig::default();
client_config.tls.disable_verify = Some(true);
client_config.version = vec!["moq-lite-05".parse().unwrap()];
let client = client_config.init().expect("failed to init client");
let url: url::Url = format!("{scheme}://localhost:{}", addr.port()).parse().unwrap();
let server_handle = tokio::spawn(async move {
let request = server.accept().await.expect("no incoming connection");
let session = request.with_publisher(&pub_origin).ok().await?;
let _broadcast = broadcast;
let _track = track;
let _ = session.closed().await;
Ok::<_, anyhow::Error>(())
});
let client = client.with_subscriber(sub_origin);
let session = tokio::time::timeout(TIMEOUT, client.connect(url))
.await
.expect("client connect timed out")
.expect("client connect failed");
let moq_native::moq_net::announce::Update { path, broadcast: bc } =
tokio::time::timeout(TIMEOUT, announcements.next())
.await
.expect("announce timed out")
.expect("origin closed");
assert_eq!(path.as_str(), "test");
let bc = bc.expect("expected announce, got unannounce");
let mut group_sub = tokio::time::timeout(TIMEOUT, async { bc.track("video").unwrap().fetch_group(0, None).await })
.await
.expect("fetch timed out")
.expect("fetch failed");
assert_eq!(group_sub.sequence, 0);
for &expected_us in &frames {
let mut frame_sub = tokio::time::timeout(TIMEOUT, group_sub.next_frame())
.await
.expect("next_frame timed out")
.expect("next_frame failed")
.expect("group closed prematurely");
let ts = frame_sub.timestamp;
assert_eq!(ts.scale(), Timescale::MICRO);
assert_eq!(ts.value(), expected_us);
let payload = frame_sub.read_all().await.expect("failed to read frame");
assert_eq!(payload, bytes::Bytes::from(format!("frame@{expected_us}")));
}
let end = tokio::time::timeout(TIMEOUT, group_sub.next_frame())
.await
.expect("next_frame timed out")
.expect("next_frame failed");
assert!(end.is_none(), "group should finish after its frames");
drop(session);
server_handle
.await
.expect("server task panicked")
.expect("server task failed");
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_moq_lite_05_fetch_webtransport() {
lite05_fetch_roundtrip("https").await;
}
async fn lite05_fetch_during_subscribe(scheme: &str) {
use moq_native::moq_net::{Timescale, Timestamp};
fn timestamped_frame(us: u64, payload: &str) -> moq_net::frame::Info {
moq_net::frame::Info {
size: payload.len() as u64,
timestamp: Timestamp::new(us, Timescale::MICRO).unwrap(),
}
}
let pub_origin = Origin::random().produce();
let mut broadcast = pub_origin
.create_broadcast("test", moq_net::broadcast::Route::new().with_announce(true))
.expect("failed to create broadcast");
let mut track = broadcast
.create_track(
"video",
moq_net::track::Info::default().with_timescale(Timescale::MICRO),
)
.expect("failed to create track");
let mut group0 = track.append_group().expect("append group 0"); let mut w = group0.create_frame(timestamped_frame(10_000, "old")).expect("frame 0");
w.write(bytes::Bytes::from_static(b"old")).expect("write 0");
w.finish().expect("finish frame 0");
group0.finish().expect("finish group 0");
let mut group1 = track.append_group().expect("append group 1"); let mut w = group1.create_frame(timestamped_frame(20_000, "new")).expect("frame 1");
w.write(bytes::Bytes::from_static(b"new")).expect("write 1");
w.finish().expect("finish frame 1");
group1.finish().expect("finish group 1");
let mut server_config = moq_native::ServerConfig::default();
server_config.bind = Some("[::]:0".to_string());
server_config.tls.generate = vec!["localhost".into()];
server_config.version = vec!["moq-lite-05".parse().unwrap()];
let mut server = server_config.init().expect("failed to init server");
let addr = server.local_addr().expect("failed to get local addr");
let sub_origin = Origin::random().produce();
let mut announcements = sub_origin.consume().announced();
let mut client_config = moq_native::ClientConfig::default();
client_config.tls.disable_verify = Some(true);
client_config.version = vec!["moq-lite-05".parse().unwrap()];
let client = client_config.init().expect("failed to init client");
let url: url::Url = format!("{scheme}://localhost:{}", addr.port()).parse().unwrap();
let server_handle = tokio::spawn(async move {
let request = server.accept().await.expect("no incoming connection");
let session = request.with_publisher(&pub_origin).ok().await?;
let _broadcast = broadcast;
let _track = track;
let _ = session.closed().await;
Ok::<_, anyhow::Error>(())
});
let client = client.with_subscriber(sub_origin);
let session = tokio::time::timeout(TIMEOUT, client.connect(url))
.await
.expect("client connect timed out")
.expect("client connect failed");
let moq_native::moq_net::announce::Update { path, broadcast: bc } =
tokio::time::timeout(TIMEOUT, announcements.next())
.await
.expect("announce timed out")
.expect("origin closed");
assert_eq!(path.as_str(), "test");
let bc = bc.expect("expected announce, got unannounce");
let mut track_sub = tokio::time::timeout(TIMEOUT, async { bc.track("video").unwrap().subscribe(None).await })
.await
.expect("subscribe timed out")
.expect("subscribe failed");
let mut live = tokio::time::timeout(TIMEOUT, track_sub.recv_group())
.await
.expect("recv_group timed out")
.expect("recv_group failed")
.expect("track closed prematurely");
assert_eq!(live.sequence, 1);
let frame = tokio::time::timeout(TIMEOUT, live.read_frame())
.await
.expect("read_frame timed out")
.expect("read_frame failed")
.expect("group closed prematurely");
assert_eq!(&frame.payload[..], b"new");
let mut fetched = tokio::time::timeout(TIMEOUT, async { bc.track("video").unwrap().fetch_group(0, None).await })
.await
.expect("fetch timed out")
.expect("fetch failed");
assert_eq!(fetched.sequence, 0);
let frame = tokio::time::timeout(TIMEOUT, fetched.read_frame())
.await
.expect("fetch read_frame timed out")
.expect("fetch read_frame failed")
.expect("fetched group closed prematurely");
assert_eq!(&frame.payload[..], b"old");
drop(session);
server_handle
.await
.expect("server task panicked")
.expect("server task failed");
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_moq_lite_05_fetch_during_subscribe_webtransport() {
lite05_fetch_during_subscribe("https").await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_moq_lite_05_default_timescale() {
use moq_native::moq_net::Timescale;
let pub_origin = Origin::random().produce();
let mut broadcast = pub_origin
.create_broadcast("test", moq_net::broadcast::Route::new().with_announce(true))
.expect("create broadcast");
let mut track = broadcast.create_track("video", None).expect("create track");
let mut group = track.append_group().expect("append group");
group
.write_frame(moq_native::moq_net::Timestamp::ZERO, b"hello".as_ref())
.expect("write frame");
group.finish().expect("finish group");
let mut server_config = moq_native::ServerConfig::default();
server_config.bind = Some("[::]:0".to_string());
server_config.tls.generate = vec!["localhost".into()];
server_config.version = vec!["moq-lite-05".parse().unwrap()];
let mut server = server_config.init().expect("init server");
let addr = server.local_addr().expect("local addr");
let sub_origin = Origin::random().produce();
let mut announcements = sub_origin.consume().announced();
let mut client_config = moq_native::ClientConfig::default();
client_config.tls.disable_verify = Some(true);
client_config.version = vec!["moq-lite-05".parse().unwrap()];
let client = client_config.init().expect("init client");
let url: url::Url = format!("https://localhost:{}", addr.port()).parse().unwrap();
let server_handle = tokio::spawn(async move {
let request = server.accept().await.expect("accept");
let session = request.with_publisher(&pub_origin).ok().await?;
let _broadcast = broadcast;
let _track = track;
let _ = session.closed().await;
Ok::<_, anyhow::Error>(())
});
let client = client.with_subscriber(sub_origin);
let session = tokio::time::timeout(TIMEOUT, client.connect(url))
.await
.expect("connect timeout")
.expect("connect failed");
let moq_native::moq_net::announce::Update { broadcast: bc, .. } =
tokio::time::timeout(TIMEOUT, announcements.next())
.await
.expect("announce timeout")
.expect("origin closed");
let bc = bc.expect("expected announce");
let mut track_sub = bc
.track("video")
.unwrap()
.subscribe(None)
.await
.expect("consume_track failed");
let mut group_sub = tokio::time::timeout(TIMEOUT, track_sub.recv_group())
.await
.expect("recv_group timeout")
.expect("recv_group failed")
.expect("track closed");
let frame_sub = tokio::time::timeout(TIMEOUT, group_sub.next_frame())
.await
.expect("next_frame timeout")
.expect("next_frame failed")
.expect("group closed");
let ts = frame_sub.timestamp;
assert_eq!(ts.scale(), Timescale::MILLI, "default timescale is milliseconds");
drop(session);
server_handle
.await
.expect("server task panicked")
.expect("server task failed");
}
async fn next_announce(announcements: &mut moq_net::announce::Consumer) -> moq_net::announce::Update {
tokio::time::timeout(TIMEOUT, announcements.next())
.await
.expect("announce timeout")
.expect("origin closed")
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_moq_lite_06_announce_lifecycle() {
let pub_origin = Origin::random().produce();
let mut first = pub_origin
.create_broadcast("first", moq_net::broadcast::Route::new().with_announce(true))
.expect("create broadcast");
let mut server_config = moq_native::ServerConfig::default();
server_config.bind = Some("[::]:0".to_string());
server_config.tls.generate = vec!["localhost".into()];
server_config.version = vec!["moq-lite-06-wip".parse().unwrap()];
let mut server = server_config.init().expect("init server");
let addr = server.local_addr().expect("local addr");
let sub_origin = Origin::random().produce();
let mut announcements = sub_origin.consume().announced();
let mut client_config = moq_native::ClientConfig::default();
client_config.tls.disable_verify = Some(true);
client_config.version = vec!["moq-lite-06-wip".parse().unwrap()];
let client = client_config.init().expect("init client");
let url: url::Url = format!("moqt://localhost:{}", addr.port()).parse().unwrap();
let server_origin = pub_origin.clone();
let server_handle = tokio::spawn(async move {
let request = server.accept().await.expect("accept");
let session = request.with_publisher(&server_origin).ok().await?;
let _ = session.closed().await;
Ok::<_, anyhow::Error>(())
});
let client = client.with_subscriber(sub_origin);
let session = tokio::time::timeout(TIMEOUT, client.connect(url))
.await
.expect("connect timeout")
.expect("connect failed");
let moq_net::announce::Update { path, broadcast } = next_announce(&mut announcements).await;
assert_eq!(path.as_str(), "first");
assert!(broadcast.is_some(), "expected initial announce");
let mut second = pub_origin
.create_broadcast("second", moq_net::broadcast::Route::new().with_announce(true))
.expect("create broadcast");
let moq_net::announce::Update { path, broadcast } = next_announce(&mut announcements).await;
assert_eq!(path.as_str(), "second");
assert!(broadcast.is_some(), "expected live announce");
second.finish();
let moq_net::announce::Update { path, broadcast } = next_announce(&mut announcements).await;
assert_eq!(path.as_str(), "second");
assert!(broadcast.is_none(), "expected unannounce");
let _second = pub_origin
.create_broadcast("second", moq_net::broadcast::Route::new().with_announce(true))
.expect("create broadcast");
let moq_net::announce::Update { path, broadcast } = next_announce(&mut announcements).await;
assert_eq!(path.as_str(), "second");
assert!(broadcast.is_some(), "expected re-announce");
first.finish();
let moq_net::announce::Update { path, broadcast } = next_announce(&mut announcements).await;
assert_eq!(path.as_str(), "first");
assert!(broadcast.is_none(), "expected the replaced unannounce");
let _replacement = pub_origin
.create_broadcast("first", moq_net::broadcast::Route::new().with_announce(true))
.expect("create replacement");
let moq_net::announce::Update { path, broadcast } = next_announce(&mut announcements).await;
assert_eq!(path.as_str(), "first");
assert!(broadcast.is_some(), "expected the replacement announce");
let _sentinel = pub_origin
.create_broadcast("sentinel", moq_net::broadcast::Route::new().with_announce(true))
.expect("create broadcast");
let moq_net::announce::Update { path, broadcast } = next_announce(&mut announcements).await;
assert_eq!(path.as_str(), "sentinel");
assert!(broadcast.is_some(), "expected sentinel announce");
drop(session);
server_handle
.await
.expect("server task panicked")
.expect("server task failed");
}
async fn read_payloads(sub: &mut moq_net::track::Subscriber, count: usize) -> Vec<String> {
let mut payloads = Vec::new();
for _ in 0..count {
let mut group = tokio::time::timeout(TIMEOUT, sub.recv_group())
.await
.expect("recv_group timeout")
.expect("recv_group failed")
.expect("track closed");
let frame = tokio::time::timeout(TIMEOUT, group.read_frame())
.await
.expect("read_frame timeout")
.expect("read_frame failed")
.expect("group empty");
payloads.push(String::from_utf8(frame.payload.to_vec()).unwrap());
}
payloads.sort();
payloads
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_route_migration() {
use moq_net::Timestamp;
let publisher = Origin::new(0x42).unwrap();
let origin_a = Origin::random().produce();
let mut hops_a = moq_net::OriginList::new();
hops_a.push(publisher).unwrap();
let mut broadcast_a = origin_a
.create_broadcast(
"test",
moq_net::broadcast::Route::new().with_hops(hops_a).with_announce(true),
)
.expect("create broadcast");
let mut track_a = broadcast_a.create_track("video", None).expect("create track");
for sequence in 0..2u64 {
let mut group = track_a
.create_group(moq_net::group::Info { sequence })
.expect("create group");
group
.write_frame(Timestamp::ZERO, format!("a{sequence}").into_bytes())
.expect("write frame");
group.finish().expect("finish group");
}
let origin_b = Origin::random().produce();
let mut hops_b = moq_net::OriginList::new();
hops_b.push(publisher).unwrap();
hops_b.push(Origin::new(0x1234).unwrap()).unwrap();
let mut broadcast_b = origin_b
.create_broadcast(
"test",
moq_net::broadcast::Route::new().with_hops(hops_b).with_announce(true),
)
.expect("create broadcast");
let mut track_b = broadcast_b.create_track("video", None).expect("create track");
for sequence in 2..4u64 {
let mut group = track_b
.create_group(moq_net::group::Info { sequence })
.expect("create group");
group
.write_frame(Timestamp::ZERO, format!("b{sequence}").into_bytes())
.expect("write frame");
group.finish().expect("finish group");
}
let mut server_a = {
let mut config = moq_native::ServerConfig::default();
config.bind = Some("[::]:0".to_string());
config.tls.generate = vec!["localhost".into()];
config.init().expect("init server a")
};
let mut server_b = {
let mut config = moq_native::ServerConfig::default();
config.bind = Some("[::]:0".to_string());
config.tls.generate = vec!["localhost".into()];
config.init().expect("init server b")
};
let addr_a = server_a.local_addr().expect("local addr");
let addr_b = server_b.local_addr().expect("local addr");
let handle_a = tokio::spawn(async move {
let request = server_a.accept().await.expect("accept");
let session = request.with_publisher(&origin_a).ok().await?;
let _broadcast = broadcast_a;
let _track = track_a;
let _ = session.closed().await;
Ok::<_, anyhow::Error>(())
});
let handle_b = tokio::spawn(async move {
let request = server_b.accept().await.expect("accept");
let session = request.with_publisher(&origin_b).ok().await?;
let _broadcast = broadcast_b;
let _track = track_b;
let _ = session.closed().await;
Ok::<_, anyhow::Error>(())
});
let sub_origin = Origin::random().produce();
let mut announcements = sub_origin.consume().announced();
let connect = |port: u16, sub: moq_net::origin::Producer| {
let mut config = moq_native::ClientConfig::default();
config.tls.disable_verify = Some(true);
let client = config.init().expect("init client");
let url: url::Url = format!("moqt://localhost:{port}").parse().unwrap();
async move {
tokio::time::timeout(TIMEOUT, client.with_subscriber(sub).connect(url))
.await
.expect("connect timeout")
.expect("connect failed")
}
};
let session_a = connect(addr_a.port(), sub_origin.clone()).await;
let _session_b = connect(addr_b.port(), sub_origin.clone()).await;
let moq_net::announce::Update { path, broadcast } = next_announce(&mut announcements).await;
assert_eq!(path.as_str(), "test");
let broadcast = broadcast.expect("expected announce");
let subscription = moq_net::track::Subscription::default().with_latency_max(Duration::from_secs(10));
let mut sub = broadcast
.track("video")
.unwrap()
.subscribe(subscription)
.await
.expect("subscribe failed");
assert_eq!(read_payloads(&mut sub, 1).await, ["a1"]);
drop(session_a);
assert_eq!(read_payloads(&mut sub, 2).await, ["b2", "b3"]);
assert!(
announcements.try_next().is_none(),
"route migration must not emit announce events"
);
handle_a.await.expect("server a panicked").expect("server a failed");
drop(_session_b);
handle_b.await.expect("server b panicked").expect("server b failed");
}
async fn route_reannounce_test(version: Option<&str>) {
use moq_net::Timestamp;
let version: Option<moq_net::Version> = version.map(|v| v.parse().expect("invalid version"));
let origin = Origin::random().produce();
let publisher_hop = Origin::new(0x4444).unwrap();
let mut initial_hops = moq_net::OriginList::new();
initial_hops.push(publisher_hop).unwrap();
let mut producer = origin
.create_broadcast(
"test",
moq_net::broadcast::Route::new()
.with_hops(initial_hops)
.with_announce(true),
)
.expect("create broadcast");
let mut track = producer.create_track("video", None).expect("create track");
{
let mut group = track
.create_group(moq_net::group::Info { sequence: 0 })
.expect("create group");
group.write_frame(Timestamp::ZERO, b"g0".as_ref()).expect("write frame");
group.finish().expect("finish group");
}
let mut server_config = moq_native::ServerConfig::default();
server_config.bind = Some("[::]:0".to_string());
server_config.tls.generate = vec!["localhost".into()];
if let Some(v) = version {
server_config.version = vec![v];
}
let mut server = server_config.init().expect("init server");
let addr = server.local_addr().expect("local addr");
let mut route_producer = producer.clone();
let handle = tokio::spawn(async move {
let request = server.accept().await.expect("accept");
let session = request.with_publisher(&origin).ok().await?;
let _producer = producer;
let _ = session.closed().await;
Ok::<_, anyhow::Error>(())
});
let sub_origin = Origin::random().produce();
let mut announcements = sub_origin.consume().announced();
let mut client_config = moq_native::ClientConfig::default();
client_config.tls.disable_verify = Some(true);
if let Some(v) = version {
client_config.version = vec![v];
}
let client = client_config.init().expect("init client");
let url: url::Url = format!("moqt://localhost:{}", addr.port()).parse().unwrap();
let session = tokio::time::timeout(TIMEOUT, client.with_subscriber(sub_origin).connect(url))
.await
.expect("connect timeout")
.expect("connect failed");
let moq_net::announce::Update { path, broadcast } = next_announce(&mut announcements).await;
assert_eq!(path.as_str(), "test");
let broadcast = broadcast.expect("expected announce");
let mut watch = broadcast.clone();
let initial = tokio::time::timeout(TIMEOUT, watch.route_changed())
.await
.expect("route timeout")
.expect("route dropped");
let mut sub = broadcast
.track("video")
.unwrap()
.subscribe(None)
.await
.expect("subscribe failed");
assert_eq!(read_payloads(&mut sub, 1).await, ["g0"]);
let mut hops = moq_net::OriginList::new();
hops.push(publisher_hop).unwrap();
hops.push(Origin::new(0x5555).unwrap()).unwrap();
route_producer
.set_route(moq_net::broadcast::Route::new().with_hops(hops).with_announce(true))
.expect("update route");
let updated = tokio::time::timeout(TIMEOUT, watch.route_changed())
.await
.expect("route update timeout")
.expect("route dropped");
assert_ne!(initial, updated, "route must change");
assert!(
updated.hops.iter().any(|h| h.id() == 0x5555),
"the new chain must carry the added hop"
);
assert!(
announcements.try_next().is_none(),
"a route change must not emit announce events"
);
{
let mut group = track
.create_group(moq_net::group::Info { sequence: 1 })
.expect("create group");
group.write_frame(Timestamp::ZERO, b"g1".as_ref()).expect("write frame");
group.finish().expect("finish group");
}
assert_eq!(read_payloads(&mut sub, 1).await, ["g1"]);
drop(session);
handle.await.expect("server panicked").expect("server failed");
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_route_reannounce() {
route_reannounce_test(None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_route_reannounce_lite_06() {
route_reannounce_test(Some("moq-lite-06-wip")).await;
}
async fn route_replaced_test(version: Option<&str>) {
use moq_net::Timestamp;
let version: Option<moq_net::Version> = version.map(|v| v.parse().expect("invalid version"));
let origin = Origin::random().produce();
let mut hops_a = moq_net::OriginList::new();
hops_a.push(Origin::new(0x1111).unwrap()).unwrap();
let mut producer = origin
.create_broadcast(
"test",
moq_net::broadcast::Route::new().with_hops(hops_a).with_announce(true),
)
.expect("create broadcast");
let mut track = producer.create_track("video", None).expect("create track");
{
let mut group = track
.create_group(moq_net::group::Info { sequence: 0 })
.expect("create group");
group.write_frame(Timestamp::ZERO, b"g0".as_ref()).expect("write frame");
group.finish().expect("finish group");
}
let mut server_config = moq_native::ServerConfig::default();
server_config.bind = Some("[::]:0".to_string());
server_config.tls.generate = vec!["localhost".into()];
if let Some(v) = version {
server_config.version = vec![v];
}
let mut server = server_config.init().expect("init server");
let addr = server.local_addr().expect("local addr");
let mut route_producer = producer.clone();
let handle = tokio::spawn(async move {
let request = server.accept().await.expect("accept");
let session = request.with_publisher(&origin).ok().await?;
let _producer = producer;
let _ = session.closed().await;
Ok::<_, anyhow::Error>(())
});
let sub_origin = Origin::random().produce();
let mut announcements = sub_origin.consume().announced();
let mut client_config = moq_native::ClientConfig::default();
client_config.tls.disable_verify = Some(true);
if let Some(v) = version {
client_config.version = vec![v];
}
let client = client_config.init().expect("init client");
let url: url::Url = format!("moqt://localhost:{}", addr.port()).parse().unwrap();
let session = tokio::time::timeout(TIMEOUT, client.with_subscriber(sub_origin).connect(url))
.await
.expect("connect timeout")
.expect("connect failed");
let moq_net::announce::Update { path, broadcast } = next_announce(&mut announcements).await;
assert_eq!(path.as_str(), "test");
let broadcast = broadcast.expect("expected announce");
let mut sub = broadcast
.track("video")
.unwrap()
.subscribe(moq_net::track::Subscription::default().with_latency_max(Duration::from_secs(10)))
.await
.expect("subscribe failed");
assert_eq!(read_payloads(&mut sub, 1).await, ["g0"]);
let mut hops_b = moq_net::OriginList::new();
hops_b.push(Origin::new(0x2222).unwrap()).unwrap();
route_producer
.set_route(moq_net::broadcast::Route::new().with_hops(hops_b).with_announce(true))
.expect("update route");
let moq_net::announce::Update { path, broadcast } = next_announce(&mut announcements).await;
assert_eq!(path.as_str(), "test");
assert!(broadcast.is_none(), "expected the old broadcast to end");
let moq_net::announce::Update { path, broadcast } = next_announce(&mut announcements).await;
assert_eq!(path.as_str(), "test");
let replacement = broadcast.expect("expected the replacement announce");
let mut sub = replacement
.track("video")
.unwrap()
.subscribe(moq_net::track::Subscription::default().with_latency_max(Duration::from_secs(10)))
.await
.expect("subscribe to the replacement failed");
assert_eq!(read_payloads(&mut sub, 1).await, ["g0"]);
drop(session);
handle.await.expect("server panicked").expect("server failed");
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_route_replaced() {
route_replaced_test(None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_route_replaced_lite_06() {
route_replaced_test(Some("moq-lite-06-wip")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_moq_lite_01() {
broadcast_test("moqt", Some("moq-lite-01"), Some("moq-lite-01")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_moq_lite_02() {
broadcast_test("moqt", Some("moq-lite-02"), Some("moq-lite-02")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_moq_lite_03() {
broadcast_test("moqt", Some("moq-lite-03"), Some("moq-lite-03")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_moq_lite_06() {
broadcast_test("moqt", Some("moq-lite-06-wip"), Some("moq-lite-06-wip")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_moq_transport_14() {
broadcast_test("moqt", Some("moq-transport-14"), Some("moq-transport-14")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_moq_transport_15() {
broadcast_test("moqt", Some("moq-transport-15"), Some("moq-transport-15")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_moq_transport_16() {
broadcast_test("moqt", Some("moq-transport-16"), Some("moq-transport-16")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_moq_transport_17() {
broadcast_test("moqt", Some("moq-transport-17"), Some("moq-transport-17")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_moq_transport_18() {
broadcast_test("moqt", Some("moq-transport-18"), Some("moq-transport-18")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_moq_transport_19() {
broadcast_test("moqt", Some("moq-transport-19"), Some("moq-transport-19")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_negotiate_server_all_client_lite_01() {
broadcast_test("moqt", Some("moq-lite-01"), None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_negotiate_server_all_client_lite_02() {
broadcast_test("moqt", Some("moq-lite-02"), None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_negotiate_server_all_client_lite_03() {
broadcast_test("moqt", Some("moq-lite-03"), None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_negotiate_server_all_client_transport_14() {
broadcast_test("moqt", Some("moq-transport-14"), None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_negotiate_server_all_client_transport_15() {
broadcast_test("moqt", Some("moq-transport-15"), None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_negotiate_server_all_client_transport_16() {
broadcast_test("moqt", Some("moq-transport-16"), None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_negotiate_server_all_client_transport_17() {
broadcast_test("moqt", Some("moq-transport-17"), None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_negotiate_server_all_client_transport_18() {
broadcast_test("moqt", Some("moq-transport-18"), None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_negotiate_server_all_client_transport_19() {
broadcast_test("moqt", Some("moq-transport-19"), None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_negotiate_client_all_server_lite_01() {
broadcast_test("moqt", None, Some("moq-lite-01")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_negotiate_client_all_server_lite_02() {
broadcast_test("moqt", None, Some("moq-lite-02")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_negotiate_client_all_server_lite_03() {
broadcast_test("moqt", None, Some("moq-lite-03")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_negotiate_client_all_server_transport_14() {
broadcast_test("moqt", None, Some("moq-transport-14")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_negotiate_client_all_server_transport_15() {
broadcast_test("moqt", None, Some("moq-transport-15")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_negotiate_client_all_server_transport_16() {
broadcast_test("moqt", None, Some("moq-transport-16")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_negotiate_client_all_server_transport_17() {
broadcast_test("moqt", None, Some("moq-transport-17")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_negotiate_client_all_server_transport_18() {
broadcast_test("moqt", None, Some("moq-transport-18")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_negotiate_client_all_server_transport_19() {
broadcast_test("moqt", None, Some("moq-transport-19")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport() {
broadcast_test("https", None, None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_moq_lite_01() {
broadcast_test("https", Some("moq-lite-01"), Some("moq-lite-01")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_moq_lite_02() {
broadcast_test("https", Some("moq-lite-02"), Some("moq-lite-02")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_moq_lite_03() {
broadcast_test("https", Some("moq-lite-03"), Some("moq-lite-03")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_moq_transport_14() {
broadcast_test("https", Some("moq-transport-14"), Some("moq-transport-14")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_moq_transport_15() {
broadcast_test("https", Some("moq-transport-15"), Some("moq-transport-15")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_moq_transport_16() {
broadcast_test("https", Some("moq-transport-16"), Some("moq-transport-16")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_moq_transport_17() {
broadcast_test("https", Some("moq-transport-17"), Some("moq-transport-17")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_moq_transport_18() {
broadcast_test("https", Some("moq-transport-18"), Some("moq-transport-18")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_moq_transport_19() {
broadcast_test("https", Some("moq-transport-19"), Some("moq-transport-19")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_negotiate_server_all_client_lite_01() {
broadcast_test("https", Some("moq-lite-01"), None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_negotiate_server_all_client_lite_02() {
broadcast_test("https", Some("moq-lite-02"), None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_negotiate_server_all_client_lite_03() {
broadcast_test("https", Some("moq-lite-03"), None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_negotiate_server_all_client_transport_14() {
broadcast_test("https", Some("moq-transport-14"), None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_negotiate_server_all_client_transport_15() {
broadcast_test("https", Some("moq-transport-15"), None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_negotiate_server_all_client_transport_16() {
broadcast_test("https", Some("moq-transport-16"), None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_negotiate_server_all_client_transport_17() {
broadcast_test("https", Some("moq-transport-17"), None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_negotiate_server_all_client_transport_18() {
broadcast_test("https", Some("moq-transport-18"), None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_negotiate_server_all_client_transport_19() {
broadcast_test("https", Some("moq-transport-19"), None).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_negotiate_client_all_server_lite_01() {
broadcast_test("https", None, Some("moq-lite-01")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_negotiate_client_all_server_lite_02() {
broadcast_test("https", None, Some("moq-lite-02")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_negotiate_client_all_server_lite_03() {
broadcast_test("https", None, Some("moq-lite-03")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_negotiate_client_all_server_transport_14() {
broadcast_test("https", None, Some("moq-transport-14")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_negotiate_client_all_server_transport_15() {
broadcast_test("https", None, Some("moq-transport-15")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_negotiate_client_all_server_transport_16() {
broadcast_test("https", None, Some("moq-transport-16")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_negotiate_client_all_server_transport_17() {
broadcast_test("https", None, Some("moq-transport-17")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_negotiate_client_all_server_transport_18() {
broadcast_test("https", None, Some("moq-transport-18")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_webtransport_negotiate_client_all_server_transport_19() {
broadcast_test("https", None, Some("moq-transport-19")).await;
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_websocket() {
use moq_native::moq_net::Origin;
let pub_origin = Origin::random().produce();
let mut broadcast = pub_origin
.create_broadcast("test", moq_net::broadcast::Route::new().with_announce(true))
.expect("failed to create broadcast");
let mut track = broadcast.create_track("video", None).expect("failed to create track");
let mut group = track.append_group().expect("failed to append group");
group
.write_frame(moq_native::moq_net::Timestamp::ZERO, b"hello".as_ref())
.expect("failed to write frame");
group.finish().expect("failed to finish group");
let mut server_config = moq_native::ServerConfig::default();
server_config.bind = Some("[::]:0".to_string());
server_config.tls.generate = vec!["localhost".into()];
let ws_listener = moq_native::websocket::Listener::bind("[::]:0".parse().unwrap())
.await
.expect("failed to bind WebSocket listener");
let ws_addr = ws_listener.local_addr().expect("failed to get ws addr");
let mut server = server_config
.init()
.expect("failed to init server")
.with_websocket(ws_listener);
let sub_origin = Origin::random().produce();
let mut announcements = sub_origin.consume().announced();
let mut client_config = moq_native::ClientConfig::default();
client_config.tls.disable_verify = Some(true);
client_config.websocket.delay = None;
let client = client_config.init().expect("failed to init client");
let url: url::Url = format!("ws://localhost:{}", ws_addr.port()).parse().unwrap();
let server_handle = tokio::spawn(async move {
let request = server.accept().await.expect("no incoming connection");
assert_eq!(request.transport(), moq_native::Transport::WebSocket);
let session = request.with_publisher(&pub_origin).ok().await?;
let _broadcast = broadcast;
let _track = track;
let _ = session.closed().await;
Ok::<_, anyhow::Error>(())
});
let client = client.with_subscriber(sub_origin);
let session = tokio::time::timeout(TIMEOUT, client.connect(url))
.await
.expect("client connect timed out")
.expect("client connect failed");
let moq_native::moq_net::announce::Update { path, broadcast: bc } =
tokio::time::timeout(TIMEOUT, announcements.next())
.await
.expect("announce timed out")
.expect("origin closed");
assert_eq!(path.as_str(), "test");
let bc = bc.expect("expected announce, got unannounce");
let mut track_sub = bc
.track("video")
.unwrap()
.subscribe(None)
.await
.expect("consume_track failed");
let mut group_sub = tokio::time::timeout(TIMEOUT, track_sub.recv_group())
.await
.expect("recv_group timed out")
.expect("recv_group failed")
.expect("track closed prematurely");
let frame = tokio::time::timeout(TIMEOUT, group_sub.read_frame())
.await
.expect("read_frame timed out")
.expect("read_frame failed")
.expect("group closed prematurely");
assert_eq!(&frame.payload[..], b"hello");
drop(session);
server_handle
.await
.expect("server task panicked")
.expect("server task failed");
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_websocket_fallback() {
use moq_native::moq_net::Origin;
let pub_origin = Origin::random().produce();
let mut broadcast = pub_origin
.create_broadcast("test", moq_net::broadcast::Route::new().with_announce(true))
.expect("failed to create broadcast");
let mut track = broadcast.create_track("video", None).expect("failed to create track");
let mut group = track.append_group().expect("failed to append group");
group
.write_frame(moq_native::moq_net::Timestamp::ZERO, b"hello".as_ref())
.expect("failed to write frame");
group.finish().expect("failed to finish group");
let mut server_config = moq_native::ServerConfig::default();
server_config.bind = Some("[::]:0".to_string());
server_config.tls.generate = vec!["localhost".into()];
let ws_listener = moq_native::websocket::Listener::bind("[::]:0".parse().unwrap())
.await
.expect("failed to bind WebSocket listener");
let ws_addr = ws_listener.local_addr().expect("failed to get ws addr");
let mut server = server_config
.init()
.expect("failed to init server")
.with_websocket(ws_listener);
let sub_origin = Origin::random().produce();
let mut announcements = sub_origin.consume().announced();
let mut client_config = moq_native::ClientConfig::default();
client_config.tls.disable_verify = Some(true);
client_config.websocket.delay = None;
let client = client_config.init().expect("failed to init client");
let url: url::Url = format!("http://localhost:{}", ws_addr.port()).parse().unwrap();
let server_handle = tokio::spawn(async move {
let request = server.accept().await.expect("no incoming connection");
assert_eq!(request.transport(), moq_native::Transport::WebSocket);
let session = request.with_publisher(&pub_origin).ok().await?;
let _broadcast = broadcast;
let _track = track;
let _ = session.closed().await;
Ok::<_, anyhow::Error>(())
});
let client = client.with_subscriber(sub_origin);
let session = tokio::time::timeout(TIMEOUT, client.connect(url))
.await
.expect("client connect timed out")
.expect("client connect failed");
let moq_native::moq_net::announce::Update { path, broadcast: bc } =
tokio::time::timeout(TIMEOUT, announcements.next())
.await
.expect("announce timed out")
.expect("origin closed");
assert_eq!(path.as_str(), "test");
let bc = bc.expect("expected announce, got unannounce");
let mut track_sub = bc
.track("video")
.unwrap()
.subscribe(None)
.await
.expect("consume_track failed");
let mut group_sub = tokio::time::timeout(TIMEOUT, track_sub.recv_group())
.await
.expect("recv_group timed out")
.expect("recv_group failed")
.expect("track closed prematurely");
let frame = tokio::time::timeout(TIMEOUT, group_sub.read_frame())
.await
.expect("read_frame timed out")
.expect("read_frame failed")
.expect("group closed prematurely");
assert_eq!(&frame.payload[..], b"hello");
drop(session);
server_handle
.await
.expect("server task panicked")
.expect("server task failed");
}
const NEWEST_LITE: &str = "moq-lite-05";
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_websocket_uses_newest_version() {
let pub_origin = Origin::random().produce();
let mut broadcast = pub_origin
.create_broadcast("test", moq_net::broadcast::Route::new().with_announce(true))
.expect("failed to create broadcast");
let mut track = broadcast.create_track("video", None).expect("failed to create track");
let mut group = track.append_group().expect("failed to append group");
group
.write_frame(moq_native::moq_net::Timestamp::ZERO, b"hello".as_ref())
.expect("failed to write frame");
group.finish().expect("failed to finish group");
let mut server_config = moq_native::ServerConfig::default();
server_config.bind = Some("[::]:0".to_string());
server_config.tls.generate = vec!["localhost".into()];
let ws_listener = moq_native::websocket::Listener::bind("[::]:0".parse().unwrap())
.await
.expect("failed to bind WebSocket listener");
let ws_addr = ws_listener.local_addr().expect("failed to get ws addr");
let mut server = server_config
.init()
.expect("failed to init server")
.with_websocket(ws_listener);
let sub_origin = Origin::random().produce();
let mut client_config = moq_native::ClientConfig::default();
client_config.tls.disable_verify = Some(true);
client_config.websocket.delay = None;
let client = client_config.init().expect("failed to init client");
let url: url::Url = format!("ws://localhost:{}", ws_addr.port()).parse().unwrap();
let expected_version: moq_net::Version = NEWEST_LITE.parse().expect("invalid version");
let server_handle = tokio::spawn(async move {
let request = server.accept().await.expect("no incoming connection");
assert_eq!(request.transport(), moq_native::Transport::WebSocket);
let session = request.with_publisher(&pub_origin).ok().await?;
assert_eq!(session.version(), expected_version, "server negotiated stale version");
let _broadcast = broadcast;
let _track = track;
let _ = session.closed().await;
Ok::<_, anyhow::Error>(())
});
let client = client.with_subscriber(sub_origin);
let cs = tokio::time::timeout(TIMEOUT, client.connect(url))
.await
.expect("client connect timed out")
.expect("client connect failed");
assert_eq!(cs.version(), expected_version, "client negotiated stale version");
drop(cs);
server_handle
.await
.expect("server task panicked")
.expect("server task failed");
}
#[tracing_test::traced_test]
#[tokio::test]
async fn broadcast_race_quic_wins() {
let pub_origin = Origin::random().produce();
let mut broadcast = pub_origin
.create_broadcast("test", moq_net::broadcast::Route::new().with_announce(true))
.expect("failed to create broadcast");
let mut track = broadcast.create_track("video", None).expect("failed to create track");
let mut group = track.append_group().expect("failed to append group");
group
.write_frame(moq_native::moq_net::Timestamp::ZERO, b"hello".as_ref())
.expect("failed to write frame");
group.finish().expect("failed to finish group");
let ws_listener = moq_native::websocket::Listener::bind("[::]:0".parse().unwrap())
.await
.expect("failed to bind WebSocket listener");
let port = ws_listener.local_addr().expect("failed to get ws addr").port();
let mut server_config = moq_native::ServerConfig::default();
server_config.bind = Some(format!("[::]:{port}"));
server_config.tls.generate = vec!["localhost".into()];
let mut server = server_config
.init()
.expect("failed to init server")
.with_websocket(ws_listener);
let sub_origin = Origin::random().produce();
let mut client_config = moq_native::ClientConfig::default();
client_config.tls.disable_verify = Some(true);
client_config.websocket.delay = None;
let client = client_config.init().expect("failed to init client");
let url: url::Url = format!("https://localhost:{port}").parse().unwrap();
let expected_version: moq_net::Version = NEWEST_LITE.parse().expect("invalid version");
let server_handle = tokio::spawn(async move {
let request = server.accept().await.expect("no incoming connection");
assert_eq!(
request.transport(),
moq_native::Transport::Quic,
"QUIC lost the race to WebSocket with both reachable",
);
let session = request.with_publisher(&pub_origin).ok().await?;
assert_eq!(session.version(), expected_version, "server negotiated stale version");
let _broadcast = broadcast;
let _track = track;
let _ = session.closed().await;
Ok::<_, anyhow::Error>(())
});
let client = client.with_subscriber(sub_origin);
let cs = tokio::time::timeout(TIMEOUT, client.connect(url))
.await
.expect("client connect timed out")
.expect("client connect failed");
assert_eq!(cs.version(), expected_version, "client negotiated stale version");
drop(cs);
server_handle
.await
.expect("server task panicked")
.expect("server task failed");
}
#[tokio::test]
async fn resubscribe_keeps_flowing_moq_lite_03() {
let pub_origin = Origin::random().produce();
let mut broadcast = pub_origin
.create_broadcast("test", moq_net::broadcast::Route::new().with_announce(true))
.expect("create broadcast");
let mut track = broadcast.create_track("video", None).expect("create track");
let mut group0 = track.append_group().expect("append group 0");
group0
.write_frame(moq_native::moq_net::Timestamp::ZERO, b"a".as_ref())
.expect("write frame 0");
group0.finish().expect("finish group 0");
let mut server_config = moq_native::ServerConfig::default();
server_config.bind = Some("[::]:0".to_string());
server_config.tls.generate = vec!["localhost".into()];
server_config.version = vec!["moq-lite-03".parse().unwrap()];
let mut server = server_config.init().expect("init server");
let addr = server.local_addr().expect("server addr");
let sub_origin = Origin::random().produce();
let mut announcements = sub_origin.consume().announced();
let mut client_config = moq_native::ClientConfig::default();
client_config.tls.disable_verify = Some(true);
client_config.version = vec!["moq-lite-03".parse().unwrap()];
let client = client_config.init().expect("init client");
let url: url::Url = format!("moqt://localhost:{}", addr.port()).parse().unwrap();
let server_handle = tokio::spawn(async move {
let request = server.accept().await.expect("accept");
let session = request.with_publisher(&pub_origin).ok().await?;
let _ = session.closed().await;
Ok::<_, anyhow::Error>(())
});
let client = client.with_subscriber(sub_origin);
let session = tokio::time::timeout(TIMEOUT, client.connect(url))
.await
.expect("connect timeout")
.expect("connect failed");
let moq_native::moq_net::announce::Update { path, broadcast: bc } =
tokio::time::timeout(TIMEOUT, announcements.next())
.await
.expect("announce timeout")
.expect("origin closed");
assert_eq!(path.as_str(), "test");
let bc = bc.expect("expected announce");
let mut sub1 = bc.track("video").unwrap().subscribe(None).await.expect("subscribe1");
let mut g = tokio::time::timeout(TIMEOUT, sub1.recv_group())
.await
.expect("recv group 0 timeout")
.expect("recv group 0 failed")
.expect("track closed early");
assert_eq!(g.sequence, 0);
let frame = tokio::time::timeout(TIMEOUT, g.read_frame())
.await
.expect("read frame 0 timeout")
.expect("read frame 0 failed")
.expect("group closed early");
assert_eq!(&frame.payload[..], b"a");
drop(g);
drop(sub1);
tokio::time::sleep(Duration::from_millis(20)).await;
let mut sub2 = bc.track("video").unwrap().subscribe(None).await.expect("subscribe2");
let mut group1 = track.append_group().expect("append group 1");
group1
.write_frame(moq_native::moq_net::Timestamp::ZERO, b"b".as_ref())
.expect("write frame 1");
group1.finish().expect("finish group 1");
let mut saw_group1 = false;
for _ in 0..2 {
let mut next = tokio::time::timeout(TIMEOUT, sub2.recv_group())
.await
.expect("recv group timeout")
.expect("recv group failed")
.expect("track closed early");
if next.sequence == 1 {
let frame = tokio::time::timeout(TIMEOUT, next.read_frame())
.await
.expect("read frame 1 timeout")
.expect("read frame 1 failed")
.expect("group closed early on resume");
assert_eq!(&frame.payload[..], b"b");
saw_group1 = true;
break;
}
}
assert!(
saw_group1,
"expected group 1 to be delivered to the resubscribed consumer"
);
drop(session);
server_handle
.await
.expect("server task panicked")
.expect("server task failed");
}
fn active_viewers(registry: &moq_net::stats::Registry) -> u64 {
registry
.snapshot()
.traffic()
.into_iter()
.filter(|(_, role, _)| matches!(role, moq_net::stats::Role::Publisher))
.map(|(_, _, traffic)| traffic.active_broadcasts())
.sum()
}
#[tokio::test]
async fn idle_subscription_releases_the_viewer_count() {
let pub_origin = Origin::random().produce();
let mut broadcast = pub_origin
.create_broadcast("test", moq_net::broadcast::Route::new().with_announce(true))
.expect("create broadcast");
let mut track = broadcast.create_track("video", None).expect("create track");
let mut group = track.append_group().expect("append group");
group
.write_frame(moq_native::moq_net::Timestamp::ZERO, b"hello".as_ref())
.expect("write frame");
group.finish().expect("finish group");
let mut server_config = moq_native::ServerConfig::default();
server_config.bind = Some("[::]:0".to_string());
server_config.tls.generate = vec!["localhost".into()];
let mut server = server_config.init().expect("init server");
let addr = server.local_addr().expect("server addr");
let registry = moq_net::stats::Registry::new(moq_net::stats::Config::new());
let stats = registry.tier(moq_net::stats::Tier::default()).session("");
let sub_origin = Origin::random().produce();
let mut announcements = sub_origin.consume().announced();
let mut client_config = moq_native::ClientConfig::default();
client_config.tls.disable_verify = Some(true);
let client = client_config.init().expect("init client");
let url: url::Url = format!("moqt://localhost:{}", addr.port()).parse().unwrap();
let server_handle = tokio::spawn(async move {
let request = server.accept().await.expect("accept");
let session = request.with_publisher(&pub_origin).with_stats(stats).ok().await?;
let _broadcast = broadcast;
let _track = track;
let _ = session.closed().await;
Ok::<_, anyhow::Error>(())
});
let client = client.with_subscriber(sub_origin);
let session = tokio::time::timeout(TIMEOUT, client.connect(url))
.await
.expect("connect timeout")
.expect("connect failed");
let moq_native::moq_net::announce::Update { broadcast: bc, .. } =
tokio::time::timeout(TIMEOUT, announcements.next())
.await
.expect("announce timeout")
.expect("origin closed");
let bc = bc.expect("expected announce");
let mut sub = bc.track("video").unwrap().subscribe(None).await.expect("subscribe");
let mut g = tokio::time::timeout(TIMEOUT, sub.recv_group())
.await
.expect("recv group timeout")
.expect("recv group failed")
.expect("track closed early");
let frame = tokio::time::timeout(TIMEOUT, g.read_frame())
.await
.expect("read frame timeout")
.expect("read frame failed")
.expect("group closed early");
assert_eq!(&frame.payload[..], b"hello");
assert_eq!(active_viewers(®istry), 1, "the live consumer must count as a viewer");
drop(g);
drop(sub);
let released = tokio::time::timeout(TIMEOUT, async {
while active_viewers(®istry) != 0 {
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await;
assert!(
released.is_ok(),
"viewer count stuck at {} after the last consumer left",
active_viewers(®istry)
);
drop(session);
server_handle
.await
.expect("server task panicked")
.expect("server task failed");
}
#[tracing_test::traced_test]
#[tokio::test]
async fn websocket_unauthorized_handshake_is_explicit() {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("failed to bind TCP listener");
let addr = listener.local_addr().expect("failed to get local addr");
let server_handle = tokio::spawn(async move {
let (mut stream, _) = listener.accept().await?;
let mut buf = [0; 1024];
let _ = stream.read(&mut buf).await?;
stream
.write_all(b"HTTP/1.1 401 Unauthorized\r\nContent-Length: 0\r\nConnection: close\r\n\r\n")
.await?;
Ok::<_, anyhow::Error>(())
});
let mut client_config = moq_native::ClientConfig::default();
client_config.websocket.delay = None;
let client = client_config.init().expect("failed to init client");
let url: url::Url = format!("ws://{addr}").parse().unwrap();
let err = tokio::time::timeout(TIMEOUT, client.connect(url))
.await
.expect("client connect timed out");
let err = expect_connect_err(err);
assert_connect_error(&err, moq_native::ConnectError::Unauthorized);
server_handle
.await
.expect("server task panicked")
.expect("server task failed");
}
#[tracing_test::traced_test]
#[tokio::test]
async fn reconnect_stops_on_websocket_unauthorized() {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("failed to bind TCP listener");
let addr = listener.local_addr().expect("failed to get local addr");
let server_handle = tokio::spawn(async move {
let (mut stream, _) = listener.accept().await?;
let mut buf = [0; 1024];
let _ = stream.read(&mut buf).await?;
stream
.write_all(b"HTTP/1.1 401 Unauthorized\r\nContent-Length: 0\r\nConnection: close\r\n\r\n")
.await?;
Ok::<_, anyhow::Error>(())
});
let mut client_config = moq_native::ClientConfig::default();
client_config.websocket.delay = None;
let client = client_config.init().expect("failed to init client");
let url: url::Url = format!("ws://{addr}").parse().unwrap();
let reconnect = client.reconnect(url);
let err = tokio::time::timeout(TIMEOUT, reconnect.closed())
.await
.expect("reconnect close timed out")
.expect_err("reconnect unexpectedly succeeded");
assert_connect_error(&err, moq_native::ConnectError::Unauthorized);
server_handle
.await
.expect("server task panicked")
.expect("server task failed");
}
#[tracing_test::traced_test]
#[tokio::test]
async fn announce_interest_unauthorized_keeps_session_alive() {
use moq_native::moq_net::Origin;
let pub_origin = Origin::random().produce();
let mut broadcast = pub_origin
.create_broadcast("allowed/test", moq_net::broadcast::Route::new().with_announce(true))
.expect("failed to create broadcast");
let mut track = broadcast.create_track("video", None).expect("failed to create track");
let mut group = track.append_group().expect("failed to append group");
group
.write_frame(moq_native::moq_net::Timestamp::ZERO, b"hello".as_ref())
.expect("failed to write frame");
group.finish().expect("failed to finish group");
let publish = pub_origin
.consume()
.scope(&["allowed".into()])
.expect("failed to scope publish origin");
let (mut server, addr) = test_server();
let sub_origin = Origin::random().produce();
let consume = sub_origin
.scope(&["allowed".into(), "denied".into()])
.expect("failed to scope consume origin");
let mut announcements = consume.consume().announced();
let client = test_client();
let url: url::Url = format!("https://localhost:{}", addr.port()).parse().unwrap();
let server_handle = tokio::spawn(async move {
let request = server.accept().await.expect("no incoming connection");
let session = request.with_publisher(publish).ok().await?;
let _broadcast = broadcast;
let _track = track;
let _ = session.closed().await;
Ok::<_, anyhow::Error>(())
});
let client = client.with_subscriber(consume);
let session = tokio::time::timeout(TIMEOUT, client.connect(url))
.await
.expect("client connect timed out")
.expect("client connect failed");
let moq_native::moq_net::announce::Update { path, broadcast: bc } =
tokio::time::timeout(TIMEOUT, announcements.next())
.await
.expect("announce timed out")
.expect("origin closed");
assert_eq!(path.as_str(), "allowed/test");
assert!(bc.is_some(), "expected announce, got unannounce");
assert!(
tokio::time::timeout(Duration::from_millis(200), session.closed())
.await
.is_err(),
"session closed after unauthorized announce interest",
);
drop(session);
server_handle
.await
.expect("server task panicked")
.expect("server task failed");
}
#[tracing_test::traced_test]
#[tokio::test]
async fn publish_only_client_to_subscribe_only_server() {
use moq_native::moq_net::Origin;
let sub_origin = Origin::random().produce();
let consume = sub_origin
.scope(&["allowed".into(), "denied".into()])
.expect("failed to scope consume origin");
let mut announcements = consume.consume().announced();
let (mut server, addr) = test_server();
let url: url::Url = format!("https://localhost:{}", addr.port()).parse().unwrap();
let server_handle = tokio::spawn(async move {
let session = server
.accept()
.await
.expect("no incoming connection")
.with_subscriber(consume)
.ok()
.await?;
let moq_native::moq_net::announce::Update { path, broadcast: bc } =
tokio::time::timeout(TIMEOUT, announcements.next())
.await
.expect("announce timed out")
.expect("origin closed");
assert_eq!(path.as_str(), "allowed/test");
let bc = bc.expect("expected announce, got unannounce");
let mut track_sub = bc
.track("video")
.unwrap()
.subscribe(None)
.await
.expect("consume_track failed");
let mut group_sub = tokio::time::timeout(TIMEOUT, track_sub.recv_group())
.await
.expect("recv_group timed out")
.expect("recv_group failed")
.expect("track closed prematurely");
let frame = tokio::time::timeout(TIMEOUT, group_sub.read_frame())
.await
.expect("read_frame timed out")
.expect("read_frame failed")
.expect("group closed prematurely");
assert_eq!(&frame.payload[..], b"hello");
assert!(
tokio::time::timeout(Duration::from_millis(200), session.closed())
.await
.is_err(),
"server session closed after unauthorized announce interest",
);
Ok::<_, anyhow::Error>(())
});
let pub_origin = Origin::random().produce();
let mut broadcast = pub_origin
.create_broadcast("allowed/test", moq_net::broadcast::Route::new().with_announce(true))
.expect("failed to create broadcast");
let mut track = broadcast.create_track("video", None).expect("failed to create track");
let mut group = track.append_group().expect("failed to append group");
group
.write_frame(moq_native::moq_net::Timestamp::ZERO, b"hello".as_ref())
.expect("failed to write frame");
group.finish().expect("failed to finish group");
let publish = pub_origin
.consume()
.scope(&["allowed".into()])
.expect("failed to scope publish origin");
let session = tokio::time::timeout(TIMEOUT, test_client().with_publisher(publish).connect(url))
.await
.expect("client connect timed out")
.expect("client connect failed");
server_handle
.await
.expect("server task panicked")
.expect("server task failed");
drop(session);
drop(track);
drop(broadcast);
}
fn test_server() -> (moq_native::Server, std::net::SocketAddr) {
let mut config = moq_native::ServerConfig::default();
config.bind = Some("[::]:0".to_string());
config.tls.generate = vec!["localhost".into()];
let server = config.init().expect("failed to init server");
let addr = server.local_addr().expect("failed to get local addr");
(server, addr)
}
fn test_client() -> moq_native::Client {
let mut config = moq_native::ClientConfig::default();
config.tls.disable_verify = Some(true);
config.init().expect("failed to init client")
}
fn assert_connect_error(err: &moq_native::Error, expected: moq_native::ConnectError) {
assert_eq!(err.connect_error(), Some(expected), "unexpected error: {err}",);
}
fn expect_connect_err(result: moq_native::Result<moq_net::Session>) -> moq_native::Error {
match result {
Ok(_) => panic!("client connect unexpectedly succeeded"),
Err(err) => err,
}
}