use anyhow::Context as _;
use std::ffi::{OsStr, OsString};
use std::time::Duration;
use crate::publish::PublishFormat;
use crate::subscribe::{CatalogFormatArg, SubscribeFormat};
#[derive(usage::Cli, Clone)]
#[usage(unknown_flags = "error", args_override_self = false)]
#[usage(name = "moq", version = env!("CARGO_PKG_VERSION"))]
#[usage(completion, settings)]
#[usage(after_help = "Separate additional import/export stages with `--`; they share one \
connection and one Origin. Every `--` starts a stage, so it is not an \
end-of-options marker: write a path starting with `-` as `./-name`.")]
pub struct Cli {
#[usage(flatten)]
pub log: moq_tokio::Log,
#[usage(flatten)]
pub moq: MoqSide,
#[usage(subcommand)]
pub command: Command,
}
#[derive(usage::Cli, Clone)]
#[usage(unknown_flags = "error", args_override_self = false)]
#[usage(name = "moq")]
#[usage(completion, settings)]
pub struct Stage {
#[usage(subcommand)]
pub command: Command,
}
pub struct Invocation {
pub log: moq_tokio::Log,
pub moq: MoqSide,
pub typed: MoqSide,
pub stages: Vec<Command>,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum ParseErrorKind {
UnknownArgument,
MissingSubcommand,
ValueValidation,
DisplayHelp,
DisplayVersion,
Other,
}
#[derive(Debug)]
pub struct ParseError {
kind: ParseErrorKind,
message: String,
}
impl ParseError {
#[cfg_attr(not(test), allow(dead_code))]
pub fn kind(&self) -> ParseErrorKind {
self.kind
}
fn new(kind: ParseErrorKind, message: impl Into<String>) -> Self {
Self {
kind,
message: message.into(),
}
}
fn exit(self) -> ! {
if matches!(self.kind, ParseErrorKind::DisplayHelp | ParseErrorKind::DisplayVersion) {
print!("{}", self.message);
std::process::exit(0);
}
eprint!("{}", self.message);
std::process::exit(2)
}
}
impl std::fmt::Display for ParseError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.message)
}
}
impl std::error::Error for ParseError {}
impl Invocation {
pub async fn parse() -> Self {
let args: Vec<OsString> = std::env::args_os().collect();
if let Some(reply) = crate::complete::answer(args.get(1..).unwrap_or_default()).await {
print!("{reply}");
std::process::exit(0);
}
match Self::try_parse_from(args) {
Ok(parsed) => parsed,
Err(err) => err.exit(),
}
}
pub fn reject(&self, command: &str) -> anyhow::Result<()> {
self.typed.reject(command)
}
pub fn dial_only(&self, command: &str) -> anyhow::Result<()> {
self.typed.dial_only(command)
}
pub fn try_parse_from<I, T>(argv: I) -> Result<Self, ParseError>
where
I: IntoIterator<Item = T>,
T: Into<OsString>,
{
let argv: Vec<OsString> = argv.into_iter().map(Into::into).collect();
let mut chunks = argv.split(|arg| arg == OsStr::new("--"));
let first = chunks.next().unwrap_or_default();
let first = first.iter().skip(1).map(OsString::as_os_str).collect::<Vec<_>>();
let cli = Cli::parse_from(&first).map_err(|err| parse_error(Cli::spec(), Cli::command(), &first, err))?;
let typed = MoqSide::from_argv(&first, Environment::Ignore).unwrap_or_else(|| cli.moq.clone());
let mut deprecated = cli.moq.deprecated();
deprecated.extend(cli.command.deprecated());
let mut stages = vec![cli.command.current()];
for chunk in chunks {
if chunk.is_empty() {
return Err(ParseError::new(
ParseErrorKind::MissingSubcommand,
"error: `--` starts another stage, so it must be followed by `import` or `export`\n",
));
}
let chunk = chunk.iter().map(OsString::as_os_str).collect::<Vec<_>>();
let command = Stage::parse_from(&chunk)
.map_err(|err| parse_error(Stage::spec(), Stage::command(), &chunk, err))?
.command;
deprecated.extend(command.deprecated());
stages.push(command.current());
}
if !deprecated.is_empty() {
return Err(ParseError::new(
ParseErrorKind::ValueValidation,
format!("error: {deprecated}\n"),
));
}
Ok(Self {
log: cli.log,
moq: cli.moq,
typed,
stages,
})
}
pub fn validate(&self) -> anyhow::Result<()> {
if self.stages.len() == 1 {
return Ok(());
}
if let Some(command) = self.stages.iter().find(|command| !command.is_stageable()) {
anyhow::bail!(
"`{}` must be the only verb; it can't share a process with another `--` stage",
command.name()
);
}
Ok(())
}
}
fn parse_error(
spec: &usage::argv::spec::Spec<'_>,
root: &usage::Command<'_>,
argv: &[&OsStr],
err: usage::Error<'_, '_>,
) -> ParseError {
let kind = match &err {
usage::Error::Help { .. } | usage::Error::HelpAll { .. } => ParseErrorKind::DisplayHelp,
usage::Error::Version { .. } => ParseErrorKind::DisplayVersion,
usage::Error::UnknownFlag { .. } | usage::Error::UnexpectedArg { .. } => ParseErrorKind::UnknownArgument,
usage::Error::MissingSubcommand | usage::Error::MissingArgsHelp { .. } => ParseErrorKind::MissingSubcommand,
usage::Error::InvalidValue(_) | usage::Error::InvalidChoice { .. } => ParseErrorKind::ValueValidation,
_ => ParseErrorKind::Other,
};
ParseError::new(kind, moq_tokio::cli::answer(spec, root, argv, err).message())
}
#[derive(Clone, Copy, Eq, PartialEq)]
pub(crate) enum Environment {
Read,
Ignore,
}
#[derive(usage::Args, Clone)]
#[usage(unknown_flags = "error", args_override_self = false)]
pub struct MoqSide {
#[usage(long, help_heading = "MoQ")]
pub broadcast: Option<String>,
#[usage(long = "name", hide = true)]
name: Option<String>,
#[usage(long, env = "MOQ_HOP", help_heading = "MoQ")]
pub hop: Option<u64>,
#[usage(name = "origin", long = "origin", env = "MOQ_ORIGIN", hide = true)]
pub origin: Option<u64>,
#[usage(flatten)]
pub client: moq_tokio::connect::Config,
#[usage(flatten)]
pub quic: moq_tokio::quic::Config,
#[usage(flatten)]
pub server: moq_tokio::listen::Config,
#[cfg(feature = "iroh")]
#[usage(flatten)]
pub iroh: moq_tokio::iroh::Config,
#[usage(flatten)]
pub cluster: moq_relay::cluster::Config,
#[usage(flatten)]
pub auth: moq_relay::auth::Config,
}
impl MoqSide {
fn deprecated(&self) -> moq_tokio::cli::Deprecated {
let mut found = self.client.deprecated();
found.extend(self.quic.deprecated());
found.extend(self.server.deprecated());
found.extend(self.cluster.deprecated());
if self.origin.is_some() {
found.flag("--origin", Some("MOQ_ORIGIN"), "--hop / MOQ_HOP");
}
if self.name.is_some() {
found.flag("--name", None, "--broadcast");
}
found
}
pub fn cluster(&self) -> anyhow::Result<moq_relay::cluster::Cluster> {
let mut config = self.cluster.clone();
match (config.id, self.hop) {
(None, Some(hop)) => config.id = Some(hop),
(Some(id), Some(hop)) if id != hop => {
anyhow::bail!("--hop {hop} and --cluster-id {id} must agree")
}
_ => {}
}
moq_relay::cluster::Cluster::new(moq_relay::cluster::Options::new(config))
}
pub fn lan(&self) -> bool {
#[cfg(feature = "cluster-lan")]
return self.cluster.lan.enabled;
#[cfg(not(feature = "cluster-lan"))]
false
}
pub fn server_config(&self) -> moq_tokio::listen::Config {
let mut config = self.server.clone();
if self.lan() {
config
.bind
.get_or_insert_with(|| moq_tokio::listen::Bind::Addr("[::]:0".parse().unwrap()));
if config.tls.generate.is_empty() && config.tls.cert.is_empty() {
config.tls.generate = vec!["moq-cluster-lan".to_string()];
}
}
config
}
pub fn serves(&self) -> bool {
self.server.has_explicit_bind() || self.lan()
}
fn cluster_dials(&self) -> bool {
!self.cluster.connect.is_empty() || self.cluster.connect_api.is_some()
}
pub fn validate(&self) -> anyhow::Result<()> {
anyhow::ensure!(
self.client.url.is_some() || self.serves() || self.cluster_dials(),
"a MoQ side is required: pass --connect <url> to dial a relay, a --listen option to self-host, --cluster-lan to mesh over the LAN, or --cluster-connect to join a cluster"
);
#[cfg(feature = "cluster-lan")]
{
self.cluster.lan.validate()?;
if self.lan() {
moq_relay::cluster::Cluster::validate_lan_versions(&self.client, &self.server_config())?;
}
}
if self.server.has_explicit_bind() {
self.auth
.validate()
.context("--listen needs --auth-url or --auth-public")?;
} else if self.auth.url.is_some() || self.auth_public() {
self.auth.validate()?;
}
Ok(())
}
fn auth_public(&self) -> bool {
!(self.auth.public.is_empty() && self.auth.public_subscribe.is_empty() && self.auth.public_publish.is_empty())
}
pub(crate) fn from_argv(argv: &[&OsStr], environment: Environment) -> Option<Self> {
use usage::spec::CommandArgs;
let mut partial = <Self as CommandArgs>::start();
let mut parser = usage::Parser::new(Cli::command(), argv);
while let Some(event) = parser.next_event() {
match event {
Ok(event) => {
<Self as CommandArgs>::apply(&mut partial, &event);
}
Err(_) => break,
}
}
if environment == Environment::Read {
<Self as CommandArgs>::apply_env(&mut partial);
}
<Self as CommandArgs>::apply_defaults(&mut partial);
<Self as CommandArgs>::build(partial).ok()
}
fn reject(&self, command: &str) -> anyhow::Result<()> {
if let Some(flag) = self.given().next() {
anyhow::bail!("`{command}` runs locally and takes no MoQ side; drop {flag}");
}
Ok(())
}
fn dial_only(&self, command: &str) -> anyhow::Result<()> {
if let Some(flag) = self.given().find(|flag| !matches!(*flag, "--connect" | "--broadcast")) {
anyhow::bail!("`{command}` only dials a relay with --connect; drop {flag}");
}
Ok(())
}
fn given(&self) -> impl Iterator<Item = &'static str> {
#[cfg(feature = "cluster-lan")]
let cluster_secret = self.cluster.lan.secret.is_some();
#[cfg(not(feature = "cluster-lan"))]
let cluster_secret = false;
#[cfg(feature = "cluster-lan")]
let cluster_app = self.cluster.lan.app.is_some();
#[cfg(not(feature = "cluster-lan"))]
let cluster_app = false;
let flags = [
("--connect", self.client.url.is_some()),
("--listen", self.server.bind.is_some()),
("--listen-tcp-bind", self.server.tcp.bind.is_some()),
("--cluster-lan", self.lan()),
("--cluster-lan-secret", cluster_secret),
("--cluster-lan-app", cluster_app),
("--cluster-connect", !self.cluster.connect.is_empty()),
("--cluster-connect-api", self.cluster.connect_api.is_some()),
("--cluster-node", self.cluster.node.is_some()),
("--cluster-mesh", self.cluster.mesh.is_some()),
("--cluster-token", self.cluster.token.is_some()),
("--cluster-id", self.cluster.id.is_some()),
("--cluster-tier", self.cluster.tier.is_some()),
("--auth-url", self.auth.url.is_some()),
("--auth-public", self.auth_public()),
("--broadcast", self.broadcast.is_some()),
("--hop", self.hop.is_some()),
];
#[cfg(unix)]
let unix = {
let allow = &self.server.unix.allow;
[
("--listen-unix-bind", self.server.unix.bind.is_some()),
("--listen-unix-allow-uid", !allow.uid.is_empty()),
("--listen-unix-allow-gid", !allow.gid.is_empty()),
("--listen-unix-allow-pid", !allow.pid.is_empty()),
]
};
#[cfg(not(unix))]
let unix = [];
flags
.into_iter()
.chain(unix)
.filter(|(_, given)| *given)
.map(|(flag, _)| flag)
}
}
#[derive(usage::Subcommands, Clone)]
pub enum Command {
Import(Import),
Export(Export),
#[usage(hide = true)]
Publish(Import),
#[usage(hide = true)]
Subscribe(Export),
Fetch(crate::fetch::Args),
#[cfg(feature = "play")]
Play(crate::play::Args),
#[cfg(feature = "transcode")]
Transcode(crate::transcode::Args),
Auth(crate::auth::Args),
Completion(crate::complete::Args),
#[cfg(feature = "capture")]
Devices,
}
impl Command {
fn deprecated(&self) -> moq_tokio::cli::Deprecated {
match self {
Self::Import(import) => import.deprecated(),
Self::Publish(import) => {
let mut found = import.deprecated();
found.flag("publish", None, "import");
found
}
Self::Export(export) => export.deprecated(),
Self::Subscribe(export) => {
let mut found = export.deprecated();
found.flag("subscribe", None, "export");
found
}
_ => moq_tokio::cli::Deprecated::default(),
}
}
fn current(self) -> Self {
match self {
Self::Publish(import) => Self::Import(import),
Self::Subscribe(export) => Self::Export(export),
other => other,
}
}
pub fn name(&self) -> &'static str {
match self {
Self::Import(_) | Self::Publish(_) => "import",
Self::Export(_) | Self::Subscribe(_) => "export",
Self::Fetch(_) => "fetch",
#[cfg(feature = "play")]
Self::Play(_) => "play",
#[cfg(feature = "transcode")]
Self::Transcode(_) => "transcode",
Self::Auth(_) => "auth",
Self::Completion(_) => "completion",
#[cfg(feature = "capture")]
Self::Devices => "devices",
}
}
pub fn is_stageable(&self) -> bool {
matches!(
self,
Self::Import(_) | Self::Export(_) | Self::Publish(_) | Self::Subscribe(_)
)
}
pub fn broadcast(&self, moq: &MoqSide) -> String {
let stage = match self {
Self::Import(import) | Self::Publish(import) => import.broadcast.as_deref(),
Self::Export(export) | Self::Subscribe(export) => export.broadcast.as_deref(),
_ => None,
};
stage.or(moq.broadcast.as_deref()).unwrap_or_default().to_string()
}
}
#[derive(usage::Args, Clone)]
#[usage(unknown_flags = "error", args_override_self = false)]
pub struct Import {
#[usage(long)]
pub broadcast: Option<String>,
#[usage(long = "name", hide = true)]
name: Option<String>,
#[usage(long)]
pub max_age: Option<crate::duration::Duration>,
#[usage(long = "latency-max", hide = true)]
latency_max: Option<crate::duration::Duration>,
#[usage(subcommand)]
pub source: ImportSource,
}
impl Import {
fn deprecated(&self) -> moq_tokio::cli::Deprecated {
let mut found = moq_tokio::cli::Deprecated::default();
if self.name.is_some() {
found.flag("--name", None, "--broadcast");
}
if self.latency_max.is_some() {
found.flag("--latency-max", None, "--max-age");
}
found
}
}
#[derive(usage::Subcommands, Clone)]
pub enum ImportSource {
Avc3,
Fmp4,
Ts,
Flv,
Hls(crate::hls::ImportArgs),
Rtmp(crate::rtmp::Args),
Srt(crate::srt::Args),
Rtc(crate::rtc::Args),
#[cfg(feature = "capture")]
Capture(crate::publish::CaptureArgs),
}
impl ImportSource {
pub fn stdin_format(&self) -> Option<PublishFormat> {
Some(match self {
Self::Avc3 => PublishFormat::Avc3,
Self::Fmp4 => PublishFormat::Fmp4,
Self::Ts => PublishFormat::Ts,
Self::Flv => PublishFormat::Flv,
_ => return None,
})
}
}
#[derive(usage::Args, Clone)]
#[usage(unknown_flags = "error", args_override_self = false)]
pub struct Export {
#[usage(long)]
pub broadcast: Option<String>,
#[usage(long = "name", hide = true)]
name: Option<String>,
#[usage(long = "catalog-format", value_enum)]
pub catalog_format: Option<CatalogFormatArg>,
#[usage(flatten)]
pub select: crate::subscribe::SelectArgs,
#[usage(subcommand)]
pub sink: ExportSink,
}
impl Export {
fn deprecated(&self) -> moq_tokio::cli::Deprecated {
let mut found = moq_tokio::cli::Deprecated::default();
if self.name.is_some() {
found.flag("--name", None, "--broadcast");
}
match &self.sink {
ExportSink::Fmp4(args) | ExportSink::Mkv(args) => found.extend(args.container.deprecated()),
ExportSink::Ts(args) => found.extend(args.container.deprecated()),
ExportSink::Flv(args) | ExportSink::H264(args) | ExportSink::H265(args) => found.extend(args.deprecated()),
ExportSink::Hls(hls) => found.extend(hls.tls.deprecated()),
ExportSink::Rtmp(rtmp) if rtmp.latency_max.is_some() => {
found.flag("--latency-max", None, "--max-age");
}
_ => {}
}
found
}
}
#[derive(usage::Subcommands, Clone)]
pub enum ExportSink {
Fmp4(Fragmented),
Mkv(Fragmented),
Ts(Transport),
Flv(Container),
H264(Container),
H265(Container),
Hls(crate::hls::ExportArgs),
Rtmp(crate::rtmp::ExportArgs),
Srt(crate::srt::Args),
Rtc(crate::rtc::Args),
}
impl ExportSink {
pub fn is_stdout(&self) -> bool {
self.stdout().is_some()
}
pub fn stdout(&self) -> Option<Stdout> {
let container = |format, container: &Container| Stdout {
format,
max_age: container.max_age.into_std(),
fragment_duration: None,
mux_rate: None,
};
Some(match self {
Self::Fmp4(args) => Stdout {
fragment_duration: args.fragment_duration.map(crate::duration::Duration::into_std),
..container(SubscribeFormat::Fmp4, &args.container)
},
Self::Mkv(args) => Stdout {
fragment_duration: args.fragment_duration.map(crate::duration::Duration::into_std),
..container(SubscribeFormat::Mkv, &args.container)
},
Self::Ts(args) => Stdout {
mux_rate: args.mux_rate,
..container(SubscribeFormat::Ts, &args.container)
},
Self::Flv(args) => container(SubscribeFormat::Flv, args),
Self::H264(args) => container(SubscribeFormat::H264, args),
Self::H265(args) => container(SubscribeFormat::H265, args),
_ => return None,
})
}
}
pub struct Stdout {
pub format: SubscribeFormat,
pub max_age: Duration,
pub fragment_duration: Option<Duration>,
pub mux_rate: Option<u64>,
}
#[derive(usage::Args, Clone)]
#[usage(unknown_flags = "error", args_override_self = false)]
pub struct Container {
#[usage(long, default = "500ms")]
pub max_age: crate::duration::Duration,
#[usage(long = "latency-max", hide = true)]
latency_max: Option<crate::duration::Duration>,
}
impl Container {
fn deprecated(&self) -> moq_tokio::cli::Deprecated {
let mut found = moq_tokio::cli::Deprecated::default();
if self.latency_max.is_some() {
found.flag("--latency-max", None, "--max-age");
}
found
}
}
#[derive(usage::Args, Clone)]
#[usage(unknown_flags = "error", args_override_self = false)]
pub struct Transport {
#[usage(flatten)]
pub container: Container,
#[usage(long)]
pub mux_rate: Option<u64>,
}
#[derive(usage::Args, Clone)]
#[usage(unknown_flags = "error", args_override_self = false)]
pub struct Fragmented {
#[usage(flatten)]
pub container: Container,
#[usage(long)]
pub fragment_duration: Option<crate::duration::Duration>,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn valid() {
let _ = Cli::to_kdl();
}
#[test]
fn valid_stage() {
let _ = Stage::to_kdl();
}
#[test]
fn single_stage() {
let cli = Invocation::try_parse_from(["moq", "--connect", "http://relay", "import", "ts"]).unwrap();
assert_eq!(cli.stages.len(), 1);
assert_eq!(cli.stages[0].name(), "import");
assert!(cli.validate().is_ok());
}
#[test]
fn released_spellings_are_refused_with_a_migration() {
let Err(err) = Invocation::try_parse_from([
"moq",
"--client-connect",
"http://relay/anon",
"--client-connect-timeout",
"9s",
"--client-tls-fingerprint",
"abcd1234",
"--client-quic-gso=false",
"--server-bind",
"[::]:4443",
"--server-tcp-bind",
"127.0.0.1:4444",
"export",
"ts",
]) else {
panic!("a released spelling must not start a run");
};
let reported = err.to_string();
for line in [
"--client-connect / MOQ_CLIENT_CONNECT -> --connect / MOQ_CONNECT",
"--client-connect-timeout / MOQ_CLIENT_CONNECT_TIMEOUT -> --connect-timeout / MOQ_CONNECT_TIMEOUT",
"--client-tls-fingerprint / MOQ_CLIENT_TLS_FINGERPRINT -> --connect-tls-fingerprint / MOQ_CONNECT_TLS_FINGERPRINT",
"--client-quic-gso / MOQ_CLIENT_QUIC_GSO -> --quic-gso / MOQ_QUIC_GSO",
"--server-bind / MOQ_SERVER_BIND -> --listen / MOQ_LISTEN",
"--server-tcp-bind / MOQ_SERVER_TCP_BIND -> --listen-tcp-bind / MOQ_LISTEN_TCP_BIND",
] {
assert!(reported.contains(line), "missing {line:?} from {reported}");
}
}
#[test]
fn the_released_origin_spelling_is_refused_with_a_migration() {
let Err(err) = Invocation::try_parse_from([
"moq",
"--origin",
"42",
"--connect",
"http://relay/anon",
"export",
"ts",
]) else {
panic!("--origin must not start a run");
};
let reported = err.to_string();
assert!(
reported.contains("--origin / MOQ_ORIGIN -> --hop / MOQ_HOP"),
"missing the migration from {reported}"
);
}
#[test]
fn a_stage_local_released_spelling_is_refused() {
let Err(err) = Invocation::try_parse_from([
"moq",
"--connect",
"http://relay/anon",
"--broadcast",
"room",
"export",
"hls",
"--tls-cert",
"/tmp/cert.pem",
"--server-tls-root",
"/tmp/ca.pem",
]) else {
panic!("a released spelling on a stage must not start a run");
};
let reported = err.to_string();
assert!(reported.contains("--tls-cert"), "{reported}");
assert!(reported.contains("--listen-tls-cert"), "{reported}");
assert!(reported.contains("--listen-tls-root"), "{reported}");
}
#[test]
fn one_released_spelling_is_enough_to_refuse() {
let Err(err) = Invocation::try_parse_from([
"moq",
"--connect",
"http://relay/anon",
"--client-quic-gso=false",
"export",
"ts",
]) else {
panic!("a current spelling alongside a released one must not excuse it");
};
assert!(err.to_string().contains("--quic-gso"), "{err}");
}
#[test]
fn tcp_only_listener_is_a_moq_side() {
let cli = Invocation::try_parse_from([
"moq",
"--listen-tcp-bind",
"127.0.0.1:0",
"--auth-public",
"**",
"import",
"ts",
])
.expect("parse");
assert!(cli.moq.validate().is_ok());
assert!(cli.moq.serves());
assert_eq!(cli.moq.server_config().tcp.bind, Some("127.0.0.1:0".parse().unwrap()));
}
#[cfg(unix)]
#[test]
fn unix_only_listener_is_a_moq_side() {
let cli = Invocation::try_parse_from([
"moq",
"--listen-unix-bind",
"/tmp/moq-cli.sock",
"--auth-public",
"**",
"export",
"ts",
])
.expect("parse");
assert!(cli.moq.validate().is_ok());
assert!(cli.moq.serves());
assert_eq!(
cli.moq.server_config().unix.bind.as_deref(),
Some(std::path::Path::new("/tmp/moq-cli.sock"))
);
}
#[test]
fn a_listener_needs_exactly_one_auth_source() {
let cli =
Invocation::try_parse_from(["moq", "--listen-tcp-bind", "127.0.0.1:0", "import", "ts"]).expect("parse");
let err = cli.moq.validate().unwrap_err().to_string();
assert!(err.contains("--auth-url or --auth-public"), "{err}");
let cli = Invocation::try_parse_from([
"moq",
"--listen-tcp-bind",
"127.0.0.1:0",
"--auth-url",
"http://127.0.0.1:4440/",
"--auth-public",
"**",
"import",
"ts",
])
.expect("parse");
assert!(cli.moq.validate().is_err());
let cli = Invocation::try_parse_from([
"moq",
"--listen-tcp-bind",
"127.0.0.1:0",
"--auth-url",
"http://127.0.0.1:4440/",
"import",
"ts",
])
.expect("parse");
assert!(cli.moq.validate().is_ok());
let cli = Invocation::try_parse_from(["moq", "--auth-public", "**", "auth", "generate"]).expect("parse");
assert!(
cli.moq
.reject("auth")
.unwrap_err()
.to_string()
.contains("--auth-public")
);
}
#[test]
fn multiple_stages() {
let cli = Invocation::try_parse_from([
"moq",
"--connect",
"http://localhost:4444/event",
"import",
"--broadcast",
"cam1.hang",
"rtmp",
"--listen",
"0.0.0.0:1935",
"--",
"import",
"--broadcast",
"cam2.hang",
"rtmp",
"--listen",
"0.0.0.0:1936",
"--",
"export",
"--broadcast",
"cam1.hang",
"hls",
"--listen",
"0.0.0.0:8080",
])
.unwrap();
assert!(cli.validate().is_ok());
assert_eq!(cli.stages.len(), 3);
assert_eq!(
cli.moq.client.url.as_ref().map(ToString::to_string).as_deref(),
Some("http://localhost:4444/event")
);
let names: Vec<String> = cli.stages.iter().map(|stage| stage.broadcast(&cli.moq)).collect();
assert_eq!(names, ["cam1.hang", "cam2.hang", "cam1.hang"]);
assert_eq!(cli.stages[2].name(), "export");
}
#[test]
fn broadcast_falls_back_to_the_global() {
let cli = Invocation::try_parse_from([
"moq",
"--connect",
"http://relay",
"--broadcast",
"room.hang",
"import",
"ts",
"--",
"export",
"--broadcast",
"other.hang",
"fmp4",
])
.unwrap();
assert_eq!(cli.stages[0].broadcast(&cli.moq), "room.hang");
assert_eq!(cli.stages[1].broadcast(&cli.moq), "other.hang");
}
#[test]
fn broadcast_defaults_to_root() {
let cli = Invocation::try_parse_from(["moq", "--connect", "http://relay", "import", "ts"]).unwrap();
assert_eq!(cli.stages[0].broadcast(&cli.moq), "");
}
#[test]
fn rejects_unstageable_verbs() {
let cli = Invocation::try_parse_from([
"moq",
"--connect",
"http://relay",
"import",
"ts",
"--",
"auth",
"generate",
"--algorithm",
"ES256",
])
.unwrap();
let err = cli.validate().unwrap_err().to_string();
assert!(err.contains("auth"), "{err}");
}
#[test]
fn stage_errors_are_parse_errors() {
let Err(err) = Invocation::try_parse_from([
"moq",
"--connect",
"http://relay",
"import",
"ts",
"--",
"import",
"rtmp",
"--bogus",
]) else {
panic!("expected a parse error")
};
assert_eq!(err.kind(), ParseErrorKind::UnknownArgument);
}
#[test]
fn a_dash_prefixed_path_is_written_relative() {
let cli =
Invocation::try_parse_from(["moq", "--connect", "http://relay", "import", "hls", "./-odd.m3u8"]).unwrap();
let Command::Import(import) = &cli.stages[0] else {
panic!("expected import")
};
let ImportSource::Hls(hls) = &import.source else {
panic!("expected hls")
};
assert_eq!(hls.playlist, "./-odd.m3u8");
}
#[test]
fn rejects_an_empty_stage() {
for argv in [
vec!["moq", "--connect", "http://relay", "import", "ts", "--"],
vec![
"moq",
"--connect",
"http://relay",
"import",
"ts",
"--",
"--",
"export",
"fmp4",
],
] {
let Err(err) = Invocation::try_parse_from(argv.clone()) else {
panic!("expected a parse error for {argv:?}")
};
assert_eq!(err.kind(), ParseErrorKind::MissingSubcommand);
assert!(err.to_string().contains("must be followed by"), "{err}");
}
}
#[test]
fn stages_reject_globals() {
let Err(err) = Invocation::try_parse_from([
"moq",
"--connect",
"http://relay",
"import",
"ts",
"--",
"--connect",
"http://other",
"import",
"fmp4",
]) else {
panic!("expected a parse error")
};
assert_eq!(err.kind(), ParseErrorKind::UnknownArgument);
}
#[test]
fn imports_without_rate_control_can_share_a_connection() {
let cli = Invocation::try_parse_from([
"moq",
"--connect",
"http://relay",
"import",
"--broadcast",
"a.hang",
"rtmp",
"--listen",
"127.0.0.1:1935",
"--",
"import",
"--broadcast",
"b.hang",
"srt",
"--listen",
"127.0.0.1:9000",
])
.unwrap();
assert!(cli.validate().is_ok());
}
#[cfg(feature = "capture")]
#[test]
fn two_encoding_stages_may_share_a_connection() {
let cli = Invocation::try_parse_from([
"moq",
"--connect",
"http://relay",
"import",
"capture",
"--",
"import",
"capture",
])
.unwrap();
assert!(cli.validate().is_ok());
}
#[test]
fn max_age_is_unset_unless_asked_for() {
let cli = Invocation::try_parse_from(["moq", "import", "ts"]).unwrap();
let Command::Import(import) = &cli.stages[0] else {
panic!("expected import")
};
assert_eq!(import.max_age, None);
let cli = Invocation::try_parse_from(["moq", "import", "--max-age", "5s", "ts"]).unwrap();
let Command::Import(import) = &cli.stages[0] else {
panic!("expected import")
};
assert_eq!(import.max_age, Some(Duration::from_secs(5).into()));
let cli =
Invocation::try_parse_from(["moq", "import", "--max-age", "5s", "rtmp", "--listen", "127.0.0.1:1935"])
.unwrap();
let Command::Import(import) = &cli.stages[0] else {
panic!("expected import")
};
assert_eq!(import.max_age, Some(Duration::from_secs(5).into()));
}
#[test]
fn latency_max_is_refused_with_a_migration() {
let Err(err) = Invocation::try_parse_from(["moq", "import", "--latency-max", "5s", "ts"]) else {
panic!("--latency-max must not start a run");
};
assert!(err.to_string().contains("--latency-max -> --max-age"), "{}", err);
let Err(err) = Invocation::try_parse_from(["moq", "export", "--broadcast", "b", "ts", "--latency-max", "1s"])
else {
panic!("--latency-max must not start a run");
};
assert!(err.to_string().contains("--latency-max -> --max-age"), "{}", err);
}
#[test]
fn the_released_name_and_verb_spellings_are_refused() {
let Err(err) = Invocation::try_parse_from(["moq", "--name", "room", "import", "ts"]) else {
panic!("--name must not start a run");
};
assert!(err.to_string().contains("--name -> --broadcast"), "{err}");
let Err(err) = Invocation::try_parse_from(["moq", "publish", "ts"]) else {
panic!("publish must not start a run");
};
assert!(err.to_string().contains("publish -> import"), "{err}");
let Err(err) = Invocation::try_parse_from(["moq", "subscribe", "ts"]) else {
panic!("subscribe must not start a run");
};
assert!(err.to_string().contains("subscribe -> export"), "{err}");
}
#[test]
fn auth_verb() {
let cli = Invocation::try_parse_from(["moq", "auth", "generate", "--algorithm", "ES256"]).unwrap();
assert!(matches!(cli.stages[0], Command::Auth(_)));
assert!(cli.moq.validate().is_err());
assert!(cli.moq.reject("auth").is_ok());
for (flag, value, reported) in [
("--connect", "https://relay.example.com", "--connect"),
("--listen-tcp-bind", "127.0.0.1:0", "--listen-tcp-bind"),
("--broadcast", "room", "--broadcast"),
("--cluster-connect", "https://relay.example", "--cluster-connect"),
(
"--cluster-connect-api",
"https://api.example/peers",
"--cluster-connect-api",
),
("--cluster-node", "https://self.example", "--cluster-node"),
("--cluster-token", "cluster.jwt", "--cluster-token"),
("--cluster-id", "1", "--cluster-id"),
("--cluster-tier", "internal", "--cluster-tier"),
] {
let cli = Invocation::try_parse_from(["moq", flag, value, "auth", "generate"]).unwrap();
let err = cli.moq.reject("auth").unwrap_err().to_string();
assert!(err.contains(reported), "{err}");
}
let cli = Invocation::try_parse_from(["moq", "--cluster-mesh", "auth", "generate"]).unwrap();
let err = cli.moq.reject("auth").unwrap_err().to_string();
assert!(err.contains("--cluster-mesh"), "{err}");
#[cfg(unix)]
{
for (flag, value, reported) in [
("--listen-unix-bind", "/tmp/moq-cli.sock", "--listen-unix-bind"),
("--listen-unix-allow-uid", "1000", "--listen-unix-allow-uid"),
("--listen-unix-allow-gid", "1000", "--listen-unix-allow-gid"),
("--listen-unix-allow-pid", "1000", "--listen-unix-allow-pid"),
] {
let cli = Invocation::try_parse_from(["moq", flag, value, "auth", "generate"]).unwrap();
let err = cli.moq.reject("auth").unwrap_err().to_string();
assert!(err.contains(reported), "{err}");
}
}
#[cfg(feature = "cluster-lan")]
{
let cli = Invocation::try_parse_from(["moq", "--cluster-lan", "auth", "generate"]).unwrap();
let err = cli.moq.reject("auth").unwrap_err().to_string();
assert!(err.contains("--cluster-lan"), "{err}");
let cli = Invocation::try_parse_from([
"moq",
"--cluster-lan=false",
"--cluster-lan-secret",
"cluster.key",
"auth",
"generate",
])
.unwrap();
let err = cli.moq.reject("auth").unwrap_err().to_string();
assert!(err.contains("--cluster-lan-secret"), "{err}");
let cli = Invocation::try_parse_from([
"moq",
"--cluster-lan=false",
"--cluster-lan-app",
"custom",
"auth",
"generate",
])
.unwrap();
let err = cli.moq.reject("auth").unwrap_err().to_string();
assert!(err.contains("--cluster-lan-app"), "{err}");
}
}
#[test]
fn cluster_connect_is_a_moq_side() {
let cli = Invocation::try_parse_from(["moq", "--cluster-connect", "https://relay.example", "import", "ts"])
.expect("parse");
assert!(cli.moq.validate().is_ok(), "a cluster dial is a MoQ side on its own");
let cli = Invocation::try_parse_from([
"moq",
"--cluster-connect-api",
"https://api.example/peers",
"import",
"ts",
])
.expect("parse");
assert!(cli.moq.validate().is_ok(), "a cluster API is a MoQ side on its own");
let cli = Invocation::try_parse_from(["moq", "--cluster-node", "https://self.example", "import", "ts"])
.expect("parse");
let err = cli.moq.validate().unwrap_err().to_string();
assert!(err.contains("MoQ side"), "{err}");
}
#[cfg(feature = "cluster-lan")]
#[test]
fn cluster_lan_is_a_moq_side_and_fills_in_a_listener() {
let cli = Invocation::try_parse_from(["moq", "--cluster-lan", "import", "ts"]).expect("parse");
assert!(cli.moq.lan());
assert!(cli.moq.validate().is_ok(), "the LAN mesh is a MoQ side on its own");
let server = cli.moq.server_config();
assert_eq!(
server.bind.as_ref().map(ToString::to_string).as_deref(),
Some("[::]:0"),
"an ephemeral port"
);
assert_eq!(server.tls.generate, ["moq-cluster-lan"], "a generated certificate");
let cli = Invocation::try_parse_from([
"moq",
"--cluster-lan",
"--listen",
"[::]:4443",
"--listen-tls-generate",
"localhost",
"import",
"ts",
])
.expect("parse");
let server = cli.moq.server_config();
assert_eq!(
server.bind.as_ref().map(ToString::to_string).as_deref(),
Some("[::]:4443")
);
assert_eq!(server.tls.generate, ["localhost"]);
let cli = Invocation::try_parse_from(["moq", "--connect", "https://relay.example.com", "import", "ts"])
.expect("parse");
assert!(!cli.moq.lan());
assert_eq!(cli.moq.server_config().bind, None);
}
#[cfg(feature = "cluster-lan")]
#[test]
fn cluster_lan_secret_requires_the_mesh() {
for lan in [None, Some("--cluster-lan=false")] {
let mut argv = vec!["moq"];
if let Some(lan) = lan {
argv.push(lan);
}
argv.extend([
"--cluster-lan-secret",
"cluster.key",
"--connect",
"https://relay.example.com",
"import",
"ts",
]);
let cli = Invocation::try_parse_from(argv).expect("parse");
let err = cli.moq.validate().unwrap_err().to_string();
assert!(err.contains("--cluster-lan=true"), "{err}");
}
let cli = Invocation::try_parse_from([
"moq",
"--cluster-lan",
"--cluster-lan-secret",
"cluster.key",
"import",
"ts",
])
.expect("parse");
assert!(cli.moq.validate().is_ok());
assert_eq!(cli.moq.cluster.lan.secret.as_deref(), Some("cluster.key"));
}
#[cfg(feature = "cluster-lan")]
#[test]
fn cluster_lan_app_requires_the_mesh_and_defaults() {
for lan in [None, Some("--cluster-lan=false")] {
let mut argv = vec!["moq"];
if let Some(lan) = lan {
argv.push(lan);
}
argv.extend([
"--cluster-lan-app",
"custom",
"--connect",
"https://relay.example.com",
"import",
"ts",
]);
let cli = Invocation::try_parse_from(argv).expect("parse");
let err = cli.moq.validate().unwrap_err().to_string();
assert!(err.contains("--cluster-lan=true"), "{err}");
}
let cli = Invocation::try_parse_from(["moq", "--cluster-lan", "import", "ts"]).expect("parse");
assert_eq!(cli.moq.cluster.lan.app.clone().unwrap_or_default().as_str(), "default");
let cli = Invocation::try_parse_from(["moq", "--cluster-lan", "--cluster-lan-app", "custom", "import", "ts"])
.expect("parse");
assert!(cli.moq.validate().is_ok());
assert_eq!(
cli.moq.cluster.lan.app.as_ref().map(ToString::to_string).as_deref(),
Some("custom")
);
let err = Invocation::try_parse_from(["moq", "--cluster-lan", "--cluster-lan-app", "Default", "import", "ts"])
.err()
.expect("uppercase must not parse")
.to_string();
assert!(err.contains("app") || err.contains("Default"), "{err}");
}
#[cfg(feature = "cluster-lan")]
#[test]
fn cluster_lan_requires_a_path_capable_version() {
for (flag, reported) in [
("--connect-version", "--connect-version"),
("--listen-version", "--listen-version"),
] {
let cli = Invocation::try_parse_from(["moq", "--cluster-lan", flag, "moq-lite-04", "import", "ts"])
.expect("parse");
let err = cli.moq.validate().unwrap_err().to_string();
assert!(err.contains(reported), "{flag}: {err}");
}
let cli = Invocation::try_parse_from([
"moq",
"--cluster-lan",
"--connect-version",
"moq-lite-04",
"--connect-version",
"moq-lite-05",
"--listen-version",
"moq-lite-05",
"import",
"ts",
])
.expect("parse");
assert!(cli.moq.validate().is_ok());
}
#[cfg(feature = "play")]
#[test]
fn play_verb() {
let cli = Invocation::try_parse_from([
"moq",
"--connect",
"https://relay.example.com/anon",
"--broadcast",
"room.hang",
"play",
"--video-name",
"hd",
"--audio-codec",
"aac",
])
.unwrap();
let Command::Play(play) = &cli.stages[0] else {
panic!("expected play")
};
assert_eq!(play.delay, Duration::from_millis(100).into());
assert_eq!(play.select.video_name.as_deref(), Some("hd"));
assert!(cli.moq.validate().is_ok());
assert!(play.validate().is_ok());
}
#[cfg(feature = "play")]
#[test]
fn play_rejects_undecodable_codecs() {
for codec in ["vp8", "vp9"] {
let cli = Invocation::try_parse_from([
"moq",
"--connect",
"https://relay.example.com/anon",
"play",
"--video-codec",
codec,
])
.unwrap();
let Command::Play(play) = &cli.stages[0] else {
panic!("expected play")
};
let err = play.validate().unwrap_err().to_string();
assert!(err.contains(codec), "{err}");
}
let cli = Invocation::try_parse_from([
"moq",
"--connect",
"https://relay.example.com/anon",
"play",
"--audio-codec",
"aac",
])
.unwrap();
let Command::Play(play) = &cli.stages[0] else {
panic!("expected play")
};
assert!(play.validate().is_ok());
}
#[test]
fn help_and_version_render_their_output() {
for args in [
vec!["moq", "--help"],
vec!["moq", "-h"],
vec!["moq", "--version"],
vec!["moq", "-V"],
vec!["moq", "publish", "--help"],
] {
let Err(err) = Invocation::try_parse_from(args.clone()) else {
panic!("{args:?} parsed instead of asking a question")
};
assert!(
matches!(err.kind(), ParseErrorKind::DisplayHelp | ParseErrorKind::DisplayVersion),
"{args:?} produced {:?}",
err.kind()
);
assert!(!err.to_string().trim().is_empty(), "{args:?} printed nothing");
}
}
#[test]
fn stage_help_renders() {
for args in [
vec![
"moq",
"--connect",
"http://localhost:4444/x",
"import",
"fmp4",
"--",
"export",
"--help",
],
vec![
"moq",
"--connect",
"http://localhost:4444/x",
"import",
"fmp4",
"--",
"export",
"fmp4",
"--help",
],
] {
let Err(err) = Invocation::try_parse_from(args.clone()) else {
panic!("{args:?} parsed instead of asking a question")
};
assert_eq!(err.kind(), ParseErrorKind::DisplayHelp, "{args:?}");
assert!(
err.to_string().contains("Usage:"),
"{args:?} rendered no help page: {err}"
);
}
}
#[test]
fn version_choices_match_the_parser() {
fn walk<'a>(cmd: &'a usage::argv::spec::CommandMeta<'a>, found: &mut Vec<(&'a str, Vec<&'a str>)>) {
for flag in cmd.flags {
let Some(long) = flag.flag.longs.first() else {
continue;
};
if long.ends_with("version") && !flag.choices.is_empty() {
found.push((long, flag.choices.to_vec()));
}
}
for sub in cmd.subcommands {
walk(sub, found);
}
}
let mut expected: Vec<&str> = hang::moq_net::Version::names().collect();
expected.sort_unstable();
let mut found = Vec::new();
walk(Cli::spec().root, &mut found);
assert!(!found.is_empty(), "no version flag carried a choice list");
for (long, choices) in &mut found {
choices.sort_unstable();
assert_eq!(choices, &expected, "--{long} is out of step with Version::names()");
}
}
}