use std::cell::RefCell;
use std::collections::BTreeSet;
use std::ffi::{OsStr, OsString};
#[cfg(feature = "capture")]
use std::future::Future;
use std::time::Duration;
use anyhow::Context;
use hang::moq_net;
use moq_mux::catalog::{CatalogFormat, Stream};
use tokio::time::{Instant, timeout_at};
use usage::complete::{Candidate, CompleteCtx, CompletionFuture, CompletionOverlay, CompletionRequest, Shell, render};
use usage::spec::{CommandArgs, ValueEnum};
use crate::args::{Cli, Environment, Export, MoqSide, Stage};
use crate::subscribe::CatalogFormatArg;
const BUDGET: Duration = Duration::from_millis(500);
const CEILING: Duration = Duration::from_millis(1_500);
const SETTLE: Duration = Duration::from_millis(30);
static NETWORK: &[CompletionOverlay<'static>] = &[
CompletionOverlay::async_any("BROADCAST", broadcasts),
CompletionOverlay::async_any("VIDEO_NAME", video_names),
CompletionOverlay::async_any("AUDIO_NAME", audio_names),
];
#[cfg(feature = "capture")]
const CAPTURE_PATH: &str = "import capture";
#[cfg(feature = "capture")]
static CAPTURE: &[CompletionOverlay<'static>] = &[
CompletionOverlay::asynchronous(CAPTURE_PATH, "CAMERA", cameras),
CompletionOverlay::asynchronous(CAPTURE_PATH, "DISPLAY", displays),
CompletionOverlay::asynchronous(CAPTURE_PATH, "WINDOW", windows),
CompletionOverlay::asynchronous(CAPTURE_PATH, "APP", apps),
CompletionOverlay::asynchronous(CAPTURE_PATH, "MICROPHONE", microphones),
];
#[cfg(not(feature = "capture"))]
static CAPTURE: &[CompletionOverlay<'static>] = &[];
fn overlays() -> Vec<CompletionOverlay<'static>> {
NETWORK.iter().chain(CAPTURE).copied().collect()
}
thread_local! {
static GLOBALS: RefCell<Option<MoqSide>> = const { RefCell::new(None) };
}
struct Globals;
impl Globals {
fn set(side: Option<MoqSide>) -> Self {
GLOBALS.with_borrow_mut(|slot| *slot = side);
Self
}
fn get() -> Option<MoqSide> {
GLOBALS.with_borrow(Clone::clone)
}
}
impl Drop for Globals {
fn drop(&mut self) {
GLOBALS.with_borrow_mut(|slot| *slot = None);
}
}
pub async fn answer(argv: &[OsString]) -> Option<String> {
let request = CompletionRequest::parse(argv)?;
let _globals = Globals::set(globals(&request));
let overlays = overlays();
let words = request.split.walked();
let staged = words[..request.split.cword]
.iter()
.rposition(|word| word == "--")
.map(|at| at + 1);
let answer = timeout_at(Instant::now() + CEILING, complete(&request, staged, &overlays))
.await
.unwrap_or_default();
Some(render(&answer, request.shell))
}
async fn complete<'a>(
request: &CompletionRequest,
staged: Option<usize>,
overlays: &'a [CompletionOverlay<'a>],
) -> usage::complete::Completions<'a> {
let words = request.split.walked();
match staged {
None => {
Cli::app()
.completion_app()
.completions(overlays)
.complete_request(request)
.await
}
Some(start) => {
let mut chunk = vec![words.first().cloned().unwrap_or_default()];
chunk.extend_from_slice(&words[start..]);
let mut request = request.clone();
request.split.cword = chunk.len() - 1;
request.split.words = chunk;
Stage::app()
.completion_app()
.completions(overlays)
.complete_request(&request)
.await
}
}
}
fn globals(request: &CompletionRequest) -> Option<MoqSide> {
let chunk = request.split.argv().split(|word| word == "--").next()?;
let argv: Vec<&OsStr> = chunk.iter().map(OsStr::new).collect();
let typed = MoqSide::from_argv(&argv, Environment::Ignore)?;
let mut side = MoqSide::from_argv(&argv, Environment::Read)?;
if typed.client.url.is_none() {
side.client.url = None;
}
Some(side)
}
fn partial<T: CommandArgs>(ctx: &CompleteCtx<'_>) -> Option<T::Partial> {
let (command, words) = ctx.command_for(T::COMMAND)?;
let argv: Vec<&OsStr> = words.iter().map(OsStr::new).collect();
let mut partial = T::start();
let mut parser = usage::Parser::new(command, &argv);
while let Some(event) = parser.next_event() {
match event {
Ok(event) => {
T::apply(&mut partial, &event);
}
Err(_) => break,
}
}
Some(partial)
}
fn given(value: &Option<Vec<u8>>) -> Option<&str> {
std::str::from_utf8(value.as_deref()?).ok()
}
fn choice<T: ValueEnum>(value: &Option<Vec<u8>>) -> Option<T> {
T::from_choice(given(value)?)
}
fn catalog_format(ctx: &CompleteCtx<'_>, export: Option<&<Export as CommandArgs>::Partial>) -> Option<CatalogFormat> {
let named = export.and_then(|export| choice::<CatalogFormatArg>(&export.catalog_format));
#[cfg(feature = "play")]
let named = named.or_else(|| {
partial::<crate::play::Args>(ctx).and_then(|play| choice::<CatalogFormatArg>(&play.catalog_format))
});
#[cfg(not(feature = "play"))]
let _ = ctx;
named.map(Into::into)
}
#[derive(usage::ValueEnum, Clone, Copy)]
pub enum ShellArg {
Bash,
Elvish,
Fish,
Nu,
#[usage(name = "powershell")]
PowerShell,
Zsh,
}
impl From<ShellArg> for Shell {
fn from(shell: ShellArg) -> Self {
match shell {
ShellArg::Bash => Self::Bash,
ShellArg::Elvish => Self::Elvish,
ShellArg::Fish => Self::Fish,
ShellArg::Nu => Self::Nu,
ShellArg::PowerShell => Self::PowerShell,
ShellArg::Zsh => Self::Zsh,
}
}
}
#[derive(usage::Args, Clone)]
#[usage(unknown_flags = "error", args_override_self = false)]
pub struct Args {
#[usage(arg, value_enum)]
pub shell: ShellArg,
#[usage(long)]
pub install: bool,
#[usage(long, requires = "--install")]
pub force: bool,
}
impl Args {
pub fn run(self) -> anyhow::Result<()> {
let shell = self.shell.into();
if !self.install {
print!("{}", Cli::completion_script(shell));
return Ok(());
}
let on_foreign = match self.force {
true => usage::install::OnForeign::Overwrite,
false => usage::install::OnForeign::Refuse,
};
let installed = Cli::install_completion(shell, &usage::install::Env::from_process(), on_foreign)
.context("failed to install the completion script")?;
let wrote = match installed.wrote {
usage::install::Wrote::Created => "wrote",
usage::install::Wrote::Unchanged => "already current",
usage::install::Wrote::Updated => "updated",
usage::install::Wrote::Replaced => "replaced",
_ => "installed",
};
eprintln!("{wrote} {}", installed.plan.path.display());
if let usage::install::Loading::Manual { line, file, why } = &installed.plan.loading {
eprintln!("\nadd this to {file}:");
for line in line.lines() {
eprintln!(" {line}");
}
eprintln!("\n{why}");
}
Ok(())
}
}
#[cfg(feature = "capture")]
async fn sources<T, E>(
found: impl Future<Output = Result<Vec<T>, E>>,
describe: impl Fn(&T) -> Candidate<'static>,
) -> Vec<Candidate<'static>> {
match timeout_at(Instant::now() + BUDGET, found).await {
Ok(Ok(items)) => items.iter().map(describe).collect(),
Ok(Err(_)) | Err(_) => Vec::new(),
}
}
#[cfg(feature = "capture")]
fn cameras(_ctx: CompleteCtx<'_>) -> CompletionFuture<'static> {
Box::pin(async move {
sources(moq_video::capture::cameras(), |camera| {
Candidate::described(camera.id.clone(), camera.name.clone())
})
.await
})
}
#[cfg(feature = "capture")]
fn displays(_ctx: CompleteCtx<'_>) -> CompletionFuture<'static> {
Box::pin(async move {
sources(moq_video::capture::displays(), |display| {
Candidate::described(
display.id.clone(),
format!("{} ({}x{})", display.name, display.width, display.height),
)
})
.await
})
}
#[cfg(feature = "capture")]
fn windows(_ctx: CompleteCtx<'_>) -> CompletionFuture<'static> {
Box::pin(async move {
sources(moq_video::capture::windows(), |window| {
let title = if window.title.is_empty() {
"(untitled)"
} else {
&window.title
};
Candidate::described(window.id.clone(), format!("{} - {title}", window.app))
})
.await
})
}
#[cfg(feature = "capture")]
fn apps(_ctx: CompleteCtx<'_>) -> CompletionFuture<'static> {
Box::pin(async move {
sources(moq_video::capture::apps(), |app| {
Candidate::described(app.id.clone(), app.name.clone())
})
.await
})
}
#[cfg(feature = "capture")]
fn microphones(_ctx: CompleteCtx<'_>) -> CompletionFuture<'static> {
Box::pin(async move {
sources(moq_audio::capture::devices(), |device| match device.default {
true => Candidate::described(device.id.clone(), "the default input"),
false => Candidate::new(device.id.clone()),
})
.await
})
}
fn broadcasts(_ctx: CompleteCtx<'_>) -> CompletionFuture<'static> {
Box::pin(async move {
let Some(side) = Globals::get() else {
return Vec::new();
};
let deadline = Instant::now() + BUDGET;
let Some((origin, connection)) = dial(&side, deadline).await else {
return Vec::new();
};
let mut announced = origin.consume().announced();
let mut live = BTreeSet::new();
let mut until = deadline;
while let Ok(Some(update)) = timeout_at(until, announced.next()).await {
until = deadline.min(Instant::now() + SETTLE);
let path = update.prefix.to_string();
if path.is_empty() {
continue;
}
match update.kind.is_active() {
true => live.insert(path),
false => live.remove(&path),
};
}
drop(connection);
live.into_iter().map(Candidate::new).collect()
})
}
fn video_names(ctx: CompleteCtx<'_>) -> CompletionFuture<'_> {
Box::pin(async move {
let Some(catalog) = renditions(&ctx).await else {
return Vec::new();
};
catalog
.video
.renditions
.iter()
.map(|(name, config)| {
let size = match (config.coded_width, config.coded_height) {
(Some(width), Some(height)) => format!(" {width}x{height}"),
_ => String::new(),
};
Candidate::described(name.clone(), format!("{}{size}", config.codec))
})
.collect()
})
}
fn audio_names(ctx: CompleteCtx<'_>) -> CompletionFuture<'_> {
Box::pin(async move {
let Some(catalog) = renditions(&ctx).await else {
return Vec::new();
};
catalog
.audio
.renditions
.iter()
.map(|(name, config)| {
Candidate::described(
name.clone(),
format!("{} {} Hz {}ch", config.codec, config.sample_rate, config.channel_count),
)
})
.collect()
})
}
async fn renditions(ctx: &CompleteCtx<'_>) -> Option<moq_mux::catalog::hang::Catalog> {
let side = Globals::get()?;
let export = partial::<Export>(ctx);
let export = export.as_ref();
let path = export
.and_then(|export| given(&export.broadcast))
.or(side.broadcast.as_deref())
.unwrap_or_default()
.to_string();
let format = catalog_format(ctx, export)
.or_else(|| CatalogFormat::detect(&path))
.unwrap_or_default();
let deadline = Instant::now() + BUDGET;
let (origin, connection) = dial(&side, deadline).await?;
let catalog = timeout_at(deadline, catalog(&origin, &path, format))
.await
.ok()
.flatten();
drop(connection);
catalog
}
async fn catalog(
origin: &moq_net::origin::Producer,
path: &str,
format: CatalogFormat,
) -> Option<moq_mux::catalog::hang::Catalog> {
let consumer = origin.consume();
consumer.routed(path).await?;
let mut stream = moq_mux::Source::new(consumer, path).catalog(format).await.ok()?;
stream.next().await.ok()?
}
async fn dial(side: &MoqSide, deadline: Instant) -> Option<(moq_net::origin::Producer, moq_tokio::Connection)> {
let url = side.client.url.clone()?;
let origin = moq_tokio::origin::spawn();
let (connect, quic) = (side.client.clone(), side.quic.clone());
let client = timeout_at(deadline, tokio::task::spawn_blocking(move || connect.init(quic)))
.await
.ok()?
.ok()?
.ok()?
.with_subscriber(origin.clone())
.with_reconnect(false);
let connection = timeout_at(deadline, client.connect(url).established())
.await
.ok()?
.ok()?;
Some((origin, connection))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::test_env::EnvGuard;
async fn complete(line: &str) -> Vec<String> {
let argv: Vec<OsString> = ["__complete_word__", "--shell", "bash", "--line", line, "--cursor"]
.iter()
.map(OsString::from)
.chain(std::iter::once(OsString::from(line.len().to_string())))
.collect();
answer(&argv)
.await
.unwrap_or_default()
.lines()
.map(str::to_string)
.collect()
}
fn relay(origin: &moq_net::origin::Producer) -> String {
let _ = moq_tokio::crypto::install_default();
let mut config = moq_tokio::listen::Config::default();
config.bind = Some("127.0.0.1:0".parse().unwrap());
config.tls.generate = vec!["localhost".to_string()];
let server = config.init(Default::default()).expect("failed to bind listener");
let port = server.local_addr().expect("no local addr").port();
tokio::spawn(server.serve_publish(origin.consume()));
format!("--connect moqt://127.0.0.1:{port} --connect-tls-insecure")
}
#[tokio::test]
async fn retargets_to_the_active_stage() {
let _env = EnvGuard::clear(&["MOQ_CONNECT"]);
let staged = complete("moq --connect http://x/y import fmp4 -- export fmp4 --").await;
assert!(!staged.is_empty(), "a later stage completed nothing");
for global in ["--connect", "--hop", "--broadcast"] {
assert!(
!staged.iter().any(|candidate| candidate == global),
"{global} leaked into a stage that refuses it: {staged:?}"
);
}
let root = complete("moq --conn").await;
assert!(
root.iter().any(|candidate| candidate == "--connect"),
"root lost its globals: {root:?}"
);
assert!(complete("moq import fmp4 --").await.is_empty());
}
#[test]
fn every_overlay_matches_only_what_it_answers_for() {
let renditions = match cfg!(feature = "play") {
true => vec!["export", "play"],
false => vec!["export"],
};
let shared = [
("BROADCAST", vec!["", "import", "export"]),
("VIDEO_NAME", renditions.clone()),
("AUDIO_NAME", renditions),
];
fn declaring(command: &usage::spec::CommandMeta<'_>, value: &str, at: &str, found: &mut Vec<String>) {
if command.hide {
return;
}
let declares = command
.flags
.iter()
.any(|field| field.value_name.unwrap_or(field.flag.name).eq_ignore_ascii_case(value))
|| command
.args
.iter()
.any(|field| field.arg.name.eq_ignore_ascii_case(value));
if declares {
found.push(at.to_string());
}
for sub in command.subcommands {
let deeper = match at.is_empty() {
true => sub.cmd.name.to_string(),
false => format!("{at} {}", sub.cmd.name),
};
declaring(sub, value, &deeper, found);
}
}
for overlay in overlays() {
let mut found = Vec::new();
declaring(Cli::spec().root, overlay.value, "", &mut found);
assert!(
!found.is_empty(),
"no flag or argument takes a value named `{}`, so its completer never runs",
overlay.value
);
let Some((_, allowed)) = shared.iter().find(|(name, _)| *name == overlay.value) else {
continue;
};
found.sort();
let mut allowed: Vec<String> = allowed.iter().map(|path| path.to_string()).collect();
allowed.sort();
assert_eq!(
found, allowed,
"`{}` is answered everywhere it appears, and the set of commands declaring it changed",
overlay.value
);
}
}
#[tokio::test]
async fn the_environment_cannot_ask_for_a_moq_side() {
let origin = moq_tokio::origin::spawn();
let _alpha = origin.create_broadcast("alpha").expect("alpha");
_alpha.announce(Default::default()).expect("alpha");
let connect = relay(&origin);
let url = connect.split_whitespace().nth(1).expect("a --connect url").to_string();
let _env = EnvGuard::set(&[("MOQ_CONNECT", &url)]);
assert!(
complete("moq --connect-tls-insecure --broadcast ").await.is_empty(),
"MOQ_CONNECT authorized a dial the line never asked for"
);
assert_eq!(
complete(&format!("moq {connect} --broadcast ")).await,
["alpha"],
"a typed --connect stopped working"
);
let ambient = crate::args::Invocation::try_parse_from(["moq", "auth", "generate"]).expect("parse");
assert!(
ambient.moq.client.url.is_some(),
"the resolved side should still pick the variable up"
);
assert!(
ambient.reject("auth").is_ok(),
"an exported MOQ_CONNECT was treated as a request"
);
let typed =
crate::args::Invocation::try_parse_from(["moq", "--connect", &url, "auth", "generate"]).expect("parse");
assert!(typed.reject("auth").is_err(), "a typed --connect stopped being refused");
}
#[tokio::test]
async fn no_relay_on_the_line_means_no_dial() {
let _env = EnvGuard::clear(&["MOQ_CONNECT"]);
for line in [
"moq --broadcast ",
"moq export --broadcast ",
"moq export --video-name ",
"moq export --audio-name ",
] {
assert!(complete(line).await.is_empty(), "{line:?} completed without a relay");
}
}
#[test]
fn every_shell_choice_names_a_real_shell() {
for choice in <ShellArg as ValueEnum>::CHOICES {
let arg = ShellArg::from_choice(choice).expect("a declared choice");
assert_eq!(
Shell::from(arg).as_str(),
*choice,
"`{choice}` is not what Usage calls it"
);
}
}
#[test]
fn the_script_registers_this_binary() {
let script = Cli::completion_script(Shell::Zsh);
assert!(
script.starts_with("#compdef moq"),
"{}",
&script[..40.min(script.len())]
);
assert!(script.contains("__complete_word__"), "the script asks nothing");
}
#[tokio::test]
async fn a_relay_on_the_line_answers_broadcast() {
let _env = EnvGuard::clear(&["MOQ_CONNECT"]);
let origin = moq_tokio::origin::spawn();
let _alpha = origin.create_broadcast("alpha").expect("alpha");
_alpha.announce(Default::default()).expect("alpha");
let _nested = origin.create_broadcast("room/beta").expect("beta");
_nested.announce(Default::default()).expect("beta");
let connect = relay(&origin);
let found = complete(&format!("moq {connect} --broadcast ")).await;
assert!(found.contains(&"alpha".to_string()), "{found:?}");
assert!(found.contains(&"room/beta".to_string()), "{found:?}");
let staged = complete(&format!("moq {connect} import fmp4 -- export --broadcast ")).await;
assert_eq!(staged, found, "a later stage lost the relay the globals named");
}
#[tokio::test]
async fn a_stage_broadcast_picks_the_catalog_to_read() {
let _env = EnvGuard::clear(&["MOQ_CONNECT"]);
use hang::catalog::{AudioCodec, AudioConfig, H264, VideoConfig};
let origin = moq_tokio::origin::spawn();
let mut keep = Vec::new();
for (path, video, audio) in [("wanted", "hd", "stereo"), ("other", "sd", "mono")] {
let mut broadcast = origin.create_broadcast(path).expect("broadcast");
broadcast.announce(Default::default()).expect("broadcast");
let mut catalog =
moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).expect("catalog");
let mut edit = catalog.modify().unwrap();
edit.video.renditions.insert(
video.to_string(),
VideoConfig::new(H264 {
profile: 0x42,
constraints: 0,
level: 0x1e,
inline: false,
}),
);
edit.audio
.renditions
.insert(audio.to_string(), AudioConfig::new(AudioCodec::Opus, 48_000, 2));
edit.commit().expect("publish the catalog");
keep.push((broadcast, catalog));
}
let connect = relay(&origin);
let line = format!("moq {connect} --broadcast other export --broadcast wanted");
assert_eq!(complete(&format!("{line} --video-name ")).await, ["hd"]);
assert_eq!(complete(&format!("{line} --audio-name ")).await, ["stereo"]);
}
#[tokio::test]
async fn the_catalog_format_on_the_line_is_honored() {
let _env = EnvGuard::clear(&["MOQ_CONNECT"]);
let origin = moq_tokio::origin::spawn();
let broadcast = origin.create_broadcast("room").expect("broadcast");
broadcast.announce(Default::default()).expect("broadcast");
let track = broadcast
.create_track(moq_msf::DEFAULT_NAME, moq_net::track::Info::default())
.expect("msf track");
let mut msf = moq_msf::Track::new("hd", moq_msf::Packaging::Loc);
msf.role = Some(moq_msf::Role::Video);
msf.codec = Some("avc1.42001e".to_string());
let catalog = moq_msf::Catalog::new(vec![msf]).to_json().expect("msf json");
let mut group = track.append_group().expect("group");
group.write_frame(moq_net::Timestamp::now(), catalog).expect("frame");
let connect = relay(&origin);
let line = format!("moq {connect} --broadcast room");
assert!(complete(&format!("{line} export --video-name ")).await.is_empty());
assert_eq!(
complete(&format!("{line} export --catalog-format msf --video-name ")).await,
["hd"],
"export ignored its own --catalog-format"
);
#[cfg(feature = "play")]
{
assert!(complete(&format!("{line} play --video-name ")).await.is_empty());
assert_eq!(
complete(&format!("{line} play --catalog-format msf --video-name ")).await,
["hd"],
"play ignored its own --catalog-format"
);
}
}
}