use std::net::SocketAddr;
use std::time::Duration;
use anyhow::Context;
use hang::moq_net;
use moq_rtmp::{Client, Request, Server};
use url::Url;
use crate::moq::{ImportTarget, notify_ready};
#[derive(usage::Args, Clone)]
#[usage(unknown_flags = "error", args_override_self = false)]
#[usage(group("rtmp-mode", required))]
pub struct Args {
#[usage(name = "rtmp-connect", long = "connect", value_name = "URL", group = "rtmp-mode")]
pub connect: Option<Url>,
#[usage(name = "rtmp-listen", long = "listen", value_name = "ADDR", group = "rtmp-mode")]
pub listen: Option<SocketAddr>,
}
#[derive(usage::Args, Clone)]
#[usage(unknown_flags = "error", args_override_self = false)]
pub struct ExportArgs {
#[usage(flatten)]
pub endpoint: Args,
#[usage(long, default = "500ms")]
pub max_age: crate::duration::Duration,
#[usage(long = "latency-max", hide = true)]
pub(crate) latency_max: Option<crate::duration::Duration>,
}
pub async fn listen_import(target: ImportTarget, addr: SocketAddr) -> anyhow::Result<()> {
let ImportTarget {
origin,
name,
max_age,
bandwidth,
} = target;
let mut server = Server::bind(addr).await?;
tracing::info!(%addr, %name, "RTMP listening (import)");
notify_ready();
while let Some(request) = server.accept().await {
match request {
Request::Publish(publish) => {
let origin = origin.clone();
let name = name.clone();
let bandwidth = bandwidth.clone();
tokio::spawn(async move {
if let Err(err) = publish
.with_max_age(max_age)
.with_bandwidth(bandwidth)
.accept(&origin, &name)
.await
{
tracing::warn!(%name, %err, "RTMP ingest ended with error");
}
});
}
Request::Play(play) => {
tokio::spawn(async move {
let _ = play.reject("this is an import listener; it does not serve plays").await;
});
}
_ => {}
}
}
Ok(())
}
pub async fn listen_export(
origin: moq_net::origin::Consumer,
addr: SocketAddr,
name: String,
max_age: Duration,
) -> anyhow::Result<()> {
let mut server = Server::bind(addr).await?;
tracing::info!(%addr, %name, "RTMP listening (export)");
notify_ready();
while let Some(request) = server.accept().await {
match request {
Request::Play(play) => {
let origin = origin.clone();
let name = name.clone();
tokio::spawn(async move {
if let Err(err) = play.with_max_age(max_age).accept(&origin, &name).await {
tracing::warn!(%name, %err, "RTMP play ended with error");
}
});
}
Request::Publish(publish) => {
tokio::spawn(async move {
let _ = publish
.reject("this is an export listener; it does not accept publishes")
.await;
});
}
_ => {}
}
}
Ok(())
}
pub async fn connect_import(target: ImportTarget, url: Url) -> anyhow::Result<()> {
let (addr, app, key) = parse_url(&url).await?;
let name = &target.name;
tracing::info!(%addr, %app, %name, "RTMP client pulling");
notify_ready();
let client = Client::connect(addr, &app)
.await?
.with_import_max_age(target.max_age)
.with_import_bandwidth(target.bandwidth);
Ok(client.pull(&key, &target.origin, name).await?)
}
pub async fn connect_export(
origin: moq_net::origin::Consumer,
url: Url,
name: String,
max_age: Duration,
) -> anyhow::Result<()> {
let (addr, app, key) = parse_url(&url).await?;
origin
.routed(&name)
.await
.with_context(|| format!("origin closed before broadcast `{name}` was announced"))?;
tracing::info!(%addr, %app, %name, "RTMP client pushing");
notify_ready();
let client = Client::connect(addr, &app).await?.with_export_max_age(max_age);
Ok(client.publish(&key, origin, &name).await?)
}
async fn parse_url(url: &Url) -> anyhow::Result<(SocketAddr, String, String)> {
anyhow::ensure!(url.scheme() == "rtmp", "rtmp url must use the rtmp scheme: {url}");
let host = url
.host_str()
.with_context(|| format!("rtmp url missing host: {url}"))?;
let port = url.port().unwrap_or(1935);
let addr = tokio::net::lookup_host((host, port))
.await?
.next()
.with_context(|| format!("could not resolve {host}:{port}"))?;
let mut segments = url.path().trim_matches('/').splitn(2, '/');
let app = segments.next().unwrap_or_default().to_string();
let key = segments.next().unwrap_or_default().to_string();
anyhow::ensure!(
!app.is_empty() && !key.is_empty(),
"rtmp url must include an app and stream key: rtmp://host/<app>/<key>"
);
Ok((addr, app, key))
}
#[cfg(test)]
mod tests {
use super::*;
async fn parse(url: &str) -> anyhow::Result<(SocketAddr, String, String)> {
parse_url(&Url::parse(url).unwrap()).await
}
#[tokio::test]
async fn ok_default_port() {
let (addr, app, key) = parse("rtmp://127.0.0.1/live/cam0").await.unwrap();
assert_eq!(addr.port(), 1935);
assert_eq!((app.as_str(), key.as_str()), ("live", "cam0"));
}
#[tokio::test]
async fn ok_explicit_port() {
let (addr, _, _) = parse("rtmp://127.0.0.1:1936/live/cam0").await.unwrap();
assert_eq!(addr.port(), 1936);
}
#[tokio::test]
async fn rejects_non_rtmp_scheme() {
assert!(parse("http://127.0.0.1/live/cam0").await.is_err());
}
#[tokio::test]
async fn requires_app_and_key() {
assert!(parse("rtmp://127.0.0.1/live").await.is_err());
assert!(parse("rtmp://127.0.0.1/").await.is_err());
}
}