Skip to main content

deepshrink_ffmpeg/
run.rs

1//! Run a single ffmpeg pass, streaming progress and surfacing failures — and
2//! stop it on request ([`CancelToken`]).
3
4use std::ffi::OsStr;
5use std::io::{BufRead, BufReader};
6use std::path::Path;
7use std::process::{Command, Stdio};
8use std::sync::atomic::{AtomicBool, Ordering};
9use std::sync::mpsc::{self, RecvTimeoutError};
10use std::sync::Arc;
11use std::time::Duration;
12
13use crate::progress::{self, Progress};
14use crate::FfmpegError;
15
16/// How often a running pass looks at its cancel token.
17const CANCEL_POLL: Duration = Duration::from_millis(100);
18
19/// A shared "stop" flag: [`CancelToken::cancel`] from any thread stops every
20/// pass run with (a clone of) this token — the ffmpeg process is killed within
21/// ~0.1 s and the pass returns [`FfmpegError::Cancelled`].
22#[derive(Debug, Clone, Default)]
23pub struct CancelToken(Arc<AtomicBool>);
24
25impl CancelToken {
26    pub fn new() -> Self {
27        Self::default()
28    }
29
30    pub fn cancel(&self) {
31        self.0.store(true, Ordering::SeqCst);
32    }
33
34    pub fn is_cancelled(&self) -> bool {
35        self.0.load(Ordering::SeqCst)
36    }
37}
38
39/// Run `ffmpeg` with `args`, invoking `on_progress` with a 0.0..=1.0 fraction as
40/// the pass proceeds. `total_secs` is the source duration (for the fraction).
41///
42/// The caller is expected to have included `-progress pipe:1 -nostats` in `args`
43/// so progress is emitted on stdout. Not cancellable — see
44/// [`run_pass_cancellable`].
45pub fn run_pass<S: AsRef<OsStr>>(
46    ffmpeg: &Path,
47    args: &[S],
48    total_secs: f64,
49    on_progress: &mut dyn FnMut(f64),
50) -> Result<(), FfmpegError> {
51    run_pass_cancellable(ffmpeg, args, total_secs, on_progress, &CancelToken::new())
52}
53
54/// [`run_pass`] that stops when `cancel` is set: before ffmpeg starts, or while
55/// it runs (the process is killed and reaped — no orphan keeps encoding). The
56/// partial output is the caller's to remove.
57pub fn run_pass_cancellable<S: AsRef<OsStr>>(
58    ffmpeg: &Path,
59    args: &[S],
60    total_secs: f64,
61    on_progress: &mut dyn FnMut(f64),
62    cancel: &CancelToken,
63) -> Result<(), FfmpegError> {
64    if cancel.is_cancelled() {
65        return Err(FfmpegError::Cancelled);
66    }
67    let mut child = Command::new(ffmpeg)
68        .args(args)
69        .stdin(Stdio::null())
70        .stdout(Stdio::piped())
71        .stderr(Stdio::piped())
72        .spawn()
73        .map_err(|source| FfmpegError::Spawn {
74            tool: "ffmpeg",
75            source,
76        })?;
77
78    // Drain stderr on a separate thread so a chatty encoder can't deadlock us
79    // while we read stdout for progress.
80    let stderr = child.stderr.take();
81    let stderr_handle = std::thread::spawn(move || {
82        let mut buf = String::new();
83        if let Some(stderr) = stderr {
84            use std::io::Read;
85            let _ = BufReader::new(stderr).read_to_string(&mut buf);
86        }
87        buf
88    });
89
90    // Progress is read on its own thread too, so this one can wake up every
91    // `CANCEL_POLL` to check the token even while ffmpeg prints nothing.
92    let (tx, rx) = mpsc::channel::<f64>();
93    let stdout = child.stdout.take();
94    let stdout_handle = std::thread::spawn(move || {
95        if let Some(stdout) = stdout {
96            for line in BufReader::new(stdout).lines().map_while(Result::ok) {
97                let f = match progress::parse_line(&line) {
98                    Some(Progress::OutTimeUs(us)) => progress::fraction(us, total_secs),
99                    Some(Progress::End) => 1.0,
100                    _ => continue,
101                };
102                if tx.send(f).is_err() {
103                    break;
104                }
105            }
106        }
107    });
108
109    loop {
110        match rx.recv_timeout(CANCEL_POLL) {
111            Ok(f) => on_progress(f),
112            Err(RecvTimeoutError::Timeout) => {}
113            Err(RecvTimeoutError::Disconnected) => break, // stdout closed: ffmpeg is done
114        }
115        if cancel.is_cancelled() {
116            let _ = child.kill();
117            let _ = child.wait();
118            let _ = stdout_handle.join();
119            let _ = stderr_handle.join();
120            return Err(FfmpegError::Cancelled);
121        }
122    }
123    let _ = stdout_handle.join();
124
125    let status = child.wait().map_err(|source| FfmpegError::Spawn {
126        tool: "ffmpeg",
127        source,
128    })?;
129    let stderr = stderr_handle.join().unwrap_or_default();
130
131    if !status.success() {
132        return Err(FfmpegError::CommandFailed {
133            tool: "ffmpeg",
134            status: status.to_string(),
135            stderr: stderr.trim().to_string(),
136        });
137    }
138    Ok(())
139}
140
141#[cfg(test)]
142mod tests {
143    use super::*;
144    use std::time::Instant;
145
146    fn ffmpeg() -> Option<std::path::PathBuf> {
147        crate::locate().ok().map(|t| t.ffmpeg)
148    }
149
150    /// A long synthetic encode to stop half-way (~20 s of work if left alone).
151    fn long_args(out: &Path) -> Vec<std::ffi::OsString> {
152        let mut a: Vec<std::ffi::OsString> = [
153            "-hide_banner",
154            "-y",
155            "-loglevel",
156            "error",
157            "-progress",
158            "pipe:1",
159            "-nostats",
160            "-f",
161            "lavfi",
162            "-i",
163            "testsrc2=size=1920x1080:rate=30:duration=600",
164            "-c:v",
165            "libx264",
166            "-preset",
167            "veryslow",
168        ]
169        .iter()
170        .map(Into::into)
171        .collect();
172        a.push(out.as_os_str().to_owned());
173        a
174    }
175
176    #[test]
177    fn a_cancelled_pass_kills_ffmpeg_promptly() {
178        let Some(ffmpeg) = ffmpeg() else {
179            eprintln!("skipping: ffmpeg not found");
180            return;
181        };
182        let out =
183            std::env::temp_dir().join(format!("deepshrink-cancel-{}.mp4", std::process::id()));
184        let token = CancelToken::new();
185        let stopper = token.clone();
186        std::thread::spawn(move || {
187            std::thread::sleep(Duration::from_millis(2500));
188            stopper.cancel();
189        });
190        let started = Instant::now();
191        let r = run_pass_cancellable(&ffmpeg, &long_args(&out), 600.0, &mut |_| {}, &token);
192        assert!(matches!(r, Err(FfmpegError::Cancelled)), "{r:?}");
193        assert!(
194            started.elapsed() < Duration::from_secs(5),
195            "{:?}",
196            started.elapsed()
197        );
198        let _ = std::fs::remove_file(&out);
199    }
200
201    #[test]
202    fn an_already_cancelled_token_never_starts_ffmpeg() {
203        let token = CancelToken::new();
204        token.cancel();
205        let r = run_pass_cancellable(
206            Path::new("/nonexistent/ffmpeg"),
207            &["-version"],
208            1.0,
209            &mut |_| {},
210            &token,
211        );
212        assert!(matches!(r, Err(FfmpegError::Cancelled)));
213    }
214}