use moq_native::moq_net::{self, Origin};
use std::time::Duration;
const TIMEOUT: Duration = Duration::from_secs(10);
#[cfg(any(feature = "quinn", feature = "quiche", feature = "noq"))]
struct ConnectTest<'a> {
scheme: &'a str,
bind: &'a str,
client_bind: Option<&'a str>,
authority: &'a str,
backend: moq_native::QuicBackend,
qlog: Option<&'a std::path::Path>,
}
#[cfg(any(feature = "quinn", feature = "quiche", feature = "noq"))]
async fn backend_test(scheme: &str, backend: moq_native::QuicBackend) {
connect_test(ConnectTest {
scheme,
bind: "[::]:0",
client_bind: None,
authority: "localhost",
backend,
qlog: None,
})
.await;
}
#[cfg(any(feature = "quinn", feature = "noq"))]
async fn no_sni_test(scheme: &str, backend: moq_native::QuicBackend) {
connect_test(ConnectTest {
scheme,
bind: "127.0.0.1:0",
client_bind: None,
authority: "127.0.0.1",
backend,
qlog: None,
})
.await;
}
#[cfg(any(feature = "quinn", feature = "quiche", feature = "noq"))]
async fn connect_test(config: ConnectTest<'_>) {
let ConnectTest {
scheme,
bind,
client_bind,
authority,
backend,
qlog,
} = config;
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(bind.to_string());
server_config.tls.generate = vec!["localhost".into()];
server_config.backend = Some(backend.clone());
server_config.quic.qlog = qlog.map(Into::into);
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.backend = Some(backend);
client_config.quic.qlog = qlog.map(Into::into);
client_config.bind = client_bind.unwrap_or(bind).parse().expect("invalid bind address");
let client = client_config.init().expect("failed to init client");
let url: url::Url = format!("{scheme}://{authority}:{}", addr.port()).parse().unwrap();
let server_handle = tokio::spawn(async move {
let request = server.accept().await.expect("no incoming connection");
assert_eq!(request.role(), Some(moq_native::moq_net::Role::Subscriber));
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");
}
#[cfg(any(feature = "quinn", feature = "quiche", feature = "noq"))]
fn generate_mtls_certs() -> (tempfile::TempDir, MtlsPaths) {
use rcgen::{
BasicConstraints, CertificateParams, DnType, ExtendedKeyUsagePurpose, IsCa, Issuer, KeyPair, KeyUsagePurpose,
};
use std::io::Write;
let dir = tempfile::tempdir().expect("failed to create tempdir");
let ca_key = KeyPair::generate().expect("ca key");
let mut ca_params = CertificateParams::new(Vec::new()).expect("ca params");
ca_params.is_ca = IsCa::Ca(BasicConstraints::Unconstrained);
ca_params.distinguished_name.push(DnType::CommonName, "moq test CA");
ca_params.key_usages = vec![KeyUsagePurpose::KeyCertSign, KeyUsagePurpose::CrlSign];
let ca_cert = ca_params.self_signed(&ca_key).expect("ca cert");
let issuer = Issuer::from_params(&ca_params, &ca_key);
let server_key = KeyPair::generate().expect("server key");
let mut server_params = CertificateParams::new(vec!["localhost".to_string()]).expect("server params");
server_params.distinguished_name.push(DnType::CommonName, "localhost");
server_params.extended_key_usages = vec![ExtendedKeyUsagePurpose::ServerAuth];
let server_cert = server_params.signed_by(&server_key, &issuer).expect("server cert");
let client_key = KeyPair::generate().expect("client key");
let mut client_params = CertificateParams::new(vec!["client.example".to_string()]).expect("client params");
client_params
.distinguished_name
.push(DnType::CommonName, "client.example");
client_params.extended_key_usages = vec![ExtendedKeyUsagePurpose::ClientAuth];
let client_cert = client_params.signed_by(&client_key, &issuer).expect("client cert");
let write = |name: &str, contents: String| {
let path = dir.path().join(name);
let mut file = std::fs::File::create(&path).expect("create pem file");
file.write_all(contents.as_bytes()).expect("write pem file");
path
};
let paths = MtlsPaths {
ca: write("ca.pem", ca_cert.pem()),
server_cert: write("server.pem", server_cert.pem()),
server_key: write("server.key", server_key.serialize_pem()),
client_cert: write("client.pem", client_cert.pem()),
client_key: write("client.key", client_key.serialize_pem()),
};
(dir, paths)
}
#[cfg(any(feature = "quinn", feature = "quiche", feature = "noq"))]
struct MtlsPaths {
ca: std::path::PathBuf,
server_cert: std::path::PathBuf,
server_key: std::path::PathBuf,
client_cert: std::path::PathBuf,
client_key: std::path::PathBuf,
}
#[cfg(any(feature = "quinn", feature = "quiche", feature = "noq"))]
async fn mtls_test(scheme: &str, backend: moq_native::QuicBackend, reject: bool) {
let (_dir, paths) = generate_mtls_certs();
let pub_origin = Origin::random().produce();
let mut server_config = moq_native::ServerConfig::default();
server_config.bind = Some("127.0.0.1:0".to_string());
server_config.tls.cert = vec![paths.server_cert.clone()];
server_config.tls.key = vec![paths.server_key.clone()];
server_config.tls.root = vec![paths.ca.clone()];
server_config.backend = Some(backend.clone());
server_config.quic.gso = Some(false);
server_config.quic.keep_alive = Some(Duration::from_secs(1));
let mut server = server_config.init().expect("failed to init server");
let addr = server.local_addr().expect("failed to get local addr");
let mut client_config = moq_native::ClientConfig::default();
client_config.tls.root = vec![paths.ca.clone()];
client_config.tls.system_roots = Some(false);
client_config.tls.cert = Some(paths.client_cert.clone());
client_config.tls.key = Some(paths.client_key.clone());
client_config.tls.host_name = Some("localhost".to_string());
client_config.backend = Some(backend);
client_config.bind = "0.0.0.0:0".parse().unwrap();
client_config.quic.gso = Some(false);
client_config.quic.keep_alive = Some(Duration::from_secs(1));
let client = client_config.init().expect("failed to init client");
let url: url::Url = format!("{scheme}://127.0.0.1:{}", addr.port()).parse().unwrap();
let (identity_tx, identity_rx) = tokio::sync::oneshot::channel();
let server_handle = tokio::spawn(async move {
let request = server.accept().await.expect("no incoming connection");
let has_cert = request.peer_identity().is_some();
let _ = identity_tx.send(has_cert);
if reject {
request.close(403).await?;
return Ok::<_, anyhow::Error>(has_cert);
}
let session = request.with_publisher(pub_origin.consume()).ok().await?;
let _ = session.closed().await;
Ok::<_, anyhow::Error>(has_cert)
});
let session = tokio::time::timeout(TIMEOUT, client.connect(url))
.await
.expect("client connect timed out");
let has_cert = tokio::time::timeout(TIMEOUT, identity_rx)
.await
.expect("identity inspection timed out")
.expect("server dropped identity result");
assert!(has_cert, "server did not observe the client certificate");
if !reject {
session.as_ref().expect("client connect failed");
}
drop(session);
if reject {
server_handle.abort();
} else {
tokio::time::timeout(TIMEOUT, server_handle)
.await
.expect("server task timed out")
.expect("server task panicked")
.expect("server task failed");
}
}
#[cfg(feature = "quinn")]
#[tracing_test::traced_test]
#[tokio::test]
async fn quinn_raw_quic() {
backend_test("moqt", moq_native::QuicBackend::Quinn).await;
}
#[cfg(feature = "quinn")]
#[tracing_test::traced_test]
#[tokio::test]
async fn quinn_raw_quic_no_sni() {
no_sni_test("moqt", moq_native::QuicBackend::Quinn).await;
}
#[cfg(feature = "quinn")]
#[tracing_test::traced_test]
#[tokio::test]
async fn quinn_mtls() {
mtls_test("https", moq_native::QuicBackend::Quinn, false).await;
}
#[cfg(feature = "quinn")]
#[tracing_test::traced_test]
#[tokio::test]
async fn quinn_webtransport() {
backend_test("https", moq_native::QuicBackend::Quinn).await;
}
#[cfg(feature = "quiche")]
#[tracing_test::traced_test]
#[tokio::test]
async fn quiche_raw_quic() {
backend_test("moqt", moq_native::QuicBackend::Quiche).await;
}
#[cfg(feature = "quiche")]
#[tracing_test::traced_test]
#[tokio::test]
async fn quiche_dual_stack_ipv4() {
connect_test(ConnectTest {
scheme: "moqt",
bind: "[::]:0",
client_bind: Some("0.0.0.0:0"),
authority: "127.0.0.1",
backend: moq_native::QuicBackend::Quiche,
qlog: None,
})
.await;
}
#[cfg(feature = "quiche")]
#[tracing_test::traced_test]
#[tokio::test]
async fn quiche_mtls() {
mtls_test("https", moq_native::QuicBackend::Quiche, true).await;
}
#[cfg(feature = "quiche")]
#[tracing_test::traced_test]
#[tokio::test]
#[ignore = "web-transport-quiche teardown bug. Two symptoms: the lite-05 TRACK stream, dropped \
after reading TRACK_INFO (closed early by design), surfaces as a connection-level \
`quiche error: Done` that tears down the whole session; and on CI it intermittently \
aborts with a SIGSEGV in the boring/quiche C stack (seen on aarch64 runners), inside \
web-transport-quiche / tokio-quiche / BoringSSL, not our code. lite-04 (no TRACK \
stream) and the quinn backend both work. Re-enable once the quiche backend handles the \
early stream drop."]
async fn quiche_webtransport() {
backend_test("https", moq_native::QuicBackend::Quiche).await;
}
#[cfg(feature = "iroh")]
#[tracing_test::traced_test]
#[tokio::test]
async fn iroh_connect() {
use moq_native::iroh::EndpointConfig;
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_iroh_config = EndpointConfig::default();
server_iroh_config.enabled = Some(true);
let server_endpoint = server_iroh_config
.bind(&moq_native::quic::Client::default())
.await
.expect("failed to bind server iroh endpoint")
.expect("server iroh endpoint not enabled");
let server_addr = server_endpoint.addr();
let server_addrs: Vec<std::net::SocketAddr> = server_addr.ip_addrs().copied().collect();
let server_endpoint_id = server_endpoint.id();
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("failed to init server")
.with_iroh(server_endpoint);
let sub_origin = Origin::random().produce();
let mut announcements = sub_origin.consume().announced();
let mut client_iroh_config = EndpointConfig::default();
client_iroh_config.enabled = Some(true);
let client_endpoint = client_iroh_config
.bind(&moq_native::quic::Client::default())
.await
.expect("failed to bind client iroh endpoint")
.expect("client iroh endpoint not enabled");
let mut client_config = moq_native::ClientConfig::default();
client_config.tls.disable_verify = Some(true);
let client = client_config
.init()
.expect("failed to init client")
.with_iroh(client_endpoint)
.with_iroh_addrs(server_addrs);
let url: url::Url = format!("iroh://{server_endpoint_id}").parse().unwrap();
let server_handle = tokio::spawn(async move {
let request = server.accept().await.expect("no incoming connection");
assert_eq!(request.role(), Some(moq_native::moq_net::Role::Subscriber));
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");
}
#[cfg(feature = "noq")]
#[tracing_test::traced_test]
#[tokio::test]
async fn noq_raw_quic() {
backend_test("moqt", moq_native::QuicBackend::Noq).await;
}
#[cfg(feature = "noq")]
#[tracing_test::traced_test]
#[tokio::test]
async fn noq_raw_quic_no_sni() {
no_sni_test("moqt", moq_native::QuicBackend::Noq).await;
}
#[cfg(feature = "noq")]
#[tracing_test::traced_test]
#[tokio::test]
async fn noq_webtransport() {
backend_test("https", moq_native::QuicBackend::Noq).await;
}
#[cfg(feature = "noq")]
#[tracing_test::traced_test]
#[tokio::test]
async fn noq_mtls() {
mtls_test("https", moq_native::QuicBackend::Noq, false).await;
}
#[cfg(all(feature = "qlog", any(feature = "quinn", feature = "quiche", feature = "noq")))]
async fn qlog_test(scheme: &str, backend: moq_native::QuicBackend) -> Vec<std::path::PathBuf> {
let dir = tempfile::tempdir().expect("failed to create tempdir");
connect_test(ConnectTest {
scheme,
bind: "[::]:0",
client_bind: None,
authority: "localhost",
backend,
qlog: Some(dir.path()),
})
.await;
let traces: Vec<_> = std::fs::read_dir(dir.path())
.expect("failed to read qlog dir")
.map(|entry| entry.expect("failed to read qlog entry").path())
.collect();
for trace in &traces {
let raw = std::fs::read(trace).expect("failed to read trace");
let header = raw.split(|&b| b == b'\n').next().unwrap_or_default();
let header = String::from_utf8_lossy(header);
assert!(
header.starts_with('\u{1e}'),
"qlog trace {} is not JSON-SEQ: {header:?}",
trace.display()
);
assert!(
header.contains("JSON-SEQ"),
"qlog trace {} has no qlog header: {header:?}",
trace.display()
);
assert!(
raw.iter().filter(|&&b| b == 0x1e).count() > 1,
"qlog trace {} has a header but no events",
trace.display()
);
}
traces
}
#[cfg(all(feature = "qlog", feature = "quinn"))]
#[tokio::test]
async fn quinn_qlog() {
let traces = qlog_test("moqt", moq_native::QuicBackend::Quinn).await;
assert_eq!(traces.len(), 2, "expected one trace per endpoint, got {traces:?}");
}
#[cfg(all(feature = "qlog", feature = "noq"))]
#[tokio::test]
async fn noq_qlog() {
let traces = qlog_test("moqt", moq_native::QuicBackend::Noq).await;
assert!(!traces.is_empty(), "expected at least one trace");
}
#[cfg(all(feature = "qlog", feature = "quiche"))]
#[tokio::test]
#[ignore = "shares the connect path that quiche_webtransport is ignored for; the qlog plumbing itself is verified by running this locally"]
async fn quiche_qlog() {
let traces = qlog_test("https", moq_native::QuicBackend::Quiche).await;
assert!(!traces.is_empty(), "expected at least one trace");
}