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; }