hyf

Context-aware query service for Radroots
git clone https://radroots.dev/git/hyf.git
Log | Files | Refs | README | LICENSE

commit 6b59a8ec82cd7754605b71a5aefe72ba069bb724
parent f3bf6c0dc429e9fca1095fb8206bf90fa32ce0e6
Author: triesap <tyson@radroots.org>
Date:   Tue, 22 Sep 2026 21:59:52 +0000

test(hyf): close C002C review findings on status, readiness and diagnostics

- Cache an observed status without consuming the child so repeated status(),
  terminate() and reap() stay consistent and reap still recovers the report
  and evaluates exit/accounting truth after status inspection
- Route readiness through parse_ready_or_cleanup so malformed ready lines
  terminate the exact owned child and report the cause plus cleanup result
- Propagate non-timeout completion-probe errors as bounded io_error instead of
  a false success, and classify a diagnostics read failure apart from EOF
- Add executed controls: status-before-reap, malformed-ready cleanup, ready
  grammar rejections, planted socket census, stderr flood overflow and
  diagnostics read failure; 60/60 provider-helper and 7/7 repo-local-process

Diffstat:
Mtests/jev_provider_helper.mojo | 35+++++++++++++++++++----------------
Mtests/max_local_process_helper.mojo | 40++++++++++++++++++++++++----------------
Mtests/parent_lifecycle.mojo | 26++++++++++++++++++++++++++
Mtests/stdio_process_helper.mojo | 4+++-
Atests/stdio_stderr_flood_entrypoint.mojo | 17+++++++++++++++++
Mtests/strict_fixture.mojo | 28++++++++++------------------
Mtests/test_provider_helpers.mojo | 103++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mtests/test_repo_local_process_contract.mojo | 22++++++++++++++++++++++
8 files changed, 223 insertions(+), 52 deletions(-)

diff --git a/tests/jev_provider_helper.mojo b/tests/jev_provider_helper.mojo @@ -22,10 +22,11 @@ from parent_lifecycle import ( dup2_fd, fork_pid, make_pipe, - parse_ready_line, + parse_ready_or_cleanup, set_alarm, terminate_owned, wait_bounded, + wait_nohang, write_raw, ) from strict_fixture import ( @@ -319,16 +320,26 @@ struct SpawnedJevStub(Movable): + String(self.state.connections) ) - def status(self) -> ProcessStatus: + def status(mut self) -> ProcessStatus: + """Observe child status without losing ownership or report truth.""" if self.state.reaped: return self.state.status.copy() - return wait_bounded(self.pid, 0) + if self.state.observed_valid: + return self.state.observed.copy() + var st = wait_nohang(self.pid) + if st.state == "running" or st.state == "interrupted": + return st^ + self.state.observed = st.copy() + self.state.observed_valid = True + return st^ def reap(mut self): """Strictly reap the owned child and decode its bounded report.""" if self.state.reaped: return var status = wait_bounded(self.pid, self.state.deadline_ms) + if self.state.observed_valid: + status = self.state.observed.copy() self.state.status = status.copy() if status.state == "running" or status.state == "interrupted": var term = terminate_owned(self.pid, TERMINATION_GRACE_MS) @@ -467,7 +478,7 @@ struct SpawnedJevStubView(Movable): def describe(self) -> String: return self.target[].describe() - def status(self) -> ProcessStatus: + def status(mut self) -> ProcessStatus: return self.target[].status() def reap(mut self): @@ -566,6 +577,8 @@ def _spawn_child_or_cleanup( connections=0, cleanup_error="", status=ProcessStatus("pending", False, -1, 0, 0, ""), + observed=ProcessStatus("pending", False, -1, 0, 0, ""), + observed_valid=False, ) var ready_line = "" try: @@ -584,21 +597,11 @@ def _spawn_child_or_cleanup( ) var reported_port = 0 try: - reported_port = parse_ready_line(ready_line, 256) + reported_port = parse_ready_or_cleanup(pid, ready_line, 256) except e: - var st = terminate_owned(pid, TERMINATION_GRACE_MS) - state.status = st.copy() state.reaped = True state.close_reader() - raise Error( - "jev stub malformed readiness (" - + String(e) - + " / report=" - + ready_line - + " / " - + st.describe() - + ")" - ) + raise Error("jev stub malformed readiness (" + String(e) + ")") _ = mode return SpawnedJevStubAuto( port=reported_port, stub=SpawnedJevStub(pid, state^) diff --git a/tests/max_local_process_helper.mojo b/tests/max_local_process_helper.mojo @@ -24,10 +24,11 @@ from parent_lifecycle import ( dup2_fd, fork_pid, make_pipe, - parse_ready_line, + parse_ready_or_cleanup, set_alarm, terminate_owned, wait_bounded, + wait_nohang, write_raw, ) from strict_fixture import ( @@ -413,10 +414,23 @@ struct SpawnedMaxLocalStub(Movable): + String(self.state.connections) ) - def status(self) -> ProcessStatus: + def status(mut self) -> ProcessStatus: + """Observe child status without losing ownership or report truth. + + If the child has already exited, the observed reap status is cached in + ``state.observed`` rather than claimed as a completed reap, so a later + ``reap()`` still reads the report and evaluates exit/accounting truth. + """ if self.state.reaped: return self.state.status.copy() - return wait_bounded(self.pid, 0) + if self.state.observed_valid: + return self.state.observed.copy() + var st = wait_nohang(self.pid) + if st.state == "running" or st.state == "interrupted": + return st^ + self.state.observed = st.copy() + self.state.observed_valid = True + return st^ def reap(mut self): """Reap the owned child and strictly decode its bounded report. @@ -429,6 +443,8 @@ struct SpawnedMaxLocalStub(Movable): if self.state.reaped: return var status = wait_bounded(self.pid, self.state.deadline_ms) + if self.state.observed_valid: + status = self.state.observed.copy() self.state.status = status.copy() if status.state == "running" or status.state == "interrupted": var term = terminate_owned(self.pid, TERMINATION_GRACE_MS) @@ -572,7 +588,7 @@ struct SpawnedMaxLocalView(Movable): def describe(self) -> String: return self.target[].describe() - def status(self) -> ProcessStatus: + def status(mut self) -> ProcessStatus: return self.target[].status() def reap(mut self): @@ -685,6 +701,8 @@ def _spawn_max_local( connections=0, cleanup_error="", status=ProcessStatus("pending", False, -1, 0, 0, ""), + observed=ProcessStatus("pending", False, -1, 0, 0, ""), + observed_valid=False, ) var ready_line = "" try: @@ -703,19 +721,9 @@ def _spawn_max_local( ) var reported_port = 0 try: - reported_port = parse_ready_line(ready_line, 256) + reported_port = parse_ready_or_cleanup(pid, ready_line, 256) except e: - var st = terminate_owned(pid, TERMINATION_GRACE_MS) - state.status = st.copy() state.reaped = True state.close_reader() - raise Error( - "max_local stub malformed readiness (" - + String(e) - + " / report=" - + ready_line - + " / " - + st.describe() - + ")" - ) + raise Error("max_local stub malformed readiness (" + String(e) + ")") return SpawnedMaxLocalStub(pid, reported_port, state^) diff --git a/tests/parent_lifecycle.mojo b/tests/parent_lifecycle.mojo @@ -584,6 +584,8 @@ struct PipedChildState(Movable): var connections: Int var cleanup_error: String var status: ProcessStatus + var observed: ProcessStatus + var observed_valid: Bool def store( mut self, @@ -695,6 +697,30 @@ def parse_ready_line(line: String, max_bytes: Int) raises -> Int: # ── Non-opening descriptor census ─────────────────────────────────────────── +def parse_ready_or_cleanup( + pid: Int, line: String, max_bytes: Int +) raises -> Int: + """Parse exact readiness or terminate and prove the owned child is gone. + + Gives malformed readiness a real cleanup path with a bounded cause carrying + the offending line and the owned child's reap result. + """ + try: + return parse_ready_line(line, max_bytes) + except e: + var status = terminate_owned(pid, TERMINATION_GRACE_MS) + raise Error( + "ready_invalid:" + + String(e) + + " report=" + + line + + " child=" + + status.describe() + + " cleanup=" + + ("proved" if status.cleanup_proved() else "unreaped") + ) + + def descriptor_census(limit: Int) -> Int: """Count open descriptors numerically via ``fcntl(F_GETFD)``. diff --git a/tests/stdio_process_helper.mojo b/tests/stdio_process_helper.mojo @@ -90,7 +90,9 @@ def drain_ready( return DrainOutcome(False, "") var buf = InlineArray[Byte, 4096](fill=0) var n = read_fd(fd, buf.unsafe_ptr(), 4096) - if n <= 0: + if n < 0: + return DrainOutcome(True, "read_error") + if n == 0: return DrainOutcome(True, "") if len(out) + n > cap: return DrainOutcome(True, "stream_overflow") diff --git a/tests/stdio_stderr_flood_entrypoint.mojo b/tests/stdio_stderr_flood_entrypoint.mojo @@ -0,0 +1,17 @@ +"""Test-only stdio entrypoint that floods stderr before reading stdin. + +Exercises LC04: the parent must drain stderr concurrently with writing the +request instead of deadlocking on a full stderr pipe, and must bound the +diagnostics it retains with a distinct overflow cause. +""" + +from std.sys import stderr + + +def main(): + var chunk = String("") + for _ in range(1000): + chunk += "e" + for _ in range(100): + print(chunk, file=stderr) + print('{"ok":true}') diff --git a/tests/strict_fixture.mojo b/tests/strict_fixture.mojo @@ -423,31 +423,23 @@ struct ConnectionReader(Movable): outcome.total_bytes = expected_total return outcome^ - def probe_completion(mut self, grace_ms: Int) -> String: + def probe_completion(mut self, grace_ms: Int) raises -> String: """Bounded completion handshake after the final expected exchange. Any already-buffered or subsequently received bytes mean the client - sent an extra exchange. A timeout with no bytes is success; this never - waits indefinitely for hypothetical future requests. + sent an extra exchange. A bounded timeout with no bytes is success; any + other I/O/setup error propagates (it is never a successful completion) + and is recorded by the serve loop as a bounded io_error. """ if len(self._buffer) > 0: return "extra_exchange_after_completion" + self._stream.set_recv_timeout(grace_ms) try: - self._stream.set_recv_timeout(grace_ms) - except: - return "completion_probe_config_error" - try: - try: - var n = self._read_more() - if n > 0: - return "extra_exchange_after_completion" - except Timeout: - # A bounded grace elapsed with no extra bytes: success. - return "" - except e: - # An unexpected I/O error is not a successful completion. - _ = String(e) - return "completion_probe_error" + var n = self._read_more() + if n > 0: + return "extra_exchange_after_completion" + except Timeout: + return "" return "" diff --git a/tests/test_provider_helpers.mojo b/tests/test_provider_helpers.mojo @@ -6,7 +6,7 @@ exception, signal or unrelated startup failure. """ from std.testing import TestSuite, assert_true, assert_equal -from std.ffi import ErrNo +from std.ffi import ErrNo, c_int, external_call from flare.net import SocketAddr from flare.tcp import TcpListener, TcpStream @@ -22,9 +22,12 @@ from parent_lifecycle import ( fork_pid, make_pipe, open_fd_count, + parse_ready_line, + parse_ready_or_cleanup, pid_not_waitable, read_all_bounded, read_line_bounded, + sleep_ms, wait_nohang, write_fd_bounded, write_raw, @@ -946,6 +949,8 @@ def _owned_report_child( connections=0, cleanup_error="", status=ProcessStatus("pending", False, -1, 0, 0, ""), + observed=ProcessStatus("pending", False, -1, 0, 0, ""), + observed_valid=False, ) return SpawnedMaxLocalStub(pid, 0, state^) @@ -1166,6 +1171,8 @@ def test_coalesced_ready_and_report_lines_retain_surplus() raises: connections=0, cleanup_error="", status=ProcessStatus("pending", False, -1, 0, 0, ""), + observed=ProcessStatus("pending", False, -1, 0, 0, ""), + observed_valid=False, ) var ready = state.read_line(2048, 500) var report = state.read_line(2048, 500) @@ -1241,6 +1248,100 @@ def test_descriptor_census_detects_planted_high_fd() raises: assert_true(after <= with_pipe) +def test_descriptor_census_detects_planted_socket() raises: + # The census must count sockets, not only regular files. + var before = open_fd_count() + var sock = Int(external_call["socket", c_int](c_int(2), c_int(1), c_int(0))) + assert_true(sock >= 0) + var with_socket = open_fd_count() + assert_true(with_socket > before) + close_fd(sock) + assert_true(open_fd_count() <= with_socket) + + +def test_ready_grammar_rejections() raises: + var valid = 0 + var raised = "" + try: + valid = parse_ready_line("ready 65535", 256) + except e: + raised = String(e) + assert_equal(valid, 65535) + assert_equal(raised, "") + var cases = List[String]() + cases.append("") + cases.append("ready") + cases.append("ready ") + cases.append("ready abc") + cases.append("ready 12345x") + cases.append("ready 70000") + cases.append("ready 0") + cases.append("not-ready") + for probe in cases: + var reason = "" + try: + _ = parse_ready_line(probe, 256) + except e: + reason = String(e) + assert_true(reason != "") + + +def test_malformed_ready_line_terminates_owned_child() raises: + # LC03: malformed readiness must clean up the exact owned child. + var pipe = make_pipe() + var pid = fork_pid() + if pid == 0: + close_fd(pipe.read_fd) + _ = write_raw(pipe.write_fd, "garbage-not-ready\n") + for _ in range(400): + sleep_ms(50) + child_exit(0) + close_fd(pipe.write_fd) + var line = read_line_bounded(pipe.read_fd, 256, 500) + close_fd(pipe.read_fd) + var message = "" + try: + _ = parse_ready_or_cleanup(pid, line, 256) + except e: + message = String(e) + assert_true(message.find("ready_invalid:ready_grammar") >= 0) + assert_true(message.find("cleanup=proved") >= 0) + assert_true(pid_not_waitable(pid)) + + +def test_status_observation_preserves_ownership_and_report() raises: + # LC01/LC02: observing an exited child must not consume the report. + var stub = _owned_report_child( + 0, + "result ok phase=complete case=- reason=ok requests=1 connections=1\n", + ) + var observed = "" + for _ in range(200): + observed = stub.status().state + if observed != "running" and observed != "interrupted": + break + sleep_ms(20) + assert_equal(observed, "reaped") + var second = stub.status() + assert_equal(second.state, "reaped") + stub.reap() + assert_true(stub.ok()) + assert_equal(stub.request_count(), 1) + assert_true(stub.status().cleanup_proved()) + + +def test_status_observation_then_terminate_is_safe() raises: + var stub = _owned_report_child(0, "") + sleep_ms(100) + _ = stub.status() + var owned_pid = stub.pid + stub.terminate() + stub.terminate() + stub.reap() + assert_true(pid_not_waitable(owned_pid)) + assert_true(not stub.ok()) + + # ── LC05: coalesced header cap accounting ─────────────────────────────────── diff --git a/tests/test_repo_local_process_contract.mojo b/tests/test_repo_local_process_contract.mojo @@ -164,5 +164,27 @@ def test_diagnostics_overflow_is_bounded_with_cause() raises: assert_true(len(out) <= 65536) +def test_diagnostics_read_failure_is_distinct_from_eof() raises: + # LC04: an unavailable descriptor is a read failure, not a clean EOF. + var out = List[UInt8]() + var d = drain_ready(-1, out, 16, POLLIN) + assert_true(d.eof) + assert_equal(d.reason, "read_error") + + +def test_run_stdio_entrypoint_drains_stderr_concurrently_and_bounds_it() raises: + # LC04: a child that floods stderr before producing output must not + # deadlock the parent, and the overflow is reported with its own cause. + var message = "" + try: + _ = run_stdio_entrypoint_with_deadline( + "tests/stdio_stderr_flood_entrypoint.mojo", "{}", "", "", 30000 + ) + except e: + message = String(e) + assert_true(message.find("stderr_stream_overflow") >= 0) + assert_true(message.find("timeout") < 0) + + def main() raises: TestSuite.discover_tests[__functions_in_module()]().run()