Report a merge worker that dies rather than leaving the page on Stop
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.
This commit is contained in:
+88
-8
@@ -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::<String>().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<MergeEvent>) -> Vec<MergeEvent> {
|
||||
/// 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<MergeEvent>) -> (Vec<MergeEvent>, 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());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -487,11 +487,21 @@ fn chip_projection(i: i32) -> Option<dr_pano::Projection> {
|
||||
|
||||
/// Take everything the job has said and reflect it on the page.
|
||||
fn drain(window: &AppWindow, ctl: &Rc<MergeController>, on_done: &Rc<impl Fn(&AppWindow)>) {
|
||||
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;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user