1use 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
16const CANCEL_POLL: Duration = Duration::from_millis(100);
18
19#[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
39pub 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
54pub 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 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 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, }
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 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}