#![doc = include_str!("../README.md")]
use std::ffi::OsString;
use std::fmt;
use std::io::{BufRead, BufReader, Read};
use std::path::{Path, PathBuf};
use std::process::{Child, Command, Stdio};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::{self, Receiver, Sender};
use std::sync::Arc;
use std::thread;
use std::time::{Duration, Instant};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ConversionMode {
ExtractAudio,
TranscodeVideo,
TranscodeAudio,
}
impl ConversionMode {
pub const ALL: [ConversionMode; 3] = [
ConversionMode::ExtractAudio,
ConversionMode::TranscodeVideo,
ConversionMode::TranscodeAudio,
];
pub fn title(self) -> &'static str {
match self {
ConversionMode::ExtractAudio => "Video -> Audio",
ConversionMode::TranscodeVideo => "Video -> Video",
ConversionMode::TranscodeAudio => "Audio -> Audio",
}
}
pub fn description(self) -> &'static str {
match self {
ConversionMode::ExtractAudio => "Strip video and encode the audio track only.",
ConversionMode::TranscodeVideo => "Re-encode video into another container/codec route.",
ConversionMode::TranscodeAudio => "Convert one audio file into another audio format.",
}
}
pub fn source_label(self) -> &'static str {
match self {
ConversionMode::ExtractAudio | ConversionMode::TranscodeVideo => "VIDEO SOURCE",
ConversionMode::TranscodeAudio => "AUDIO SOURCE",
}
}
pub fn note_text(self) -> &'static str {
match self {
ConversionMode::ExtractAudio => {
"EXTRACT ROUTE // SOURCE VIDEO -> AUDIO CODEC -> OUTPUT CONTAINER"
}
ConversionMode::TranscodeVideo => {
"VIDEO ROUTE // SOURCE STREAMS -> TRANSCODE MATRIX -> OUTPUT CONTAINER"
}
ConversionMode::TranscodeAudio => {
"AUDIO ROUTE // SOURCE AUDIO -> CODEC CONVERSION -> OUTPUT FILE"
}
}
}
pub fn expects_video_input(self) -> bool {
matches!(
self,
ConversionMode::ExtractAudio | ConversionMode::TranscodeVideo
)
}
pub fn output_extension(self, audio: AudioFormat, video: VideoFormat) -> &'static str {
match self {
ConversionMode::ExtractAudio | ConversionMode::TranscodeAudio => audio.extension(),
ConversionMode::TranscodeVideo => video.extension(),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AudioFormat {
Mp3,
M4a,
Flac,
Wav,
Ogg,
Opus,
}
impl AudioFormat {
pub const ALL: [AudioFormat; 6] = [
AudioFormat::Mp3,
AudioFormat::M4a,
AudioFormat::Flac,
AudioFormat::Wav,
AudioFormat::Ogg,
AudioFormat::Opus,
];
pub fn label(self) -> &'static str {
match self {
AudioFormat::Mp3 => "MP3",
AudioFormat::M4a => "M4A/AAC",
AudioFormat::Flac => "FLAC",
AudioFormat::Wav => "WAV",
AudioFormat::Ogg => "OGG",
AudioFormat::Opus => "OPUS",
}
}
pub fn extension(self) -> &'static str {
match self {
AudioFormat::Mp3 => "mp3",
AudioFormat::M4a => "m4a",
AudioFormat::Flac => "flac",
AudioFormat::Wav => "wav",
AudioFormat::Ogg => "ogg",
AudioFormat::Opus => "opus",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum VideoFormat {
Mp4,
Mkv,
Webm,
Mov,
}
impl VideoFormat {
pub const ALL: [VideoFormat; 4] = [
VideoFormat::Mp4,
VideoFormat::Mkv,
VideoFormat::Webm,
VideoFormat::Mov,
];
pub fn label(self) -> &'static str {
match self {
VideoFormat::Mp4 => "MP4",
VideoFormat::Mkv => "MKV",
VideoFormat::Webm => "WEBM",
VideoFormat::Mov => "MOV",
}
}
pub fn extension(self) -> &'static str {
match self {
VideoFormat::Mp4 => "mp4",
VideoFormat::Mkv => "mkv",
VideoFormat::Webm => "webm",
VideoFormat::Mov => "mov",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum QualityPreset {
Compact,
Balanced,
Archive,
}
impl QualityPreset {
pub const ALL: [QualityPreset; 3] = [
QualityPreset::Compact,
QualityPreset::Balanced,
QualityPreset::Archive,
];
pub fn label(self) -> &'static str {
match self {
QualityPreset::Compact => "Compact",
QualityPreset::Balanced => "Balanced",
QualityPreset::Archive => "Archive",
}
}
pub fn description(self) -> &'static str {
match self {
QualityPreset::Compact => "smaller file",
QualityPreset::Balanced => "general use",
QualityPreset::Archive => "higher fidelity",
}
}
fn x264_crf(self) -> &'static str {
match self {
QualityPreset::Compact => "28",
QualityPreset::Balanced => "23",
QualityPreset::Archive => "18",
}
}
fn vp9_crf(self) -> &'static str {
match self {
QualityPreset::Compact => "38",
QualityPreset::Balanced => "32",
QualityPreset::Archive => "24",
}
}
fn audio_bitrate(self) -> &'static str {
match self {
QualityPreset::Compact => "128k",
QualityPreset::Balanced => "192k",
QualityPreset::Archive => "256k",
}
}
fn vorbis_quality(self) -> &'static str {
match self {
QualityPreset::Compact => "3",
QualityPreset::Balanced => "5",
QualityPreset::Archive => "8",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ConversionRequest {
pub mode: ConversionMode,
pub input_path: PathBuf,
pub output_path: PathBuf,
pub audio_format: AudioFormat,
pub video_format: VideoFormat,
pub quality: QualityPreset,
pub overwrite: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum ConversionRoute {
ExtractAudio(AudioFormat),
TranscodeVideo(VideoFormat),
TranscodeAudio(AudioFormat),
}
impl ConversionRoute {
pub fn mode(self) -> ConversionMode {
match self {
ConversionRoute::ExtractAudio(_) => ConversionMode::ExtractAudio,
ConversionRoute::TranscodeVideo(_) => ConversionMode::TranscodeVideo,
ConversionRoute::TranscodeAudio(_) => ConversionMode::TranscodeAudio,
}
}
pub fn audio_format(self) -> Option<AudioFormat> {
match self {
ConversionRoute::ExtractAudio(format) | ConversionRoute::TranscodeAudio(format) => {
Some(format)
}
ConversionRoute::TranscodeVideo(_) => None,
}
}
pub fn video_format(self) -> Option<VideoFormat> {
match self {
ConversionRoute::TranscodeVideo(format) => Some(format),
ConversionRoute::ExtractAudio(_) | ConversionRoute::TranscodeAudio(_) => None,
}
}
pub fn output_extension(self) -> &'static str {
match self {
ConversionRoute::ExtractAudio(format) | ConversionRoute::TranscodeAudio(format) => {
format.extension()
}
ConversionRoute::TranscodeVideo(format) => format.extension(),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
#[non_exhaustive]
pub enum RouteInferenceError {
#[error("unsupported input extension for {0}")]
UnsupportedInput(PathBuf),
#[error("unsupported output extension for {0}")]
UnsupportedOutput(PathBuf),
#[error("audio input cannot be converted into a video output")]
AudioToVideo,
}
pub fn infer_route(input: &Path, output: &Path) -> Result<ConversionRoute, RouteInferenceError> {
let input_is_audio = supported_audio_extension(input);
let input_is_video = supported_video_extension(input);
if !input_is_audio && !input_is_video {
return Err(RouteInferenceError::UnsupportedInput(input.to_path_buf()));
}
if let Some(format) = audio_format_from_path(output) {
return if input_is_audio {
Ok(ConversionRoute::TranscodeAudio(format))
} else {
Ok(ConversionRoute::ExtractAudio(format))
};
}
if let Some(format) = video_format_from_path(output) {
return if input_is_audio {
Err(RouteInferenceError::AudioToVideo)
} else {
Ok(ConversionRoute::TranscodeVideo(format))
};
}
Err(RouteInferenceError::UnsupportedOutput(output.to_path_buf()))
}
pub fn request_from_paths(
input: impl AsRef<Path>,
output: impl AsRef<Path>,
) -> Result<ConversionRequest, RouteInferenceError> {
let input = input.as_ref();
let output = output.as_ref();
let route = infer_route(input, output)?;
let (audio_format, video_format) = match route {
ConversionRoute::ExtractAudio(format) | ConversionRoute::TranscodeAudio(format) => {
(format, VideoFormat::Mp4)
}
ConversionRoute::TranscodeVideo(format) => (AudioFormat::Mp3, format),
};
Ok(ConversionRequest {
mode: route.mode(),
input_path: input.to_path_buf(),
output_path: output.to_path_buf(),
audio_format,
video_format,
quality: QualityPreset::Balanced,
overwrite: false,
})
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CommandSpec {
pub program: String,
pub args: Vec<String>,
}
impl fmt::Display for CommandSpec {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(formatter, "{}", self.program)?;
for arg in &self.args {
write!(formatter, " {}", shell_quote(arg))?;
}
Ok(())
}
}
#[derive(Debug, thiserror::Error)]
pub enum ConversionError {
#[error("select a source media file first")]
MissingInput,
#[error("source file does not exist: {0}")]
MissingInputFile(String),
#[error("source path must be a regular file: {0}")]
InputNotFile(String),
#[error("selected route expects a video source: {0}")]
ExpectedVideoInput(String),
#[error("selected route expects an audio source: {0}")]
ExpectedAudioInput(String),
#[error("choose an output file path first")]
MissingOutput,
#[error("output parent folder does not exist: {0}")]
MissingOutputFolder(String),
#[error("output cannot be the same file as the source")]
SameInputOutput,
#[error("output extension should be .{expected} for the selected route")]
OutputExtensionMismatch { expected: &'static str },
#[error("output already exists; enable overwrite or choose another path: {0}")]
OutputExists(String),
}
pub fn build_command(request: &ConversionRequest) -> Result<CommandSpec, ConversionError> {
validate_request(request)?;
let args = ffmpeg_arguments(request)
.into_iter()
.map(|arg| arg.to_string_lossy().into_owned())
.collect();
Ok(CommandSpec {
program: "ffmpeg".to_owned(),
args,
})
}
fn ffmpeg_arguments(request: &ConversionRequest) -> Vec<OsString> {
let mut args = [
"-hide_banner",
"-nostdin",
"-progress",
"pipe:1",
"-nostats",
if request.overwrite { "-y" } else { "-n" },
"-i",
]
.into_iter()
.map(OsString::from)
.collect::<Vec<_>>();
args.push(request.input_path.as_os_str().to_owned());
match request.mode {
ConversionMode::ExtractAudio | ConversionMode::TranscodeAudio => {
push_args(&mut args, ["-map", "0:a:0", "-vn"]);
apply_audio_args(&mut args, request.audio_format, request.quality);
}
ConversionMode::TranscodeVideo => {
push_args(&mut args, ["-map", "0:v:0", "-map", "0:a?", "-sn"]);
apply_video_args(&mut args, request.video_format, request.quality);
}
}
args.push(request.output_path.as_os_str().to_owned());
args
}
pub fn validate_request(request: &ConversionRequest) -> Result<(), ConversionError> {
if request.input_path.as_os_str().is_empty() {
return Err(ConversionError::MissingInput);
}
let input_display = request.input_path.display().to_string();
let metadata = std::fs::metadata(&request.input_path)
.map_err(|_| ConversionError::MissingInputFile(input_display.clone()))?;
if !metadata.is_file() {
return Err(ConversionError::InputNotFile(input_display));
}
if request.mode.expects_video_input() && !supported_video_extension(&request.input_path) {
return Err(ConversionError::ExpectedVideoInput(
request.input_path.display().to_string(),
));
}
if matches!(request.mode, ConversionMode::TranscodeAudio)
&& !supported_audio_extension(&request.input_path)
{
return Err(ConversionError::ExpectedAudioInput(
request.input_path.display().to_string(),
));
}
if request.output_path.as_os_str().is_empty() {
return Err(ConversionError::MissingOutput);
}
if request.input_path == request.output_path
|| canonical_match(&request.input_path, &request.output_path)
{
return Err(ConversionError::SameInputOutput);
}
if let Some(parent) = request.output_path.parent() {
if !parent.as_os_str().is_empty() && !parent.is_dir() {
return Err(ConversionError::MissingOutputFolder(
parent.display().to_string(),
));
}
}
let expected = request
.mode
.output_extension(request.audio_format, request.video_format);
if path_extension_lower(&request.output_path).as_deref() != Some(expected) {
return Err(ConversionError::OutputExtensionMismatch { expected });
}
if request.output_path.exists() && !request.overwrite {
return Err(ConversionError::OutputExists(
request.output_path.display().to_string(),
));
}
Ok(())
}
fn canonical_match(left: &Path, right: &Path) -> bool {
match (std::fs::canonicalize(left), std::fs::canonicalize(right)) {
(Ok(left), Ok(right)) => left == right,
_ => false,
}
}
fn apply_audio_args(args: &mut Vec<OsString>, format: AudioFormat, quality: QualityPreset) {
match format {
AudioFormat::Mp3 => {
push_args(
args,
["-c:a", "libmp3lame", "-b:a", quality.audio_bitrate()],
);
}
AudioFormat::M4a => {
push_args(args, ["-c:a", "aac", "-b:a", quality.audio_bitrate()]);
push_args(args, ["-movflags", "+faststart"]);
}
AudioFormat::Flac => push_args(args, ["-c:a", "flac"]),
AudioFormat::Wav => push_args(args, ["-c:a", "pcm_s16le"]),
AudioFormat::Ogg => {
push_args(
args,
["-c:a", "libvorbis", "-q:a", quality.vorbis_quality()],
);
}
AudioFormat::Opus => {
push_args(args, ["-c:a", "libopus", "-b:a", quality.audio_bitrate()]);
}
}
}
fn apply_video_args(args: &mut Vec<OsString>, format: VideoFormat, quality: QualityPreset) {
match format {
VideoFormat::Mp4 | VideoFormat::Mkv | VideoFormat::Mov => {
push_args(
args,
[
"-c:v",
"libx264",
"-preset",
"medium",
"-crf",
quality.x264_crf(),
"-pix_fmt",
"yuv420p",
"-c:a",
"aac",
"-b:a",
quality.audio_bitrate(),
],
);
if matches!(format, VideoFormat::Mp4 | VideoFormat::Mov) {
push_args(args, ["-movflags", "+faststart"]);
}
}
VideoFormat::Webm => {
push_args(
args,
[
"-c:v",
"libvpx-vp9",
"-b:v",
"0",
"-crf",
quality.vp9_crf(),
"-c:a",
"libopus",
"-b:a",
quality.audio_bitrate(),
],
);
}
}
}
fn push_args<const N: usize>(args: &mut Vec<OsString>, values: [&str; N]) {
args.extend(values.into_iter().map(OsString::from));
}
pub fn probe_duration(path: &Path) -> Option<f64> {
let output = Command::new("ffprobe")
.args([
"-v",
"error",
"-show_entries",
"format=duration",
"-of",
"default=noprint_wrappers=1:nokey=1",
])
.arg(path)
.output()
.ok()?;
if !output.status.success() {
return None;
}
let text = String::from_utf8_lossy(&output.stdout);
text.lines()
.find_map(|line| line.trim().parse::<f64>().ok())
.filter(|duration| *duration > 0.0 && duration.is_finite())
}
#[derive(Debug, Clone, PartialEq)]
pub enum OperationEvent {
Started(String),
Log(String),
Progress {
current_seconds: f64,
total_seconds: f64,
},
Finished(Result<(), String>),
}
pub fn start_operation(
command: CommandSpec,
expected_duration: Option<f64>,
) -> Receiver<OperationEvent> {
let (tx, rx) = mpsc::channel();
thread::spawn(move || run_operation(command, expected_duration, tx));
rx
}
pub fn start_conversion(
request: &ConversionRequest,
) -> Result<Receiver<OperationEvent>, ConversionError> {
let command = build_command(request)?;
let expected_duration = probe_duration(&request.input_path);
Ok(start_operation(command, expected_duration))
}
const DEFAULT_TOOL_CHECK_TIMEOUT: Duration = Duration::from_secs(2);
const DEFAULT_PROBE_TIMEOUT: Duration = Duration::from_secs(10);
const PROCESS_POLL_INTERVAL: Duration = Duration::from_millis(20);
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Converter {
ffmpeg_path: PathBuf,
ffprobe_path: PathBuf,
tool_check_timeout: Duration,
probe_timeout: Duration,
}
impl Default for Converter {
fn default() -> Self {
Self {
ffmpeg_path: PathBuf::from("ffmpeg"),
ffprobe_path: PathBuf::from("ffprobe"),
tool_check_timeout: DEFAULT_TOOL_CHECK_TIMEOUT,
probe_timeout: DEFAULT_PROBE_TIMEOUT,
}
}
}
impl Converter {
pub fn new() -> Self {
Self::default()
}
pub fn with_programs(
ffmpeg_path: impl Into<PathBuf>,
ffprobe_path: impl Into<PathBuf>,
) -> Self {
Self {
ffmpeg_path: ffmpeg_path.into(),
ffprobe_path: ffprobe_path.into(),
..Self::default()
}
}
pub fn with_tool_check_timeout(mut self, timeout: Duration) -> Self {
self.tool_check_timeout = timeout;
self
}
pub fn with_probe_timeout(mut self, timeout: Duration) -> Self {
self.probe_timeout = timeout;
self
}
pub fn ffmpeg_path(&self) -> &Path {
&self.ffmpeg_path
}
pub fn ffprobe_path(&self) -> &Path {
&self.ffprobe_path
}
pub fn check_tools(&self) -> ToolAvailability {
ToolAvailability {
ffmpeg: check_tool(&self.ffmpeg_path, self.tool_check_timeout),
ffprobe: check_tool(&self.ffprobe_path, self.tool_check_timeout),
}
}
pub fn spawn(&self, request: ConversionRequest) -> ConversionJob {
let (tx, events) = mpsc::channel();
let cancellation = Arc::new(AtomicBool::new(false));
let worker_cancellation = Arc::clone(&cancellation);
let converter = self.clone();
let input_path = request.input_path.clone();
let output_path = request.output_path.clone();
thread::spawn(move || {
run_conversion_job(converter, request, worker_cancellation, tx);
});
ConversionJob {
events,
cancellation,
input_path,
output_path,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum ToolStatus {
Available,
Unavailable(String),
TimedOut,
}
impl ToolStatus {
pub fn is_available(&self) -> bool {
matches!(self, ToolStatus::Available)
}
}
impl fmt::Display for ToolStatus {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
ToolStatus::Available => formatter.write_str("available"),
ToolStatus::Unavailable(error) => write!(formatter, "unavailable: {error}"),
ToolStatus::TimedOut => formatter.write_str("timed out"),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ToolAvailability {
ffmpeg: ToolStatus,
ffprobe: ToolStatus,
}
impl ToolAvailability {
pub fn ffmpeg(&self) -> &ToolStatus {
&self.ffmpeg
}
pub fn ffprobe(&self) -> &ToolStatus {
&self.ffprobe
}
pub fn conversion_ready(&self) -> bool {
self.ffmpeg.is_available()
}
pub fn duration_progress_available(&self) -> bool {
self.ffprobe.is_available()
}
}
pub fn check_tools() -> ToolAvailability {
Converter::new().check_tools()
}
fn check_tool(program: &Path, timeout: Duration) -> ToolStatus {
let mut command = Command::new(program);
command
.arg("-version")
.stdin(Stdio::null())
.stdout(Stdio::null())
.stderr(Stdio::null());
configure_process_group(&mut command);
let mut child = match command.spawn() {
Ok(child) => child,
Err(error) => return ToolStatus::Unavailable(error.to_string()),
};
let started = Instant::now();
loop {
match child.try_wait() {
Ok(Some(status)) if status.success() => return ToolStatus::Available,
Ok(Some(status)) => return ToolStatus::Unavailable(format!("exited with {status}")),
Ok(None) if started.elapsed() >= timeout => {
terminate_child(&mut child);
return ToolStatus::TimedOut;
}
Ok(None) => thread::sleep(PROCESS_POLL_INTERVAL),
Err(error) => {
terminate_child(&mut child);
return ToolStatus::Unavailable(error.to_string());
}
}
}
}
fn configure_process_group(command: &mut Command) {
#[cfg(unix)]
{
use std::os::unix::process::CommandExt;
command.process_group(0);
}
}
fn terminate_child(child: &mut Child) {
#[cfg(unix)]
if let Ok(process_group) = i32::try_from(child.id()) {
unsafe {
libc::kill(-process_group, libc::SIGKILL);
}
}
let _ = child.kill();
let _ = child.wait();
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum ConversionOutcome {
Succeeded,
Cancelled,
Failed(String),
}
impl ConversionOutcome {
pub fn is_success(&self) -> bool {
matches!(self, ConversionOutcome::Succeeded)
}
pub fn is_cancelled(&self) -> bool {
matches!(self, ConversionOutcome::Cancelled)
}
}
#[derive(Debug, Clone, PartialEq)]
#[non_exhaustive]
pub enum ConversionEvent {
Started(String),
Log(String),
#[non_exhaustive]
Progress {
current_seconds: f64,
total_seconds: Option<f64>,
},
Finished(ConversionOutcome),
}
#[must_use = "dropping a conversion job requests cancellation"]
pub struct ConversionJob {
events: Receiver<ConversionEvent>,
cancellation: Arc<AtomicBool>,
input_path: PathBuf,
output_path: PathBuf,
}
impl ConversionJob {
pub fn input_path(&self) -> &Path {
&self.input_path
}
pub fn output_path(&self) -> &Path {
&self.output_path
}
pub fn try_recv(&self) -> Result<ConversionEvent, mpsc::TryRecvError> {
self.events.try_recv()
}
pub fn recv(&self) -> Result<ConversionEvent, mpsc::RecvError> {
self.events.recv()
}
pub fn recv_timeout(
&self,
timeout: Duration,
) -> Result<ConversionEvent, mpsc::RecvTimeoutError> {
self.events.recv_timeout(timeout)
}
pub fn cancel(&self) {
self.cancellation.store(true, Ordering::Release);
}
pub fn cancellation_requested(&self) -> bool {
self.cancellation.load(Ordering::Acquire)
}
}
impl Drop for ConversionJob {
fn drop(&mut self) {
self.cancel();
}
}
pub fn spawn_conversion(request: ConversionRequest) -> ConversionJob {
Converter::new().spawn(request)
}
fn run_conversion_job(
converter: Converter,
request: ConversionRequest,
cancellation: Arc<AtomicBool>,
tx: Sender<ConversionEvent>,
) {
if cancellation.load(Ordering::Acquire) {
let _ = tx.send(ConversionEvent::Finished(ConversionOutcome::Cancelled));
return;
}
if let Err(error) = validate_request(&request) {
let outcome = if cancellation.load(Ordering::Acquire) {
ConversionOutcome::Cancelled
} else {
ConversionOutcome::Failed(error.to_string())
};
let _ = tx.send(ConversionEvent::Finished(outcome));
return;
}
let expected_duration = match probe_duration_for_job(
&converter.ffprobe_path,
&request.input_path,
converter.probe_timeout,
&cancellation,
) {
ProbeOutcome::Duration(duration) => duration,
ProbeOutcome::Cancelled => {
let _ = tx.send(ConversionEvent::Finished(ConversionOutcome::Cancelled));
return;
}
};
if cancellation.load(Ordering::Acquire) {
let _ = tx.send(ConversionEvent::Finished(ConversionOutcome::Cancelled));
return;
}
run_job_operation(
&converter.ffmpeg_path,
ffmpeg_arguments(&request),
expected_duration,
cancellation,
tx,
);
}
enum ProbeOutcome {
Duration(Option<f64>),
Cancelled,
}
fn probe_duration_for_job(
program: &Path,
path: &Path,
timeout: Duration,
cancellation: &AtomicBool,
) -> ProbeOutcome {
if cancellation.load(Ordering::Acquire) {
return ProbeOutcome::Cancelled;
}
let mut command = Command::new(program);
command
.args([
"-v",
"error",
"-show_entries",
"format=duration",
"-of",
"default=noprint_wrappers=1:nokey=1",
])
.arg(path)
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::null());
configure_process_group(&mut command);
let mut child = match command.spawn() {
Ok(child) => child,
Err(_) => return ProbeOutcome::Duration(None),
};
let started = Instant::now();
let status = loop {
match child.try_wait() {
Ok(Some(status)) => break Some(status),
Ok(None) if cancellation.load(Ordering::Acquire) => {
terminate_child(&mut child);
return ProbeOutcome::Cancelled;
}
Ok(None) if started.elapsed() >= timeout => {
terminate_child(&mut child);
break None;
}
Ok(None) => thread::sleep(PROCESS_POLL_INTERVAL),
Err(_) => {
terminate_child(&mut child);
break None;
}
}
};
if !matches!(status, Some(status) if status.success()) {
return ProbeOutcome::Duration(None);
}
let mut stdout = String::new();
let duration = child
.stdout
.take()
.and_then(|mut stream| stream.read_to_string(&mut stdout).ok())
.and_then(|_| {
stdout
.lines()
.find_map(|line| line.trim().parse::<f64>().ok())
})
.filter(|duration| *duration > 0.0 && duration.is_finite());
ProbeOutcome::Duration(duration)
}
fn run_job_operation(
program: &Path,
args: Vec<OsString>,
expected_duration: Option<f64>,
cancellation: Arc<AtomicBool>,
tx: Sender<ConversionEvent>,
) {
if cancellation.load(Ordering::Acquire) {
let _ = tx.send(ConversionEvent::Finished(ConversionOutcome::Cancelled));
return;
}
let mut command = Command::new(program);
command
.args(&args)
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::piped());
configure_process_group(&mut command);
let mut child = match command.spawn() {
Ok(child) => child,
Err(error) => {
let _ = tx.send(ConversionEvent::Finished(ConversionOutcome::Failed(
format!("failed to start ffmpeg: {error}"),
)));
return;
}
};
if cancellation.load(Ordering::Acquire) {
terminate_child(&mut child);
let _ = tx.send(ConversionEvent::Finished(ConversionOutcome::Cancelled));
return;
}
let _ = tx.send(ConversionEvent::Started(display_os_command(program, &args)));
let stdout = child.stdout.take();
let stderr = child.stderr.take();
let stdout_thread =
stdout.map(|stream| forward_conversion_stream(stream, tx.clone(), expected_duration));
let stderr_thread =
stderr.map(|stream| forward_conversion_stream(stream, tx.clone(), expected_duration));
let outcome = loop {
match child.try_wait() {
Ok(Some(status)) if status.success() => break ConversionOutcome::Succeeded,
Ok(Some(status)) => {
break ConversionOutcome::Failed(format!("ffmpeg exited with status {status}"));
}
Ok(None) if cancellation.load(Ordering::Acquire) => {
terminate_child(&mut child);
break ConversionOutcome::Cancelled;
}
Ok(None) => thread::sleep(PROCESS_POLL_INTERVAL),
Err(error) => {
terminate_child(&mut child);
break ConversionOutcome::Failed(format!("failed to wait for ffmpeg: {error}"));
}
}
};
if let Some(thread) = stdout_thread {
let _ = thread.join();
}
if let Some(thread) = stderr_thread {
let _ = thread.join();
}
let _ = tx.send(ConversionEvent::Finished(outcome));
}
fn display_os_command(program: &Path, args: &[OsString]) -> String {
let command = CommandSpec {
program: program.display().to_string(),
args: args
.iter()
.map(|arg| arg.to_string_lossy().into_owned())
.collect(),
};
command.to_string()
}
fn forward_conversion_stream<R>(
stream: R,
tx: Sender<ConversionEvent>,
expected_duration: Option<f64>,
) -> thread::JoinHandle<()>
where
R: Read + Send + 'static,
{
thread::spawn(move || {
let reader = BufReader::new(stream);
for line in reader.lines() {
let line = match line {
Ok(line) => line.trim().to_owned(),
Err(error) => {
let _ = tx.send(ConversionEvent::Log(format!("stream read failed: {error}")));
break;
}
};
if line.is_empty() {
continue;
}
if let Some(current_seconds) = parse_ffmpeg_progress_seconds(&line) {
let current_seconds = expected_duration
.map(|total| current_seconds.min(total))
.unwrap_or(current_seconds);
let _ = tx.send(ConversionEvent::Progress {
current_seconds,
total_seconds: expected_duration,
});
continue;
}
if is_ffmpeg_progress_metadata(&line) {
continue;
}
let _ = tx.send(ConversionEvent::Log(line));
}
})
}
fn run_operation(command: CommandSpec, expected_duration: Option<f64>, tx: Sender<OperationEvent>) {
let _ = tx.send(OperationEvent::Started(command.to_string()));
let mut child = match Command::new(&command.program)
.args(&command.args)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
{
Ok(child) => child,
Err(error) => {
let _ = tx.send(OperationEvent::Finished(Err(format!(
"failed to start ffmpeg: {error}"
))));
return;
}
};
let stdout = child.stdout.take();
let stderr = child.stderr.take();
let stdout_thread = stdout.map(|stream| forward_stream(stream, tx.clone(), expected_duration));
let stderr_thread = stderr.map(|stream| forward_stream(stream, tx.clone(), expected_duration));
let result = match child.wait() {
Ok(status) if status.success() => Ok(()),
Ok(status) => Err(format!("ffmpeg exited with status {status}")),
Err(error) => Err(format!("failed to wait for ffmpeg: {error}")),
};
if let Some(thread) = stdout_thread {
let _ = thread.join();
}
if let Some(thread) = stderr_thread {
let _ = thread.join();
}
let _ = tx.send(OperationEvent::Finished(result));
}
fn forward_stream<R>(
stream: R,
tx: Sender<OperationEvent>,
expected_duration: Option<f64>,
) -> thread::JoinHandle<()>
where
R: Read + Send + 'static,
{
thread::spawn(move || {
let reader = BufReader::new(stream);
for line in reader.lines() {
let line = match line {
Ok(line) => line.trim().to_owned(),
Err(error) => {
let _ = tx.send(OperationEvent::Log(format!("stream read failed: {error}")));
break;
}
};
if line.is_empty() {
continue;
}
if let Some(current_seconds) = parse_ffmpeg_progress_seconds(&line) {
if let Some(total_seconds) = expected_duration {
let _ = tx.send(OperationEvent::Progress {
current_seconds: current_seconds.min(total_seconds),
total_seconds,
});
}
continue;
}
if is_ffmpeg_progress_metadata(&line) {
continue;
}
let _ = tx.send(OperationEvent::Log(line));
}
})
}
pub fn parse_ffmpeg_progress_seconds(line: &str) -> Option<f64> {
let (key, value) = line.trim().split_once('=')?;
match key {
"out_time_us" | "out_time_ms" => value
.trim()
.parse::<f64>()
.ok()
.map(|microseconds| microseconds / 1_000_000.0),
"out_time" => parse_ffmpeg_timestamp(value.trim()),
_ => None,
}
}
fn parse_ffmpeg_timestamp(value: &str) -> Option<f64> {
let mut parts = value.split(':');
let hours = parts.next()?.parse::<f64>().ok()?;
let minutes = parts.next()?.parse::<f64>().ok()?;
let seconds = parts.next()?.parse::<f64>().ok()?;
if parts.next().is_some() {
return None;
}
Some(hours * 3600.0 + minutes * 60.0 + seconds)
}
fn is_ffmpeg_progress_metadata(line: &str) -> bool {
let Some((key, _)) = line.split_once('=') else {
return false;
};
matches!(
key,
"bitrate"
| "drop_frames"
| "dup_frames"
| "fps"
| "frame"
| "progress"
| "speed"
| "stream_0_0_q"
| "total_size"
)
}
pub fn duration_progress_label(current_seconds: f64, total_seconds: f64) -> String {
let percent = if total_seconds <= 0.0 {
0.0
} else {
current_seconds / total_seconds * 100.0
}
.clamp(0.0, 100.0);
format!(
"{} / {} ({percent:.1}%)",
format_duration(current_seconds),
format_duration(total_seconds)
)
}
pub fn format_duration(seconds: f64) -> String {
let seconds = seconds.max(0.0).round() as u64;
let hours = seconds / 3600;
let minutes = (seconds % 3600) / 60;
let seconds = seconds % 60;
if hours > 0 {
format!("{hours}:{minutes:02}:{seconds:02}")
} else {
format!("{minutes}:{seconds:02}")
}
}
pub fn suggest_output_path(
input_path: &Path,
mode: ConversionMode,
audio_format: AudioFormat,
video_format: VideoFormat,
) -> Option<PathBuf> {
if input_path.as_os_str().is_empty() {
return None;
}
let mut filename = input_path.file_stem()?.to_os_string();
let extension = mode.output_extension(audio_format, video_format);
filename.push("-caery.");
filename.push(extension);
let mut output = input_path.parent().map(PathBuf::from).unwrap_or_default();
output.push(filename);
Some(output)
}
pub fn supported_media_extension(path: &Path) -> bool {
supported_audio_extension(path) || supported_video_extension(path)
}
pub fn audio_format_from_path(path: &Path) -> Option<AudioFormat> {
match path_extension_lower(path).as_deref()? {
"mp3" => Some(AudioFormat::Mp3),
"m4a" => Some(AudioFormat::M4a),
"flac" => Some(AudioFormat::Flac),
"wav" => Some(AudioFormat::Wav),
"ogg" => Some(AudioFormat::Ogg),
"opus" => Some(AudioFormat::Opus),
_ => None,
}
}
pub fn video_format_from_path(path: &Path) -> Option<VideoFormat> {
match path_extension_lower(path).as_deref()? {
"mp4" => Some(VideoFormat::Mp4),
"mkv" => Some(VideoFormat::Mkv),
"webm" => Some(VideoFormat::Webm),
"mov" => Some(VideoFormat::Mov),
_ => None,
}
}
pub fn supported_audio_extension(path: &Path) -> bool {
matches!(
path_extension_lower(path).as_deref(),
Some(
"aac"
| "aif"
| "aiff"
| "alac"
| "flac"
| "m4a"
| "mp3"
| "oga"
| "ogg"
| "opus"
| "wav"
| "wma"
)
)
}
pub fn supported_video_extension(path: &Path) -> bool {
matches!(
path_extension_lower(path).as_deref(),
Some(
"3gp"
| "avi"
| "flv"
| "m2ts"
| "m4v"
| "mkv"
| "mov"
| "mp4"
| "mpeg"
| "mpg"
| "ogv"
| "ts"
| "webm"
| "wmv"
)
)
}
pub fn media_kind_label(path: &Path) -> &'static str {
if supported_video_extension(path) {
"VIDEO"
} else if supported_audio_extension(path) {
"AUDIO"
} else {
"FILE"
}
}
pub fn path_extension_lower(path: &Path) -> Option<String> {
path.extension()
.and_then(|extension| extension.to_str())
.map(|extension| extension.to_ascii_lowercase())
}
pub fn human_size(bytes: u64) -> String {
const UNITS: [&str; 5] = ["B", "KiB", "MiB", "GiB", "TiB"];
let mut value = bytes as f64;
let mut unit = 0;
while value >= 1024.0 && unit < UNITS.len() - 1 {
value /= 1024.0;
unit += 1;
}
if unit == 0 {
format!("{} {}", bytes, UNITS[unit])
} else {
format!("{value:.1} {}", UNITS[unit])
}
}
pub fn shell_quote(value: &str) -> String {
if value.is_empty() {
return "''".to_owned();
}
format!("'{}'", value.replace('\'', "'\"'\"'"))
}
#[cfg(test)]
mod tests {
use super::*;
use std::fs;
use std::io::Write;
use tempfile::{tempdir, NamedTempFile};
fn request(input: &Path, output: &Path, mode: ConversionMode) -> ConversionRequest {
ConversionRequest {
mode,
input_path: input.to_path_buf(),
output_path: output.to_path_buf(),
audio_format: AudioFormat::Mp3,
video_format: VideoFormat::Mp4,
quality: QualityPreset::Balanced,
overwrite: true,
}
}
#[cfg(unix)]
fn executable_script(directory: &Path, name: &str, body: &str) -> PathBuf {
use std::os::unix::fs::PermissionsExt;
let path = directory.join(name);
fs::write(&path, format!("#!/bin/sh\n{body}\n")).expect("write script");
let mut permissions = fs::metadata(&path).expect("script metadata").permissions();
permissions.set_mode(0o755);
fs::set_permissions(&path, permissions).expect("make script executable");
path
}
fn terminal_outcome(job: &ConversionJob) -> ConversionOutcome {
loop {
let event = job
.recv_timeout(Duration::from_secs(2))
.expect("conversion event");
if let ConversionEvent::Finished(outcome) = event {
return outcome;
}
}
}
#[test]
fn infers_routes_without_accessing_the_filesystem() {
assert_eq!(
infer_route(Path::new("song.WAV"), Path::new("song.flac")),
Ok(ConversionRoute::TranscodeAudio(AudioFormat::Flac))
);
assert_eq!(
infer_route(Path::new("clip.mp4"), Path::new("clip.opus")),
Ok(ConversionRoute::ExtractAudio(AudioFormat::Opus))
);
assert_eq!(
infer_route(Path::new("clip.mkv"), Path::new("clip.WEBM")),
Ok(ConversionRoute::TranscodeVideo(VideoFormat::Webm))
);
}
#[test]
fn infers_every_selectable_output_format() {
for format in AudioFormat::ALL {
let output = PathBuf::from(format!("output.{}", format.extension()));
assert_eq!(
infer_route(Path::new("song.wav"), &output),
Ok(ConversionRoute::TranscodeAudio(format))
);
assert_eq!(
infer_route(Path::new("clip.mp4"), &output),
Ok(ConversionRoute::ExtractAudio(format))
);
}
for format in VideoFormat::ALL {
let output = PathBuf::from(format!("output.{}", format.extension()));
assert_eq!(
infer_route(Path::new("clip.mp4"), &output),
Ok(ConversionRoute::TranscodeVideo(format))
);
}
}
#[test]
fn route_inference_rejects_unsupported_pairs_and_extensions() {
assert_eq!(
infer_route(Path::new("song.wav"), Path::new("song.mp4")),
Err(RouteInferenceError::AudioToVideo)
);
assert!(matches!(
infer_route(Path::new("notes.txt"), Path::new("notes.mp3")),
Err(RouteInferenceError::UnsupportedInput(_))
));
assert!(matches!(
infer_route(Path::new("song.wav"), Path::new("song.aac")),
Err(RouteInferenceError::UnsupportedOutput(_))
));
}
#[test]
fn inferred_request_uses_embedding_defaults() {
let request = request_from_paths("clip.mp4", "clip.flac").expect("infer request");
assert_eq!(request.mode, ConversionMode::ExtractAudio);
assert_eq!(request.audio_format, AudioFormat::Flac);
assert_eq!(request.video_format, VideoFormat::Mp4);
assert_eq!(request.quality, QualityPreset::Balanced);
assert!(!request.overwrite);
}
#[cfg(unix)]
#[test]
fn tool_checks_report_available_missing_and_timed_out_programs() {
let directory = tempdir().expect("temp dir");
let available = executable_script(directory.path(), "available", "exit 0");
let slow = executable_script(directory.path(), "slow", "sleep 5");
let missing = directory.path().join("missing");
let available_tools = Converter::with_programs(&available, &available).check_tools();
assert!(available_tools.conversion_ready());
assert!(available_tools.duration_progress_available());
let unavailable_tools = Converter::with_programs(&missing, &slow)
.with_tool_check_timeout(Duration::from_millis(100))
.check_tools();
assert!(matches!(
unavailable_tools.ffmpeg(),
ToolStatus::Unavailable(_)
));
assert_eq!(unavailable_tools.ffprobe(), &ToolStatus::TimedOut);
}
#[cfg(unix)]
#[test]
fn spawned_conversion_returns_immediately_and_cancels_during_probe() {
let directory = tempdir().expect("temp dir");
let input = directory.path().join("clip.mp4");
let output = directory.path().join("clip.mp3");
fs::write(&input, "fake video").expect("write input");
let ffmpeg = executable_script(directory.path(), "ffmpeg", "exit 0");
let ffprobe = executable_script(directory.path(), "ffprobe", "sleep 5");
let request = request(&input, &output, ConversionMode::ExtractAudio);
let converter =
Converter::with_programs(ffmpeg, ffprobe).with_probe_timeout(Duration::from_secs(5));
let started = Instant::now();
let job = converter.spawn(request);
assert!(started.elapsed() < Duration::from_millis(250));
assert_eq!(job.input_path(), input);
assert_eq!(job.output_path(), output);
thread::sleep(Duration::from_millis(80));
job.cancel();
assert_eq!(terminal_outcome(&job), ConversionOutcome::Cancelled);
}
#[cfg(unix)]
#[test]
fn spawned_conversion_cancels_running_ffmpeg() {
let directory = tempdir().expect("temp dir");
let input = directory.path().join("clip.mp4");
let output = directory.path().join("clip.mp3");
fs::write(&input, "fake video").expect("write input");
let ffprobe = executable_script(directory.path(), "ffprobe", "printf '10.0\\n'");
let ffmpeg = executable_script(
directory.path(),
"ffmpeg",
"printf 'out_time_us=1000000\\n'\nsleep 5",
);
let request = request(&input, &output, ConversionMode::ExtractAudio);
let job = Converter::with_programs(ffmpeg, ffprobe).spawn(request);
loop {
if matches!(
job.recv_timeout(Duration::from_secs(2))
.expect("conversion event"),
ConversionEvent::Started(_)
) {
break;
}
}
job.cancel();
assert_eq!(terminal_outcome(&job), ConversionOutcome::Cancelled);
}
#[cfg(unix)]
#[test]
fn missing_ffprobe_still_reports_timestamp_progress() {
let directory = tempdir().expect("temp dir");
let input = directory.path().join("clip.mp4");
let output = directory.path().join("clip.mp3");
fs::write(&input, "fake video").expect("write input");
let ffmpeg = executable_script(
directory.path(),
"ffmpeg",
"printf 'out_time_us=1250000\\n'",
);
let request = request(&input, &output, ConversionMode::ExtractAudio);
let job = Converter::with_programs(ffmpeg, directory.path().join("missing")).spawn(request);
let mut progress = None;
loop {
match job
.recv_timeout(Duration::from_secs(2))
.expect("conversion event")
{
ConversionEvent::Progress {
current_seconds,
total_seconds,
} => progress = Some((current_seconds, total_seconds)),
ConversionEvent::Finished(outcome) => {
assert_eq!(outcome, ConversionOutcome::Succeeded);
break;
}
_ => {}
}
}
assert_eq!(progress, Some((1.25, None)));
}
#[cfg(unix)]
#[test]
fn embedded_runner_preserves_non_utf8_paths() {
use std::os::unix::ffi::OsStringExt;
let directory = tempdir().expect("temp dir");
let input = directory
.path()
.join(OsString::from_vec(b"clip-\xff.mp4".to_vec()));
let output = directory
.path()
.join(OsString::from_vec(b"clip-\xfe.mp3".to_vec()));
fs::write(&input, "fake video").expect("write input");
let ffmpeg = executable_script(
directory.path(),
"ffmpeg",
"for arg in \"$@\"; do last=$arg; done\nprintf converted > \"$last\"",
);
let request = request(&input, &output, ConversionMode::ExtractAudio);
let job = Converter::with_programs(ffmpeg, directory.path().join("missing")).spawn(request);
assert_eq!(terminal_outcome(&job), ConversionOutcome::Succeeded);
assert_eq!(
fs::read_to_string(output).expect("read output"),
"converted"
);
}
#[test]
fn quotes_shell_values_safely() {
assert_eq!(shell_quote(""), "''");
assert_eq!(shell_quote("clip.mp4"), "'clip.mp4'");
assert_eq!(
shell_quote("artist's clip.mp4"),
"'artist'\"'\"'s clip.mp4'"
);
}
#[test]
fn suggests_caery_output_name() {
let path = Path::new("/tmp/source.video.mp4");
let output = suggest_output_path(
path,
ConversionMode::ExtractAudio,
AudioFormat::Flac,
VideoFormat::Mp4,
)
.expect("suggest output");
assert_eq!(output, PathBuf::from("/tmp/source.video-caery.flac"));
}
#[test]
fn builds_audio_extraction_command() {
let mut input = NamedTempFile::with_suffix(".mp4").expect("temp video");
writeln!(input, "fake video").expect("write temp video");
let output = input.path().with_file_name("clip-caery.mp3");
let request = request(input.path(), &output, ConversionMode::ExtractAudio);
let command = build_command(&request).expect("build command");
assert_eq!(command.program, "ffmpeg");
assert!(command.args.contains(&"-vn".to_owned()));
assert!(command.args.contains(&"libmp3lame".to_owned()));
assert!(command.args.contains(&output.display().to_string()));
}
#[test]
fn builds_webm_transcode_command() {
let mut input = NamedTempFile::with_suffix(".mkv").expect("temp video");
writeln!(input, "fake video").expect("write temp video");
let output = input.path().with_file_name("clip-caery.webm");
let mut request = request(input.path(), &output, ConversionMode::TranscodeVideo);
request.video_format = VideoFormat::Webm;
let command = build_command(&request).expect("build command");
assert!(command.args.contains(&"libvpx-vp9".to_owned()));
assert!(command.args.contains(&"libopus".to_owned()));
assert!(command.args.contains(&"0:a?".to_owned()));
}
#[test]
fn high_level_conversion_rejects_invalid_requests_before_spawning() {
let request = ConversionRequest {
mode: ConversionMode::ExtractAudio,
input_path: PathBuf::new(),
output_path: PathBuf::from("output.mp3"),
audio_format: AudioFormat::Mp3,
video_format: VideoFormat::Mp4,
quality: QualityPreset::Balanced,
overwrite: false,
};
let error = start_conversion(&request).expect_err("reject invalid request");
assert!(matches!(error, ConversionError::MissingInput));
}
#[test]
fn rejects_existing_output_without_overwrite() {
let mut input = NamedTempFile::with_suffix(".wav").expect("temp audio");
writeln!(input, "fake audio").expect("write temp audio");
let output = NamedTempFile::with_suffix(".mp3").expect("temp output");
let mut request = request(input.path(), output.path(), ConversionMode::TranscodeAudio);
request.overwrite = false;
let error = validate_request(&request).expect_err("reject existing output");
assert!(matches!(error, ConversionError::OutputExists(_)));
}
#[test]
fn parses_ffmpeg_progress_lines() {
assert_eq!(
parse_ffmpeg_progress_seconds("out_time_us=1250000"),
Some(1.25)
);
assert_eq!(
parse_ffmpeg_progress_seconds("out_time_ms=2500000"),
Some(2.5)
);
assert_eq!(
parse_ffmpeg_progress_seconds("out_time=00:01:02.500000"),
Some(62.5)
);
assert_eq!(parse_ffmpeg_progress_seconds("frame=20"), None);
}
#[test]
fn labels_duration_progress() {
assert_eq!(duration_progress_label(30.0, 120.0), "0:30 / 2:00 (25.0%)");
}
#[test]
fn formats_human_sizes() {
assert_eq!(human_size(0), "0 B");
assert_eq!(human_size(1024), "1.0 KiB");
assert_eq!(human_size(4 * 1024 * 1024 * 1024), "4.0 GiB");
}
}