hyf

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

commit 5f82d3d097b7afa5f40a0d945f7ab0da36c72ce8
parent 951b345c7aff3c88fb5d376302e00538641a1a1a
Author: triesap <tyson@radroots.org>
Date:   Wed, 23 Sep 2026 02:50:08 +0000

test(hyf): complete report truth and retained cleanup ownership (C002D)

ADR-0017 D37 PC01-PC05 and the inherited LC/FX obligations, repairing the
guards the period-8 probes falsified. Test-only: no product src, schema,
dependency or lock change.

PC01 complete report stream/semantics: the owned-child reader retains surplus
as undecoded bytes (no premature UTF-8 decode), records newline termination,
drains every byte after the report through EOF so aligned/split/coalesced
duplicate and trailing reports fail regardless of chunk boundary, and rejects
unterminated, malformed, empty-field and inconsistent success reports on both
provider reap paths. Successful emitted reports must be phase=complete,
reason=ok with expected accounting.

PC02 retained ownership: cleanup/terminate only release ownership once the
child is provably collected; an uncertain wait stays retryable and records the
failure in a caller-owned CleanupLedger that stays observable after scope exit.

PC03 bounded process I/O: stdio fails when intended request bytes were not all
delivered (early HUP/error/peer exit) even with valid JSON and exit 0, uses the
declared finite budget across write/read/wait, and distinguishes read errors
from EOF/EOF-timeouts with cause-specific polling.

PC04 cause-specific evidence: the completion probe now propagates a real
socket read error after successful timeout setup (the Mojo error model cannot
distinguish the struct payloads by handler type) instead of swallowing it as a
timeout; the descriptor census fails explicitly when complete coverage inside
the admitted ceiling is unavailable.

PC05 truthful closure: 74/74 provider-helper and 8/8 repo-local-process tests
pass on final source; test-stdio remains 56/32/24/0 with identical D16 names
and reasons; format/architecture/build checks green.

Diffstat:
Mtests/jev_provider_helper.mojo | 134++++++++++++++++++++++++++++++++++++++++++++++++-------------------------------
Mtests/max_local_process_helper.mojo | 107++++++++++++++++++++++++++++++++++++++++++++-----------------------------------
Mtests/parent_lifecycle.mojo | 265+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------------
Atests/stdio_no_read_entrypoint.mojo | 10++++++++++
Mtests/stdio_process_helper.mojo | 15+++++++++++++--
Mtests/strict_fixture.mojo | 27++++++++++++++++++++-------
Mtests/test_provider_helpers.mojo | 332++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------
Mtests/test_repo_local_process_contract.mojo | 19+++++++++++++++++++
8 files changed, 685 insertions(+), 224 deletions(-)

diff --git a/tests/jev_provider_helper.mojo b/tests/jev_provider_helper.mojo @@ -14,6 +14,7 @@ from flare.utils import usleep from parent_lifecycle import ( FIXTURE_DEFAULT_DEADLINE_MS, TERMINATION_GRACE_MS, + CleanupLedger, PipedChildState, PipeFds, ProcessStatus, @@ -22,7 +23,9 @@ from parent_lifecycle import ( dup2_fd, fork_owned_or_close, make_pipe, + now_ms, parse_ready_or_cleanup, + piped_child_state, set_alarm, terminate_owned, wait_bounded, @@ -275,15 +278,25 @@ struct SpawnedJevStub(Movable): self.cleanup() def cleanup(mut self): - """Fast, non-raising owned cleanup for assertion/error/early return.""" + """Fast, non-raising owned cleanup for assertion/error/early return. + + Ownership is released only once the child is provably collected. An + uncertain wait keeps the handle retryable and records the failure in + the caller-owned ledger so it stays observable after scope exit + instead of being silently marked complete. + """ if self.state.reaped: return var status = terminate_owned(self.pid, TERMINATION_GRACE_MS) self.state.status = status.copy() - self.state.reaped = True - self.state.close_reader() - if not status.cleanup_proved(): - self.state.cleanup_error = "unreaped:" + status.describe() + if status.cleanup_proved(): + self.state.reaped = True + self.state.close_reader() + return + self.state.cleanup_error = "unreaped:" + status.describe() + self.state.ledger.record( + "owned-child cleanup unproved " + status.describe() + ) def ok(self) -> Bool: return self.state.ok @@ -337,7 +350,12 @@ struct SpawnedJevStub(Movable): """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) + var remaining = self.state.deadline_ms - ( + now_ms() - self.state.spawn_ms + ) + if remaining < 1: + remaining = 1 + var status = wait_bounded(self.pid, remaining) if self.state.observed_valid: status = self.state.observed.copy() self.state.status = status.copy() @@ -359,13 +377,23 @@ struct SpawnedJevStub(Movable): self.state.close_reader() return var report_text = "" - var read_error = "" + var report_error = "" try: report_text = self.state.read_line(STRICT_MAX_REPORT_BYTES, 1000) + if self.state.last_terminated: + var surplus = self.state.drain_surplus( + STRICT_MAX_REPORT_BYTES, 1000 + ) + if surplus > 0: + report_error = "duplicate_report" + elif report_text.byte_length() > 0: + # Bytes at EOF without a terminating newline are a truncated + # report, never a complete one. + report_error = "unterminated_report" except e: - read_error = String(e) + report_error = String(e) self.state.close_reader() - if report_text == "": + if report_text == "" and report_error == "": if status.exited and status.exit_code == 0: self.state.store(False, "startup", "-", "missing_report", 0, 0) elif status.exited: @@ -386,15 +414,15 @@ struct SpawnedJevStub(Movable): 0, 0, ) - if read_error != "": - self.state.cleanup_error = "report_read:" + read_error + self.state.reaped = True + return + if report_error != "": + self.state.store(False, "parse", "-", report_error, -1, -1) self.state.reaped = True return var parsed = parse_report(report_text) if parsed.phase == "parse": self.state.store(False, "parse", "-", parsed.reason, -1, -1) - if read_error != "": - self.state.cleanup_error = "report_read:" + read_error self.state.reaped = True return self.state.store( @@ -405,13 +433,7 @@ struct SpawnedJevStub(Movable): parsed.requests, parsed.connections, ) - if self.state.pending != "": - # A second report line after the first is a duplicate/malformed - # report, never a success. - self.state.ok = False - self.state.phase = "parse" - self.state.reason = "duplicate_report" - elif not report_status_matches_exit( + if not report_status_matches_exit( status.exited, status.exit_code, parsed.ok ): self.state.ok = False @@ -440,12 +462,15 @@ struct SpawnedJevStub(Movable): return var status = terminate_owned(self.pid, TERMINATION_GRACE_MS) self.state.status = status.copy() - self.state.reaped = True - self.state.close_reader() - if not status.cleanup_proved(): - raise Error( - "lifecycle: owned child not reaped: " + status.describe() - ) + if status.cleanup_proved(): + self.state.reaped = True + self.state.close_reader() + return + self.state.cleanup_error = "unreaped:" + status.describe() + self.state.ledger.record( + "owned-child cleanup unproved " + status.describe() + ) + raise Error("lifecycle: owned child not reaped: " + status.describe()) struct SpawnedJevStubView(Movable): @@ -531,9 +556,12 @@ def reserve_jev_port() raises -> Int: def spawn_jev_stub_auto( - mode: String, requests: Int, deadline_ms: Int = FIXTURE_DEFAULT_DEADLINE_MS + mode: String, + requests: Int, + deadline_ms: Int = FIXTURE_DEFAULT_DEADLINE_MS, + ledger: CleanupLedger = CleanupLedger(), ) raises -> SpawnedJevStubAuto: - return _spawn_jev_stub(0, mode, requests, deadline_ms) + return _spawn_jev_stub(0, mode, requests, deadline_ms, ledger) def spawn_jev_stub( @@ -558,33 +586,22 @@ def serve_jev_scripted( def spawn_jev_scripted_auto( var scripts: List[ExchangeScript], deadline_ms: Int = FIXTURE_DEFAULT_DEADLINE_MS, + ledger: CleanupLedger = CleanupLedger(), ) raises -> SpawnedJevStubAuto: - return _spawn_jev_scripted(0, scripts^, deadline_ms) + return _spawn_jev_scripted(0, scripts^, deadline_ms, ledger) def _spawn_child_or_cleanup( - pipe: PipeFds, pid: Int, mode: String, deadline_ms: Int, requests: Int + pipe: PipeFds, + pid: Int, + mode: String, + deadline_ms: Int, + requests: Int, + ledger: CleanupLedger, ) raises -> SpawnedJevStubAuto: """Build the owned state, read exact readiness, or clean up and raise.""" - var state = PipedChildState( - pid=pid, - report_fd=pipe.read_fd, - pending="", - eof=False, - closed=False, - deadline_ms=deadline_ms, - expected_requests=requests, - reaped=False, - ok=False, - phase="pending", - case_label="-", - reason="not_reaped", - requests=0, - connections=0, - cleanup_error="", - status=ProcessStatus("pending", False, -1, 0, 0, ""), - observed=ProcessStatus("pending", False, -1, 0, 0, ""), - observed_valid=False, + var state = piped_child_state( + pid, pipe.read_fd, deadline_ms, requests, ledger ) var ready_line = "" try: @@ -615,7 +632,10 @@ def _spawn_child_or_cleanup( def _spawn_jev_scripted( - port: Int, var scripts: List[ExchangeScript], deadline_ms: Int + port: Int, + var scripts: List[ExchangeScript], + deadline_ms: Int, + ledger: CleanupLedger = CleanupLedger(), ) raises -> SpawnedJevStubAuto: var total = len(scripts) var pipe = make_pipe() @@ -637,11 +657,17 @@ def _spawn_jev_scripted( write_raw(1, report_line(failed) + "\n") child_exit(125) close_fd(pipe.write_fd) - return _spawn_child_or_cleanup(pipe, pid, "scripted", deadline_ms, total) + return _spawn_child_or_cleanup( + pipe, pid, "scripted", deadline_ms, total, ledger + ) def _spawn_jev_stub( - port: Int, mode: String, requests: Int, deadline_ms: Int + port: Int, + mode: String, + requests: Int, + deadline_ms: Int, + ledger: CleanupLedger = CleanupLedger(), ) raises -> SpawnedJevStubAuto: var pipe = make_pipe() var pid = fork_owned_or_close(pipe.copy()) @@ -662,4 +688,6 @@ def _spawn_jev_stub( write_raw(1, report_line(failed) + "\n") child_exit(125) close_fd(pipe.write_fd) - return _spawn_child_or_cleanup(pipe, pid, mode, deadline_ms, requests) + return _spawn_child_or_cleanup( + pipe, pid, mode, deadline_ms, requests, ledger + ) diff --git a/tests/max_local_process_helper.mojo b/tests/max_local_process_helper.mojo @@ -17,6 +17,7 @@ from flare.utils import usleep from parent_lifecycle import ( FIXTURE_DEFAULT_DEADLINE_MS, TERMINATION_GRACE_MS, + CleanupLedger, PipedChildState, ProcessStatus, child_exit, @@ -24,7 +25,9 @@ from parent_lifecycle import ( dup2_fd, fork_owned_or_close, make_pipe, + now_ms, parse_ready_or_cleanup, + piped_child_state, set_alarm, terminate_owned, wait_bounded, @@ -369,15 +372,25 @@ struct SpawnedMaxLocalStub(Movable): self.cleanup() def cleanup(mut self): - """Fast, non-raising owned cleanup for assertion/error/early return.""" + """Fast, non-raising owned cleanup for assertion/error/early return. + + Ownership is released only once the child is provably collected. An + uncertain wait keeps the handle retryable and records the failure in + the caller-owned ledger so it stays observable after scope exit + instead of being silently marked complete. + """ if self.state.reaped: return var status = terminate_owned(self.pid, TERMINATION_GRACE_MS) self.state.status = status.copy() - self.state.reaped = True - self.state.close_reader() - if not status.cleanup_proved(): - self.state.cleanup_error = "unreaped:" + status.describe() + if status.cleanup_proved(): + self.state.reaped = True + self.state.close_reader() + return + self.state.cleanup_error = "unreaped:" + status.describe() + self.state.ledger.record( + "owned-child cleanup unproved " + status.describe() + ) def ok(self) -> Bool: return self.state.ok @@ -442,7 +455,12 @@ struct SpawnedMaxLocalStub(Movable): """ if self.state.reaped: return - var status = wait_bounded(self.pid, self.state.deadline_ms) + var remaining = self.state.deadline_ms - ( + now_ms() - self.state.spawn_ms + ) + if remaining < 1: + remaining = 1 + var status = wait_bounded(self.pid, remaining) if self.state.observed_valid: status = self.state.observed.copy() self.state.status = status.copy() @@ -464,13 +482,23 @@ struct SpawnedMaxLocalStub(Movable): self.state.close_reader() return var report_text = "" - var read_error = "" + var report_error = "" try: report_text = self.state.read_line(STRICT_MAX_REPORT_BYTES, 1000) + if self.state.last_terminated: + var surplus = self.state.drain_surplus( + STRICT_MAX_REPORT_BYTES, 1000 + ) + if surplus > 0: + report_error = "duplicate_report" + elif report_text.byte_length() > 0: + # Bytes at EOF without a terminating newline are a truncated + # report, never a complete one. + report_error = "unterminated_report" except e: - read_error = String(e) + report_error = String(e) self.state.close_reader() - if report_text == "": + if report_text == "" and report_error == "": if status.exited and status.exit_code == 0: self.state.store(False, "startup", "-", "missing_report", 0, 0) elif status.exited: @@ -491,15 +519,15 @@ struct SpawnedMaxLocalStub(Movable): 0, 0, ) - if read_error != "": - self.state.cleanup_error = "report_read:" + read_error + self.state.reaped = True + return + if report_error != "": + self.state.store(False, "parse", "-", report_error, -1, -1) self.state.reaped = True return var parsed = parse_report(report_text) if parsed.phase == "parse": self.state.store(False, "parse", "-", parsed.reason, -1, -1) - if read_error != "": - self.state.cleanup_error = "report_read:" + read_error self.state.reaped = True return self.state.store( @@ -510,13 +538,7 @@ struct SpawnedMaxLocalStub(Movable): parsed.requests, parsed.connections, ) - if self.state.pending != "": - # A second report line after the first is a duplicate/malformed - # report, never a success. - self.state.ok = False - self.state.phase = "parse" - self.state.reason = "duplicate_report" - elif not report_status_matches_exit( + if not report_status_matches_exit( status.exited, status.exit_code, parsed.ok ): self.state.ok = False @@ -545,12 +567,15 @@ struct SpawnedMaxLocalStub(Movable): return var status = terminate_owned(self.pid, TERMINATION_GRACE_MS) self.state.status = status.copy() - self.state.reaped = True - self.state.close_reader() - if not status.cleanup_proved(): - raise Error( - "lifecycle: owned child not reaped: " + status.describe() - ) + if status.cleanup_proved(): + self.state.reaped = True + self.state.close_reader() + return + self.state.cleanup_error = "unreaped:" + status.describe() + self.state.ledger.record( + "owned-child cleanup unproved " + status.describe() + ) + raise Error("lifecycle: owned child not reaped: " + status.describe()) struct SpawnedMaxLocalView(Movable): @@ -640,18 +665,22 @@ def spawn_max_local_stub( mode: String, requests: Int, deadline_ms: Int = FIXTURE_DEFAULT_DEADLINE_MS, + ledger: CleanupLedger = CleanupLedger(), ) raises -> SpawnedMaxLocalStub: var scripts = List[ExchangeScript]() - return _spawn_max_local(port, scripts^, mode, requests, False, deadline_ms) + return _spawn_max_local( + port, scripts^, mode, requests, False, deadline_ms, ledger + ) def spawn_max_local_scripted( port: Int, var scripts: List[ExchangeScript], deadline_ms: Int = FIXTURE_DEFAULT_DEADLINE_MS, + ledger: CleanupLedger = CleanupLedger(), ) raises -> SpawnedMaxLocalStub: return _spawn_max_local( - port, scripts^, "scripted", len(scripts), True, deadline_ms + port, scripts^, "scripted", len(scripts), True, deadline_ms, ledger ) @@ -662,6 +691,7 @@ def _spawn_max_local( requests: Int, scripted: Bool, deadline_ms: Int, + ledger: CleanupLedger, ) raises -> SpawnedMaxLocalStub: var pipe = make_pipe() var pid = fork_owned_or_close(pipe.copy()) @@ -684,25 +714,8 @@ def _spawn_max_local( write_raw(1, report_line(failed) + "\n") child_exit(125) close_fd(pipe.write_fd) - var state = PipedChildState( - pid=pid, - report_fd=pipe.read_fd, - pending="", - eof=False, - closed=False, - deadline_ms=deadline_ms, - expected_requests=requests, - reaped=False, - ok=False, - phase="pending", - case_label="-", - reason="not_reaped", - requests=0, - connections=0, - cleanup_error="", - status=ProcessStatus("pending", False, -1, 0, 0, ""), - observed=ProcessStatus("pending", False, -1, 0, 0, ""), - observed_valid=False, + var state = piped_child_state( + pid, pipe.read_fd, deadline_ms, requests, ledger ) var ready_line = "" try: diff --git a/tests/parent_lifecycle.mojo b/tests/parent_lifecycle.mojo @@ -53,7 +53,7 @@ comptime SIGKILL: Int = 9 comptime SIGTERM: Int = 15 comptime SIG_IGN: Int = 1 comptime F_GETFD: Int = 1 -comptime CENSUS_MAX_FDS: Int = 65536 +comptime CENSUS_MAX_FDS: Int = 1048576 comptime FIXTURE_DEFAULT_DEADLINE_MS: Int = 20000 comptime TERMINATION_GRACE_MS: Int = 2000 @@ -209,6 +209,12 @@ def sleep_ms(ms: Int): def poll_fd(fd: Int, events: Int, timeout_ms: Int) -> Int: + """Poll one descriptor. + + Returns the ``revents`` mask, ``0`` on timeout and ``-1`` on a real + ``poll(2)`` error so callers can distinguish a read-phase error from an + ordinary timeout instead of collapsing both to ``0``. + """ var cell = InlineArray[Int32, 2](fill=0) cell[0] = Int32(fd) cell[1] = Int32(events) @@ -217,7 +223,9 @@ def poll_fd(fd: Int, events: Int, timeout_ms: Int) -> Int: cell.unsafe_ptr(), c_uint(1), c_int(timeout_ms) ) ) - if n <= 0: + if n < 0: + return -1 + if n == 0: return 0 return (Int(cell[1]) >> 16) & 0xFFFF @@ -290,6 +298,13 @@ def write_raw(fd: Int, text: String) -> Int: return n +def write_raw_bytes(fd: Int, bytes: List[UInt8]) -> Int: + """Best-effort blocking write of raw bytes (test-only split controls).""" + if len(bytes) == 0: + return 0 + return _write_fd(fd, bytes.unsafe_ptr(), len(bytes)) + + comptime WRITE_CHUNK_BYTES: Int = 512 @@ -308,6 +323,8 @@ def write_fd_bounded(fd: Int, data: String, deadline_ms: Int) -> String: if now_ms() - start >= deadline_ms: return "write_deadline_expired" var ev = poll_fd(fd, POLLOUT, LIFECYCLE_POLL_SLICE_MS) + if ev < 0: + return "write_poll_error" if ev == 0: continue if (ev & (POLLERR | POLLHUP | POLLNVAL)) != 0: @@ -428,11 +445,15 @@ struct BoundedLineReader(Movable): if now_ms() - start >= deadline_ms: raise Error("read_deadline_expired") var ev = poll_fd(self.fd, POLLIN, LIFECYCLE_POLL_SLICE_MS) + if ev < 0: + raise Error("read_error") if ev == 0: continue var buf = InlineArray[Byte, 512](fill=0) var n = read_fd(self.fd, buf.unsafe_ptr(), 512) - if n <= 0: + if n < 0: + raise Error("read_error") + if n == 0: self._eof = True continue for index in range(n): @@ -462,10 +483,14 @@ def read_all_bounded( if now_ms() - start >= deadline_ms: raise Error("read_deadline_expired") var ev = poll_fd(fd, POLLIN, LIFECYCLE_POLL_SLICE_MS) + if ev < 0: + raise Error("read_error") if ev == 0: continue var n = read_fd(fd, buf.unsafe_ptr(), 4096) - if n <= 0: + if n < 0: + raise Error("read_error") + if n == 0: break if len(out) + n > max_bytes: raise Error("stdout_overflow") @@ -480,6 +505,48 @@ def bytes_to_string(bytes: List[UInt8]) raises -> String: return String(from_utf8=Span(ptr=bytes.unsafe_ptr(), length=len(bytes))) +def _utf8_line(bytes: List[UInt8]) raises -> String: + """Decode a complete line as UTF-8, or fail with a bounded cause.""" + try: + return bytes_to_string(bytes) + except: + raise Error("invalid_utf8") + + +struct CleanupLedger(Copyable, Movable): + """Caller-owned record of owned-child cleanup failures. + + The ledger is deliberately *not* owned by the handle: an unproved cleanup + stays observable after the owning ``with`` block ends and the handle is + destroyed, because the calling test keeps the ``List`` and can assert it + is empty. A default-constructed ``CleanupLedger`` is inert. + """ + + var errors: Optional[UnsafePointer[List[String], MutAnyOrigin]] + + def __init__(out self): + self.errors = None + + def __init__(out self, errors: UnsafePointer[List[String], MutAnyOrigin]): + self.errors = errors + + def record(self, text: String): + if self.errors: + self.errors.value()[].append(text) + + def count(self) -> Int: + if self.errors: + return len(self.errors.value()[]) + return -1 + + def last(self) -> String: + if self.errors: + var recorded = self.errors.value()[] + if len(recorded) > 0: + return String(recorded[len(recorded) - 1]) + return "" + + # ── Child lifecycle ───────────────────────────────────────────────────────── @@ -621,14 +688,16 @@ def pid_not_waitable(pid: Int) -> Bool: struct PipedChildState(Movable): """Single mutable lifecycle record shared by every copy of one handle. - Holds no ``List``: in-place container mutation through a shared reference - is avoided so the record stays safe to share across handle copies. The - retained line surplus is an immutable ``String`` reassigned in place. + Retained read surplus is an undecoded byte buffer, so a chunk that splits a + multi-byte character is preserved without raising on an incomplete UTF-8 + fragment. ``last_terminated`` records whether the most recent line ended + with a newline, so a truncated report can never be read as complete. + ``ledger`` is the caller-owned record of cleanup failures. """ var pid: Int var report_fd: Int - var pending: String + var pending: List[UInt8] var eof: Bool var closed: Bool var deadline_ms: Int @@ -644,6 +713,9 @@ struct PipedChildState(Movable): var status: ProcessStatus var observed: ProcessStatus var observed_valid: Bool + var last_terminated: Bool + var spawn_ms: Int + var ledger: CleanupLedger def store( mut self, @@ -667,66 +739,125 @@ struct PipedChildState(Movable): self.closed = True def read_line(mut self, max_bytes: Int, deadline_ms: Int) raises -> String: - """Bounded line read that retains surplus after each newline. + """Bounded line read that retains surplus as undecoded bytes. The byte cap is enforced inside every read chunk (including a newline - in the same chunk), and bytes after the returned newline are retained - for the next consumer. Raises ``ready_output_overflow`` past the cap - and ``read_deadline_expired`` when the deadline elapses first. + in the same chunk) and bytes after the returned newline stay buffered + for the next consumer. ``last_terminated`` reports whether the returned + line ended with a newline. Raises ``ready_output_overflow`` past the cap, + ``read_deadline_expired`` on the read deadline, ``read_error`` for a + real read/poll failure (never conflated with EOF) and ``invalid_utf8`` + for a line that is not valid UTF-8. """ - var out = List[UInt8]() + var line = List[UInt8]() var start = now_ms() + self.last_terminated = False while True: - var nl = self.pending.find("\n") - if nl >= 0: - var line = String(self.pending[byte=0:nl]) - var rest = String(self.pending[byte = nl + 1 :]) - if rest.byte_length() > max_bytes: - raise Error("ready_output_overflow") + var found = -1 + for index in range(len(self.pending)): + if Int(self.pending[index]) == 10: + found = index + break + if found >= 0: + for index in range(found): + line.append(self.pending[index]) + var rest = List[UInt8]() + for index in range(found + 1, len(self.pending)): + rest.append(self.pending[index]) self.pending = rest^ - if line.byte_length() > max_bytes: + if len(line) > max_bytes or len(self.pending) > max_bytes: raise Error("ready_output_overflow") - return line^ - for byte in self.pending.as_bytes(): - out.append(UInt8(Int(byte))) - self.pending = "" - if len(out) > max_bytes: + self.last_terminated = True + return _utf8_line(line^) + for index in range(len(self.pending)): + line.append(self.pending[index]) + self.pending = List[UInt8]() + if len(line) > max_bytes: raise Error("ready_output_overflow") if self.eof: - return bytes_to_string(out) + return _utf8_line(line^) if now_ms() - start >= deadline_ms: raise Error("read_deadline_expired") var ev = poll_fd(self.report_fd, POLLIN, LIFECYCLE_POLL_SLICE_MS) + if ev < 0: + raise Error("read_error") if ev == 0: continue var buf = InlineArray[Byte, 512](fill=0) var n = read_fd(self.report_fd, buf.unsafe_ptr(), 512) - if n <= 0: + if n < 0: + raise Error("read_error") + if n == 0: self.eof = True continue - var newline_at = -1 for index in range(n): - if Int(buf[index]) == 10: - newline_at = index - break - if newline_at < 0: - for index in range(n): - out.append(UInt8(Int(buf[index]))) - if len(out) > max_bytes: - raise Error("ready_output_overflow") + self.pending.append(UInt8(Int(buf[index]))) + + def drain_surplus(mut self, max_bytes: Int, deadline_ms: Int) raises -> Int: + """Consume every byte after the last returned line, through EOF. + + A real report is exactly one newline-terminated line, so any surplus + byte is a duplicate or trailing report regardless of chunk alignment. + A read/poll failure is a distinct cause and never a clean end of + stream. + """ + var total = len(self.pending) + self.pending = List[UInt8]() + if total > max_bytes: + raise Error("ready_output_overflow") + var start = now_ms() + while not self.eof: + if now_ms() - start >= deadline_ms: + raise Error("read_deadline_expired") + var ev = poll_fd(self.report_fd, POLLIN, LIFECYCLE_POLL_SLICE_MS) + if ev < 0: + raise Error("read_error") + if ev == 0: continue - for index in range(newline_at): - out.append(UInt8(Int(buf[index]))) - if len(out) > max_bytes: - raise Error("ready_output_overflow") - var rest = List[UInt8]() - for index in range(newline_at + 1, n): - rest.append(UInt8(Int(buf[index]))) - var surplus = bytes_to_string(rest) - if surplus.byte_length() > max_bytes: + var buf = InlineArray[Byte, 1024](fill=0) + var n = read_fd(self.report_fd, buf.unsafe_ptr(), 1024) + if n < 0: + raise Error("read_error") + if n == 0: + self.eof = True + continue + total += n + if total > max_bytes: raise Error("ready_output_overflow") - self.pending = surplus^ - return bytes_to_string(out) + return total + + +def piped_child_state( + pid: Int, + report_fd: Int, + deadline_ms: Int, + expected_requests: Int, + ledger: CleanupLedger, +) -> PipedChildState: + """Build one owned-child lifecycle record with explicit ownership truth.""" + return PipedChildState( + pid=pid, + report_fd=report_fd, + pending=List[UInt8](), + eof=False, + closed=False, + deadline_ms=deadline_ms, + expected_requests=expected_requests, + reaped=False, + ok=False, + phase="pending", + case_label="-", + reason="not_reaped", + requests=0, + connections=0, + cleanup_error="", + status=ProcessStatus("pending", False, -1, 0, 0, ""), + observed=ProcessStatus("pending", False, -1, 0, 0, ""), + observed_valid=False, + last_terminated=False, + spawn_ms=now_ms(), + ledger=ledger.copy(), + ) def parse_ready_line(line: String, max_bytes: Int) raises -> Int: @@ -791,29 +922,55 @@ def descriptor_census(limit: Int) -> Int: return -1 var count = 0 for fd in range(0, limit): - if Int(external_call["fcntl", c_int](c_int(fd), c_int(F_GETFD))) >= 0: - count += 1 + while True: + var rc = Int( + external_call["fcntl", c_int](c_int(fd), c_int(F_GETFD)) + ) + if rc >= 0: + count += 1 + break + if get_errno() == ErrNo.EINTR: + continue + break if count == 0: return -1 return count def fd_scan_limit() -> Int: + """Return the OS descriptor-table size, or -1 when unavailable. + + The raw size is returned; callers decide whether complete coverage inside + the admitted range is possible rather than silently truncating the scan. + """ var n = Int(external_call["getdtablesize", c_int]()) if n <= 0: return -1 - if n > CENSUS_MAX_FDS: - n = CENSUS_MAX_FDS return n def open_fd_count() -> Int: """Numeric open-descriptor census for this process (-1 if unavailable).""" - return descriptor_census(fd_scan_limit()) + var limit = fd_scan_limit() + if limit <= 0 or limit > CENSUS_MAX_FDS: + return -1 + return descriptor_census(limit) -def open_fd_count_checked() raises -> Int: - var count = open_fd_count() +def open_fd_count_checked(limit: Int = -1) raises -> Int: + """Checked census that fails explicitly when coverage cannot be complete. + + ``limit`` defaults to the OS descriptor-table size. A nonpositive, + oversized or otherwise unavailable census raises + ``descriptor_census_unavailable`` rather than returning a partial count + that a caller could read as a pass. + """ + var effective = limit + if effective < 0: + effective = fd_scan_limit() + if effective <= 0 or effective > CENSUS_MAX_FDS: + raise Error("descriptor_census_unavailable") + var count = descriptor_census(effective) if count < 0: raise Error("descriptor_census_unavailable") return count diff --git a/tests/stdio_no_read_entrypoint.mojo b/tests/stdio_no_read_entrypoint.mojo @@ -0,0 +1,10 @@ +"""Test-only stdio entrypoint that ignores stdin, emits valid JSON and exits 0. + +Exercises PC03: the parent must reject a run whose intended request bytes were +never delivered even though the child's stdout is valid JSON and its exit code +is zero. +""" + + +def main(): + print('{"ok":true}') diff --git a/tests/stdio_process_helper.mojo b/tests/stdio_process_helper.mojo @@ -256,6 +256,10 @@ def run_stdio_entrypoint_with_deadline( continue if not stdin_done: if (pr.r0 & (POLLERR | POLLHUP | POLLNVAL)) != 0: + # An early peer close or error is a cause-specific failure even + # when the child later exits 0 with valid stdout: the intended + # request bytes were not delivered. + write_reason = "write_pipe_closed" stdin_done = True elif (pr.r0 & POLLOUT) != 0: var cw = write_fd_chunk(stdin_write_fd, request, sent) @@ -285,6 +289,10 @@ def run_stdio_entrypoint_with_deadline( read_reason = "stderr_" + d.reason break + if write_reason == "" and sent < request.byte_length(): + # The loop only ends with stdin finished; guard any path that would + # otherwise leave intended request bytes unwritten. + write_reason = "write_incomplete" close_fd(stdin_write_fd) if write_reason != "": _terminate_and_raise( @@ -299,8 +307,11 @@ def run_stdio_entrypoint_with_deadline( pid, -1, stdout_read_fd, stderr_read_fd, read_reason ) - var st = wait_bounded(pid, TERMINATION_GRACE_MS) - if not st.reaped(): + var remaining = budget - (now_ms() - start) + if remaining < 1: + remaining = 1 + var st = wait_bounded(pid, remaining) + if not st.cleanup_proved(): _terminate_and_raise(pid, -1, stdout_read_fd, stderr_read_fd, "timeout") close_fd(stdout_read_fd) close_fd(stderr_read_fd) diff --git a/tests/strict_fixture.mojo b/tests/strict_fixture.mojo @@ -427,20 +427,27 @@ struct ConnectionReader(Movable): """Bounded completion handshake after the final expected exchange. Any already-buffered or subsequently received bytes mean the client - 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. + sent an extra exchange. A bounded timeout with no bytes is success; a + real read error is propagated as the exact read cause (never read as + success or as a timeout) and is recorded by the serve loop as a + bounded io_error. The Mojo error model cannot discriminate these + struct payloads by handler type, so the timeout is identified from the + rendered cause. """ if len(self._buffer) > 0: return "extra_exchange_after_completion" self._stream.set_recv_timeout(grace_ms) + var outcome = "" try: var n = self._read_more() if n > 0: - return "extra_exchange_after_completion" - except Timeout: - return "" - return "" + outcome = "extra_exchange_after_completion" + except e: + var text = String(e) + if text.startswith("Timeout"): + return "" + raise Error("probe_completion_read_error:" + text) + return outcome^ # ── Scripted exchanges (FX01/FX02/FX04/FX05) ──────────────────────────────── @@ -812,6 +819,8 @@ def parse_report(line: String) -> ServeReport: return _report_failure("malformed_field") var key = String(field[byte=0:eq]) var value = String(field[byte = eq + 1 :]) + if value.byte_length() == 0: + return _report_failure("empty_field") if key == "phase": if have_phase: return _report_failure("duplicate_field") @@ -851,6 +860,10 @@ def parse_report(line: String) -> ServeReport: and have_connections ): return _report_failure("missing_field") + if status == "ok" and (phase != "complete" or reason != "ok"): + # A claimed success must carry the emitted success phase/reason; any + # other pairing is an inconsistent report, never a success. + return _report_failure("inconsistent_status") return ServeReport( status == "ok", phase, case_label, reason, requests, connections ) diff --git a/tests/test_provider_helpers.mojo b/tests/test_provider_helpers.mojo @@ -13,6 +13,8 @@ from flare.net.socket import RawSocket from flare.tcp import TcpListener, TcpStream from parent_lifecycle import ( + CENSUS_MAX_FDS, + CleanupLedger, PipedChildState, ProcessStatus, child_exit, @@ -25,17 +27,20 @@ from parent_lifecycle import ( fork_pid, make_pipe, make_three_pipes, + now_ms, open_fd_count, open_fd_count_checked, parse_ready_line, parse_ready_or_cleanup, pid_not_waitable, + piped_child_state, read_all_bounded, read_line_bounded, sleep_ms, wait_nohang, write_fd_bounded, write_raw, + write_raw_bytes, ) from strict_fixture import ( ConnectionReader, @@ -55,6 +60,7 @@ from max_local_process_helper import ( spawn_max_local_stub, ) from jev_provider_helper import ( + SpawnedJevStub, spawn_jev_scripted_auto, spawn_jev_stub_auto, ) @@ -938,29 +944,29 @@ def _owned_report_child( _ = write_raw(1, report) child_exit(exit_code) close_fd(pipe.write_fd) - var state = PipedChildState( - pid=pid, - report_fd=pipe.read_fd, - pending="", - eof=False, - closed=False, - deadline_ms=2000, - expected_requests=1, - reaped=False, - ok=False, - phase="pending", - case_label="-", - reason="not_reaped", - requests=0, - connections=0, - cleanup_error="", - status=ProcessStatus("pending", False, -1, 0, 0, ""), - observed=ProcessStatus("pending", False, -1, 0, 0, ""), - observed_valid=False, - ) + var state = piped_child_state(pid, pipe.read_fd, 2000, 1, CleanupLedger()) return SpawnedMaxLocalStub(pid, 0, state^) +def _owned_jev_report_child( + exit_code: Int, report: String +) raises -> SpawnedJevStub: + """Same controlled report child, reaped through the Jev provider path.""" + var pipe = make_pipe() + var pid = fork_pid() + if pid == 0: + if dup2_fd(pipe.write_fd, 1) < 0: + child_exit(126) + close_fd(pipe.read_fd) + close_fd(pipe.write_fd) + if report != "": + _ = write_raw(1, report) + child_exit(exit_code) + close_fd(pipe.write_fd) + var state = piped_child_state(pid, pipe.read_fd, 2000, 1, CleanupLedger()) + return SpawnedJevStub(pid, state^) + + # ── LC01: automatic scope ownership ───────────────────────────────────────── @@ -1161,25 +1167,12 @@ def test_coalesced_ready_and_report_lines_retain_surplus() raises: " connections=1\n" ), ) - var state = PipedChildState( + var state = piped_child_state( pid=0, report_fd=pipe.read_fd, - pending="", - eof=False, - closed=False, deadline_ms=500, expected_requests=1, - reaped=False, - ok=False, - phase="pending", - case_label="-", - reason="not_reaped", - requests=0, - connections=0, - cleanup_error="", - status=ProcessStatus("pending", False, -1, 0, 0, ""), - observed=ProcessStatus("pending", False, -1, 0, 0, ""), - observed_valid=False, + ledger=CleanupLedger(), ) var ready = state.read_line(2048, 500) var report = state.read_line(2048, 500) @@ -1308,31 +1301,27 @@ def test_result_truth_rejects_duplicate_report_line() raises: def test_cleanup_failure_is_observable() raises: - # LC01/D36: cleanup failure must be observable, never silently swallowed. - var state = PipedChildState( + # LC01/PC02: cleanup failure must be observable and must not claim the + # child was collected or discard retryable ownership. + var recorded = List[String]() + var ledger = CleanupLedger(UnsafePointer(to=recorded)) + var state = piped_child_state( pid=0, report_fd=-1, - pending="", - eof=False, - closed=False, deadline_ms=100, expected_requests=1, - reaped=False, - ok=False, - phase="pending", - case_label="-", - reason="not_reaped", - requests=0, - connections=0, - cleanup_error="", - status=ProcessStatus("pending", False, -1, 0, 0, ""), - observed=ProcessStatus("pending", False, -1, 0, 0, ""), - observed_valid=False, + ledger=ledger, ) var stub = SpawnedMaxLocalStub(0, 0, state^) stub.cleanup() assert_true(stub.cleanup_error().find("unreaped") >= 0) assert_true(not stub.status().cleanup_proved()) + assert_equal(len(recorded), 1) + assert_true(recorded[0].find("unproved") >= 0) + # A second cleanup still retries the same owned identity rather than + # short-circuiting on a false "reaped" flag. + stub.cleanup() + assert_equal(len(recorded), 2) def test_result_truth_rejects_wrong_request_count() raises: @@ -1355,21 +1344,57 @@ def test_result_truth_rejects_invalid_connection_count() raises: assert_equal(stub.reason(), "connection_count_invalid") -def test_completion_probe_error_is_not_success() raises: - # LC05: a non-timeout completion-probe I/O/setup error must not be read as - # a successful completion. - var pipe = make_pipe() - var sock = RawSocket(c_int(pipe.read_fd), c_int(2), c_int(1), True) +def test_completion_probe_read_error_is_distinct_from_setup_and_timeout() raises: + # LC05/PC04: a real socket read error after a successful timeout setup must + # propagate as the exact read cause, not be read as success, a setup + # failure or a timeout. + # + # 1. Successful setup on a real (unconnected) socket, then a real read + # error (ENOTCONN): the probe must raise and never return success. + var sock = RawSocket(c_int(2), c_int(1)) var stream = TcpStream(sock^, SocketAddr.localhost(UInt16(1))) - var reader = ConnectionReader(stream^) - var raised = False + var setup_ok = False try: - _ = reader.probe_completion(20) + stream.set_recv_timeout(20) + setup_ok = True except e: - raised = True _ = String(e) + assert_true(setup_ok) + var reader = ConnectionReader(stream^) + var read_message = "" + var read_result = "" + try: + read_result = reader.probe_completion(20) + except e: + read_message = String(e) + assert_true(read_result == "") + assert_true(read_message.find("recv") >= 0) + assert_true(read_message.find("timeout") < 0) + + # 2. A non-socket descriptor fails at setup, a distinct cause. + var pipe = make_pipe() + var pipe_sock = RawSocket(c_int(pipe.read_fd), c_int(2), c_int(1), True) + var pipe_stream = TcpStream(pipe_sock^, SocketAddr.localhost(UInt16(1))) + var pipe_reader = ConnectionReader(pipe_stream^) + var setup_failure = "" + try: + _ = pipe_reader.probe_completion(20) + except e: + setup_failure = String(e) close_fd(pipe.write_fd) - assert_true(raised) + assert_true(setup_failure != "") + assert_true(setup_failure.find("setsockopt") >= 0) + + # 3. A quiet connected socket times out, which is the probe's success path + # and stays distinct from the read error above. + var listener = TcpListener.bind(SocketAddr.localhost(0)) + var port = Int(listener.local_addr().port) + var client = TcpStream.connect(SocketAddr.localhost(UInt16(port))) + var server = listener.accept() + var server_reader = ConnectionReader(server^) + var quiet = server_reader.probe_completion(20) + client.close() + assert_equal(quiet, "") def test_descriptor_census_detects_planted_socket() raises: @@ -1493,5 +1518,190 @@ def test_coalesced_large_body_does_not_charge_header_cap() raises: stub.wait() +# ── PC01/PC02/PC04: complete report truth and retained ownership ───────────── + + +comptime VALID_REPORT = ( + "result ok phase=complete case=- reason=ok requests=1 connections=1\n" +) + + +def _report_controls( + report: String, exit_code: Int, expected_reason: String +) raises: + """Every report-framing control must fail for its cause on BOTH providers. + """ + var stub = _owned_report_child(exit_code, report) + var owned = stub.pid + stub.reap() + assert_true(not stub.ok()) + assert_equal(stub.reason(), expected_reason) + assert_true(pid_not_waitable(owned)) + var jev_stub = _owned_jev_report_child(exit_code, report) + var jev_owned = jev_stub.pid + jev_stub.reap() + assert_true(not jev_stub.ok()) + assert_equal(jev_stub.reason(), expected_reason) + assert_true(pid_not_waitable(jev_owned)) + + +def test_report_stream_controls_both_providers() raises: + # PC01: the complete bounded report stream is validated through EOF; a + # positive control and the aligned/split/coalesced duplicate, no-LF, + # malformed and inconsistent-field controls all execute on both providers. + var stub = _owned_report_child(0, VALID_REPORT) + stub.reap() + assert_true(stub.ok()) + assert_equal(stub.phase(), "complete") + assert_equal(stub.request_count(), 1) + var jev_stub = _owned_jev_report_child(0, VALID_REPORT) + jev_stub.reap() + assert_true(jev_stub.ok()) + assert_equal(jev_stub.request_count(), 1) + + var first = "result ok phase=complete case=" + var tail = " reason=ok requests=1 connections=1\n" + while first.byte_length() + tail.byte_length() < 512: + first += "x" + first += tail + assert_equal(first.byte_length(), 512) + _report_controls(first + VALID_REPORT, 0, "duplicate_report") + _report_controls(VALID_REPORT + VALID_REPORT, 0, "duplicate_report") + var unterminated = String( + VALID_REPORT[byte = 0 : VALID_REPORT.byte_length() - 1] + ) + _report_controls(unterminated, 0, "unterminated_report") + _report_controls( + ( + "result ok phase=read case=- reason=io_error requests=1" + " connections=1\n" + ), + 0, + "inconsistent_status", + ) + _report_controls( + "result ok phase= case=- reason=ok requests=1 connections=1\n", + 0, + "empty_field", + ) + _report_controls( + "result ok phase=complete case=- reason=ok requests=1\n", + 0, + "missing_field", + ) + + +def test_descriptor_read_error_is_distinct_from_eof() raises: + # PC01/PC03: an unavailable descriptor is a bounded read error, never an + # EOF/empty success, for both the shared reader and the owned-child state. + var pipe = make_pipe() + var closed_fd = pipe.read_fd + close_fd(pipe.read_fd) + close_fd(pipe.write_fd) + var line_message = "" + try: + _ = read_line_bounded(closed_fd, 64, 200) + except e: + line_message = String(e) + assert_equal(line_message, "read_error") + var state = piped_child_state(closed_fd, closed_fd, 200, 1, CleanupLedger()) + var state_message = "" + try: + _ = state.read_line(64, 200) + except e: + state_message = String(e) + assert_equal(state_message, "read_error") + assert_true(not state.last_terminated) + + +def test_multibyte_surplus_is_not_decoded_prematurely() raises: + # PC01/LC03: a chunk that splits a multi-byte character after a newline must + # be retained as bytes instead of raising a premature UTF-8 decode error. + var pipe = make_pipe() + _ = write_raw(pipe.write_fd, "ready 4242\n") + var lead = List[UInt8]() + lead.append(UInt8(0xC3)) + _ = write_raw_bytes(pipe.write_fd, lead) + var state = piped_child_state( + pid=0, + report_fd=pipe.read_fd, + deadline_ms=500, + expected_requests=1, + ledger=CleanupLedger(), + ) + var ready = state.read_line(64, 500) + assert_equal(ready, "ready 4242") + assert_true(state.last_terminated) + var trail = List[UInt8]() + trail.append(UInt8(0xA9)) + trail.append(UInt8(10)) + _ = write_raw_bytes(pipe.write_fd, trail) + var letter = state.read_line(64, 500) + state.close_reader() + close_fd(pipe.write_fd) + assert_equal(letter, "\u00e9") + assert_true(state.last_terminated) + + +def test_cleanup_failure_preserves_retryable_ownership() raises: + # PC02: a controlled wait failure on a real owned child must not mark it + # reaped or discard ownership; the retry with the restored exact identity + # still collects it, and the failure stays recorded in the caller ledger. + var recorded = List[String]() + var ledger = CleanupLedger(UnsafePointer(to=recorded)) + var stub = spawn_max_local_stub(0, "count_requests", 1, 2000, ledger) + var actual = stub.pid + stub.pid = 0 + stub.cleanup() + assert_true(stub.cleanup_error().find("unreaped") >= 0) + assert_equal(len(recorded), 1) + stub.pid = actual + stub.cleanup() + assert_true(stub.status().cleanup_proved()) + assert_true(pid_not_waitable(actual)) + + var jev_stub = spawn_jev_stub_auto("ok", 1, 2000, ledger) + var jev_actual = jev_stub.stub.pid + jev_stub.stub.pid = 0 + jev_stub.stub.cleanup() + assert_true(jev_stub.stub.cleanup_error().find("unreaped") >= 0) + assert_equal(len(recorded), 2) + jev_stub.stub.pid = jev_actual + jev_stub.stub.cleanup() + assert_true(jev_stub.stub.status().cleanup_proved()) + assert_true(pid_not_waitable(jev_actual)) + + +def test_cleanup_failure_survives_scope_exit() raises: + # PC02: cleanup failure must remain observable after the owning handle is + # destroyed at scope exit, not merely stored in an inaccessible object. + var recorded = List[String]() + var ledger = CleanupLedger(UnsafePointer(to=recorded)) + var state = piped_child_state(0, -1, 100, 1, ledger) + with SpawnedMaxLocalStub(0, 0, state^) as holder: + _ = holder + assert_equal(len(recorded), 1) + assert_true(recorded[0].find("unproved") >= 0) + + +def test_descriptor_census_unavailable_propagates() raises: + # PC04: an unavailable or incomplete-range census must fail explicitly + # through the checked caller rather than look like a small passing count. + assert_equal(descriptor_census(0), -1) + var zero_reason = "" + try: + _ = open_fd_count_checked(0) + except e: + zero_reason = String(e) + assert_equal(zero_reason, "descriptor_census_unavailable") + var ceiling_reason = "" + try: + _ = open_fd_count_checked(CENSUS_MAX_FDS + 1) + except e: + ceiling_reason = String(e) + assert_equal(ceiling_reason, "descriptor_census_unavailable") + assert_true(open_fd_count_checked() > 0) + + def main() raises: TestSuite.discover_tests[__functions_in_module()]().run() diff --git a/tests/test_repo_local_process_contract.mojo b/tests/test_repo_local_process_contract.mojo @@ -105,6 +105,25 @@ def test_run_stdio_entrypoint_reaps_stalled_child_under_deadline() raises: message = String(e) assert_true(message.find("stdio-entrypoint") >= 0) assert_true(message.find("signal=") >= 0 or message.find("exited=") >= 0) + assert_true(message.find("cleanup_error=") >= 0) + + +def test_run_stdio_entrypoint_rejects_unread_request() raises: + # PC03: an incomplete request write is a cause-specific failure even though + # the child emits valid JSON and exits zero, and the failure still exposes + # the owned child's cleanup truth. + var request = String("") + for _ in range(150000): + request += "r" + var message = "" + try: + _ = run_stdio_entrypoint_with_deadline( + "tests/stdio_no_read_entrypoint.mojo", request, "", "", 10000 + ) + except e: + message = String(e) + assert_true(message.find("write_") >= 0) + assert_true(message.find("cleanup_error=") >= 0) def test_run_stdio_entrypoint_classifies_loader_failure() raises: