From f1760436326902e8f1df8646029b0355a37f78bd Mon Sep 17 00:00:00 2001 From: Duncan Tourolle Date: Sun, 20 Sep 2026 00:12:36 +0200 Subject: [PATCH] Report a merge worker that dies rather than leaving the page on Stop MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit wgpu reports a device out of memory by panicking, and a twelve-frame merge on a GPU another process is using is where that happens. The panic unwound the worker, the sender went with it, and the page sat on "Stop" with every control disabled and nothing to say why — the crash record on disk was the only sign. Catch the panic and send it as a failure, and treat a closed channel with no final event as a dead worker too. --- ui/dr-ui/src/merge.rs | 96 ++++++++++++++++++++++++++++++++++++---- ui/dr-ui/src/merge_ui.rs | 12 ++++- 2 files changed, 99 insertions(+), 9 deletions(-) diff --git a/ui/dr-ui/src/merge.rs b/ui/dr-ui/src/merge.rs index 370d5a8..a660c4c 100644 --- a/ui/dr-ui/src/merge.rs +++ b/ui/dr-ui/src/merge.rs @@ -218,6 +218,17 @@ pub enum MergeEvent { Cancelled, } +impl MergeEvent { + /// Whether this is the job's last word: after one of these the worker + /// has nothing more to say, so its channel closing is expected. + pub fn is_final(&self) -> bool { + matches!( + self, + MergeEvent::Done { .. } | MergeEvent::Failed(_) | MergeEvent::Cancelled + ) + } +} + /// The alignment, described for a panel. #[derive(Debug, Clone)] pub struct AlignmentReport { @@ -255,10 +266,29 @@ pub fn run( let send = |e: MergeEvent| { let _ = events.send(e); }; - match run_inner(&ctx, &request, &events, &decision, &cancel) { - Ok(Some(done)) => send(done), - Ok(None) => send(MergeEvent::Cancelled), - Err(e) => send(MergeEvent::Failed(e)), + // Caught rather than allowed to unwind the thread: wgpu reports a device + // that has run out of memory by panicking, and a twelve-frame merge on a + // GPU another process is using is exactly where that happens. Uncaught, + // the thread died, the sender went with it, and the page sat on "Stop" + // with every control disabled and nothing to say why — the crash record + // on disk was the only sign. The panic hook still writes that record; + // this is what puts the reason on the page. + let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + run_inner(&ctx, &request, &events, &decision, &cancel) + })); + match result { + Ok(Ok(Some(done))) => send(done), + Ok(Ok(None)) => send(MergeEvent::Cancelled), + Ok(Err(e)) => send(MergeEvent::Failed(e)), + Err(panic) => { + let detail = panic + .downcast_ref::<&str>() + .map(|s| (*s).to_string()) + .or_else(|| panic.downcast_ref::().cloned()) + .unwrap_or_else(|| "panicked with a non-string payload".to_string()); + log::error!("merge worker panicked: {detail}"); + send(MergeEvent::Failed(format!("internal error: {detail}"))); + } } } @@ -1272,12 +1302,18 @@ fn unused_name(dir: &Path, name: &str) -> PathBuf { } /// Drain everything a job has said so far. -pub fn drain(rx: &Receiver) -> Vec { +/// Everything the worker has said since the last call, and whether it has +/// hung up — a closed channel after nothing [`MergeEvent::is_final`] is a +/// worker that died mid-job. +pub fn drain(rx: &Receiver) -> (Vec, bool) { let mut out = Vec::new(); - while let Ok(e) = rx.try_recv() { - out.push(e); + loop { + match rx.try_recv() { + Ok(e) => out.push(e), + Err(std::sync::mpsc::TryRecvError::Empty) => return (out, false), + Err(std::sync::mpsc::TryRecvError::Disconnected) => return (out, true), + } } - out } /// The filler's input as the debugging example reads it: `input.ppm`, the @@ -1300,3 +1336,47 @@ fn dump_fill_input( pgm.extend(known.iter().map(|&k| if k { 255u8 } else { 0 })); std::fs::write(dir.join("known.pgm"), pgm) } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn a_worker_that_hangs_up_mid_job_is_reported_as_gone() { + // The failure the page could not see: a panic drops the sender with + // no final event, and `try_recv`'s Disconnected used to be folded + // into "nothing new". + let (tx, rx) = std::sync::mpsc::channel(); + tx.send(MergeEvent::Progress { + stage: "Reading", + done: 1, + total: 12, + }) + .unwrap(); + let (events, gone) = drain(&rx); + assert_eq!(events.len(), 1); + assert!(!gone, "the sender is still alive"); + + drop(tx); + let (events, gone) = drain(&rx); + assert!(events.is_empty()); + assert!(gone, "and now it is not"); + } + + #[test] + fn a_job_that_said_its_last_word_is_not_a_dead_worker() { + let (tx, rx) = std::sync::mpsc::channel(); + tx.send(MergeEvent::Failed("out of memory".into())).unwrap(); + drop(tx); + let (events, gone) = drain(&rx); + assert!(gone); + assert!(events.iter().any(MergeEvent::is_final), "Failed is final"); + assert!(MergeEvent::Cancelled.is_final()); + assert!(!MergeEvent::Progress { + stage: "x", + done: 0, + total: 0 + } + .is_final()); + } +} diff --git a/ui/dr-ui/src/merge_ui.rs b/ui/dr-ui/src/merge_ui.rs index 60452b3..e5b4e5e 100644 --- a/ui/dr-ui/src/merge_ui.rs +++ b/ui/dr-ui/src/merge_ui.rs @@ -487,11 +487,21 @@ fn chip_projection(i: i32) -> Option { /// Take everything the job has said and reflect it on the page. fn drain(window: &AppWindow, ctl: &Rc, on_done: &Rc) { - let events = { + let (events, gone) = { let job = ctl.job.borrow(); let Some(job) = job.as_ref() else { return }; merge::drain(&job.rx) }; + // The worker is the only sender, so a closed channel with nothing final + // said is a worker that died without saying anything — the one way + // `merge::run` cannot report, since reporting is what it was doing. + // Without this the page stayed on "Stop" for ever. + let mut events = events; + if gone && !events.iter().any(MergeEvent::is_final) { + events.push(MergeEvent::Failed( + "the merge stopped unexpectedly; the log has the reason".into(), + )); + } if events.is_empty() { return; }