use crate::cache;
use crate::cli::Cli;
use crate::commands::{
convert_format, format_to_provider_format, language_to_str, output_dry_run,
output_error_traced, output_success, parse_video_id_from_url, with_deadline, DeliveredLanguage,
ResolvedTarget, TargetSource,
};
use crate::error::{AppError, AppResult};
use crate::io;
use crate::provider::{ProviderAttempt, ProviderChain};
use futures::stream::{self, StreamExt};
use std::collections::{BTreeMap, BTreeSet};
use std::path::PathBuf;
use std::process::ExitCode;
use std::time::Instant;
use tracing::instrument;
const DEFAULT_MAX_JOBS: usize = 8;
fn max_jobs() -> usize {
crate::config::tuning_usize_in_range("cli.max_jobs", DEFAULT_MAX_JOBS, 1, 64)
}
fn effective_jobs(cli: &Cli, total: usize) -> usize {
let requested = if cli.jobs > 0 {
usize::try_from(cli.jobs).unwrap_or(1)
} else {
std::thread::available_parallelism().map_or(1, std::num::NonZeroUsize::get)
};
requested.clamp(1, max_jobs()).min(total.max(1))
}
enum ItemOutput {
Success {
provider: &'static str,
video_id: String,
converted: String,
source: String,
duration_ms: u64,
delivered_language: DeliveredLanguage,
},
DryRun {
video_id: String,
},
Failed {
error: AppError,
video_id: Option<String>,
attempts: Vec<ProviderAttempt>,
},
}
struct ItemResult {
idx: usize,
target: ResolvedTarget,
outputs: Vec<ItemOutput>,
}
#[instrument(skip(cli, chain), fields(total))]
pub async fn run(cli: &Cli, chain: &ProviderChain) -> AppResult<ExitCode> {
let urls = io::read_urls_from_stdin(cli.no_input).await?;
tracing::Span::current().record("total", urls.len());
let unique = dedupe_preserving_order(urls);
let mut progress = BatchProgress::open(cli.resume);
let unique: Vec<String> = unique
.into_iter()
.filter(|url| !progress.already_done(url))
.collect();
let total = unique.len();
let jobs = effective_jobs(cli, total);
tracing::info!(target: "events", event = "batch_started", total = total, jobs = jobs);
let mut tally = Tally::default();
let mut parked: BTreeMap<usize, ItemResult> = BTreeMap::new();
let mut next = 0usize;
let mut stream = stream::iter(unique.iter().enumerate())
.map(|(idx, url)| process_item(cli, chain, idx, total, url))
.buffer_unordered(jobs);
while let Some(result) = stream.next().await {
parked.insert(result.idx, result);
while let Some(item) = parked.remove(&next) {
emit_result(cli, total, item, &mut tally, &mut progress).await;
next += 1;
}
}
debug_assert!(
parked.is_empty(),
"{} items never reached the emit window",
parked.len()
);
let Tally {
last_err,
succeeded,
..
} = tally;
let lost = total.saturating_sub(succeeded as usize);
if !cli.quiet {
if lost > 0 {
tracing::warn!(
target: "events",
event = "batch_completed",
succeeded,
total,
lost,
);
} else {
tracing::info!(target: "events", event = "batch_completed", succeeded, total);
}
}
match last_err {
Some(e) if lost > 0 => Ok(ExitCode::from(e.exit_code())),
_ => Ok(ExitCode::SUCCESS),
}
}
#[derive(Default)]
struct Tally {
last_err: Option<AppError>,
succeeded: u32,
written: usize,
}
struct BatchProgress {
path: Option<PathBuf>,
done: std::collections::HashSet<String>,
}
impl BatchProgress {
fn open(resume: bool) -> Self {
Self::at(
crate::config::state_dir()
.ok()
.map(|d| d.join("batch-progress.txt")),
resume,
)
}
fn at(path: Option<PathBuf>, resume: bool) -> Self {
let mut done = std::collections::HashSet::new();
match (&path, resume) {
(Some(p), true) => {
if let Ok(text) = std::fs::read_to_string(p) {
done.extend(text.lines().map(str::to_owned).filter(|l| !l.is_empty()));
}
}
(Some(p), false) => {
let _ = std::fs::remove_file(p);
}
(None, _) => {}
}
Self { path, done }
}
fn already_done(&self, url: &str) -> bool {
self.done.contains(url)
}
fn record(&mut self, url: &str) {
if !self.done.insert(url.to_owned()) {
return;
}
let Some(path) = &self.path else {
return;
};
if let Some(parent) = path.parent() {
let _ = std::fs::create_dir_all(parent);
}
if let Ok(mut file) = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(path)
{
use std::io::Write as _;
let _ = writeln!(file, "{url}");
let _ = file.flush();
}
}
}
async fn emit_result(
cli: &Cli,
total: usize,
result: ItemResult,
tally: &mut Tally,
progress: &mut BatchProgress,
) {
for output in result.outputs {
match output {
ItemOutput::Success {
provider,
video_id,
converted,
source,
duration_ms,
delivered_language,
} => {
match write_item(
cli,
tally.written,
total,
provider,
&video_id,
&result.target,
&converted,
&source,
duration_ms,
&delivered_language,
)
.await
{
Ok(()) => {
tally.succeeded += 1;
tally.written += 1;
progress.record(&result.target.value);
}
Err(e) => {
output_error_traced(cli, &e, Some(&result.target), Some(&video_id), &[])
.await
.ok();
tally.last_err = Some(e);
}
}
}
ItemOutput::DryRun { video_id } => {
tracing::info!(target: "events", event = "dry_run_cache_miss", video_id = %video_id);
let _ = output_dry_run(cli, &video_id, &result.target, true).await;
}
ItemOutput::Failed {
error,
video_id,
attempts,
} => {
output_error_traced(
cli,
&error,
Some(&result.target),
video_id.as_deref(),
&attempts,
)
.await
.ok();
tally.last_err = Some(error);
}
}
}
}
async fn process_item(
cli: &Cli,
chain: &ProviderChain,
idx: usize,
total: usize,
url: &str,
) -> ItemResult {
if !cli.quiet && !cli.no_progress {
tracing::info!(target: "events", event = "progress", index = idx + 1, total, url = %url);
}
let target = ResolvedTarget {
value: url.to_string(),
source: TargetSource::BatchFile,
};
let mut outputs = Vec::new();
let video_id = match parse_video_id_from_url(cli, url) {
Ok(id) => id,
Err(e) => {
if !cli.json {
tracing::error!(target: "events", event = "failed", url = %url, error = %e);
}
outputs.push(ItemOutput::Failed {
error: e,
video_id: None,
attempts: Vec::new(),
});
return ItemResult {
idx,
target,
outputs,
};
}
};
let lang = language_to_str(cli.lang);
let format = format_to_provider_format(cli.format);
let cache_ttl = cli.cache_ttl_duration();
let started = Instant::now();
if !cli.no_cache {
if let Ok(path) = cache::cache_path(&video_id, lang, format.extension(), cache_ttl) {
if let Ok(Some((bytes, format_hint))) =
cache::read_cache_with_hint(&path, cache_ttl).await
{
let duration_ms = started.elapsed().as_millis() as u64;
match convert_format(&bytes, cli.format, format_hint) {
Ok(converted) => {
outputs.push(ItemOutput::Success {
provider: "cache",
video_id,
converted,
source: "cache".to_string(),
duration_ms,
delivered_language: DeliveredLanguage::unknown(),
});
return ItemResult {
idx,
target,
outputs,
};
}
Err(e) => outputs.push(ItemOutput::Failed {
error: e,
video_id: Some(video_id.clone()),
attempts: Vec::new(),
}),
}
}
}
}
if cli.dry_run {
outputs.push(ItemOutput::DryRun { video_id });
return ItemResult {
idx,
target,
outputs,
};
}
let traced = with_deadline(cli, async {
Ok(chain.fetch_subtitle_traced(&video_id, lang, format).await)
})
.await;
let (result, attempts) = match traced {
Ok(pair) => pair,
Err(e) => (Err(e), Vec::new()),
};
match result {
Ok((info, content)) => {
tracing::info!(target: "events", event = "fetched", video_id = %video_id, source = %info.source_url);
if !cli.no_cache {
if let Ok(path) = cache::cache_path(&video_id, lang, format.extension(), cache_ttl)
{
let _ = cache::write_cache_with_hint(&path, &content, info.format_hint).await;
}
}
match convert_format(&content, cli.format, info.format_hint) {
Ok(converted) => outputs.push(ItemOutput::Success {
provider: info.provider,
video_id,
converted,
source: info.source_url.clone(),
duration_ms: started.elapsed().as_millis() as u64,
delivered_language: DeliveredLanguage::observed(&info),
}),
Err(e) => outputs.push(ItemOutput::Failed {
error: e,
video_id: Some(video_id),
attempts,
}),
}
}
Err(e) => {
if !cli.json {
tracing::error!(target: "events", event = "failed", video_id = %video_id, error = %e);
}
outputs.push(ItemOutput::Failed {
error: e,
video_id: Some(video_id),
attempts,
});
}
}
ItemResult {
idx,
target,
outputs,
}
}
fn dedupe_preserving_order(urls: Vec<String>) -> Vec<String> {
let keep: Vec<bool> = {
let mut seen: BTreeSet<&str> = BTreeSet::new();
urls.iter().map(|u| seen.insert(u.as_str())).collect()
};
let mut keep_iter = keep.into_iter();
let mut unique = urls;
unique.retain(|_| keep_iter.next().unwrap_or(false));
unique
}
#[allow(clippy::too_many_arguments)]
async fn write_item(
cli: &Cli,
idx: usize,
total: usize,
provider: &'static str,
video_id: &str,
target: &ResolvedTarget,
converted: &str,
source: &str,
duration_ms: u64,
delivered_language: &DeliveredLanguage,
) -> AppResult<()> {
if cli.json {
output_success(
cli,
provider,
video_id,
target,
converted,
source,
duration_ms,
delivered_language,
)
.await?;
} else {
let separator = if idx > 0 { "\n---\n" } else { "" };
let mut buf = Vec::with_capacity(separator.len() + converted.len() + 1);
buf.extend_from_slice(separator.as_bytes());
buf.extend_from_slice(converted.as_bytes());
buf.push(b'\n');
io::write_subtitle_to_stdout(&buf).await?;
}
let _ = total;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::cli::Cli;
use clap::Parser;
#[test]
fn dry_run_envelope_uses_stable_event_field() {
let cli = Cli::parse_from(["youtube-legend-cli", "--batch", "--dry-run", "--json"]);
let payload = serde_json::json!({
"event": "dry_run_cache_miss",
"video_id": "dQw4w9WgXcQ",
"language": crate::commands::language_to_str(cli.lang),
"format": crate::commands::format_to_str(cli.format),
"would_fetch": true,
});
let json = payload.to_string();
assert!(
json.contains("\"event\":\"dry_run_cache_miss\""),
"envelope must carry the stable event field, got: {json}"
);
assert!(
json.contains("\"would_fetch\":true"),
"envelope must include the would_fetch flag, got: {json}"
);
}
#[test]
fn dry_run_does_not_populate_last_err() {
let cli = Cli::parse_from(["youtube-legend-cli", "--batch", "--dry-run"]);
assert!(cli.dry_run, "dry-run flag must be parsed");
assert!(cli.batch, "batch flag must be parsed");
let _ = output_dry_run;
}
#[test]
fn dedupe_keeps_first_occurrence_in_input_order() {
let urls = vec![
"https://youtu.be/zzz".to_string(),
"https://youtu.be/aaa".to_string(),
"https://youtu.be/zzz".to_string(),
"https://youtu.be/mmm".to_string(),
"https://youtu.be/aaa".to_string(),
];
let unique = dedupe_preserving_order(urls);
assert_eq!(
unique,
vec![
"https://youtu.be/zzz".to_string(),
"https://youtu.be/aaa".to_string(),
"https://youtu.be/mmm".to_string(),
]
);
}
#[test]
fn effective_jobs_respects_the_ceiling_and_the_item_count() {
let cli = Cli::parse_from(["youtube-legend-cli", "--batch", "--jobs", "4"]);
assert_eq!(effective_jobs(&cli, 10), 4);
assert_eq!(effective_jobs(&cli, 2), 2);
let greedy = Cli::parse_from(["youtube-legend-cli", "--batch", "--jobs", "9999"]);
assert!(effective_jobs(&greedy, 10_000) <= max_jobs());
}
#[test]
fn effective_jobs_derives_a_positive_default() {
let cli = Cli::parse_from(["youtube-legend-cli", "--batch"]);
assert!(effective_jobs(&cli, 100) >= 1);
assert!(effective_jobs(&cli, 0) >= 1);
}
#[test]
fn progress_is_remembered_on_resume_and_forgotten_otherwise() {
let path =
std::env::temp_dir().join(format!("ylc-batch-progress-a-{}.txt", std::process::id()));
let _ = std::fs::remove_file(&path);
let mut first = BatchProgress::at(Some(path.clone()), false);
first.record("https://youtu.be/aaaaaaaaaaa");
first.record("https://youtu.be/bbbbbbbbbbb");
let resumed = BatchProgress::at(Some(path.clone()), true);
assert!(resumed.already_done("https://youtu.be/aaaaaaaaaaa"));
assert!(resumed.already_done("https://youtu.be/bbbbbbbbbbb"));
assert!(!resumed.already_done("https://youtu.be/ccccccccccc"));
let fresh = BatchProgress::at(Some(path.clone()), false);
assert!(
!fresh.already_done("https://youtu.be/aaaaaaaaaaa"),
"a run without --resume must start clean"
);
assert!(
!path.exists(),
"a run without --resume must remove the previous file"
);
}
#[test]
fn progress_records_each_url_once() {
let path =
std::env::temp_dir().join(format!("ylc-batch-progress-b-{}.txt", std::process::id()));
let _ = std::fs::remove_file(&path);
let mut progress = BatchProgress::at(Some(path.clone()), false);
progress.record("https://youtu.be/aaaaaaaaaaa");
progress.record("https://youtu.be/aaaaaaaaaaa");
let text = std::fs::read_to_string(&path).expect("file written");
assert_eq!(text.lines().count(), 1, "duplicate line written: {text:?}");
}
#[test]
fn progress_without_a_path_never_fails() {
let mut progress = BatchProgress::at(None, true);
progress.record("https://youtu.be/aaaaaaaaaaa");
assert!(progress.already_done("https://youtu.be/aaaaaaaaaaa"));
}
#[test]
fn the_reorder_window_emits_in_input_order() {
let arrivals = [3usize, 1, 0, 4, 2];
let mut parked: BTreeMap<usize, usize> = BTreeMap::new();
let mut next = 0usize;
let mut emitted = Vec::new();
for idx in arrivals {
parked.insert(idx, idx);
while let Some(item) = parked.remove(&next) {
emitted.push(item);
next += 1;
}
}
assert!(parked.is_empty(), "the window did not drain");
assert_eq!(emitted, vec![0, 1, 2, 3, 4]);
}
#[test]
fn results_are_reordered_by_input_position() {
let mut results: Vec<ItemResult> = [2usize, 0, 1]
.into_iter()
.map(|idx| ItemResult {
idx,
target: ResolvedTarget {
value: format!("https://youtu.be/{idx}"),
source: TargetSource::BatchFile,
},
outputs: Vec::new(),
})
.collect();
results.sort_by_key(|r| r.idx);
assert_eq!(
results.iter().map(|r| r.idx).collect::<Vec<_>>(),
vec![0, 1, 2]
);
}
#[test]
fn dedupe_handles_empty_and_single_input() {
assert!(dedupe_preserving_order(Vec::new()).is_empty());
assert_eq!(
dedupe_preserving_order(vec!["https://youtu.be/aaa".to_string()]),
vec!["https://youtu.be/aaa".to_string()]
);
}
}