Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 16 additions & 4 deletions kernel/relayflowd/tests/crash_resume/agent.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,8 +17,8 @@ use super::{
start_run,
},
support::{
completed_step_count, journal_entries, kill_group, kill_process_group, read_pid,
resume_cli, wait_until,
completed_step_count, describe_stalled_resume, journal_entries, kill_group,
kill_process_group, read_pid, resume_cli, wait_until,
},
};

Expand Down Expand Up @@ -127,8 +127,20 @@ fn rung_c_sigkill_boundaries_resume_only_unfinished_steps_via_real_cli() {
let _server = fixture.server();
let mut worker = attached_worker(&fixture, "boundary-agent-stub");
let run_id = super::support::only_run_id(&fixture.data_dir);
let resume = spawn_resume(&fixture, &run_id);
let dispatch = worker.event("step.dispatch").unwrap();
let mut resume = spawn_resume(&fixture, &run_id);
// Do not `.unwrap()` this. A missing dispatch is #174, and the whole
// difficulty there has been that the failure carries no daemon-side
// state -- so capture it here rather than losing it to the unwind.
// This fires on ANY protocol error, not only a timeout; the dump is
// useful either way, since it reports whether the child is stalled or
// already gone.
let dispatch = match worker.event("step.dispatch") {
Ok(dispatch) => dispatch,
Err(error) => panic!(
"{label}: no step.dispatch after resume: {error}{}",
describe_stalled_resume(&fixture.data_dir, &mut resume)
),
};
assert!(!record_effect(&fixture, &mut worker, &dispatch).unwrap());
complete(&mut worker, &dispatch).unwrap();
let output = resume.wait_with_output().unwrap();
Expand Down
21 changes: 18 additions & 3 deletions kernel/relayflowd/tests/crash_resume/llm.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,10 @@ use super::{
llm_support::{
LlmFixture, ProtocolClient, ServerGuard, attached_worker, complete, spawn_resume, start_run,
},
support::{journal_entries, kill_group, only_run_id, read_pid, resume_cli, wait_until},
support::{
describe_stalled_resume, journal_entries, kill_group, only_run_id, read_pid,
resume_cli, wait_until,
},
};

#[test]
Expand Down Expand Up @@ -106,8 +109,20 @@ fn sigkill_sweep_covers_before_and_between_the_rung_b_steps() {
let _server = ServerGuard::start(&fixture);
let mut worker = attached_worker(&fixture, "boundary-stub");
let run_id = only_run_id(&fixture.data_dir);
let resume = spawn_resume(&fixture, &run_id);
let dispatch = worker.event("step.dispatch").unwrap();
let mut resume = spawn_resume(&fixture, &run_id);
// Do not `.unwrap()` this. A missing dispatch is #174, and the whole
// difficulty there has been that the failure carries no daemon-side
// state -- so capture it here rather than losing it to the unwind.
// This fires on ANY protocol error, not only a timeout; the dump is
// useful either way, since it reports whether the child is stalled or
// already gone.
let dispatch = match worker.event("step.dispatch") {
Ok(dispatch) => dispatch,
Err(error) => panic!(
"{label}: no step.dispatch after resume: {error}{}",
describe_stalled_resume(&fixture.data_dir, &mut resume)
),
};
complete(&mut worker, &dispatch, json!({"answer": 4})).unwrap();
let output = resume.wait_with_output().unwrap();
assert!(
Expand Down
55 changes: 55 additions & 0 deletions kernel/relayflowd/tests/crash_resume/support.rs
Original file line number Diff line number Diff line change
Expand Up @@ -156,6 +156,61 @@ pub fn journal_entries(data_dir: &Path) -> Option<Vec<JournalEntry>> {
journal.scan_segment(segment).ok()
}

/// Everything known about a stalled resume, as one panic message.
///
/// When a dispatch never arrives (#174) the useful state is all on the daemon
/// side and none of it is captured today: `wait_with_output` is never reached,
/// so the child's output is discarded when the test unwinds, and the journal is
/// never read. Four occurrences produced four test names and nothing else.
///
/// The child may be STALLED or may have ALREADY EXITED -- #174 turned out to be
/// the second, a resume that died instantly with `run_not_found` while the test
/// waited on it. Both look identical from here, which is the point: this
/// attempts termination and reaps it either way (both results are discarded
/// because "already gone" is a normal outcome, not an error), then reports its
/// output alongside the run's journal.
///
/// The journal listing covers the CURRENT segment, which is what
/// `journal_entries` scans. These crash tests do not compact, so that is every
/// entry they produce; it would not be under compaction.
pub fn describe_stalled_resume(data_dir: &Path, resume: &mut Child) -> String {
let mut report = String::new();

// Kill before read. The child is not going to finish on its own.
let _ = resume.kill();
let _ = resume.wait();

let mut stdout = String::new();
let mut stderr = String::new();
if let Some(mut handle) = resume.stdout.take() {
let _ = std::io::Read::read_to_string(&mut handle, &mut stdout);
}
if let Some(mut handle) = resume.stderr.take() {
let _ = std::io::Read::read_to_string(&mut handle, &mut stderr);
}
report.push_str(&format!(
"\n--- resume child ---\nstdout ({} bytes):\n{stdout}\nstderr ({} bytes):\n{stderr}\n",
stdout.len(),
stderr.len()
));

// The journal says how far the run actually got, which is the question a
// missing dispatch raises: did the daemon resume and stall, or never resume?
match journal_entries(data_dir) {
None => report.push_str("--- journal --- absent (no run directory)\n"),
Some(entries) => {
report.push_str(&format!("--- journal ({} entries) ---\n", entries.len()));
for entry in &entries {
report.push_str(&format!(
" seq={} type={:?} step={:?}\n",
entry.seq, entry.entry_type, entry.step_id
));
}
}
}
report
}

pub fn completed_step_count(entries: &[JournalEntry]) -> usize {
entries
.iter()
Expand Down
Loading