Robustness: hang watchdog kills stalled yt-dlp/ffmpeg jobs (2.5)
A yt-dlp process can stall with the OS process alive but no progress — a wedged TLS handshake, an unresponsive server, a dead socket mid-fragment. Before, that job hung forever, holding a concurrency slot until the user noticed and cancelled. Now each Job tracks last_activity (bumped on any stdout/stderr line or progress message). poll() runs check_watchdog(): a Running job silent for HANG_TIMEOUT (5 min) gets SIGKILLed by pid via libc::kill. The threshold sits well above any legitimate quiet gap — yt-dlp prints a line per download chunk, polls every 30s under --wait-for-video, and our longest sleep is ~30s — so 5 min of total silence is a real stall. To capture the pid, spawn_job + enqueue_transcode now spawn the child *before* the reader thread (the thread still owns it for the blocking wait()). A watchdog kill sets watchdog_killed, and drain() forces the resulting failure to NetworkError (retryable) instead of the Other that classify() would return for an output-less kill — so the existing auto-retry path re-queues it (a hang is transient). Non-Unix builds keep the flagging but can't force-kill (no portable pid-kill; Windows isn't first-class). JobState gains Debug for tests. 2 tests: watchdog kill → retryable class; normal failures still classify from the log. 98 pass. Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
This commit is contained in:
parent
111d09b130
commit
ffa864c470
1 changed files with 152 additions and 22 deletions
|
|
@ -100,7 +100,7 @@ pub fn recheck_url(ch: &crate::library::Channel) -> String {
|
||||||
|
|
||||||
|
|
||||||
/// Lifecycle state of a download job.
|
/// Lifecycle state of a download job.
|
||||||
#[derive(Clone, Copy, PartialEq, Eq)]
|
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
|
||||||
pub enum JobState {
|
pub enum JobState {
|
||||||
Running,
|
Running,
|
||||||
Done,
|
Done,
|
||||||
|
|
@ -151,6 +151,17 @@ pub struct Job {
|
||||||
/// Set once we've enqueued the convert pass for this finished job, so
|
/// Set once we've enqueued the convert pass for this finished job, so
|
||||||
/// the per-poll scan doesn't enqueue it twice.
|
/// the per-poll scan doesn't enqueue it twice.
|
||||||
convert_handled: bool,
|
convert_handled: bool,
|
||||||
|
/// OS pid of the child process, captured at spawn. Lets the hang
|
||||||
|
/// watchdog SIGKILL a stalled process (Unix). `None` for synthetic
|
||||||
|
/// jobs (preflight failures) that have no process.
|
||||||
|
child_pid: Option<u32>,
|
||||||
|
/// Last time this job produced any output/progress. The watchdog kills
|
||||||
|
/// a Running job whose `last_activity` is older than [`HANG_TIMEOUT`].
|
||||||
|
last_activity: std::time::Instant,
|
||||||
|
/// True when the watchdog killed this job (vs. a genuine yt-dlp exit).
|
||||||
|
/// On the resulting failure we force a retryable class so auto-retry
|
||||||
|
/// re-queues it — a hang is transient, worth one more try.
|
||||||
|
watchdog_killed: bool,
|
||||||
rx: Receiver<Msg>,
|
rx: Receiver<Msg>,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -214,17 +225,30 @@ const MAX_AUTO_RETRIES: u8 = 2;
|
||||||
/// rate-limit window pass; short enough that the user isn't left waiting.
|
/// rate-limit window pass; short enough that the user isn't left waiting.
|
||||||
const AUTO_RETRY_COOLDOWN: std::time::Duration = std::time::Duration::from_secs(90);
|
const AUTO_RETRY_COOLDOWN: std::time::Duration = std::time::Duration::from_secs(90);
|
||||||
|
|
||||||
|
/// A Running job that produces no output/progress for this long is
|
||||||
|
/// considered hung and gets SIGKILLed by the watchdog. Chosen well above
|
||||||
|
/// any legitimate quiet gap: yt-dlp prints a progress line per chunk
|
||||||
|
/// while downloading, polls every 30 s under `--wait-for-video`, and our
|
||||||
|
/// longest configured sleep (retry-sleep cap / adaptive throttle) is
|
||||||
|
/// ~30 s — so 5 minutes of total silence means a real stall (a stuck TLS
|
||||||
|
/// handshake or unresponsive server), not slow-but-alive progress.
|
||||||
|
const HANG_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(300);
|
||||||
|
|
||||||
impl Job {
|
impl Job {
|
||||||
fn drain(&mut self) {
|
fn drain(&mut self) {
|
||||||
while let Ok(msg) = self.rx.try_recv() {
|
while let Ok(msg) = self.rx.try_recv() {
|
||||||
match msg {
|
match msg {
|
||||||
Msg::Line(line) => {
|
Msg::Line(line) => {
|
||||||
|
self.last_activity = std::time::Instant::now();
|
||||||
self.log.push_back(line);
|
self.log.push_back(line);
|
||||||
while self.log.len() > JOB_LOG_CAP {
|
while self.log.len() > JOB_LOG_CAP {
|
||||||
self.log.pop_front();
|
self.log.pop_front();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Msg::Progress(p) => self.progress = p,
|
Msg::Progress(p) => {
|
||||||
|
self.last_activity = std::time::Instant::now();
|
||||||
|
self.progress = p;
|
||||||
|
}
|
||||||
Msg::Finished(ok) => {
|
Msg::Finished(ok) => {
|
||||||
self.state = if ok { JobState::Done } else { JobState::Failed };
|
self.state = if ok { JobState::Done } else { JobState::Failed };
|
||||||
// Classify only on the failure transition so the
|
// Classify only on the failure transition so the
|
||||||
|
|
@ -232,9 +256,15 @@ impl Job {
|
||||||
// long-finished job. The log is already in `self.log`
|
// long-finished job. The log is already in `self.log`
|
||||||
// by this point since we drained Line messages above.
|
// by this point since we drained Line messages above.
|
||||||
if !ok && self.failure_class.is_none() {
|
if !ok && self.failure_class.is_none() {
|
||||||
self.failure_class = Some(crate::error_class::classify(
|
self.failure_class = Some(if self.watchdog_killed {
|
||||||
self.log.iter().map(|s| s.as_str()),
|
// A watchdog kill leaves no telltale error line, so
|
||||||
));
|
// classify() would return Other (non-retryable).
|
||||||
|
// Force NetworkError: a hang is transient and the
|
||||||
|
// auto-retry path should give it another go.
|
||||||
|
crate::error_class::ErrorClass::NetworkError
|
||||||
|
} else {
|
||||||
|
crate::error_class::classify(self.log.iter().map(|s| s.as_str()))
|
||||||
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -880,16 +910,23 @@ impl Downloader {
|
||||||
.display()
|
.display()
|
||||||
.to_string();
|
.to_string();
|
||||||
|
|
||||||
thread::spawn(move || {
|
// Spawn the child *before* the reader thread so we can capture its
|
||||||
let mut child = match cmd.spawn() {
|
// pid for the hang watchdog (the thread then owns the Child and
|
||||||
Ok(child) => child,
|
// does the blocking wait()). On spawn failure, push a synthetic
|
||||||
Err(err) => {
|
// failed job and bail.
|
||||||
let _ = tx.send(Msg::Line(format!("could not launch yt-dlp: {err}")));
|
let mut child = match cmd.spawn() {
|
||||||
let _ = tx.send(Msg::Finished(false));
|
Ok(child) => child,
|
||||||
return;
|
Err(err) => {
|
||||||
}
|
self.push_synthetic_failure(
|
||||||
};
|
url, label, crate::error_class::ErrorClass::Other,
|
||||||
|
format!("could not launch yt-dlp: {err}"),
|
||||||
|
);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
let child_pid = Some(child.id());
|
||||||
|
|
||||||
|
thread::spawn(move || {
|
||||||
if let Some(stderr) = child.stderr.take() {
|
if let Some(stderr) = child.stderr.take() {
|
||||||
let tx = tx.clone();
|
let tx = tx.clone();
|
||||||
let cookies_abs = cookies_abs.clone();
|
let cookies_abs = cookies_abs.clone();
|
||||||
|
|
@ -949,6 +986,9 @@ impl Downloader {
|
||||||
retry_handled: false,
|
retry_handled: false,
|
||||||
convert_on_finish,
|
convert_on_finish,
|
||||||
convert_handled: false,
|
convert_handled: false,
|
||||||
|
child_pid,
|
||||||
|
last_activity: std::time::Instant::now(),
|
||||||
|
watchdog_killed: false,
|
||||||
rx,
|
rx,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
@ -986,6 +1026,9 @@ impl Downloader {
|
||||||
retry_handled: true, // never retry a synthetic preflight failure
|
retry_handled: true, // never retry a synthetic preflight failure
|
||||||
convert_on_finish: None,
|
convert_on_finish: None,
|
||||||
convert_handled: true,
|
convert_handled: true,
|
||||||
|
child_pid: None, // no process to watchdog
|
||||||
|
last_activity: std::time::Instant::now(),
|
||||||
|
watchdog_killed: false,
|
||||||
rx,
|
rx,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
@ -997,6 +1040,7 @@ impl Downloader {
|
||||||
for job in &mut self.jobs {
|
for job in &mut self.jobs {
|
||||||
job.drain();
|
job.drain();
|
||||||
}
|
}
|
||||||
|
self.check_watchdog();
|
||||||
self.schedule_auto_retries();
|
self.schedule_auto_retries();
|
||||||
self.fire_due_retries();
|
self.fire_due_retries();
|
||||||
self.schedule_conversions();
|
self.schedule_conversions();
|
||||||
|
|
@ -1009,6 +1053,28 @@ impl Downloader {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Hang watchdog: SIGKILL any Running job that has produced no output
|
||||||
|
/// or progress for [`HANG_TIMEOUT`]. Marks it `watchdog_killed` so the
|
||||||
|
/// resulting failure is classified retryable (see [`Job::drain`]) and
|
||||||
|
/// the auto-retry path gives it another go. The reader thread's
|
||||||
|
/// `child.wait()` returns once the process dies, sending `Finished`
|
||||||
|
/// normally — we don't reap here, just deliver the kill signal.
|
||||||
|
fn check_watchdog(&mut self) {
|
||||||
|
let now = std::time::Instant::now();
|
||||||
|
for job in &mut self.jobs {
|
||||||
|
if job.state != JobState::Running || job.watchdog_killed { continue; }
|
||||||
|
let Some(pid) = job.child_pid else { continue };
|
||||||
|
if now.duration_since(job.last_activity) < HANG_TIMEOUT { continue; }
|
||||||
|
kill_pid(pid);
|
||||||
|
job.watchdog_killed = true;
|
||||||
|
job.last_activity = now; // don't re-kill before wait() returns
|
||||||
|
job.log.push_back(format!(
|
||||||
|
"⏱ watchdog: no output for {}s — killed (will auto-retry)",
|
||||||
|
HANG_TIMEOUT.as_secs(),
|
||||||
|
));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Scan freshly-Done download jobs that have a convert spec; for each
|
/// Scan freshly-Done download jobs that have a convert spec; for each
|
||||||
/// `CONVERT_PATH:` line their yt-dlp run printed, enqueue an ffmpeg
|
/// `CONVERT_PATH:` line their yt-dlp run printed, enqueue an ffmpeg
|
||||||
/// transcode job. Marks the job handled so it only fires once.
|
/// transcode job. Marks the job handled so it only fires once.
|
||||||
|
|
@ -1200,15 +1266,20 @@ impl Downloader {
|
||||||
) {
|
) {
|
||||||
let (tx, rx) = channel();
|
let (tx, rx) = channel();
|
||||||
cmd.stdin(Stdio::null()).stdout(Stdio::piped()).stderr(Stdio::piped());
|
cmd.stdin(Stdio::null()).stdout(Stdio::piped()).stderr(Stdio::piped());
|
||||||
|
// Spawn before the thread so the watchdog gets a pid (ffmpeg can
|
||||||
|
// also hang on a bad input). On spawn failure, surface it.
|
||||||
|
let mut child = match cmd.spawn() {
|
||||||
|
Ok(c) => c,
|
||||||
|
Err(e) => {
|
||||||
|
self.push_synthetic_failure(
|
||||||
|
url, label, crate::error_class::ErrorClass::CodecMissing,
|
||||||
|
format!("could not launch ffmpeg: {e}"),
|
||||||
|
);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
let child_pid = Some(child.id());
|
||||||
thread::spawn(move || {
|
thread::spawn(move || {
|
||||||
let mut child = match cmd.spawn() {
|
|
||||||
Ok(c) => c,
|
|
||||||
Err(e) => {
|
|
||||||
let _ = tx.send(Msg::Line(format!("could not launch ffmpeg: {e}")));
|
|
||||||
let _ = tx.send(Msg::Finished(false));
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
};
|
|
||||||
if let Some(stderr) = child.stderr.take() {
|
if let Some(stderr) = child.stderr.take() {
|
||||||
let tx = tx.clone();
|
let tx = tx.clone();
|
||||||
thread::spawn(move || {
|
thread::spawn(move || {
|
||||||
|
|
@ -1260,6 +1331,9 @@ impl Downloader {
|
||||||
retry_handled: true,
|
retry_handled: true,
|
||||||
convert_on_finish: None, // a transcode doesn't itself convert
|
convert_on_finish: None, // a transcode doesn't itself convert
|
||||||
convert_handled: true,
|
convert_handled: true,
|
||||||
|
child_pid,
|
||||||
|
last_activity: std::time::Instant::now(),
|
||||||
|
watchdog_killed: false,
|
||||||
rx,
|
rx,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
@ -1355,6 +1429,20 @@ fn redact_sensitive(line: &str, cookies_abs: &str) -> String {
|
||||||
line.replace(cookies_abs, "cookies.txt")
|
line.replace(cookies_abs, "cookies.txt")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// SIGKILL a child process by pid (used by the hang watchdog). The reader
|
||||||
|
/// thread's `child.wait()` reaps it; we only deliver the signal here.
|
||||||
|
#[cfg(unix)]
|
||||||
|
fn kill_pid(pid: u32) {
|
||||||
|
// SAFETY: libc::kill with a valid pid + signal is sound; an invalid
|
||||||
|
// pid just returns ESRCH which we ignore.
|
||||||
|
unsafe { libc::kill(pid as libc::pid_t, libc::SIGKILL); }
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Non-Unix fallback: no portable pid-kill, so the watchdog can flag the
|
||||||
|
/// stall but can't force-terminate. (Windows isn't a first-class target.)
|
||||||
|
#[cfg(not(unix))]
|
||||||
|
fn kill_pid(_pid: u32) {}
|
||||||
|
|
||||||
/// Parse a yt-dlp `[download] 42.7% …` line into a `[0.0, 1.0]` fraction.
|
/// Parse a yt-dlp `[download] 42.7% …` line into a `[0.0, 1.0]` fraction.
|
||||||
fn parse_progress(line: &str) -> Option<f32> {
|
fn parse_progress(line: &str) -> Option<f32> {
|
||||||
let rest = line.trim_start().strip_prefix("[download]")?.trim_start();
|
let rest = line.trim_start().strip_prefix("[download]")?.trim_start();
|
||||||
|
|
@ -1483,6 +1571,48 @@ mod tests {
|
||||||
assert_eq!(s.audio_format, "mp3"); // empty global → default mp3
|
assert_eq!(s.audio_format, "mp3"); // empty global → default mp3
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ── Hang watchdog ────────────────────────────────────────────────────
|
||||||
|
#[test]
|
||||||
|
fn watchdog_killed_job_fails_retryable() {
|
||||||
|
// A job the watchdog killed leaves no error line; drain() must
|
||||||
|
// force a retryable class so auto-retry re-queues it (a hang is
|
||||||
|
// transient), rather than classify() returning Other (terminal).
|
||||||
|
let (tx, rx) = channel();
|
||||||
|
let mut job = Job {
|
||||||
|
url: "u".into(), label: "l".into(), state: JobState::Running,
|
||||||
|
progress: 0.0, log: VecDeque::new(), failure_class: None,
|
||||||
|
retry_spec: None, retry_count: 0, retry_handled: false,
|
||||||
|
convert_on_finish: None, convert_handled: false,
|
||||||
|
child_pid: Some(999_999),
|
||||||
|
last_activity: std::time::Instant::now(),
|
||||||
|
watchdog_killed: true,
|
||||||
|
rx,
|
||||||
|
};
|
||||||
|
tx.send(Msg::Finished(false)).unwrap();
|
||||||
|
job.drain();
|
||||||
|
assert_eq!(job.state, JobState::Failed);
|
||||||
|
assert_eq!(job.failure_class, Some(crate::error_class::ErrorClass::NetworkError));
|
||||||
|
assert!(Downloader::is_retryable(job.failure_class.unwrap()));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn normal_failure_still_classified_from_log() {
|
||||||
|
// Without watchdog_killed, drain() classifies from the log as before.
|
||||||
|
let (tx, rx) = channel();
|
||||||
|
let mut job = Job {
|
||||||
|
url: "u".into(), label: "l".into(), state: JobState::Running,
|
||||||
|
progress: 0.0, log: VecDeque::new(), failure_class: None,
|
||||||
|
retry_spec: None, retry_count: 0, retry_handled: false,
|
||||||
|
convert_on_finish: None, convert_handled: false,
|
||||||
|
child_pid: Some(1), last_activity: std::time::Instant::now(),
|
||||||
|
watchdog_killed: false, rx,
|
||||||
|
};
|
||||||
|
tx.send(Msg::Line("ERROR: Video unavailable. This video has been removed".into())).unwrap();
|
||||||
|
tx.send(Msg::Finished(false)).unwrap();
|
||||||
|
job.drain();
|
||||||
|
assert_eq!(job.failure_class, Some(crate::error_class::ErrorClass::NotFound));
|
||||||
|
}
|
||||||
|
|
||||||
// ── Subtitle resolver ────────────────────────────────────────────────
|
// ── Subtitle resolver ────────────────────────────────────────────────
|
||||||
fn dl_with_subs(subs: crate::config::SubtitlesSection) -> Downloader {
|
fn dl_with_subs(subs: crate::config::SubtitlesSection) -> Downloader {
|
||||||
let mut d = Downloader::new(
|
let mut d = Downloader::new(
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue