From fd34d15a010ac9a539648fe49795128ae84558cc Mon Sep 17 00:00:00 2001 From: kjgbot Date: Sat, 5 Sep 2026 14:05:52 +0200 Subject: [PATCH 1/2] 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 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 42768e78b35b48ecada1aabc68489ca33f636849 Mon Sep 17 00:00:00 2001 From: kjgbot Date: Sat, 5 Sep 2026 14:15:05 +0200 Subject: [PATCH 2/2] 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 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();