use anyhow::Context;
use crate::Net;
use crate::args::MoqSide;
use hang::moq_net;
#[derive(clap::Args, Clone)]
pub struct Args {
#[arg(long)]
pub output: Option<String>,
#[arg(long = "rung", value_parser = parse_rung)]
pub rungs: Vec<moq_transcode::Rung>,
#[arg(long, default_value = "auto")]
pub encoder: String,
#[arg(long, default_value = "auto")]
pub decoder: String,
#[arg(long, default_value = "auto", value_parser = parse_resize_acceleration)]
pub resize_acceleration: moq_video::resize::Acceleration,
}
fn parse_rung(arg: &str) -> Result<moq_transcode::Rung, String> {
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(moq_transcode::Rung::new(height, bitrate))
}
fn parse_resize_acceleration(arg: &str) -> Result<moq_video::resize::Acceleration, String> {
match arg {
"auto" => Ok(moq_video::resize::Acceleration::Auto),
"cpu" => Ok(moq_video::resize::Acceleration::Cpu),
"gpu" => Ok(moq_video::resize::Acceleration::Gpu),
_ => Err(format!("expected auto, cpu, or gpu, got `{arg}`")),
}
}
pub async fn run(moq: MoqSide, args: Args, net: Net) -> anyhow::Result<()> {
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
.connect
.clone()
.context("`transcode` requires a relay: pass --client-connect <url>")?;
let publish = moq_net::Origin::random().produce();
let remote = moq_net::Origin::random().produce();
let session = net
.client(moq.client.clone())?
.with_publisher(&publish)
.with_subscriber(remote.clone())
.reconnect(url);
let consumer = remote.consume();
tokio::select! {
announced = consumer.announced_broadcast(&source_path) => {
announced.context("origin closed before the source broadcast was announced")?;
}
closed = session.closed() => {
closed.context("session failed before the source broadcast was announced")?;
anyhow::bail!("session closed before the source broadcast was announced");
}
}
let source = consumer
.request_broadcast(&source_path)
.await
.context("source broadcast unavailable")?;
let mut config = moq_transcode::Config::default();
if !args.rungs.is_empty() {
config.rungs = args.rungs.clone();
}
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.acceleration = args.resize_acceleration;
config.source = moq_transcode::source_reference(&source_path, &output_path);
let output = publish
.create_broadcast(&output_path, moq_net::broadcast::Route::new().with_announce(true))
.context("failed to create 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?),
}
}