mod args;
mod auth;
mod complete;
#[cfg(feature = "capture")]
mod devices;
mod duration;
mod fetch;
mod hls;
mod moq;
mod play;
mod publish;
mod rtc;
mod rtmp;
mod srt;
mod subscribe;
#[cfg(test)]
mod test_env;
#[cfg(feature = "transcode")]
mod transcode;
mod web;
use args::{Command, Export, ExportSink, Import, ImportSource, Invocation, MoqSide};
use hang::moq_net;
use publish::Publish;
use subscribe::{Subscribe, SubscribeArgs};
use anyhow::Context;
use tokio::task::JoinSet;
#[cfg(feature = "jemalloc")]
#[global_allocator]
static ALLOC: moq_tokio::jemalloc::tikv_jemallocator::Jemalloc = moq_tokio::jemalloc::tikv_jemallocator::Jemalloc;
#[derive(Clone)]
struct Net {
quic: moq_tokio::quic::Config,
#[cfg(feature = "iroh")]
iroh: Option<moq_tokio::iroh::Endpoint>,
}
impl Net {
fn client(&self, config: moq_tokio::connect::Config) -> anyhow::Result<moq_tokio::Client> {
let client = config.init(self.quic.clone())?;
#[cfg(feature = "iroh")]
let client = match self.iroh.clone() {
Some(iroh) => client.with_iroh(iroh),
None => client,
};
Ok(client)
}
fn server(&self, config: moq_tokio::listen::Config) -> anyhow::Result<moq_tokio::Server> {
let mut server = moq_tokio::server::Config::default();
server.listen = config;
server.quic = self.quic.clone();
#[cfg(feature = "iroh")]
{
server.iroh = self.iroh.clone();
}
Ok(server.init()?)
}
}
async fn spawn_server(
tasks: &mut JoinSet<anyhow::Result<()>>,
moq: &MoqSide,
cluster: &moq_relay::cluster::Cluster,
net: &Net,
directions: Directions,
) -> anyhow::Result<moq_relay::cluster::Started> {
if !moq.serves() {
return cluster.clone().start().await.context("cluster failed to start");
}
let server = net.server(moq.server_config())?;
let certificates = server.certificates();
let cluster = attach_lan(cluster.clone(), moq, &server)?;
let node = moq.cluster.node.clone().unwrap_or_default();
let auth = match moq.auth.validate() {
Ok(()) => moq.auth.init(node, &moq.client.tls)?,
Err(_) => moq_relay::auth::Auth::refuse(node),
};
let started = cluster.clone().start().await.context("cluster failed to start")?;
let origin = &cluster.origin;
let listener = server.listen().await.context("failed to bind listeners")?;
if moq.lan() {
spawn_cluster_serve(
tasks,
listener,
cluster.clone(),
auth,
origin.clone(),
directions,
moq.server.bind.is_some(),
);
} else {
spawn_serve(tasks, listener, auth, origin.clone(), directions);
}
if let Some(web_bind) = moq.server.bind.clone() {
tasks.spawn(async move { web::run_web(web_bind, certificates).await });
}
Ok(started)
}
fn attach_lan(
cluster: moq_relay::cluster::Cluster,
moq: &MoqSide,
server: &moq_tokio::Server,
) -> anyhow::Result<moq_relay::cluster::Cluster> {
if !moq.lan() {
return Ok(cluster);
}
let port = server
.local_addr()
.context("--cluster-lan needs a QUIC listener")?
.port();
let mut advertise = moq_relay::cluster::LanAdvertise::new(port);
if !moq.server_config().tls.generate.is_empty()
&& let Some(fingerprint) = server.certificates().fingerprints().into_iter().next()
{
advertise = advertise.with_fingerprint(fingerprint);
}
Ok(cluster.with_advertise(advertise))
}
fn spawn_cluster_serve(
tasks: &mut JoinSet<anyhow::Result<()>>,
mut listener: moq_tokio::Listener,
cluster: moq_relay::cluster::Cluster,
auth: moq_relay::auth::Auth,
origin: moq_net::origin::Producer,
directions: Directions,
public_quic: bool,
) {
if let Ok(addr) = listener.local_addr() {
tracing::info!(%addr, "listening");
}
tasks.spawn(async move {
let mut sessions = tokio::task::JoinSet::new();
while let Some(request) = listener.accept().await {
while sessions.try_join_next().is_some() {}
if moq_relay::cluster::Cluster::is_lan_path(request.path()) {
let conn = moq_relay::Connection::new(request, cluster.clone(), auth.clone())
.with_id(cluster.next_connection_id());
sessions.spawn(async move {
if let Err(err) = conn.run().await {
tracing::warn!(%err, "LAN peer session ended");
}
});
continue;
}
if !is_public_transport(request.transport(), public_quic) {
tracing::debug!(path = %request.path(), "refusing a non-peer request on the LAN mesh listener");
request.reject(moq_tokio::server::Reject::App(404)).await.ok();
continue;
}
let auth = auth.clone();
let origin = origin.clone();
sessions.spawn(async move {
if let Err(err) = serve_client(request, &auth, &origin, directions).await {
tracing::warn!(%err, "session ended with error");
}
});
}
anyhow::bail!("the MoQ listener stopped accepting")
});
}
async fn serve_client(
request: moq_tokio::server::Request,
auth: &moq_relay::auth::Auth,
origin: &moq_net::origin::Producer,
directions: Directions,
) -> anyhow::Result<()> {
let auth_request = moq_relay::auth::request_for(auth, &request);
let lease = match auth.admit(auth_request).await {
Ok(lease) => lease,
Err(err) => {
let status = axum::http::StatusCode::from(&err);
let reject = match status {
axum::http::StatusCode::UNAUTHORIZED => moq_tokio::server::Reject::Unauthorized,
axum::http::StatusCode::FORBIDDEN => moq_tokio::server::Reject::Forbidden,
status => moq_tokio::server::Reject::App(status.as_u16()),
};
request.reject(reject).await.ok();
return Err(anyhow::Error::new(err).context("session refused"));
}
};
let token = lease.token();
let publish = directions
.publish
.then(|| origin.scope(&token.root, &token.subscribe).ok())
.flatten();
let subscribe = directions
.consume
.then(|| origin.scope(&token.root, &token.publish).ok())
.flatten();
if publish.is_none() && subscribe.is_none() {
request.reject(moq_tokio::server::Reject::Forbidden).await.ok();
anyhow::bail!("grant allows nothing this endpoint serves at {}", token.root);
}
let mut request = request;
if let Some(publish) = publish {
request = request.with_publisher(publish.consume());
}
if let Some(subscribe) = subscribe {
request = request.with_subscriber(subscribe);
}
let session = request.ok().await?;
moq_relay::supervise(session, lease, moq_relay::shutdown::Observer::disabled(), None).await
}
fn is_public_transport(transport: moq_tokio::server::Transport, public_quic: bool) -> bool {
match transport {
moq_tokio::server::Transport::Tcp | moq_tokio::server::Transport::Unix => true,
_ => public_quic,
}
}
fn spawn_serve(
tasks: &mut JoinSet<anyhow::Result<()>>,
mut listener: moq_tokio::Listener,
auth: moq_relay::auth::Auth,
origin: moq_net::origin::Producer,
directions: Directions,
) {
if let Ok(addr) = listener.local_addr() {
tracing::info!(%addr, "listening");
}
tasks.spawn(async move {
while let Some(request) = listener.accept().await {
let auth = auth.clone();
let origin = origin.clone();
tokio::spawn(async move {
if let Err(err) = serve_client(request, &auth, &origin, directions).await {
tracing::warn!(%err, "session ended with error");
}
});
}
Ok(())
});
}
#[tokio::main]
async fn main() -> anyhow::Result<()> {
moq_tokio::crypto::install_default().expect("failed to install default crypto provider");
let mut cli = Invocation::parse().await;
cli.log.init()?;
cli.validate()?;
let mut stages = std::mem::take(&mut cli.stages);
if stages.len() == 1 {
match stages.remove(0) {
Command::Auth(auth) => {
cli.reject("auth")?;
return auth.run().await;
}
Command::Completion(completion) => {
cli.reject("completion")?;
return completion.run();
}
#[cfg(feature = "capture")]
Command::Devices => {
cli.reject("devices")?;
return devices::run().await;
}
other => stages.push(other),
}
}
if let [Command::Fetch(_)] = stages.as_slice() {
cli.dial_only("fetch")?;
} else {
cli.moq.validate()?;
}
let net = Net {
quic: cli.moq.quic.clone(),
#[cfg(feature = "iroh")]
iroh: cli.moq.iroh.clone().bind(&cli.moq.quic).await?,
};
#[cfg(feature = "jemalloc")]
let jemalloc = moq_tokio::jemalloc::run();
#[cfg(not(feature = "jemalloc"))]
let jemalloc = std::future::pending::<anyhow::Result<()>>();
let run = async move {
if stages.len() == 1 && !stages[0].is_stageable() {
match stages.remove(0) {
Command::Fetch(args) => return fetch::run(cli.moq, args, net).await,
#[cfg(feature = "play")]
Command::Play(args) => return run_play(cli.moq, args, net).await,
#[cfg(feature = "transcode")]
Command::Transcode(args) => return transcode::run(cli.moq, args, net).await,
_ => unreachable!("the local verbs returned before the transport was bound"),
}
}
run_stages(cli.moq, stages, net).await
};
tokio::select! {
result = run => result,
Err(err) = jemalloc => Err(err).context("jemalloc profiler failed"),
}
}
#[derive(Clone, Copy, Default)]
pub struct Directions {
pub publish: bool,
pub consume: bool,
}
impl Directions {
fn of(stages: &[Command]) -> Self {
Self {
publish: stages.iter().any(|stage| matches!(stage, Command::Import(_))),
consume: stages.iter().any(|stage| matches!(stage, Command::Export(_))),
}
}
}
async fn spawn_moq(
moq: &MoqSide,
net: &Net,
cluster: moq_relay::cluster::Cluster,
directions: Directions,
tasks: &mut JoinSet<anyhow::Result<()>>,
) -> anyhow::Result<(moq_net::bandwidth::Allocator, moq_net::origin::Producer)> {
let mut bandwidth = moq_net::bandwidth::Allocator::unlimited();
let client = net.client(moq.client.clone())?;
let cluster = cluster
.with_client(client.clone())
.with_client_tls(moq.client.tls.build()?)
.with_connect(moq.client.clone(), moq.quic.clone());
let origin = cluster.origin.clone();
if let Some(url) = moq.client.url.clone() {
let mut client = client;
if directions.publish {
client = client.with_publisher(origin.consume());
}
if directions.consume {
client = client.with_subscriber(origin.clone());
}
let reconnect = client.connect(url);
bandwidth = moq_net::bandwidth::Allocator::new(reconnect.send_bandwidth());
tasks.spawn(async move { Ok(reconnect.closed().await?) });
}
let started =
notify_when_initialized(spawn_server(tasks, moq, &cluster, net, directions), moq::notify_ready).await?;
if !started.standalone() {
tasks.spawn(async move { started.run().await });
}
Ok((bandwidth, origin))
}
async fn notify_when_initialized<T>(
initialization: impl std::future::Future<Output = anyhow::Result<T>>,
notify_ready: impl FnOnce(),
) -> anyhow::Result<T> {
let value = initialization.await?;
notify_ready();
Ok(value)
}
#[cfg(feature = "play")]
async fn run_play(moq: MoqSide, args: play::Args, net: Net) -> anyhow::Result<()> {
args.validate()?;
let cluster = moq.cluster()?;
let name = moq.broadcast.clone().unwrap_or_default();
let mut tasks: JoinSet<anyhow::Result<()>> = JoinSet::new();
let directions = Directions {
consume: true,
..Default::default()
};
let (_, origin) = spawn_moq(&moq, &net, cluster, directions, &mut tasks).await?;
play::run(origin.consume(), name, args, tasks)
}
async fn run_stages(moq: MoqSide, stages: Vec<Command>, net: Net) -> anyhow::Result<()> {
let cluster = moq.cluster()?;
let mut tasks: JoinSet<anyhow::Result<()>> = JoinSet::new();
let mut locals: Vec<Publish> = Vec::new();
let (bandwidth, origin) = spawn_moq(&moq, &net, cluster, Directions::of(&stages), &mut tasks).await?;
let mut stdin = None;
let mut stdout = None;
for stage in stages {
let name = stage.broadcast(&moq);
match stage {
Command::Import(import) => {
if import.source.stdin_format().is_some() {
claim("stdin", &mut stdin, &name)?;
}
if let Some(publish) = spawn_import(&origin, import, name, bandwidth.clone(), &mut tasks)? {
locals.push(publish);
}
}
Command::Export(export) => {
if export.sink.is_stdout() {
claim("stdout", &mut stdout, &name)?;
}
spawn_export(&origin, export, name, &mut tasks)?;
}
other => unreachable!("`{}` is not a stage", other.name()),
}
}
if locals.is_empty() {
return drive(tasks).await;
}
let local = tokio::task::LocalSet::new();
supervise(&local, locals.into_iter().map(Publish::run), &mut tasks);
local.run_until(drive(tasks)).await
}
fn supervise<F>(
local: &tokio::task::LocalSet,
pipelines: impl IntoIterator<Item = F>,
tasks: &mut JoinSet<anyhow::Result<()>>,
) where
F: std::future::Future<Output = anyhow::Result<()>> + 'static,
{
for pipeline in pipelines {
let pipeline = local.spawn_local(pipeline);
tasks.spawn(async move { pipeline.await.context("pipeline panicked")? });
}
}
fn claim(stream: &str, held: &mut Option<String>, name: &str) -> anyhow::Result<()> {
if let Some(first) = held {
anyhow::bail!(
"only one stage can use {stream}, but both `{}` and `{}` do",
display_name(first),
display_name(name),
);
}
*held = Some(name.to_string());
Ok(())
}
fn display_name(name: &str) -> &str {
if name.is_empty() { "<root>" } else { name }
}
fn spawn_import(
origin: &moq_net::origin::Producer,
import: Import,
name: String,
bandwidth: moq_net::bandwidth::Allocator,
tasks: &mut JoinSet<anyhow::Result<()>>,
) -> anyhow::Result<Option<Publish>> {
if let ImportSource::Rtc(rtc) = &import.source
&& rtc.connect.is_some()
{
reject_listener_cors(&rtc.cors, "import rtc")?;
}
let max_age = import.max_age.map(crate::duration::Duration::into_std);
let target = |name: String| crate::moq::ImportTarget {
origin: origin.clone(),
name,
max_age,
bandwidth: bandwidth.clone(),
};
let mut local = None;
if let Some(format) = import.source.stdin_format() {
warn_if_missing_format(&name);
let broadcast = origin.create_broadcast(&name).context("failed to create broadcast")?;
let config = moq_mux::catalog::Config::default()
.with_max_age(max_age)
.with_bandwidth(bandwidth.clone());
let publish = Publish::new(broadcast, &format, config)?;
publish.announce()?;
local = Some(publish);
} else {
match import.source {
ImportSource::Hls(hls) => {
warn_if_missing_format(&name);
tasks.spawn(hls::import(target(name), hls.playlist));
}
ImportSource::Rtmp(rtmp) => {
if let Some(addr) = rtmp.listen {
let name = require_broadcast(name, "import rtmp --listen")?;
tasks.spawn(rtmp::listen_import(target(name), addr));
} else if let Some(url) = rtmp.connect {
tasks.spawn(rtmp::connect_import(target(name), url));
}
}
ImportSource::Srt(srt) => {
if let Some(addr) = srt.listen {
let name = require_broadcast(name, "import srt --listen")?;
tasks.spawn(srt::listen_import(target(name), addr, srt.latency.into_std()));
} else if let Some(url) = srt.connect {
tasks.spawn(srt::connect_import(target(name), url, srt.latency.into_std()));
}
}
ImportSource::Rtc(rtc) => {
if let Some(addr) = rtc.listen {
let name = require_broadcast(name, "import rtc --listen")?;
tasks.spawn(rtc::listen_import(
target(name),
rtc::Listen {
addr,
udp_bind: rtc.udp_bind,
public_addr: rtc.public_addr,
cors: rtc.cors,
},
));
} else if let Some(url) = rtc.connect {
tasks.spawn(rtc::connect_import(target(name), url));
}
}
#[cfg(feature = "capture")]
ImportSource::Capture(capture) => {
warn_if_missing_format(&name);
let broadcast = origin.create_broadcast(&name).context("failed to create broadcast")?;
let publish = Publish::capture(broadcast, &capture, bandwidth.clone(), max_age)?;
publish.announce()?;
local = Some(publish);
}
_ => unreachable!("container formats are handled by stdin_format above"),
}
}
Ok(local)
}
fn spawn_export(
origin: &moq_net::origin::Producer,
export: Export,
name: String,
tasks: &mut JoinSet<anyhow::Result<()>>,
) -> anyhow::Result<()> {
if let ExportSink::Rtc(rtc) = &export.sink
&& rtc.connect.is_some()
{
reject_listener_cors(&rtc.cors, "export rtc")?;
}
if let Some(stdout) = export.sink.stdout() {
let args = SubscribeArgs {
format: stdout.format,
max_age: stdout.max_age,
fragment_duration: stdout.fragment_duration,
mux_rate: stdout.mux_rate,
catalog: export.catalog_format,
select: export.select,
};
let consumer = origin.consume();
tasks.spawn(async move { run_stdout(consumer, name, args).await });
} else {
match export.sink {
ExportSink::Hls(args) => {
let name = require_broadcast(name, "export hls")?;
tasks.spawn(hls::export(origin.consume(), args, name));
}
ExportSink::Rtmp(rtmp) => {
let max_age = rtmp.max_age.into_std();
if let Some(addr) = rtmp.endpoint.listen {
let name = require_broadcast(name, "export rtmp --listen")?;
tasks.spawn(rtmp::listen_export(origin.consume(), addr, name, max_age));
} else if let Some(url) = rtmp.endpoint.connect {
tasks.spawn(rtmp::connect_export(origin.consume(), url, name, max_age));
}
}
ExportSink::Srt(srt) => {
if let Some(addr) = srt.listen {
let name = require_broadcast(name, "export srt --listen")?;
tasks.spawn(srt::listen_export(origin.consume(), addr, name, srt.latency.into_std()));
} else if let Some(url) = srt.connect {
tasks.spawn(srt::connect_export(origin.consume(), url, name, srt.latency.into_std()));
}
}
ExportSink::Rtc(rtc) => {
if let Some(addr) = rtc.listen {
let name = require_broadcast(name, "export rtc --listen")?;
tasks.spawn(rtc::listen_export(
origin.consume(),
name,
rtc::Listen {
addr,
udp_bind: rtc.udp_bind,
public_addr: rtc.public_addr,
cors: rtc.cors,
},
));
} else if let Some(url) = rtc.connect {
tasks.spawn(rtc::connect_export(origin.consume(), url, name));
}
}
_ => unreachable!("container formats are handled by stdout_format above"),
}
}
Ok(())
}
async fn run_stdout(consumer: moq_net::origin::Consumer, name: String, args: SubscribeArgs) -> anyhow::Result<()> {
let catalog = args.catalog_format(&name);
consumer
.routed(&name)
.await
.ok_or_else(|| anyhow::anyhow!("origin closed before broadcast `{name}` was announced"))?;
let source = moq_mux::Source::new(consumer, &name);
Subscribe::new(source, catalog, args).run().await
}
async fn drive(mut tasks: JoinSet<anyhow::Result<()>>) -> anyhow::Result<()> {
tasks.spawn(async {
let _ = tokio::signal::ctrl_c().await;
Ok(())
});
while let Some(res) = tasks.join_next().await {
match res {
Ok(Ok(())) => return Ok(()),
Ok(Err(err)) => return Err(err),
Err(err) if err.is_cancelled() => continue,
Err(err) => return Err(err.into()),
}
}
Ok(())
}
fn require_broadcast(name: String, endpoint: &str) -> anyhow::Result<String> {
anyhow::ensure!(
!name.is_empty(),
"`{endpoint}` requires a broadcast: pass --broadcast <name>"
);
Ok(name)
}
fn warn_if_missing_format(name: &str) {
if !name.is_empty() && moq_mux::catalog::CatalogFormat::detect(name).is_none() {
tracing::warn!(
name,
"You should append .hang to your broadcast name to make the catalog format explicit."
);
}
}
fn reject_listener_cors(cors: &crate::web::Cors, endpoint: &str) -> anyhow::Result<()> {
anyhow::ensure!(
cors.origin.is_empty(),
"`--cors-origin` only applies to `{endpoint} --listen`"
);
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use std::cell::Cell;
use std::future::Future;
use std::pin::Pin;
type Pipeline = Pin<Box<dyn Future<Output = anyhow::Result<()>>>>;
#[tokio::test]
async fn a_panicking_pipeline_ends_the_process() {
let local = tokio::task::LocalSet::new();
let mut tasks = JoinSet::new();
let pipelines: Vec<Pipeline> = vec![
Box::pin(async { panic!("pipeline died") }),
Box::pin(std::future::pending()),
];
supervise(&local, pipelines, &mut tasks);
let err = local.run_until(drive(tasks)).await.unwrap_err();
assert!(err.to_string().contains("pipeline panicked"), "{err}");
}
#[tokio::test]
async fn a_finished_pipeline_ends_the_process() {
let local = tokio::task::LocalSet::new();
let mut tasks = JoinSet::new();
let pipelines: Vec<Pipeline> = vec![Box::pin(async { Ok(()) }), Box::pin(std::future::pending())];
supervise(&local, pipelines, &mut tasks);
local.run_until(drive(tasks)).await.unwrap();
}
#[tokio::test]
async fn a_stream_bind_failure_prevents_readiness() {
let occupied = std::net::TcpListener::bind("127.0.0.1:0").expect("occupy a port");
let addr = occupied.local_addr().expect("occupied address").to_string();
let invocation = Invocation::try_parse_from(["moq", "--listen-tcp-bind", &addr, "import", "ts"])
.expect("parse stream-only invocation");
let cluster = invocation.moq.cluster().expect("create cluster");
let net = Net {
quic: invocation.moq.quic.clone(),
#[cfg(feature = "iroh")]
iroh: None,
};
let mut tasks = JoinSet::new();
let ready = Cell::new(false);
let err = match notify_when_initialized(
spawn_server(
&mut tasks,
&invocation.moq,
&cluster,
&net,
Directions {
publish: true,
consume: false,
},
),
|| ready.set(true),
)
.await
{
Ok(_) => panic!("the occupied port must fail initialization"),
Err(err) => err,
};
assert!(err.to_string().contains("failed to bind listeners"), "{err:#}");
assert!(!ready.get(), "readiness must be withheld after a bind failure");
assert!(tasks.is_empty(), "nothing should be spawned after a bind failure");
}
#[tokio::test]
async fn tcp_only_moq_side_starts_a_server() {
let invocation = Invocation::try_parse_from(["moq", "--listen-tcp-bind", "127.0.0.1:0", "import", "ts"])
.expect("parse TCP-only invocation");
let cluster = invocation.moq.cluster().expect("create cluster");
let net = Net {
quic: invocation.moq.quic.clone(),
#[cfg(feature = "iroh")]
iroh: None,
};
let mut tasks = JoinSet::new();
if let Err(err) = spawn_server(
&mut tasks,
&invocation.moq,
&cluster,
&net,
Directions {
publish: true,
consume: false,
},
)
.await
{
panic!("start TCP-only server: {err:#}");
}
assert!(
tokio::time::timeout(std::time::Duration::from_millis(50), tasks.join_next())
.await
.is_err(),
"the server task should still be accepting connections"
);
tasks.abort_all();
}
#[tokio::test]
async fn cluster_connect_api_http_attaches_client_tls() {
let _ = moq_tokio::crypto::install_default();
let invocation = Invocation::try_parse_from([
"moq",
"--cluster-connect-api",
"https://api.example/peers",
"import",
"ts",
])
.expect("parse");
assert!(invocation.moq.validate().is_ok());
let net = Net {
quic: invocation.moq.quic.clone(),
#[cfg(feature = "iroh")]
iroh: None,
};
let client = net.client(invocation.moq.client.clone()).expect("client");
let cluster = invocation
.moq
.cluster()
.expect("cluster")
.with_client(client)
.with_connect(invocation.moq.client.clone(), invocation.moq.quic.clone());
let err = cluster
.clone()
.start()
.await
.expect_err("http API without TLS")
.to_string();
assert!(err.contains("client TLS"), "{err}");
cluster
.with_client_tls(invocation.moq.client.tls.build().expect("tls"))
.start()
.await
.expect("http API with TLS");
}
#[test]
fn explicit_stream_listeners_are_public_without_exposing_mesh_quic() {
assert!(is_public_transport(moq_tokio::server::Transport::Tcp, false));
assert!(is_public_transport(moq_tokio::server::Transport::Unix, false));
assert!(!is_public_transport(moq_tokio::server::Transport::Quic, false));
assert!(is_public_transport(moq_tokio::server::Transport::Quic, true));
}
}