use anyhow::Context;
use crate::Net;
use crate::args::MoqSide;
use hang::moq_net;
#[derive(usage::Args, Clone)]
#[usage(unknown_flags = "error", args_override_self = false)]
pub struct Args {
#[usage(long)]
pub output: Option<String>,
#[usage(long = "rung")]
rungs: Vec<RungArg>,
#[usage(long, default = "auto")]
pub encoder: String,
#[usage(long, default = "auto")]
pub decoder: String,
#[usage(long, default = "native")]
frames: OutputArg,
}
#[derive(Clone)]
struct RungArg(moq_transcode::Rung);
impl std::str::FromStr for RungArg {
type Err = String;
fn from_str(arg: &str) -> Result<Self, Self::Err> {
let (height, bitrate) = arg
.split_once(':')
.ok_or_else(|| format!("expected height:bitrate, got `{arg}`"))?;
let height: u32 = height.parse().map_err(|e| format!("invalid height `{height}`: {e}"))?;
let bitrate: u64 = bitrate
.parse()
.map_err(|e| format!("invalid bitrate `{bitrate}`: {e}"))?;
Ok(Self(moq_transcode::Rung::new(
height,
moq_net::bandwidth::Rate::from_bps(bitrate),
)))
}
}
#[derive(Clone, Copy)]
struct OutputArg(moq_video::Output);
impl std::str::FromStr for OutputArg {
type Err = String;
fn from_str(arg: &str) -> Result<Self, Self::Err> {
match arg {
"native" => Ok(Self(moq_video::Output::Native)),
"cpu" => Ok(Self(moq_video::Output::Cpu)),
_ => Err(format!("expected native or cpu, got `{arg}`")),
}
}
}
pub async fn run(moq: MoqSide, args: Args, net: Net) -> anyhow::Result<()> {
let mut config = moq_transcode::Config::default();
if !args.rungs.is_empty() {
config.ladder =
moq_transcode::Ladder::new(args.rungs.iter().map(|rung| rung.0)).context("invalid --rung ladder")?;
}
config.encoder = match args.encoder.as_str() {
"auto" => moq_video::encode::Kind::Auto,
"hardware" => moq_video::encode::Kind::Hardware,
"software" => moq_video::encode::Kind::Software,
name => moq_video::encode::Kind::Named(name.to_string()),
};
config.decoder = match args.decoder.as_str() {
"auto" => moq_video::decode::Kind::Auto,
"hardware" => moq_video::decode::Kind::Hardware,
"software" => moq_video::decode::Kind::Software,
name => moq_video::decode::Kind::Named(name.to_string()),
};
config.resize.output = args.frames.0;
let source_path = moq_net::PathOwned::from(
moq.broadcast
.clone()
.context("`transcode` requires the source broadcast: pass --broadcast <name>")?,
);
if source_path.is_empty() {
anyhow::bail!("`transcode` requires the source broadcast: pass --broadcast <name>");
}
let output_path = moq_net::PathOwned::from(
args.output
.clone()
.unwrap_or_else(|| format!("{source_path}/transcode.hang")),
);
let url = moq
.client
.url
.clone()
.context("`transcode` requires a relay: pass --connect <url>")?;
let publish = moq_tokio::origin::spawn();
let remote = moq_tokio::origin::spawn();
let session = net
.client(moq.client.clone())?
.with_publisher(&publish)
.with_subscriber(remote.clone())
.connect(url);
let consumer = remote.consume();
let source = tokio::select! {
source = consumer.routed_broadcast(&source_path) => {
source.context("source broadcast unavailable")?
}
closed = session.closed() => {
closed.context("session failed before the source broadcast was announced")?;
anyhow::bail!("session closed before the source broadcast was announced");
}
};
config.source = source_path.relative(&output_path).filter(|rel| !rel.is_empty());
let output = publish
.create_broadcast(&output_path)
.context("failed to create the derivative broadcast")?;
output
.announce(Default::default())
.context("failed to announce the derivative broadcast")?;
tracing::info!(source = %source_path, output = %output_path, "transcoding");
tokio::select! {
res = moq_transcode::run(source, output, config) => Ok(res?),
res = session.closed() => Ok(res?),
}
}