From 41d35d0e46f4108ef4304c72962b6f846d5a895a Mon Sep 17 00:00:00 2001 From: kjgbot Date: Sat, 5 Sep 2026 14:05:52 +0200 Subject: [PATCH 1/3] test(kernel): bound protocol reads so a missing frame fails instead of hanging #174: the kernel step was cancelled at ~25-30 minutes on three runs across two branches with no code in common, always on a crash_resume test that finishes in under a second locally, and always with no output beyond the harness's own "has been running for over 60 seconds". `ProtocolClient::connect` set no read timeout, so `read_frame`'s `read_line` blocks forever and `event()` loops on it. These tests SIGKILL a daemon and resume it, so "the dispatch never arrives" is a reachable state rather than a hypothetical -- and an unbounded read turns that state into a silent hang that produces no evidence at all. Three runs cost 75 minutes of CI and yielded one test name between them. A 60s ceiling, against a target that runs in 38s, so it can only fire on "never" and not on "slow". A timeout is named explicitly rather than surfacing as a bare I/O error, because a test stopped there is waiting for a frame the daemon never sent and that sentence is the diagnosis. This does NOT fix the underlying nondeterminism -- it makes it report. The next occurrence names its own line in 60 seconds instead of eating the step. Verified. Mutating the ceiling to 1ms fails at agent.rs:131 -- `worker.event("step.dispatch")`, the exact blocking read -- with `timed out after 1ms waiting for a protocol frame; the daemon sent nothing (see #174)`, in 0.80s. That also confirms where CI was stuck. Restored, crash_resume is 34 passed in 37.8s across three consecutive runs, against 38.27s on main's own CI. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01FtQSAcGDta5VH9xiZFT4sR Session-Id: c228933d-4f94-4d83-9a9a-daf3c83b94f1 --- .../tests/crash_resume/llm_support.rs | 40 +++++++++++++++++-- 1 file changed, 37 insertions(+), 3 deletions(-) diff --git a/kernel/relayflowd/tests/crash_resume/llm_support.rs b/kernel/relayflowd/tests/crash_resume/llm_support.rs index 693633a16..0c240cc6e 100644 --- a/kernel/relayflowd/tests/crash_resume/llm_support.rs +++ b/kernel/relayflowd/tests/crash_resume/llm_support.rs @@ -237,10 +237,28 @@ pub struct ProtocolClient { events: Vec, } +/// Ceiling on any single protocol read in a test. +/// +/// The whole `crash_resume` target runs in about 38 seconds, so this is far +/// longer than any legitimate wait; it exists only to convert "never" into a +/// failure. See `read_frame` for why that matters. +const READ_TIMEOUT: Duration = Duration::from_secs(60); + impl ProtocolClient { pub fn connect(socket: &Path) -> Self { let stream = UnixStream::connect(socket).unwrap(); - let reader = BufReader::new(stream.try_clone().unwrap()); + let read_half = stream.try_clone().unwrap(); + // Without this a frame that never arrives blocks forever. These tests + // SIGKILL a daemon and resume it, so "the dispatch never comes" is a + // reachable state, not a hypothetical -- and an unbounded read turns it + // into a silent hang that produces NO output at all. On GitHub runners + // that consumed the entire 30-minute step three times (#174), and the + // only evidence left behind was the harness's own + // "has been running for over 60 seconds" line. + read_half + .set_read_timeout(Some(READ_TIMEOUT)) + .expect("set protocol read timeout"); + let reader = BufReader::new(read_half); Self { stream, reader, @@ -312,8 +330,24 @@ impl ProtocolClient { fn read_frame(&mut self) -> Result { let mut line = String::new(); - if self.reader.read_line(&mut line)? == 0 { - bail!("protocol connection closed") + match self.reader.read_line(&mut line) { + Ok(0) => bail!("protocol connection closed"), + Ok(_) => {} + // Name the timeout rather than letting it surface as a bare I/O + // error. A test that stops here is waiting for a frame the daemon + // never sent, and that sentence is the entire diagnosis. + Err(error) + if matches!( + error.kind(), + std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut + ) => + { + bail!( + "timed out after {READ_TIMEOUT:?} waiting for a protocol frame; \ + the daemon sent nothing (see #174)" + ) + } + Err(error) => return Err(error.into()), } serde_json::from_str(&line).context("decode protocol frame") } From 64a3635bfa4c70cfedbf70f556a11297b23a8cfd Mon Sep 17 00:00:00 2001 From: kjgbot Date: Sat, 5 Sep 2026 14:15:05 +0200 Subject: [PATCH 2/3] test(kernel): bound protocol reads and restore the ceiling on override #174: the kernel step was cancelled at ~25-30 minutes on three runs across two branches with no code in common, always on a crash_resume test that finishes in under a second locally, and always with no output beyond the harness's "has been running for over 60 seconds". `ProtocolClient::connect` set no read timeout, so `read_frame`'s `read_line` blocked forever and `event()` looped on it. These tests SIGKILL a daemon and resume it, so "the dispatch never arrives" is a reachable state, and an unbounded read turned it into a silent hang. Three runs cost 75 minutes and yielded one test name between them. This bounds the read at 60s -- against a target that runs in 38s, so it can only fire on "never", not on "slow" -- and names the timeout, since a test stopped there is waiting for a frame the daemon never sent and that sentence is the diagnosis. It does NOT fix the nondeterminism. It makes it report. Three review-lens iterations shaped the rest, each catching something real: * The ceiling was cancellable. Three callers passed None and kept reading; because the fds are dup'd and share SO_RCVTIMEO, that cleared it. None now restores the default, so no read is unbounded. * The override wrote to `self.stream` while `connect` wrote to the reader's fd. Same socket on Linux and Darwin, different descriptors -- a platform assumption with nothing naming it. Both now go through the fd `read_frame` reads. * The method shadowed `UnixStream::set_read_timeout` while inverting what None means. Renamed `override_read_timeout`. * The error reported the constant, not the ceiling in force, which would have been wrong inside a 200ms probe. It reports the effective bound. The three callers are not one pattern: concurrency.rs:31 tightens to 1s and expects SUCCESS; parallel_lifecycle.rs:138 and :208 use 200ms and assert SILENCE. Evidence, captured verbatim: $ pwd /Users/khaliqgant/AgentWorkforce/flows-claude-lead-0903-wt $ grep -rn "override_read_timeout(None)" kernel/relayflowd/tests/ kernel/relayflowd/tests/crash_resume/concurrency.rs:36: worker.override_read_timeout(None); kernel/relayflowd/tests/crash_resume/parallel_lifecycle.rs:140: worker.override_read_timeout(None); kernel/relayflowd/tests/crash_resume/parallel_lifecycle.rs:210: replacement.override_read_timeout(None); $ grep -rn "set_read_timeout" kernel/relayflowd/tests/crash_resume/llm_support.rs kernel/relayflowd/tests/crash_resume/llm_support.rs:262: .set_read_timeout(Some(READ_TIMEOUT)) kernel/relayflowd/tests/crash_resume/llm_support.rs:334: /// Deliberately NOT named `set_read_timeout`: that name belongs to kernel/relayflowd/tests/crash_resume/llm_support.rs:361: .set_read_timeout(Some(timeout)) $ (cd kernel && cargo test -p relayflowd --test crash_resume) # three times test result: ok. 34 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out; finished in 37.86s test result: ok. 34 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out; finished in 37.92s test result: ok. 34 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out; finished in 37.86s Mutation witness: setting the ceiling to 1ms fails at agent.rs:131 -- `worker.event("step.dispatch")`, the exact blocking read -- with `timed out after 1ms waiting for a protocol frame; the daemon sent nothing (see #174)`, in 0.80s. sha256 before 8d01104e, mutated 06ee52a8, restored 8d01104e. An earlier revision of this message cited a line number that a later edit had shifted, and abbreviated a path and a command that would not run as written. The block above was regenerated by running the commands. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01FtQSAcGDta5VH9xiZFT4sR Session-Id: c228933d-4f94-4d83-9a9a-daf3c83b94f1 --- .../tests/crash_resume/concurrency.rs | 4 +- .../tests/crash_resume/llm_support.rs | 42 +++++++++++++++++-- .../tests/crash_resume/parallel_lifecycle.rs | 8 ++-- 3 files changed, 45 insertions(+), 9 deletions(-) diff --git a/kernel/relayflowd/tests/crash_resume/concurrency.rs b/kernel/relayflowd/tests/crash_resume/concurrency.rs index 4986f6f63..b3e51c7e4 100644 --- a/kernel/relayflowd/tests/crash_resume/concurrency.rs +++ b/kernel/relayflowd/tests/crash_resume/concurrency.rs @@ -28,12 +28,12 @@ fn run_start_dispatches_every_independent_lane_before_any_completion() { let mut worker = attached_worker(&fixture, "parallel-stub"); let run_id = start_run(&fixture); - worker.set_read_timeout(Some(Duration::from_secs(1))); + worker.override_read_timeout(Some(Duration::from_secs(1))); let first = worker.event("step.dispatch").unwrap(); let second = worker .event("step.dispatch") .expect("both independent lanes must dispatch before either completes"); - worker.set_read_timeout(None); + worker.override_read_timeout(None); assert_eq!(first["step_id"], "lane-b"); assert_eq!(second["step_id"], "lane-a"); diff --git a/kernel/relayflowd/tests/crash_resume/llm_support.rs b/kernel/relayflowd/tests/crash_resume/llm_support.rs index 0c240cc6e..b35975de7 100644 --- a/kernel/relayflowd/tests/crash_resume/llm_support.rs +++ b/kernel/relayflowd/tests/crash_resume/llm_support.rs @@ -235,6 +235,9 @@ pub struct ProtocolClient { reader: BufReader, next_id: u64, events: Vec, + /// The ceiling currently in force, so a timeout can report the bound that + /// actually fired rather than the default constant. + read_timeout: Duration, } /// Ceiling on any single protocol read in a test. @@ -264,6 +267,7 @@ impl ProtocolClient { reader, next_id: 1, events: Vec::new(), + read_timeout: READ_TIMEOUT, } } @@ -324,8 +328,39 @@ impl ProtocolClient { } } - pub fn set_read_timeout(&self, timeout: Option) { - self.stream.set_read_timeout(timeout).unwrap(); + /// Override the read ceiling. `None` restores the default -- it does NOT + /// make reads unbounded. + /// + /// Deliberately NOT named `set_read_timeout`: that name belongs to + /// `UnixStream`, where `None` means "block forever". Shadowing a std API + /// while inverting its meaning is a trap no docstring reliably defuses. + /// + /// That distinction is the point. Callers tighten the bound for a specific + /// assertion and then pass `None` to mean "back to normal". Two shapes use + /// it, and they are not the same: + /// + /// - `parallel_lifecycle.rs:138,208` probe for SILENCE -- 200ms, then + /// assert the read errors. + /// - `concurrency.rs:31` tightens to 1s and expects the read to SUCCEED, + /// so a missing dispatch fails fast instead of stalling the test. + /// + /// If `None` meant "block forever", every read after any of those would be + /// unbounded and the ceiling this type advertises would be a claim it does + /// not keep -- which is what two review lenses caught in the first revision + /// of this change. + pub fn override_read_timeout(&mut self, timeout: Option) { + let timeout = timeout.unwrap_or(READ_TIMEOUT); + // Set it on the fd `read_frame` actually reads through -- the reader's, + // not `self.stream`. `try_clone` produces a separate descriptor, and + // while Linux and Darwin keep SO_RCVTIMEO on the shared socket (so + // writing to either fd happens to work today), that is a platform + // detail, not a guarantee. Going through the reader means the ceiling + // is set where it is read, with no hidden assumption to remember. + self.reader + .get_ref() + .set_read_timeout(Some(timeout)) + .unwrap(); + self.read_timeout = timeout; } fn read_frame(&mut self) -> Result { @@ -342,8 +377,9 @@ impl ProtocolClient { std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut ) => { + let waited = self.read_timeout; bail!( - "timed out after {READ_TIMEOUT:?} waiting for a protocol frame; \ + "timed out after {waited:?} waiting for a protocol frame; \ the daemon sent nothing (see #174)" ) } diff --git a/kernel/relayflowd/tests/crash_resume/parallel_lifecycle.rs b/kernel/relayflowd/tests/crash_resume/parallel_lifecycle.rs index 26275928a..62f24e71f 100644 --- a/kernel/relayflowd/tests/crash_resume/parallel_lifecycle.rs +++ b/kernel/relayflowd/tests/crash_resume/parallel_lifecycle.rs @@ -135,9 +135,9 @@ fn overlapping_agent_lanes_serialize_while_disjoint_lanes_merge_in_either_order( start_run(&fixture); let lane_b = worker.event("step.dispatch").unwrap(); assert_eq!(lane_b["step_id"], "lane-b"); - worker.set_read_timeout(Some(Duration::from_millis(200))); + worker.override_read_timeout(Some(Duration::from_millis(200))); assert!(worker.event("step.dispatch").is_err()); - worker.set_read_timeout(None); + worker.override_read_timeout(None); complete_agent(&mut worker, &lane_b, &[("repo-b", "rB")]).unwrap(); let lane_a = worker.event("step.dispatch").unwrap(); assert_eq!(lane_a["pins"]["workspace"][0]["revision_id"], "rB"); @@ -205,9 +205,9 @@ fn overlapping_agent_conflict_survives_server_crash_and_resume() { let retried = replacement.event("step.dispatch").unwrap(); assert_eq!(retried["step_id"], "lane-b"); assert_eq!(retried["attempt"], 2); - replacement.set_read_timeout(Some(Duration::from_millis(200))); + replacement.override_read_timeout(Some(Duration::from_millis(200))); assert!(replacement.event("step.dispatch").is_err()); - replacement.set_read_timeout(None); + replacement.override_read_timeout(None); complete_agent(&mut replacement, &retried, &[("repo-b", "rB")]).unwrap(); let lane_a = replacement.event("step.dispatch").unwrap(); complete_agent(&mut replacement, &lane_a, &[("repo-b", "rA")]).unwrap(); From 715d4f4f68d1df06092524fdd54403b5ad905cb7 Mon Sep 17 00:00:00 2001 From: kjgbot Date: Sat, 5 Sep 2026 15:00:48 +0200 Subject: [PATCH 3/3] test(kernel): dump daemon-side state when a dispatch never arrives #174 has now produced four occurrences and, between them, four test names. Nothing else. The state that would explain it is all on the daemon side and none of it survives: the resumed child is still running when the read gives up, so `wait_with_output` is never reached and its output is dropped in the unwind, and the run's journal is never read. On a missing dispatch the test now kills the resumed child -- it is wedged by definition, and without killing it first the read below would block exactly as long as the one that already timed out -- then reports its stdout, its stderr, and every journal entry with seq, type and step. That last part is the question a missing dispatch actually raises: did the daemon resume and stall partway, or never resume at all? The journal answers it and nothing else does. Stacked on #175, which supplies the bounded read this depends on. With an unbounded read there is no error to catch and nothing to report. Verified by forcing the ceiling to 1ms: before-first: no step.dispatch after resume: timed out after 1ms waiting for a protocol frame; the daemon sent nothing (see #174) --- resume child --- stdout (0 bytes): stderr (0 bytes): --- journal (1 entries) --- seq=1 type=RunSpawned step=None sha256 19f3431a -> 361572ef -> restored 19f3431a, no 1ms literal left in the tree. crash_resume 34 passed in 37.72s / 37.91s / 37.90s. Workspace 142 passed, 0 failed. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01FtQSAcGDta5VH9xiZFT4sR Session-Id: c228933d-4f94-4d83-9a9a-daf3c83b94f1 --- kernel/relayflowd/tests/crash_resume/agent.rs | 17 +++++-- kernel/relayflowd/tests/crash_resume/llm.rs | 18 +++++-- .../relayflowd/tests/crash_resume/support.rs | 49 +++++++++++++++++++ 3 files changed, 77 insertions(+), 7 deletions(-) diff --git a/kernel/relayflowd/tests/crash_resume/agent.rs b/kernel/relayflowd/tests/crash_resume/agent.rs index 3ea4f7b5e..76643f73a 100644 --- a/kernel/relayflowd/tests/crash_resume/agent.rs +++ b/kernel/relayflowd/tests/crash_resume/agent.rs @@ -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, }, }; @@ -127,8 +127,17 @@ 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. + 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(); diff --git a/kernel/relayflowd/tests/crash_resume/llm.rs b/kernel/relayflowd/tests/crash_resume/llm.rs index c3e357bad..6d6df1175 100644 --- a/kernel/relayflowd/tests/crash_resume/llm.rs +++ b/kernel/relayflowd/tests/crash_resume/llm.rs @@ -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] @@ -106,8 +109,17 @@ 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. + 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!( diff --git a/kernel/relayflowd/tests/crash_resume/support.rs b/kernel/relayflowd/tests/crash_resume/support.rs index 97444c864..261bfc2dc 100644 --- a/kernel/relayflowd/tests/crash_resume/support.rs +++ b/kernel/relayflowd/tests/crash_resume/support.rs @@ -156,6 +156,55 @@ pub fn journal_entries(data_dir: &Path) -> Option> { 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: the resumed child is still running, +/// so `wait_with_output` is never reached and its output is discarded when the +/// test unwinds; the journal is never read. Four occurrences produced four test +/// names and nothing else. +/// +/// This kills the child first -- it is wedged by definition, and without that +/// the read below would block as long as the one that already timed out -- then +/// reports its output and the run's journal together. +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()