use std::{
collections::HashMap,
path::PathBuf,
sync::{
Arc, Mutex,
atomic::{AtomicU64, Ordering},
},
};
use async_trait::async_trait;
use clap::{CommandFactory, Parser};
use indicatif::{HumanBytes, MultiProgress, ProgressBar, ProgressState, ProgressStyle};
use odl::{
Download,
config::Config,
conflict::{
FileChangedResolution, FinalFileExistsResolution, NotResumableResolution,
SameDownloadExistsResolution, SaveConflictResolver, ServerConflictResolver,
},
credentials::Credentials,
download_manager::{DownloadManager, DownloadRequest, EvaluateRequest},
engine::EnginePreference,
error::OdlError,
format::{
DefaultFormatSelector, FixedFormatSelector, FormatOffer, FormatSelector, QualityTier,
},
progress::{
AsyncReporter, CancellationToken, DownloadContext, Phase, ProgressEvent, ProgressReporter,
SAMPLE_INTERVAL,
},
};
use reqwest::Url;
use serde_json::json;
use std::io::IsTerminal;
use std::process::ExitCode;
use tokio::{self, io::AsyncBufReadExt};
mod args;
mod json;
use args::{Args, LogLevel, OutputFormat};
use json::JsonReporter;
use odl::download_manager::DownloadStatus;
use odl::engine::Engine;
use tracing::instrument;
fn init_tracing(level: LogLevel) {
use tracing_subscriber::{EnvFilter, fmt};
let default_directive = match level {
LogLevel::Off => "off",
LogLevel::Error => "error",
LogLevel::Warn => "warn",
LogLevel::Info => "info",
LogLevel::Debug => "debug",
LogLevel::Trace => "trace",
};
let filter =
EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new(default_directive));
fmt()
.with_env_filter(filter)
.with_writer(std::io::stderr)
.init();
}
fn exit_code(e: &OdlError) -> u8 {
match e {
OdlError::CliError { .. }
| OdlError::EmptyInputFile
| OdlError::UrlDecodeError { .. }
| OdlError::ConfigBuilderError(_)
| OdlError::DownloadOptionsBuilderError(_) => 2,
OdlError::Network(_) | OdlError::Ytdlp(odl::error::YtdlpError::RateLimited { .. }) => 3,
OdlError::Conflict(_) => 4,
OdlError::StdIoError { .. } => 5,
OdlError::MetadataError(_) => 6,
OdlError::Ytdlp(_) => 7,
OdlError::NotEvaluated { .. } | OdlError::InvalidRequest { .. } => 2,
OdlError::Cancelled => 130,
OdlError::Other { .. } | _ => 1,
}
}
fn error_kind(e: &OdlError) -> &'static str {
match e {
OdlError::CliError { .. } => "cli",
OdlError::EmptyInputFile => "empty_input_file",
OdlError::UrlDecodeError { .. } => "url_decode",
OdlError::ConfigBuilderError(_) | OdlError::DownloadOptionsBuilderError(_) => "config",
OdlError::Network(_) | OdlError::Ytdlp(odl::error::YtdlpError::RateLimited { .. }) => {
"network"
}
OdlError::Conflict(_) => "conflict",
OdlError::StdIoError { .. } => "io",
OdlError::MetadataError(_) => "metadata",
OdlError::Ytdlp(_) => "ytdlp",
OdlError::NotEvaluated { .. } => "not_evaluated",
OdlError::InvalidRequest { .. } => "invalid_request",
OdlError::Cancelled => "cancelled",
OdlError::Other { .. } | _ => "other",
}
}
static FAILURE_ALREADY_SHOWN: std::sync::atomic::AtomicBool =
std::sync::atomic::AtomicBool::new(false);
fn report_error(e: &OdlError, format: OutputFormat) -> ExitCode {
match format {
OutputFormat::Text if matches!(e, OdlError::Cancelled) => {}
OutputFormat::Text if FAILURE_ALREADY_SHOWN.load(std::sync::atomic::Ordering::Relaxed) => {}
OutputFormat::Text => {
eprintln!("Error: {}", e);
#[cfg(debug_assertions)]
{
eprintln!("{e:?}");
}
}
OutputFormat::Json => {
let v = json!({
"type": "error",
"kind": error_kind(e),
"message": e.to_string(),
"exit_code": exit_code(e),
});
eprintln!("{v}");
}
}
ExitCode::from(exit_code(e))
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum DownloadType {
Url(Url),
File(Box<PathBuf>),
FileAtUrl(Url),
}
#[derive(Copy, Clone)]
struct CliResolver {
file_changed: FileChangedResolution,
not_resumable: NotResumableResolution,
same_download_exists: SameDownloadExistsResolution,
final_file_exists: FinalFileExistsResolution,
}
#[async_trait]
impl ServerConflictResolver for CliResolver {
async fn resolve_file_changed(&self, _: &Download) -> FileChangedResolution {
if self.file_changed == FileChangedResolution::Restart {
FileChangedResolution::Restart
} else {
FileChangedResolution::Abort
}
}
async fn resolve_not_resumable(&self, _: &Download) -> NotResumableResolution {
if self.not_resumable == NotResumableResolution::Restart {
NotResumableResolution::Restart
} else {
NotResumableResolution::Abort
}
}
}
#[async_trait]
impl SaveConflictResolver for CliResolver {
async fn same_download_exists(&self, _: &Download) -> SameDownloadExistsResolution {
self.same_download_exists
}
async fn final_file_exists(&self, _: &Download) -> FinalFileExistsResolution {
if self.final_file_exists == FinalFileExistsResolution::ReplaceAndContinue {
FinalFileExistsResolution::ReplaceAndContinue
} else if self.final_file_exists == FinalFileExistsResolution::AddNumberToNameAndContinue {
FinalFileExistsResolution::AddNumberToNameAndContinue
} else {
FinalFileExistsResolution::Abort
}
}
}
struct InteractiveFormatSelector {
mp: Arc<MultiProgress>,
}
#[async_trait]
impl FormatSelector for InteractiveFormatSelector {
async fn select(&self, offer: &FormatOffer) -> Option<String> {
let tiers = offer.quality_tiers();
if tiers.len() < 2 {
return DefaultFormatSelector.select(offer).await;
}
let mp = Arc::clone(&self.mp);
let title = offer.title.clone();
let can_merge = offer.can_merge;
let fallback = DefaultFormatSelector.select(offer).await;
tokio::task::spawn_blocking(move || {
mp.suspend(|| prompt_for_quality(&title, &tiers, can_merge, fallback))
})
.await
.unwrap_or(None)
}
}
fn prompt_for_quality(
title: &str,
tiers: &[QualityTier],
can_merge: bool,
fallback: Option<String>,
) -> Option<String> {
use std::io::Write;
let mut err = std::io::stderr();
let _ = writeln!(err, "\nTitle: {title}");
let default = tiers.iter().position(|t| t.available)?;
let width = tiers
.iter()
.map(|t| t.quality.to_string().chars().count())
.max()
.unwrap_or(8);
let ext_width = tiers
.iter()
.map(|t| t.ext.chars().count())
.max()
.unwrap_or(4);
for (i, tier) in tiers.iter().enumerate() {
let size = match tier.size {
Some(s) if tier.size_is_approx => format!("~{}", HumanBytes(s)),
Some(s) => HumanBytes(s).to_string(),
None => "size unknown".to_owned(),
};
let note = if tier.available {
""
} else {
" — needs ffmpeg"
};
let marker = if i == default { '*' } else { ' ' };
let quality = tier.quality.to_string();
let _ = writeln!(
err,
" {marker}{:>2}) {quality:<width$} {:<ext_width$} {size}{note}",
i + 1,
tier.ext
);
}
if !can_merge && tiers.iter().any(|t| !t.available) {
let _ = writeln!(
err,
"\n Higher qualities are served as separate video and audio streams,\n \
which need ffmpeg to join. Install it with `odl tools install ffmpeg`."
);
}
let _ = write!(
err,
"\nQuality [1-{}, default {}]: ",
tiers.len(),
default + 1
);
let _ = err.flush();
let mut line = String::new();
if std::io::stdin().read_line(&mut line).is_err() {
return fallback;
}
let trimmed = line.trim();
if trimmed.is_empty() {
return Some(tiers[default].format_id.clone());
}
match trimmed.parse::<usize>() {
Ok(n) if (1..=tiers.len()).contains(&n) && tiers[n - 1].available => {
Some(tiers[n - 1].format_id.clone())
}
Ok(n) if (1..=tiers.len()).contains(&n) => {
let _ = writeln!(
err,
"That quality needs ffmpeg, which is not installed; taking {} instead.",
tiers[default].quality
);
Some(tiers[default].format_id.clone())
}
_ => {
let _ = writeln!(err, "Not a listed choice; taking the best available.");
Some(tiers[default].format_id.clone())
}
}
}
fn spawn_interrupt_handler(cancel: CancellationToken) {
tokio::spawn(async move {
if tokio::signal::ctrl_c().await.is_err() {
return;
}
eprintln!("\nInterrupted; finishing up. Press Ctrl-C again to quit immediately.");
cancel.cancel();
if tokio::signal::ctrl_c().await.is_ok() {
#[cfg(feature = "ytdlp")]
odl::ytdlp::process::kill_all_groups();
std::process::exit(130);
}
});
}
fn should_prompt_for_format(
choice: args::ChooseFormat,
format: OutputFormat,
url_count: usize,
) -> bool {
let can_ask = format == OutputFormat::Text
&& std::io::stdin().is_terminal()
&& std::io::stderr().is_terminal();
if !can_ask {
return false;
}
match choice {
args::ChooseFormat::Never => false,
args::ChooseFormat::Always => true,
args::ChooseFormat::Auto => url_count == 1,
}
}
pub const PROGRESS_CHARS: &str = "█▇▆▅▄▃▂▁";
struct BarMetrics {
speed_bits: Arc<AtomicU64>,
downloaded: Arc<AtomicU64>,
total: Arc<AtomicU64>,
}
impl BarMetrics {
fn new() -> Self {
Self {
speed_bits: Arc::new(AtomicU64::new(0)),
downloaded: Arc::new(AtomicU64::new(0)),
total: Arc::new(AtomicU64::new(0)),
}
}
}
struct PartBar {
bar: ProgressBar,
metrics: BarMetrics,
}
struct CliReporter {
mp: Arc<MultiProgress>,
parent: ProgressBar,
parts: Mutex<HashMap<String, PartBar>>,
parent_metrics: BarMetrics,
filename: Mutex<Option<String>>,
url: String,
solo: Mutex<Option<PartBar>>,
}
impl CliReporter {
fn new(
mp: Arc<MultiProgress>,
parent: ProgressBar,
parent_metrics: BarMetrics,
url: String,
) -> Self {
Self {
url,
mp,
parent,
parts: Mutex::new(HashMap::new()),
parent_metrics,
filename: Mutex::new(None),
solo: Mutex::new(None),
}
}
fn ensure_solo_row(&self) {
if !self.parts.lock().unwrap().is_empty() {
return;
}
let mut solo = self.solo.lock().unwrap();
if solo.is_some() {
return;
}
let metrics = BarMetrics::new();
let total = self.parent_metrics.total.load(Ordering::Relaxed);
let bar = if total > 0 {
ProgressBar::new(total)
} else {
ProgressBar::new_spinner()
};
let bar = self.mp.add(bar.with_style(build_solo_row_style()));
bar.enable_steady_tick(SAMPLE_INTERVAL);
*solo = Some(PartBar { bar, metrics });
}
fn clear_solo_row(&self) {
if let Some(p) = self.solo.lock().unwrap().take() {
p.bar.finish_and_clear();
}
}
}
fn phase_label(phase: Phase) -> &'static str {
match phase {
Phase::Evaluating => "Evaluating",
Phase::ResolvingConflicts => "Resolving conflicts",
Phase::Downloading => "Downloading",
Phase::PostProcessing => "Processing",
Phase::Assembling => "Assembling",
Phase::Flushing => "Flushing data to disk",
Phase::Verifying => "Verifying checksum",
_ => "Working",
}
}
impl CliReporter {
fn describe_target(&self) -> String {
self.filename
.lock()
.unwrap()
.clone()
.unwrap_or_else(|| self.url.clone())
}
fn finish_with_status(&self, template: &str, message: String) {
self.clear_solo_row();
{
let mut parts = self.parts.lock().unwrap();
for (_, p) in parts.drain() {
p.bar.finish_and_clear();
}
}
let style =
ProgressStyle::with_template(template).expect("templating a status line cannot fail");
self.parent.set_style(style);
self.parent.finish_with_message(message);
}
}
impl ProgressReporter for CliReporter {
fn on_event(&self, event: ProgressEvent) {
match event {
ProgressEvent::FilenameResolved(name) => {
*self.filename.lock().unwrap() = Some(name.clone());
self.parent.set_message(name);
}
ProgressEvent::PhaseChanged(phase) => {
match phase {
Phase::Downloading => {
if let Some(name) = self.filename.lock().unwrap().clone() {
self.parent.set_message(name);
}
self.ensure_solo_row();
}
_ => {
self.parent.set_message(phase_label(phase));
}
}
}
ProgressEvent::Progress { downloaded, total } => {
self.parent_metrics
.downloaded
.store(downloaded, Ordering::Relaxed);
if let Some(t) = total {
self.parent_metrics.total.store(t, Ordering::Relaxed);
self.parent.set_length(t);
}
self.parent.set_position(downloaded);
if let Some(p) = self.solo.lock().unwrap().as_ref() {
p.metrics.downloaded.store(downloaded, Ordering::Relaxed);
if let Some(t) = total {
p.metrics.total.store(t, Ordering::Relaxed);
p.bar.set_length(t);
}
p.bar.set_position(downloaded);
}
}
ProgressEvent::Speed { bytes_per_second } => {
self.parent_metrics
.speed_bits
.store(bytes_per_second.to_bits(), Ordering::Relaxed);
if let Some(p) = self.solo.lock().unwrap().as_ref() {
p.metrics
.speed_bits
.store(bytes_per_second.to_bits(), Ordering::Relaxed);
}
}
ProgressEvent::PartAdded { ulid, size, .. } => {
self.clear_solo_row();
let metrics = BarMetrics::new();
metrics.total.store(size, Ordering::Relaxed);
let style = build_child_style(&metrics);
let bar = self.mp.add(ProgressBar::new(size).with_style(style));
match ulid.as_str() {
odl::progress::ASSEMBLY_ULID => bar.set_message("assembling"),
odl::progress::VERIFY_ULID => bar.set_message("verifying"),
_ => {}
}
bar.enable_steady_tick(SAMPLE_INTERVAL);
self.parts
.lock()
.unwrap()
.insert(ulid, PartBar { bar, metrics });
}
ProgressEvent::PartProgress {
ulid,
downloaded,
total,
} => {
if let Some(p) = self.parts.lock().unwrap().get(&ulid) {
p.metrics.downloaded.store(downloaded, Ordering::Relaxed);
p.metrics.total.store(total, Ordering::Relaxed);
p.bar.set_length(total);
p.bar.set_position(downloaded);
}
}
ProgressEvent::PartSpeed {
ulid,
bytes_per_second,
} => {
if let Some(p) = self.parts.lock().unwrap().get(&ulid) {
p.metrics
.speed_bits
.store(bytes_per_second.to_bits(), Ordering::Relaxed);
}
}
ProgressEvent::PartFinished { ulid } => {
if let Some(p) = self.parts.lock().unwrap().remove(&ulid) {
p.bar.finish_and_clear();
}
}
ProgressEvent::PartsCleared => {
for (_, p) in self.parts.lock().unwrap().drain() {
p.bar.finish_and_clear();
}
}
ProgressEvent::PartRetrying { ulid, attempt } => {
if let Some(p) = self.parts.lock().unwrap().get(&ulid) {
p.bar.set_message(format!("retry #{attempt}"));
}
}
ProgressEvent::RetryScheduled {
ulid: Some(ulid),
attempt,
max_attempts,
delay,
server_requested,
} => {
if let Some(p) = self.parts.lock().unwrap().get(&ulid) {
let rounded = if delay >= std::time::Duration::from_secs(1) {
std::time::Duration::from_secs(delay.as_secs())
} else {
delay
};
let why = if server_requested {
" (Rate limited)"
} else {
""
};
p.bar.set_message(format!(
"retry #{attempt}/{max_attempts} in {}{why}",
humantime::format_duration(rounded)
));
}
}
ProgressEvent::Message(msg) => {
if !msg.is_empty() {
self.parent.set_message(msg);
}
}
ProgressEvent::Completed {
path,
already_complete,
} => {
self.clear_solo_row();
{
let mut parts = self.parts.lock().unwrap();
for (_, p) in parts.drain() {
p.bar.finish_and_clear();
}
}
let suffix = if already_complete {
" (already complete)"
} else {
""
};
let final_style = ProgressStyle::with_template("✓ Saved to {msg}")
.expect("templating final progress should not fail");
self.parent.set_style(final_style);
self.parent
.finish_with_message(format!("{}{}", path.display(), suffix));
}
ProgressEvent::Cancelled => {
self.finish_with_status("✕ Cancelled: {msg}", self.describe_target());
}
ProgressEvent::Failed { message } => {
self.finish_with_status(
"✕ {msg}",
format!("{}: {message}", self.describe_target()),
);
}
_ => {}
}
}
}
fn install_metric_keys(style: ProgressStyle, metrics: &BarMetrics) -> ProgressStyle {
let speed_for_key = Arc::clone(&metrics.speed_bits);
let dl_for_eta = Arc::clone(&metrics.downloaded);
let total_for_eta = Arc::clone(&metrics.total);
let speed_for_eta = Arc::clone(&metrics.speed_bits);
style
.with_key(
"fast_speed",
move |_state: &ProgressState, w: &mut dyn std::fmt::Write| {
let bytes = f64::from_bits(speed_for_key.load(Ordering::Relaxed));
let bytes = if bytes.is_finite() && bytes >= 0.0 {
bytes as u64
} else {
0
};
let _ = write!(w, "{}/s", HumanBytes(bytes));
},
)
.with_key(
"fast_eta",
move |_state: &ProgressState, w: &mut dyn std::fmt::Write| {
let total = total_for_eta.load(Ordering::Relaxed);
let downloaded = dl_for_eta.load(Ordering::Relaxed);
let speed = f64::from_bits(speed_for_eta.load(Ordering::Relaxed));
if total == 0 || downloaded >= total || !speed.is_finite() || speed <= 1.0 {
let _ = write!(w, "--");
return;
}
let remaining = (total - downloaded) as f64 / speed;
let secs = remaining as u64;
let _ = write!(
w,
"{:02}:{:02}:{:02}",
secs / 3600,
(secs / 60) % 60,
secs % 60
);
},
)
}
fn build_solo_row_style() -> ProgressStyle {
ProgressStyle::with_template(" ↳ {bar:30.cyan/blue} {msg}")
.expect("templating progress bar should not fail")
.progress_chars(PROGRESS_CHARS)
}
fn build_parent_style(metrics: &BarMetrics) -> ProgressStyle {
let style = ProgressStyle::with_template(
"{spinner} {maybe_connect} {msg:!40} {percent:>3}% {decimal_bytes:<10} / {decimal_total_bytes:<10} {fast_speed:<14} eta {fast_eta:>9} elapsed {elapsed}",
)
.expect("templating progress bar should not fail")
.progress_chars(PROGRESS_CHARS)
.with_key(
"maybe_connect",
|state: &ProgressState, w: &mut dyn std::fmt::Write| {
if state.len().is_none() || state.len().is_some_and(|x| x == 0) {
let _ = write!(w, "━");
} else {
let _ = write!(w, "┌");
}
},
);
install_metric_keys(style, metrics)
}
fn build_child_style(metrics: &BarMetrics) -> ProgressStyle {
let style = ProgressStyle::with_template(
" ↳ {spinner} {bar:30.cyan/blue} {percent:>3}% {decimal_bytes:<10} / {decimal_total_bytes:<10} {fast_speed:<14} eta {fast_eta:>9} {msg}",
)
.expect("templating progress bar should not fail")
.progress_chars(PROGRESS_CHARS);
install_metric_keys(style, metrics)
}
#[tokio::main]
async fn main() -> ExitCode {
let args: Args = Args::parse();
init_tracing(args.log_level);
let format = args.format;
match run(args).await {
Ok(()) => ExitCode::SUCCESS,
Err(e) => report_error(&e, format),
}
}
async fn run(args: Args) -> Result<(), OdlError> {
let format = args.format;
if args.command.is_none() && args.input.is_none() {
let mut cmd = Args::command();
if let Err(e) = cmd.print_help() {
eprintln!("Failed to print help: {}", e);
}
println!();
return Ok(());
}
if let Some(cmd) = &args.command {
match cmd {
args::Commands::Config {
show,
config_file,
download_dir,
max_connections,
max_concurrent_downloads,
max_retries,
wait_between_retries,
n_fixed_retries,
speed_limit,
user_agent,
randomize_user_agent,
proxy,
timeout,
use_server_time,
accept_invalid_certs,
http2,
dynamic_split,
rampup,
rampup_batch_size,
rampup_delay_min,
rampup_delay_max,
} => {
let config_path = if let Some(c) = config_file {
c.clone()
} else if let Some(c) = &args.config_file {
c.clone()
} else {
Config::default_config_file()
};
let mut cfg: Config = Config::load_from_file(&config_path)
.await
.unwrap_or_default();
if *show {
match format {
OutputFormat::Json => {
let v = json!({
"type": "config",
"path": config_path.display().to_string(),
"config": cfg,
});
println!("{v}");
}
OutputFormat::Text => {
println!("# config path: {}", config_path.display());
match toml::to_string_pretty(&cfg) {
Ok(s) => {
if s.trim().is_empty() {
println!("# config is empty")
} else {
println!("{}", s)
}
}
Err(e) => eprintln!("Failed to format config: {}", e),
}
}
}
return Ok(());
}
let mut dl_b = cfg.download().clone().into_builder();
if let Some(v) = max_connections {
dl_b.max_connections(*v);
}
if let Some(v) = max_retries {
dl_b.max_retries(*v);
}
if let Some(v) = wait_between_retries {
dl_b.wait_between_retries(*v);
}
if let Some(v) = n_fixed_retries {
dl_b.n_fixed_retries(*v);
}
if let Some(v) = user_agent {
dl_b.user_agent(Some(v.clone()));
}
if let Some(v) = randomize_user_agent {
dl_b.randomize_user_agent(*v);
}
if let Some(v) = proxy {
dl_b.proxy(Some(v.clone()));
}
if let Some(v) = use_server_time {
dl_b.use_server_time(*v);
}
if let Some(v) = accept_invalid_certs {
dl_b.accept_invalid_certs(*v);
}
if let Some(v) = speed_limit {
dl_b.speed_limit(Some(*v));
}
if let Some(v) = *timeout {
dl_b.connect_timeout(Some(v));
}
if let Some(v) = http2 {
dl_b.http2(*v);
}
if let Some(v) = dynamic_split {
dl_b.dynamic_split(*v);
}
if let Some(v) = rampup {
dl_b.rampup(*v);
}
if let Some(v) = rampup_batch_size {
dl_b.rampup_batch_size(*v);
}
if let Some(v) = rampup_delay_min {
dl_b.rampup_delay_min(*v);
}
if let Some(v) = rampup_delay_max {
dl_b.rampup_delay_max(*v);
}
let new_download = dl_b.build()?;
let mut cfg_b = cfg.into_builder();
cfg_b.download(new_download);
if let Some(v) = download_dir {
cfg_b.download_dir(v.clone());
}
if let Some(v) = max_concurrent_downloads {
cfg_b.max_concurrent_downloads(*v);
}
cfg = cfg_b.build()?;
match cfg.save_to_file(&config_path).await {
Ok(()) => match format {
OutputFormat::Json => {
let v = json!({
"type": "config_saved",
"path": config_path.display().to_string(),
});
println!("{v}");
}
OutputFormat::Text => {
println!("Saved configuration to {}", config_path.display())
}
},
Err(e) => {
return Err(OdlError::StdIoError {
e,
extra_info: Some("failed to save configuration".to_string()),
});
}
}
return Ok(());
}
args::Commands::Tools { action } => {
return run_tools(&args, action, format).await;
}
#[cfg(feature = "self-update")]
args::Commands::Update { check, yes } => {
return run_update(&args, *check, *yes, format).await;
}
args::Commands::Probe { url } => {
return run_probe(&args, url, format).await;
}
args::Commands::Status { filter } => {
return run_status(&args, filter.as_deref(), format, false).await;
}
args::Commands::List { filter } => {
return run_status(&args, filter.as_deref(), format, true).await;
}
}
}
let mp = Arc::new(MultiProgress::new());
let dlm = build_download_manager(&args).await?;
let mut download_type = determine_download_type(&args)?;
if let DownloadType::FileAtUrl(url) = &download_type {
let path = download_remote_file(&dlm, url.clone()).await?;
download_type = DownloadType::File(Box::new(path));
}
let mut user_provided_filename: Option<String> = None;
let save_dir: PathBuf = if let Some(path) = args.output.clone() {
if let DownloadType::Url(_) = &download_type {
user_provided_filename = path
.file_name()
.map(|os_str| os_str.to_string_lossy().into_owned());
path.parent()
.expect("Failed to get output's parent directory")
.to_path_buf()
} else {
path
}
} else {
std::env::current_dir()?
};
let mut urls = Vec::new();
match &download_type {
DownloadType::Url(url) => {
urls.push(url.clone());
}
DownloadType::File(path) => {
let file = tokio::fs::File::open(&**path).await?;
let reader = tokio::io::BufReader::new(file);
let mut lines = tokio::io::BufReader::new(reader).lines();
while let Some(line) = lines.next_line().await? {
let trimmed = line.trim();
if trimmed.is_empty() || trimmed.starts_with('#') || trimmed.starts_with("//") {
continue;
}
match Url::parse(trimmed) {
Ok(url) => urls.push(url),
Err(e) => {
println!("Skipping invalid URL '{}': {}", trimmed, e);
}
}
}
}
DownloadType::FileAtUrl(_) => {
panic!("FileAtUrl should have been handled already");
}
}
let mut handles = Vec::new();
let mut forwarders: Vec<Arc<AsyncReporter>> = Vec::new();
let expected_checksums = args
.checksums
.iter()
.map(|s| odl::hash::HashDigest::parse_cli(s))
.collect::<Result<Vec<_>, _>>()
.map_err(|message| OdlError::CliError { message })?;
let dlm = Arc::new(dlm);
let credentials = if let Some(user) = args.http_user.as_deref() {
Some(Credentials::new(user, args.http_password.as_deref()))
} else {
None
};
let resolver = {
let file_changed = match args.on_file_changed {
args::FileChangedAction::Abort => FileChangedResolution::Abort,
args::FileChangedAction::Restart => FileChangedResolution::Restart,
};
let not_resumable = match args.on_not_resumable {
args::NotResumableAction::Abort => NotResumableResolution::Abort,
args::NotResumableAction::Restart => NotResumableResolution::Restart,
};
let same_download_exists = match args.on_same_download_exists {
args::SameDownloadAction::Abort => SameDownloadExistsResolution::Abort,
args::SameDownloadAction::Resume => SameDownloadExistsResolution::Resume,
args::SameDownloadAction::AddNumberToNameAndContinue => {
SameDownloadExistsResolution::AddNumberToNameAndContinue
}
};
let final_file_exists = match args.on_final_file_exists {
args::FinalFileAction::Abort => FinalFileExistsResolution::Abort,
args::FinalFileAction::ReplaceAndContinue => {
FinalFileExistsResolution::ReplaceAndContinue
}
args::FinalFileAction::AddNumberToNameAndContinue => {
FinalFileExistsResolution::AddNumberToNameAndContinue
}
};
CliResolver {
file_changed,
not_resumable,
same_download_exists,
final_file_exists,
}
};
let engine_preference = match args.engine {
args::EngineChoice::Auto => EnginePreference::Auto,
args::EngineChoice::Http => EnginePreference::Engine(Engine::HttpMultipart),
args::EngineChoice::Ytdlp => EnginePreference::Engine(Engine::Ytdlp),
};
let dlm = if maybe_offer_missing_tools(&args, &urls, format, &dlm).await? {
Arc::new(build_download_manager(&args).await?)
} else {
dlm
};
let cancel = CancellationToken::new();
spawn_interrupt_handler(cancel.clone());
let reselect_format =
args.format_id.is_some() || args.choose_format == args::ChooseFormat::Always;
let format_selector: Arc<dyn FormatSelector> = match args.format_id.clone() {
Some(id) => Arc::new(FixedFormatSelector(id)),
None if should_prompt_for_format(args.choose_format, format, urls.len()) => {
Arc::new(InteractiveFormatSelector {
mp: Arc::clone(&mp),
})
}
None => Arc::new(DefaultFormatSelector),
};
for url in urls.into_iter() {
let reporter: Arc<dyn ProgressReporter> = match format {
OutputFormat::Json => Arc::new(JsonReporter::new(url.to_string())),
OutputFormat::Text => {
let parent_metrics = BarMetrics::new();
let parent_style = build_parent_style(&parent_metrics);
let parent = mp.add(ProgressBar::new_spinner().with_style(parent_style));
parent.set_message(format!("{url} (warming up)"));
parent.enable_steady_tick(SAMPLE_INTERVAL);
let cli_reporter =
CliReporter::new(Arc::clone(&mp), parent, parent_metrics, url.to_string());
let async_reporter = AsyncReporter::spawn(cli_reporter);
forwarders.push(Arc::clone(&async_reporter));
async_reporter as Arc<dyn ProgressReporter>
}
};
let ctx = DownloadContext::new()
.with_reporter(reporter)
.with_url(url.clone())
.with_cancel(cancel.clone());
let dlm = Arc::clone(&dlm);
let save_dir = save_dir.clone();
let user_provided_filename = user_provided_filename.clone();
let credentials = credentials.clone();
let expected_checksums = expected_checksums.clone();
let format_selector = Arc::clone(&format_selector);
let permit = dlm
.acquire_download_permit()
.await
.expect("didn't expect the semaphore to close at this point");
let handle = tokio::spawn(async move {
let _permit = permit;
let result: Result<PathBuf, OdlError> = async {
let mut request = EvaluateRequest::new(url, save_dir, &resolver)
.ctx(&ctx)
.engine(engine_preference)
.format_selector(&*format_selector)
.reselect_format(reselect_format);
if let Some(c) = credentials {
request = request.credentials(c);
}
let mut instruction = match dlm.evaluate(request).await {
Ok(instruction) => instruction,
Err(OdlError::Cancelled) => {
ctx.emit(ProgressEvent::Cancelled);
return Err(OdlError::Cancelled);
}
Err(e) => {
ctx.emit(ProgressEvent::Failed {
message: e.to_string(),
});
return Err(e);
}
};
if let Some(filename) = user_provided_filename {
instruction.set_filename(filename);
}
instruction.add_checksums(expected_checksums);
dlm.download(DownloadRequest::new(instruction, &resolver).ctx(&ctx))
.await
}
.await;
result
});
handles.push(handle);
}
let mut results: Vec<Result<PathBuf, OdlError>> = Vec::new();
for h in handles {
match h.await {
Ok(Ok(path)) => results.push(Ok(path)),
Ok(Err(e)) => results.push(Err(e)),
Err(join_err) => results.push(Err(OdlError::from(join_err))),
}
}
for forwarder in &forwarders {
forwarder.drained().await;
}
let first_err = results.into_iter().find_map(|r| r.err());
if let Some(e) = first_err {
if !forwarders.is_empty() && !mp.is_hidden() {
FAILURE_ALREADY_SHOWN.store(true, std::sync::atomic::Ordering::Relaxed);
}
return Err(e);
}
Ok(())
}
#[instrument(skip(args), name = "Warming up odl...")]
async fn build_download_manager(args: &Args) -> Result<DownloadManager, OdlError> {
let config_file: PathBuf = if let Some(c) = &args.config_file {
c.clone()
} else {
Config::default_config_file()
};
let cfg = Config::load_from_file(&config_file)
.await
.unwrap_or_default();
let wait_between_retries = args.wait_between_retries.and_then(|d| {
let secs = d.as_secs_f64();
if secs.is_finite() && secs >= 0.0 {
Some(d)
} else {
None
}
});
let connect_timeout = args.timeout.and_then(|d| {
let secs = d.as_secs_f64();
if secs.is_finite() && secs >= 0.0 {
Some(d)
} else {
None
}
});
let headers = if args.headers.is_empty() {
None
} else {
let mut headers_map = indexmap::IndexMap::new();
for header in &args.headers {
if let Some((key, value)) = header.split_once(':') {
let key = key.trim();
let value = value.trim();
headers_map.insert(key.to_string(), value.to_string());
} else {
return Err(OdlError::CliError {
message: format!("Header must be in KEY:VALUE format: '{}'", header),
});
}
}
Some(headers_map)
};
let mut dl_b = cfg.download().clone().into_builder();
if let Some(v) = args.max_connections {
dl_b.max_connections(v);
}
if let Some(v) = args.max_retries {
dl_b.max_retries(v);
}
if let Some(v) = wait_between_retries {
dl_b.wait_between_retries(v);
}
if let Some(v) = args.n_fixed_retries {
dl_b.n_fixed_retries(v);
}
if let Some(v) = args.user_agent.clone() {
dl_b.user_agent(Some(v));
}
if let Some(v) = args.randomize_user_agent {
dl_b.randomize_user_agent(v);
}
if let Some(v) = args.proxy.clone() {
dl_b.proxy(Some(v));
}
if let Some(v) = args.use_server_time {
dl_b.use_server_time(v);
}
if let Some(v) = args.accept_invalid_certs {
dl_b.accept_invalid_certs(v);
}
if let Some(v) = args.speed_limit {
dl_b.speed_limit(Some(v));
}
if args.no_verify_checksums {
dl_b.verify_checksums(false);
}
if args.ascii_filenames {
dl_b.ascii_filenames(true);
}
if let Some(v) = connect_timeout {
dl_b.connect_timeout(Some(v));
}
if let Some(v) = headers {
dl_b.headers(Some(v));
}
if let Some(v) = args.http2 {
dl_b.http2(v);
}
if let Some(v) = args.dynamic_split {
dl_b.dynamic_split(v);
}
if let Some(v) = args.rampup {
dl_b.rampup(v);
}
if let Some(v) = args.rampup_batch_size {
dl_b.rampup_batch_size(v);
}
if let Some(v) = args.rampup_delay_min {
dl_b.rampup_delay_min(v);
}
if let Some(v) = args.rampup_delay_max {
dl_b.rampup_delay_max(v);
}
let download = dl_b.build()?;
let mut cfg_b = cfg.into_builder();
cfg_b.download(download);
if let Some(v) = args.download_dir.clone() {
cfg_b.download_dir(v);
}
if let Some(v) = args.max_concurrent_downloads {
cfg_b.max_concurrent_downloads(v);
}
let cfg = cfg_b.build()?;
Ok(DownloadManager::new(cfg))
}
#[instrument(skip(args), name = "Determining download type")]
fn determine_download_type(args: &Args) -> Result<DownloadType, OdlError> {
let input = args.input.as_ref().ok_or(OdlError::CliError {
message: "Missing input. Provide a URL or file path, or use a subcommand like `config`."
.to_string(),
})?;
Ok(match Url::parse(input) {
Ok(url) => {
if args.remote_list {
DownloadType::FileAtUrl(url)
} else {
DownloadType::Url(url)
}
}
Err(_) => {
let path = PathBuf::from(input);
if path.try_exists()? {
if args.remote_list {
return Err(OdlError::CliError {
message: "Expected input to be a Url, found file path instead".to_string(),
});
}
DownloadType::File(Box::new(path))
} else {
return Err(OdlError::CliError {
message: "Input is not a valid Url or a valid file path. Check file permissions if file exists.".to_string(),
});
}
}
})
}
struct ForcedResolver;
#[async_trait]
impl ServerConflictResolver for ForcedResolver {
async fn resolve_file_changed(&self, _: &Download) -> FileChangedResolution {
FileChangedResolution::Restart
}
async fn resolve_not_resumable(&self, _: &Download) -> NotResumableResolution {
NotResumableResolution::Restart
}
}
#[async_trait]
impl SaveConflictResolver for ForcedResolver {
async fn same_download_exists(&self, _: &Download) -> SameDownloadExistsResolution {
SameDownloadExistsResolution::Resume
}
async fn final_file_exists(&self, _: &Download) -> FinalFileExistsResolution {
FinalFileExistsResolution::ReplaceAndContinue
}
}
#[instrument(skip(dlm), name = "Downloading remote file containing links")]
async fn download_remote_file(dlm: &DownloadManager, url: Url) -> Result<PathBuf, OdlError> {
let resolver = ForcedResolver {};
let tmpdir = tempfile::Builder::new()
.prefix("odl")
.tempdir()
.map_err(|e| OdlError::CliError {
message: format!("Failed to create temp dir: {e}"),
})?;
let save_dir = tmpdir.path().to_path_buf();
let _permit = dlm.acquire_download_permit().await?;
let instruction = dlm
.evaluate(
EvaluateRequest::new(url, save_dir, &resolver)
.engine(EnginePreference::Engine(Engine::HttpMultipart)),
)
.await?;
let path = dlm
.download(DownloadRequest::new(instruction, &resolver))
.await?;
Ok(path)
}
fn checksum_json(c: &odl::download_metadata::FileChecksum) -> serde_json::Value {
use odl::download_metadata::{ChecksumAlgorithm, ChecksumEncoding};
let algorithm = ChecksumAlgorithm::try_from(c.algorithm)
.map(|a| a.as_str_name())
.unwrap_or("unknown");
let encoding = ChecksumEncoding::try_from(c.encoding)
.map(|e| e.as_str_name())
.unwrap_or("unknown");
json!({"algorithm": algorithm, "digest": c.digest, "encoding": encoding})
}
#[cfg(feature = "ytdlp")]
async fn maybe_offer_missing_tools(
args: &Args,
urls: &[Url],
format: OutputFormat,
dlm: &DownloadManager,
) -> Result<bool, OdlError> {
use odl::ytdlp::install::{self, Tool};
if format != OutputFormat::Text
|| !std::io::stdin().is_terminal()
|| args.engine == args::EngineChoice::Http
{
return Ok(false);
}
let config_path = config_path_for(args);
let mut cfg = Config::load_from_file(&config_path).await?;
if !cfg.ytdlp().enabled() {
return Ok(false);
}
let forced = args.engine == args::EngineChoice::Ytdlp;
let wanted = forced
|| urls
.iter()
.any(|u| odl::ytdlp::should_delegate(u, cfg.ytdlp()));
if !wanted {
return Ok(false);
}
let mut tools = odl::ytdlp::tools(cfg.ytdlp()).await.ok();
let mut changed = false;
let mut installed_something = false;
if tools.is_none() && cfg.ytdlp().offer_ytdlp_install() {
let mut ytdlp_cfg = cfg.ytdlp().clone();
match offer_tool(Tool::Ytdlp, None, &config_path, cfg.download(), dlm, false).await? {
OfferOutcome::Installed(path) => {
ytdlp_cfg.set_binary_path(Some(path));
changed = true;
installed_something = true;
}
OfferOutcome::Declined => {
ytdlp_cfg.set_offer_ytdlp_install(false);
changed = true;
}
OfferOutcome::NothingToDo => {}
}
cfg = cfg.into_builder().ytdlp(ytdlp_cfg).build()?;
if changed {
tools = odl::ytdlp::tools(cfg.ytdlp()).await.ok();
}
}
if let Some(t) = &tools
&& t.ffmpeg.is_none()
&& cfg.ytdlp().offer_ffmpeg_install()
&& install::can_install(Tool::Ffmpeg)
{
let mut ytdlp_cfg = cfg.ytdlp().clone();
match offer_tool(Tool::Ffmpeg, None, &config_path, cfg.download(), dlm, false).await? {
OfferOutcome::Installed(path) => {
ytdlp_cfg.set_ffmpeg_path(Some(path));
changed = true;
installed_something = true;
}
OfferOutcome::Declined => {
ytdlp_cfg.set_offer_ffmpeg_install(false);
changed = true;
}
OfferOutcome::NothingToDo => {}
}
cfg = cfg.into_builder().ytdlp(ytdlp_cfg).build()?;
}
if changed {
cfg.save_to_file(&config_path).await?;
eprintln!("Updated {}\n", config_path.display());
}
Ok(installed_something)
}
#[cfg(not(feature = "ytdlp"))]
async fn maybe_offer_missing_tools(
_args: &Args,
_urls: &[Url],
_format: OutputFormat,
_dlm: &DownloadManager,
) -> Result<bool, OdlError> {
Ok(false)
}
#[cfg(feature = "ytdlp")]
fn config_path_for(args: &Args) -> PathBuf {
args.config_file
.clone()
.unwrap_or_else(Config::default_config_file)
}
#[cfg(feature = "ytdlp")]
fn confirm(question: &str) -> Option<bool> {
use std::io::Write;
if !std::io::stdin().is_terminal() || !std::io::stderr().is_terminal() {
return None;
}
let mut err = std::io::stderr();
let _ = write!(err, "{question} [y/N]: ");
let _ = err.flush();
let mut line = String::new();
if std::io::stdin().read_line(&mut line).is_err() {
return None;
}
Some(matches!(
line.trim().to_ascii_lowercase().as_str(),
"y" | "yes"
))
}
#[cfg(any(feature = "ytdlp", feature = "self-update"))]
fn install_client(net: &odl::config::DownloadOptions) -> Result<reqwest::Client, OdlError> {
let mut builder = reqwest::Client::builder();
if let Some(proxy) = Option::<reqwest::Proxy>::from(net) {
builder = builder.proxy(proxy);
}
if net.accept_invalid_certs() {
builder = builder.danger_accept_invalid_certs(true);
}
if let Some(timeout) = net.connect_timeout() {
builder = builder.connect_timeout(timeout);
}
builder.build().map_err(|e| OdlError::CliError {
message: format!("could not create an HTTP client: {e}"),
})
}
#[cfg(feature = "ytdlp")]
enum OfferOutcome {
Installed(PathBuf),
Declined,
NothingToDo,
}
#[cfg(feature = "ytdlp")]
async fn offer_tool(
tool: odl::ytdlp::install::Tool,
installed: Option<&std::path::Path>,
config_path: &std::path::Path,
net: &odl::config::DownloadOptions,
dlm: &DownloadManager,
assume_yes: bool,
) -> Result<OfferOutcome, OdlError> {
use odl::ytdlp::install;
if let Some(path) = installed {
eprintln!("{}: already installed at {}", tool.as_str(), path.display());
return Ok(OfferOutcome::NothingToDo);
}
if !install::can_install(tool) {
eprintln!(
"{}: not installed, and odl has no verified build for this platform.\n Install it yourself — {} — then set `{}` in {}.",
tool.as_str(),
tool.manual_instructions(),
tool.config_key(),
config_path.display(),
);
return Ok(OfferOutcome::NothingToDo);
}
let dir = install::tools_dir();
eprintln!("\n{} is not installed.", tool.as_str());
eprintln!(" {}", tool.purpose());
eprintln!(" odl can download it from {}.", tool.source_description());
eprintln!(" It will be verified against the checksums published with it, saved to");
eprintln!(
" {}, and recorded as `{}` in",
dir.display(),
tool.config_key()
);
eprintln!(" {}.", config_path.display());
eprintln!(
" Or install it yourself — {} — and set that key by hand.",
tool.manual_instructions()
);
if !assume_yes {
match confirm(&format!("Download {} now?", tool.as_str())) {
Some(true) => {}
Some(false) => {
eprintln!(
"Skipped {0}. odl will not ask again; run `odl tools install {0}` when you want it.",
tool.as_str()
);
return Ok(OfferOutcome::Declined);
}
None => {
eprintln!(
"Not installing {0}: no terminal to confirm on. Re-run with `-y` to install it without asking.",
tool.as_str()
);
return Ok(OfferOutcome::NothingToDo);
}
}
}
let client = install_client(net)?;
let plan = install::plan(&client, tool).await.map_err(OdlError::from)?;
eprintln!("Downloading {} ({})…", tool.as_str(), plan.name);
let downloaded = download_asset(dlm, &plan).await?;
let path = install::finish(tool, &downloaded, &dir)
.await
.map_err(OdlError::from)?;
let _ = tokio::fs::remove_file(&downloaded).await;
eprintln!("Installed {} to {}", tool.as_str(), path.display());
Ok(OfferOutcome::Installed(path))
}
#[cfg(feature = "ytdlp")]
async fn download_asset(
dlm: &DownloadManager,
plan: &odl::ytdlp::install::AssetPlan,
) -> Result<PathBuf, OdlError> {
use odl::ytdlp::install;
let url = Url::parse(&plan.url).map_err(|e| OdlError::CliError {
message: format!("release asset URL is not usable: {e}"),
})?;
let staging = install::staging_dir();
tokio::fs::create_dir_all(&staging).await?;
let resolver = ForcedResolver {};
let mut instruction = dlm
.evaluate(
EvaluateRequest::new(url, staging, &resolver)
.engine(EnginePreference::Engine(Engine::HttpMultipart)),
)
.await?;
instruction.set_filename(plan.name.clone());
let digest = odl::hash::HashDigest::parse_cli(&format!("sha256:{}", plan.sha256))
.map_err(|message| OdlError::CliError { message })?;
instruction.add_checksums([digest]);
dlm.download(DownloadRequest::new(instruction, &resolver))
.await
}
#[cfg(feature = "self-update")]
async fn run_update(
args: &Args,
check_only: bool,
assume_yes: bool,
format: OutputFormat,
) -> Result<(), OdlError> {
use odl::self_update;
let current = env!("CARGO_PKG_VERSION");
let exe = std::env::current_exe()
.and_then(|p| p.canonicalize())
.map_err(|e| OdlError::CliError {
message: format!("could not find the running odl binary: {e}"),
})?;
let blocked = self_update::eligibility(&exe).await.err();
if !check_only && let Some(reason) = blocked {
return Err(OdlError::CliError {
message: reason.explain(),
});
}
let cfg = Config::load_from_file(&config_path_for(args)).await?;
let client = install_client(cfg.download())?;
let Some(plan) = self_update::plan(&client, current).await? else {
match format {
OutputFormat::Json => println!(
"{}",
json!({
"type": "update",
"status": "up_to_date",
"current_version": current,
"path": exe.to_string_lossy(),
"can_install": blocked.is_none(),
"blocked_because": blocked.as_ref().map(|r| r.explain()),
})
),
OutputFormat::Text => println!("odl {current} is the latest release."),
}
return Ok(());
};
if check_only {
match format {
OutputFormat::Json => println!(
"{}",
json!({
"type": "update",
"status": "available",
"current_version": current,
"new_version": plan.version,
"tag": plan.tag,
"asset": plan.name,
"size": plan.size,
"path": exe.to_string_lossy(),
"can_install": blocked.is_none(),
"blocked_because": blocked.as_ref().map(|r| r.explain()),
})
),
OutputFormat::Text => {
println!("odl {} is available (running {current}).", plan.version);
match &blocked {
None => println!("Run `odl update` to install it over {}.", exe.display()),
Some(reason) => println!("odl cannot install it: {}", reason.explain()),
}
}
}
return Ok(());
}
if !assume_yes && format == OutputFormat::Text {
eprintln!(
"This replaces {} with odl {} from {}.",
exe.display(),
plan.version,
plan.url
);
match confirm("Update now?") {
Some(true) => {}
Some(false) => {
eprintln!("Left odl {current} in place.");
return Ok(());
}
None => {
return Err(OdlError::CliError {
message: "no terminal to confirm on; re-run `odl update -y` to update without asking"
.to_string(),
});
}
}
}
let dlm = build_download_manager(args).await?;
let archive = download_release_asset(&dlm, &plan).await?;
let replaced = self_update::finish(&archive, &exe).await?;
let _ = tokio::fs::remove_file(&archive).await;
match format {
OutputFormat::Json => println!(
"{}",
json!({
"type": "update",
"status": "updated",
"previous_version": current,
"new_version": plan.version,
"tag": plan.tag,
"path": replaced.to_string_lossy(),
})
),
OutputFormat::Text => println!(
"Updated {} from {current} to {}.",
replaced.display(),
plan.version
),
}
Ok(())
}
#[cfg(feature = "self-update")]
async fn download_release_asset(
dlm: &DownloadManager,
plan: &odl::self_update::UpdatePlan,
) -> Result<PathBuf, OdlError> {
let url = Url::parse(&plan.url).map_err(|e| OdlError::CliError {
message: format!("release asset URL is not usable: {e}"),
})?;
let staging = odl::self_update::staging_dir();
tokio::fs::create_dir_all(&staging).await?;
let resolver = ForcedResolver {};
let mut instruction = dlm
.evaluate(
EvaluateRequest::new(url, staging, &resolver)
.engine(EnginePreference::Engine(Engine::HttpMultipart)),
)
.await?;
instruction.set_filename(plan.name.clone());
let digest = odl::hash::HashDigest::parse_cli(&format!("sha256:{}", plan.sha256))
.map_err(|message| OdlError::CliError { message })?;
instruction.add_checksums([digest]);
dlm.download(DownloadRequest::new(instruction, &resolver))
.await
}
#[cfg(feature = "ytdlp")]
async fn run_tools(
args: &Args,
action: &args::ToolsAction,
format: OutputFormat,
) -> Result<(), OdlError> {
use odl::ytdlp::install::Tool;
let config_path = config_path_for(args);
let mut cfg = Config::load_from_file(&config_path).await?;
let dlm = build_download_manager(args).await?;
let found = odl::ytdlp::tools(cfg.ytdlp()).await.ok();
let ytdlp_path = found.as_ref().map(|t| t.ytdlp.clone());
let ffmpeg_path = found.as_ref().and_then(|t| t.ffmpeg.clone());
match action {
args::ToolsAction::Status => {
match format {
OutputFormat::Json => {
let v = json!({
"type": "tools",
"config_path": config_path.to_string_lossy(),
"tools_dir": odl::ytdlp::install::tools_dir().to_string_lossy(),
"yt_dlp": ytdlp_path.as_ref().map(|p| p.to_string_lossy()),
"ffmpeg": ffmpeg_path.as_ref().map(|p| p.to_string_lossy()),
"can_install_yt_dlp": odl::ytdlp::install::can_install(Tool::Ytdlp),
"can_install_ffmpeg": odl::ytdlp::install::can_install(Tool::Ffmpeg),
});
println!("{v}");
}
OutputFormat::Text => {
println!(
"yt-dlp: {}",
ytdlp_path
.as_ref()
.map(|p| p.display().to_string())
.unwrap_or_else(|| "not installed".to_owned())
);
println!(
"ffmpeg: {}",
ffmpeg_path
.as_ref()
.map(|p| p.display().to_string())
.unwrap_or_else(|| "not installed".to_owned())
);
println!("config: {}", config_path.display());
}
}
Ok(())
}
args::ToolsAction::Install { tool, yes } => {
let wanted: Vec<Tool> = match tool {
Some(args::ToolChoice::Ytdlp) => vec![Tool::Ytdlp],
Some(args::ToolChoice::Ffmpeg) => vec![Tool::Ffmpeg],
None => vec![Tool::Ytdlp, Tool::Ffmpeg],
};
let mut changed = false;
for t in wanted {
let current = match t {
Tool::Ytdlp => ytdlp_path.clone(),
Tool::Ffmpeg => ffmpeg_path.clone(),
};
let outcome = offer_tool(
t,
current.as_deref(),
&config_path,
cfg.download(),
&dlm,
*yes,
)
.await?;
let mut ytdlp_cfg = cfg.ytdlp().clone();
match (t, outcome) {
(Tool::Ytdlp, OfferOutcome::Installed(path)) => {
ytdlp_cfg.set_binary_path(Some(path));
ytdlp_cfg.set_offer_ytdlp_install(true);
changed = true;
}
(Tool::Ffmpeg, OfferOutcome::Installed(path)) => {
ytdlp_cfg.set_ffmpeg_path(Some(path));
ytdlp_cfg.set_offer_ffmpeg_install(true);
changed = true;
}
(Tool::Ytdlp, OfferOutcome::Declined) => {
ytdlp_cfg.set_offer_ytdlp_install(false);
changed = true;
}
(Tool::Ffmpeg, OfferOutcome::Declined) => {
ytdlp_cfg.set_offer_ffmpeg_install(false);
changed = true;
}
(_, OfferOutcome::NothingToDo) => {}
}
cfg = cfg.into_builder().ytdlp(ytdlp_cfg).build()?;
}
if changed {
cfg.save_to_file(&config_path).await?;
eprintln!("Updated {}", config_path.display());
}
Ok(())
}
}
}
#[cfg(not(feature = "ytdlp"))]
async fn run_tools(
_args: &Args,
_action: &args::ToolsAction,
_format: OutputFormat,
) -> Result<(), OdlError> {
Err(OdlError::CliError {
message: "this build of odl was compiled without yt-dlp support".to_owned(),
})
}
async fn run_probe(args: &Args, url_str: &str, format: OutputFormat) -> Result<(), OdlError> {
let url = Url::parse(url_str).map_err(|e| OdlError::CliError {
message: format!("Invalid URL '{url_str}': {e}"),
})?;
let dlm = build_download_manager(args).await?;
let save_dir = std::env::current_dir()?;
let resolver = ForcedResolver;
let _permit = dlm.acquire_download_permit().await?;
let instruction = dlm
.evaluate(EvaluateRequest::new(url, save_dir, &resolver))
.await?;
let checksums: Vec<serde_json::Value> = instruction
.as_metadata()
.checksums
.iter()
.map(checksum_json)
.collect();
let last_modified_rfc3339 = instruction.last_modified_as_date().map(|d| d.to_rfc3339());
let engine = instruction.engine();
let caps = engine.capabilities();
match format {
OutputFormat::Json => {
let v = json!({
"type": "probe",
"url": instruction.url().to_string(),
"filename": instruction.filename(),
"size": instruction.size(),
"size_is_approx": instruction.size_is_approx(),
"engine": engine.as_str(),
"quality": instruction.quality().map(|q| q.to_string()),
"resumable": instruction.is_resumable(),
"etag": instruction.etag(),
"last_modified": instruction.last_modified(),
"last_modified_rfc3339": last_modified_rfc3339,
"requires_auth": instruction.requires_auth(),
"requires_basic_auth": instruction.requires_basic_auth(),
"checksums": checksums,
});
println!("{v}");
}
OutputFormat::Text => {
println!("url: {}", instruction.url());
println!("filename: {}", instruction.filename());
if engine != Engine::HttpMultipart {
println!("engine: {}", engine.as_str());
}
if let Some(quality) = instruction.quality() {
println!("quality: {quality}");
}
let approx = if instruction.size_is_approx() {
"~"
} else {
""
};
match instruction.size() {
Some(s) => println!("size: {}{} ({} bytes)", approx, HumanBytes(s), s),
None => println!("size: unknown"),
}
println!("resumable: {}", instruction.is_resumable());
if caps.response_headers {
println!(
"etag: {}",
instruction.etag().as_deref().unwrap_or("-")
);
println!(
"last_modified: {}",
last_modified_rfc3339.as_deref().unwrap_or("-")
);
println!("requires_auth: {}", instruction.requires_auth());
}
for c in instruction.as_metadata().checksums.iter() {
use odl::download_metadata::{ChecksumAlgorithm, ChecksumEncoding};
let algo = ChecksumAlgorithm::try_from(c.algorithm)
.map(|a| a.as_str_name())
.unwrap_or("unknown");
let enc = ChecksumEncoding::try_from(c.encoding)
.map(|e| e.as_str_name())
.unwrap_or("unknown");
println!("checksum: {} {} ({})", algo, c.digest, enc);
}
}
}
Ok(())
}
fn percent_complete(downloaded: u64, size: Option<u64>, finished: bool) -> Option<f64> {
if finished {
return Some(100.0);
}
match size {
Some(total) if total > 0 => Some((downloaded as f64 / total as f64) * 100.0),
_ => None,
}
}
fn status_json(d: &DownloadStatus) -> serde_json::Value {
json!({
"filename": d.filename,
"url": d.url,
"save_dir": d.save_dir.to_string_lossy(),
"final_file_path": d.final_file_path.to_string_lossy(),
"final_file_exists": d.final_file_exists,
"download_dir": d.download_dir.to_string_lossy(),
"size": d.size,
"downloaded": d.downloaded,
"percent": percent_complete(d.downloaded, d.size, d.finished),
"finished": d.finished,
"resumable": d.is_resumable,
"parts_total": d.parts_total,
"parts_finished": d.parts_finished,
"engine": d.engine.as_str(),
"size_is_approx": d.size_is_approx,
"quality": d.quality.as_ref().map(|q| q.to_string()),
})
}
async fn run_status(
args: &Args,
filter: Option<&str>,
format: OutputFormat,
brief: bool,
) -> Result<(), OdlError> {
let dlm = build_download_manager(args).await?;
let mut downloads = dlm.list_downloads().await?;
if let Some(f) = filter {
downloads.retain(|d| d.url.contains(f) || d.filename.contains(f));
}
match format {
OutputFormat::Json => {
let items: Vec<serde_json::Value> = downloads.iter().map(status_json).collect();
let v = json!({"type": "status", "count": items.len(), "downloads": items});
println!("{v}");
}
OutputFormat::Text => {
if downloads.is_empty() {
println!("No tracked downloads.");
return Ok(());
}
for d in &downloads {
let pct = percent_complete(d.downloaded, d.size, d.finished)
.map(|p| format!("{p:.1}%"))
.unwrap_or_else(|| "?%".to_string());
let state = if d.finished {
"done"
} else if d.final_file_exists {
"assembled"
} else {
"partial"
};
if brief {
println!("{:>10} {:>9} {}", state, pct, d.filename);
} else {
let caps = d.engine.capabilities();
println!("{}", d.filename);
println!(" url: {}", d.url);
println!(" state: {}", state);
if d.engine != Engine::HttpMultipart {
println!(" engine: {}", d.engine.as_str());
}
if let Some(quality) = &d.quality {
println!(" quality: {quality}");
}
let approx = if d.size_is_approx { "~" } else { "" };
let downloaded = if d.finished {
d.size.unwrap_or(d.downloaded)
} else {
d.downloaded
};
match d.size {
Some(s) => println!(
" progress: {} / {}{} ({})",
HumanBytes(downloaded),
approx,
HumanBytes(s),
pct
),
None if d.finished => println!(" progress: complete"),
None => println!(" progress: {} / unknown", HumanBytes(downloaded)),
}
if caps.multipart {
println!(" parts: {}/{}", d.parts_finished, d.parts_total);
}
println!(" resumable: {}", d.is_resumable);
println!(" final file: {}", d.final_file_path.display());
}
}
}
}
Ok(())
}