use std::net::SocketAddr;
use anyhow::Context;
use axum::http::Method;
use hang::moq_net;
use moq_tokio::RedactedUrl;
use url::Url;
use crate::moq::{ImportTarget, notify_ready};
#[derive(usage::Args, Clone)]
#[usage(unknown_flags = "error", args_override_self = false)]
#[usage(group("rtc-mode", required))]
pub struct Args {
#[usage(name = "rtc-connect", long = "connect", value_name = "URL", group = "rtc-mode")]
pub connect: Option<Url>,
#[usage(name = "rtc-listen", long = "listen", value_name = "ADDR", group = "rtc-mode")]
pub listen: Option<SocketAddr>,
#[usage(long, requires = "--listen", default = "[::]:0")]
pub udp_bind: SocketAddr,
#[usage(long, requires = "rtc-listen")]
pub public_addr: Vec<SocketAddr>,
#[usage(flatten)]
pub cors: crate::web::Cors,
}
pub struct Listen {
pub addr: SocketAddr,
pub udp_bind: SocketAddr,
pub public_addr: Vec<SocketAddr>,
pub cors: crate::web::Cors,
}
pub async fn listen_import(target: ImportTarget, listen: Listen) -> anyhow::Result<()> {
let publisher = scope_producer(&target.origin, &target.name)?;
let mut config = server_config(&listen);
config.max_age = target.max_age;
config.bandwidth = target.bandwidth;
let server = moq_rtc::Server::new(config);
serve(server.publish_router(publisher), "WHIP", listen).await
}
pub async fn listen_export(origin: moq_net::origin::Consumer, name: String, listen: Listen) -> anyhow::Result<()> {
let scope = moq_net::Patterns::from(
moq_net::Pattern::subtree(&name).with_context(|| format!("invalid broadcast name `{name}`"))?,
);
let subscriber = origin
.scope("", &scope)
.with_context(|| format!("failed to scope origin to broadcast `{name}`"))?;
let server = moq_rtc::Server::new(server_config(&listen));
serve(server.subscribe_router(subscriber), "WHEP", listen).await
}
fn scope_producer(origin: &moq_net::origin::Producer, name: &str) -> anyhow::Result<moq_net::origin::Producer> {
let scope = moq_net::Patterns::from(
moq_net::Pattern::subtree(name).with_context(|| format!("invalid broadcast name `{name}`"))?,
);
origin
.scope("", &scope)
.with_context(|| format!("failed to scope origin to broadcast `{name}`"))
}
pub async fn connect_import(target: ImportTarget, url: Url) -> anyhow::Result<()> {
let name = &target.name;
let producer = target
.origin
.create_broadcast(name)
.context("failed to create broadcast")?;
producer
.announce(Default::default())
.context("failed to announce broadcast")?;
tracing::info!(url = %RedactedUrl::new(&url), %name, "WHEP client pulling");
notify_ready();
let mut config = moq_rtc::client::Config::default();
config.max_age = target.max_age;
config.bandwidth = target.bandwidth;
let client = moq_rtc::Client::new(config);
Ok(client.subscribe(url, producer).await?)
}
pub async fn connect_export(origin: moq_net::origin::Consumer, url: Url, name: String) -> anyhow::Result<()> {
origin
.routed(&name)
.await
.with_context(|| format!("origin closed before broadcast `{name}` was announced"))?;
tracing::info!(url = %RedactedUrl::new(&url), %name, "WHIP client pushing");
notify_ready();
let client = moq_rtc::Client::new(moq_rtc::client::Config::default());
Ok(client.publish(url, origin, &name).await?)
}
fn server_config(listen: &Listen) -> moq_rtc::server::Config {
let mut config = moq_rtc::server::Config::default();
config.udp_bind = listen.udp_bind;
config.ice_candidates.clone_from(&listen.public_addr);
config
}
async fn serve(router: axum::Router, role: &str, listen: Listen) -> anyhow::Result<()> {
let cors = listen
.cors
.layer([Method::POST, Method::PATCH, Method::DELETE, Method::OPTIONS])?;
let app = router.layer(cors);
let listener = moq_tokio::bind::tcp(listen.addr)?;
tracing::info!(listen = %listen.addr, role, "serving WebRTC");
notify_ready();
crate::web::serve(listener, app, None).await
}