hyf

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

commit da09203356ebfefa27bed5bd4bcffd1a20ecd32c
parent e323906b602c5ac878e4136b4cb35a271a8effc3
Author: triesap <tyson@radroots.org>
Date:   Wed, 23 Sep 2026 13:46:12 +0000

test(hyf): enforce caller-visible cleanup truth and one work budget (C002D)

- Replace the optional inert cleanup ledger with a required caller-held CleanupGuard and integrate every supported spawner and caller, so an unproved cleanup fails the owning test while body/assertion causes stay preserved.
- Add a bounded test-only LifecycleFaults seam against real owned children and retain close-once recoverable ownership; a transient wait error stays retryable, unexpected nonterminal status fails closed and terminal status is cached before further wait/signal.
- Give wait, report line and EOF drain one spawn-relative work budget with deadline-bounded EINTR retries, and prove late-report rejection, within-budget completion and a phase-isolated stdio late-exit control with separate work/cleanup timing.
- Classify descriptor-census lookup errors (EBADF closed, EINTR bounded retry, other unavailable) through the checked caller, keeping ordinary, high-FD, socket and no-leak controls.

Diffstat:
Mtests/jev_provider_helper.mojo | 221++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------------------
Mtests/max_local_process_helper.mojo | 206+++++++++++++++++++++++++++++++++++++++++++++++++++++++------------------------
Mtests/parent_lifecycle.mojo | 570++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------
Mtests/stdio_process_helper.mojo | 226+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------
Mtests/test_jev.mojo | 38+++++++++++++++++++++++++++++++-------
Mtests/test_provider_adapter.mojo | 13+++++++++++--
Mtests/test_provider_helpers.mojo | 797+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------
Mtests/test_repo_local_process_contract.mojo | 101++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----
Mtests/test_stdio_contract.mojo | 48+++++++++++++++++++++++++++++++++++++++++-------
9 files changed, 1829 insertions(+), 391 deletions(-)

diff --git a/tests/jev_provider_helper.mojo b/tests/jev_provider_helper.mojo @@ -14,7 +14,7 @@ from flare.utils import usleep from parent_lifecycle import ( FIXTURE_DEFAULT_DEADLINE_MS, TERMINATION_GRACE_MS, - CleanupLedger, + CleanupGuard, PipedChildState, PipeFds, ProcessStatus, @@ -24,13 +24,10 @@ from parent_lifecycle import ( finalize_owned_failure, fork_owned_or_close, make_pipe, - now_ms, parse_ready_or_cleanup, piped_child_state, set_alarm, - terminate_owned, - wait_bounded, - wait_nohang, + sleep_ms, write_raw, ) from strict_fixture import ( @@ -284,21 +281,26 @@ struct SpawnedJevStub(Movable): """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. + uncertain wait keeps the handle retryable, records the failure in the + required caller-held guard and retains the exact pid/report descriptor + for recovery, instead of being silently marked complete. Never raises, + so a body/assertion cause is preserved at scope exit. """ if self.state.reaped: return - var status = terminate_owned(self.pid, TERMINATION_GRACE_MS) + var report_fd = self.state.report_fd + var status = self.state.terminate_once(self.pid, TERMINATION_GRACE_MS) self.state.status = status.copy() if status.cleanup_proved(): self.state.reaped = True self.state.close_reader() + self.state.guard[].resolve_pid(self.pid) return - self.state.cleanup_error = "unreaped:" + status.describe() - self.state.ledger.record( - "owned-child cleanup unproved " + status.describe() + self.state.record_unproved( + self.pid, + report_fd, + "owned-child cleanup unproved", + "unreaped:" + status.describe(), ) def ok(self) -> Bool: @@ -337,43 +339,79 @@ struct SpawnedJevStub(Movable): ) def status(mut self) -> ProcessStatus: - """Observe child status without losing ownership or report truth.""" + """Observe child status without losing ownership or report truth. + + A terminal observation is cached so a later ``reap`` never performs a + new wait on a stale/reused identity. A transient ``wait_error`` or any + other nonterminal result is returned but deliberately *not* cached, so + it stays retryable and cannot overwrite a valid terminal result. + """ if self.state.reaped: return self.state.status.copy() 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 + var st = self.state.wait_once(self.pid) + if st.state == "reaped" or st.state == "gone": + 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.""" + """Reap the owned child within one declared finite work budget. + + Never raises (so ``__exit__`` cannot mask a body error). Startup, wait, + report line and EOF drain share ``deadline_ms`` measured from spawn: an + expired budget fails with ``work_budget_expired`` and only the bounded + cleanup allowance, and no fresh successful interval is granted. + """ if self.state.reaped: return - 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.faults.wait_delay_ms > 0: + sleep_ms(self.state.faults.wait_delay_ms) + self.state.faults.wait_delay_ms = 0 + var remaining = self.state.work_remaining_ms() + if remaining <= 0: + self.state.store( + False, "watchdog", "-", "work_budget_expired", 0, 0 + ) + var expired = self.state.terminate_once( + self.pid, TERMINATION_GRACE_MS + ) + self.state.status = expired.copy() + if expired.cleanup_proved(): + self.state.reaped = True + self.state.close_reader() + self.state.guard[].resolve_pid(self.pid) + else: + self.state.reason = "work_budget_expired_unreaped" + self.state.record_unproved( + self.pid, + self.state.report_fd, + "owned-child cleanup unproved", + "unreaped:" + expired.describe(), + ) + return + var status = ProcessStatus("pending", False, -1, 0, 0, "") if self.state.observed_valid: status = self.state.observed.copy() + else: + status = self.state.wait_until(self.pid, remaining) self.state.status = status.copy() if status.state == "running" or status.state == "interrupted": - var term = terminate_owned(self.pid, TERMINATION_GRACE_MS) + var term = self.state.terminate_once(self.pid, TERMINATION_GRACE_MS) self.state.status = term.copy() self.state.store(False, "watchdog", "-", "timeout", 0, 0) if term.cleanup_proved(): self.state.reaped = True self.state.close_reader() + self.state.guard[].resolve_pid(self.pid) else: - self.state.cleanup_error = "unreaped:" + term.describe() self.state.reason = "timeout_unreaped" - self.state.ledger.record( - "owned-child cleanup unproved " + term.describe() + self.state.record_unproved( + self.pid, + self.state.report_fd, + "owned-child cleanup unproved", + "unreaped:" + term.describe(), ) return if status.state == "gone": @@ -382,32 +420,56 @@ struct SpawnedJevStub(Movable): self.state.store(False, "watchdog", "-", "gone", 0, 0) self.state.reaped = True self.state.close_reader() + self.state.guard[].resolve_pid(self.pid) return if status.state == "wait_error": # Identity/ownership is unproved: retain it for a retry and record # the uncertainty instead of claiming the child was collected. self.state.store(False, "watchdog", "-", "wait_error", 0, 0) - self.state.cleanup_error = "wait_error:" + status.error - self.state.ledger.record( - "owned-child wait unproved " + status.describe() + self.state.record_unproved( + self.pid, + self.state.report_fd, + "owned-child wait unproved", + "wait_error:" + status.error, + ) + return + if status.state != "reaped": + # Unexpected nonterminal/unknown taxonomy: fail closed and keep the + # exact ownership so a later retry can still collect it. + self.state.store(False, "watchdog", "-", "unexpected_status", 0, 0) + self.state.record_unproved( + self.pid, + self.state.report_fd, + "owned-child unexpected status", + "unexpected:" + status.describe(), ) return var report_text = "" 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 + var report_remaining = self.state.work_remaining_ms() + if report_remaining <= 0: + report_error = "report_budget_expired" + else: + try: + report_text = self.state.read_line( + STRICT_MAX_REPORT_BYTES, report_remaining ) - 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: - report_error = String(e) + if self.state.last_terminated: + var drain_remaining = self.state.work_remaining_ms() + if drain_remaining <= 0: + report_error = "drain_deadline_expired" + else: + var surplus = self.state.drain_surplus( + STRICT_MAX_REPORT_BYTES, drain_remaining + ) + 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: + report_error = String(e) self.state.close_reader() if report_text == "" and report_error == "": if status.exited and status.exit_code == 0: @@ -431,15 +493,18 @@ struct SpawnedJevStub(Movable): 0, ) self.state.reaped = True + self.state.guard[].resolve_pid(self.pid) return if report_error != "": self.state.store(False, "parse", "-", report_error, -1, -1) self.state.reaped = True + self.state.guard[].resolve_pid(self.pid) return var parsed = parse_report(report_text) if parsed.phase == "parse": self.state.store(False, "parse", "-", parsed.reason, -1, -1) self.state.reaped = True + self.state.guard[].resolve_pid(self.pid) return self.state.store( parsed.ok, @@ -467,6 +532,7 @@ struct SpawnedJevStub(Movable): self.state.phase = "accounting" self.state.reason = "connection_count_invalid" self.state.reaped = True + self.state.guard[].resolve_pid(self.pid) def wait(mut self) raises: self.reap() @@ -476,15 +542,19 @@ struct SpawnedJevStub(Movable): def terminate(mut self) raises: if self.state.reaped: return - var status = terminate_owned(self.pid, TERMINATION_GRACE_MS) + var report_fd = self.state.report_fd + var status = self.state.terminate_once(self.pid, TERMINATION_GRACE_MS) self.state.status = status.copy() if status.cleanup_proved(): self.state.reaped = True self.state.close_reader() + self.state.guard[].resolve_pid(self.pid) return - self.state.cleanup_error = "unreaped:" + status.describe() - self.state.ledger.record( - "owned-child cleanup unproved " + status.describe() + self.state.record_unproved( + self.pid, + report_fd, + "owned-child cleanup unproved", + "unreaped:" + status.describe(), ) raise Error("lifecycle: owned child not reaped: " + status.describe()) @@ -537,6 +607,27 @@ struct SpawnedJevStubView(Movable): def terminate(mut self) raises: self.target[].terminate() + def inject_cleanup_failure(mut self): + """Test-only bounded fault: force one unproved cleanup attempt. + + The exact-owned child keeps running, so the follow-up retry exercises + the real recoverable ownership path rather than a synthetic identity. + """ + self.target[].state.faults.cleanup_failures += 1 + + def inject_wait_error(mut self): + """Test-only bounded fault: force one transient wait error.""" + self.target[].state.faults.wait_errors += 1 + + def inject_nonterminal_status(mut self): + """Test-only bounded fault: force one unexpected nonterminal status.""" + self.target[].state.faults.nonterminal += 1 + + def inject_wait_delay_ms(mut self, ms: Int): + """Test-only bounded fault: consume work-budget time in the wait phase. + """ + self.target[].state.faults.wait_delay_ms = ms + @fieldwise_init struct SpawnedJevStubAuto(Movable): @@ -574,10 +665,10 @@ def reserve_jev_port() raises -> Int: def spawn_jev_stub_auto( mode: String, requests: Int, + mut guard: CleanupGuard, deadline_ms: Int = FIXTURE_DEFAULT_DEADLINE_MS, - ledger: CleanupLedger = CleanupLedger(), ) raises -> SpawnedJevStubAuto: - var stub = _spawn_jev_stub(0, mode, requests, deadline_ms, ledger) + var stub = _spawn_jev_stub(0, mode, requests, deadline_ms, guard) return SpawnedJevStubAuto(port=stub.port, stub=stub^) @@ -585,10 +676,10 @@ def spawn_jev_stub( port: Int, mode: String, requests: Int, + mut guard: CleanupGuard, deadline_ms: Int = FIXTURE_DEFAULT_DEADLINE_MS, - ledger: CleanupLedger = CleanupLedger(), ) raises -> SpawnedJevStub: - return _spawn_jev_stub(port, mode, requests, deadline_ms, ledger) + return _spawn_jev_stub(port, mode, requests, deadline_ms, guard) def serve_jev_scripted( @@ -602,10 +693,10 @@ def serve_jev_scripted( def spawn_jev_scripted_auto( var scripts: List[ExchangeScript], + mut guard: CleanupGuard, deadline_ms: Int = FIXTURE_DEFAULT_DEADLINE_MS, - ledger: CleanupLedger = CleanupLedger(), ) raises -> SpawnedJevStubAuto: - return _spawn_jev_scripted(0, scripts^, deadline_ms, ledger) + return _spawn_jev_scripted(0, scripts^, deadline_ms, guard) def _spawn_child_or_cleanup( @@ -614,11 +705,11 @@ def _spawn_child_or_cleanup( mode: String, deadline_ms: Int, requests: Int, - ledger: CleanupLedger, + guard: UnsafePointer[CleanupGuard, MutAnyOrigin], ) raises -> SpawnedJevStub: """Build the owned state, read exact readiness, or clean up and raise.""" var state = piped_child_state( - pid, pipe.read_fd, deadline_ms, requests, ledger + pid, pipe.read_fd, deadline_ms, requests, guard ) var ready_line = "" try: @@ -662,7 +753,7 @@ def _spawn_jev_scripted( port: Int, var scripts: List[ExchangeScript], deadline_ms: Int, - ledger: CleanupLedger = CleanupLedger(), + mut guard: CleanupGuard, ) raises -> SpawnedJevStubAuto: var total = len(scripts) var pipe = make_pipe() @@ -675,17 +766,17 @@ def _spawn_jev_scripted( _ = set_alarm(STUB_ALARM_SECONDS) try: var report = serve_jev_scripted(port, scripts^) - write_raw(1, report_line(report) + "\n") + _ = write_raw(1, report_line(report) + "\n") child_exit(0 if report.ok else 125) except: var failed = ServeReport( False, "startup", "scripted", "serve_failed", 0, 0 ) - write_raw(1, report_line(failed) + "\n") + _ = write_raw(1, report_line(failed) + "\n") child_exit(125) close_fd(pipe.write_fd) var stub = _spawn_child_or_cleanup( - pipe, pid, "scripted", deadline_ms, total, ledger + pipe, pid, "scripted", deadline_ms, total, UnsafePointer(to=guard) ) return SpawnedJevStubAuto(port=stub.port, stub=stub^) @@ -695,7 +786,7 @@ def _spawn_jev_stub( mode: String, requests: Int, deadline_ms: Int, - ledger: CleanupLedger = CleanupLedger(), + mut guard: CleanupGuard, ) raises -> SpawnedJevStub: var pipe = make_pipe() var pid = fork_owned_or_close(pipe.copy()) @@ -707,15 +798,15 @@ def _spawn_jev_stub( _ = set_alarm(STUB_ALARM_SECONDS) try: var report = serve_jev(port, mode, requests) - write_raw(1, report_line(report) + "\n") + _ = write_raw(1, report_line(report) + "\n") child_exit(0 if report.ok else 125) except: var failed = ServeReport( False, "startup", mode, "serve_failed", 0, 0 ) - write_raw(1, report_line(failed) + "\n") + _ = 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, ledger + pipe, pid, mode, deadline_ms, requests, UnsafePointer(to=guard) ) diff --git a/tests/max_local_process_helper.mojo b/tests/max_local_process_helper.mojo @@ -17,7 +17,7 @@ from flare.utils import usleep from parent_lifecycle import ( FIXTURE_DEFAULT_DEADLINE_MS, TERMINATION_GRACE_MS, - CleanupLedger, + CleanupGuard, PipedChildState, ProcessStatus, child_exit, @@ -26,13 +26,10 @@ from parent_lifecycle import ( finalize_owned_failure, fork_owned_or_close, make_pipe, - now_ms, parse_ready_or_cleanup, piped_child_state, set_alarm, - terminate_owned, - wait_bounded, - wait_nohang, + sleep_ms, write_raw, ) from strict_fixture import ( @@ -376,21 +373,26 @@ struct SpawnedMaxLocalStub(Movable): """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. + uncertain wait keeps the handle retryable, records the failure in the + required caller-held guard and retains the exact pid/report descriptor + for recovery, instead of being silently marked complete. Never raises, + so a body/assertion cause is preserved at scope exit. """ if self.state.reaped: return - var status = terminate_owned(self.pid, TERMINATION_GRACE_MS) + var report_fd = self.state.report_fd + var status = self.state.terminate_once(self.pid, TERMINATION_GRACE_MS) self.state.status = status.copy() if status.cleanup_proved(): self.state.reaped = True self.state.close_reader() + self.state.guard[].resolve_pid(self.pid) return - self.state.cleanup_error = "unreaped:" + status.describe() - self.state.ledger.record( - "owned-child cleanup unproved " + status.describe() + self.state.record_unproved( + self.pid, + report_fd, + "owned-child cleanup unproved", + "unreaped:" + status.describe(), ) def ok(self) -> Bool: @@ -431,52 +433,77 @@ struct SpawnedMaxLocalStub(Movable): 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. + A terminal observation is cached so a later ``reap`` never performs a + new wait on a stale/reused identity. A transient ``wait_error`` or any + other nonterminal result is returned but deliberately *not* cached, so + it stays retryable and cannot overwrite a valid terminal result. """ if self.state.reaped: return self.state.status.copy() 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 + var st = self.state.wait_once(self.pid) + if st.state == "reaped" or st.state == "gone": + 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. + """Reap the owned child within one declared finite work budget. - Never raises (so ``__exit__`` cannot mask a body error). Success - requires a complete strictly parsed report, matching child exit status - and exact request accounting; missing/truncated/mismatched reports and - nonzero exit or signals fail explicitly. + Never raises (so ``__exit__`` cannot mask a body error). Startup, wait, + report line and EOF drain share ``deadline_ms`` measured from spawn: an + expired budget fails with ``work_budget_expired`` and only the bounded + cleanup allowance, and no fresh successful interval is granted. """ if self.state.reaped: return - 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.faults.wait_delay_ms > 0: + sleep_ms(self.state.faults.wait_delay_ms) + self.state.faults.wait_delay_ms = 0 + var remaining = self.state.work_remaining_ms() + if remaining <= 0: + self.state.store( + False, "watchdog", "-", "work_budget_expired", 0, 0 + ) + var expired = self.state.terminate_once( + self.pid, TERMINATION_GRACE_MS + ) + self.state.status = expired.copy() + if expired.cleanup_proved(): + self.state.reaped = True + self.state.close_reader() + self.state.guard[].resolve_pid(self.pid) + else: + self.state.reason = "work_budget_expired_unreaped" + self.state.record_unproved( + self.pid, + self.state.report_fd, + "owned-child cleanup unproved", + "unreaped:" + expired.describe(), + ) + return + var status = ProcessStatus("pending", False, -1, 0, 0, "") if self.state.observed_valid: status = self.state.observed.copy() + else: + status = self.state.wait_until(self.pid, remaining) self.state.status = status.copy() if status.state == "running" or status.state == "interrupted": - var term = terminate_owned(self.pid, TERMINATION_GRACE_MS) + var term = self.state.terminate_once(self.pid, TERMINATION_GRACE_MS) self.state.status = term.copy() self.state.store(False, "watchdog", "-", "timeout", 0, 0) if term.cleanup_proved(): self.state.reaped = True self.state.close_reader() + self.state.guard[].resolve_pid(self.pid) else: - self.state.cleanup_error = "unreaped:" + term.describe() self.state.reason = "timeout_unreaped" - self.state.ledger.record( - "owned-child cleanup unproved " + term.describe() + self.state.record_unproved( + self.pid, + self.state.report_fd, + "owned-child cleanup unproved", + "unreaped:" + term.describe(), ) return if status.state == "gone": @@ -485,32 +512,56 @@ struct SpawnedMaxLocalStub(Movable): self.state.store(False, "watchdog", "-", "gone", 0, 0) self.state.reaped = True self.state.close_reader() + self.state.guard[].resolve_pid(self.pid) return if status.state == "wait_error": # Identity/ownership is unproved: retain it for a retry and record # the uncertainty instead of claiming the child was collected. self.state.store(False, "watchdog", "-", "wait_error", 0, 0) - self.state.cleanup_error = "wait_error:" + status.error - self.state.ledger.record( - "owned-child wait unproved " + status.describe() + self.state.record_unproved( + self.pid, + self.state.report_fd, + "owned-child wait unproved", + "wait_error:" + status.error, + ) + return + if status.state != "reaped": + # Unexpected nonterminal/unknown taxonomy: fail closed and keep the + # exact ownership so a later retry can still collect it. + self.state.store(False, "watchdog", "-", "unexpected_status", 0, 0) + self.state.record_unproved( + self.pid, + self.state.report_fd, + "owned-child unexpected status", + "unexpected:" + status.describe(), ) return var report_text = "" 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 + var report_remaining = self.state.work_remaining_ms() + if report_remaining <= 0: + report_error = "report_budget_expired" + else: + try: + report_text = self.state.read_line( + STRICT_MAX_REPORT_BYTES, report_remaining ) - 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: - report_error = String(e) + if self.state.last_terminated: + var drain_remaining = self.state.work_remaining_ms() + if drain_remaining <= 0: + report_error = "drain_deadline_expired" + else: + var surplus = self.state.drain_surplus( + STRICT_MAX_REPORT_BYTES, drain_remaining + ) + 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: + report_error = String(e) self.state.close_reader() if report_text == "" and report_error == "": if status.exited and status.exit_code == 0: @@ -534,15 +585,18 @@ struct SpawnedMaxLocalStub(Movable): 0, ) self.state.reaped = True + self.state.guard[].resolve_pid(self.pid) return if report_error != "": self.state.store(False, "parse", "-", report_error, -1, -1) self.state.reaped = True + self.state.guard[].resolve_pid(self.pid) return var parsed = parse_report(report_text) if parsed.phase == "parse": self.state.store(False, "parse", "-", parsed.reason, -1, -1) self.state.reaped = True + self.state.guard[].resolve_pid(self.pid) return self.state.store( parsed.ok, @@ -570,6 +624,7 @@ struct SpawnedMaxLocalStub(Movable): self.state.phase = "accounting" self.state.reason = "connection_count_invalid" self.state.reaped = True + self.state.guard[].resolve_pid(self.pid) def wait(mut self) raises: self.reap() @@ -579,15 +634,19 @@ struct SpawnedMaxLocalStub(Movable): def terminate(mut self) raises: if self.state.reaped: return - var status = terminate_owned(self.pid, TERMINATION_GRACE_MS) + var report_fd = self.state.report_fd + var status = self.state.terminate_once(self.pid, TERMINATION_GRACE_MS) self.state.status = status.copy() if status.cleanup_proved(): self.state.reaped = True self.state.close_reader() + self.state.guard[].resolve_pid(self.pid) return - self.state.cleanup_error = "unreaped:" + status.describe() - self.state.ledger.record( - "owned-child cleanup unproved " + status.describe() + self.state.record_unproved( + self.pid, + report_fd, + "owned-child cleanup unproved", + "unreaped:" + status.describe(), ) raise Error("lifecycle: owned child not reaped: " + status.describe()) @@ -645,6 +704,27 @@ struct SpawnedMaxLocalView(Movable): def terminate(mut self) raises: self.target[].terminate() + def inject_cleanup_failure(mut self): + """Test-only bounded fault: force one unproved cleanup attempt. + + The exact-owned child keeps running, so the follow-up retry exercises + the real recoverable ownership path rather than a synthetic identity. + """ + self.target[].state.faults.cleanup_failures += 1 + + def inject_wait_error(mut self): + """Test-only bounded fault: force one transient wait error.""" + self.target[].state.faults.wait_errors += 1 + + def inject_nonterminal_status(mut self): + """Test-only bounded fault: force one unexpected nonterminal status.""" + self.target[].state.faults.nonterminal += 1 + + def inject_wait_delay_ms(mut self, ms: Int): + """Test-only bounded fault: consume work-budget time in the wait phase. + """ + self.target[].state.faults.wait_delay_ms = ms + def reserve_loopback_port() raises -> Int: var listener = TcpListener.bind(SocketAddr.localhost(0)) @@ -678,23 +758,23 @@ def spawn_max_local_stub( port: Int, mode: String, requests: Int, + mut guard: CleanupGuard, 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, ledger + port, scripts^, mode, requests, False, deadline_ms, guard ) def spawn_max_local_scripted( port: Int, var scripts: List[ExchangeScript], + mut guard: CleanupGuard, deadline_ms: Int = FIXTURE_DEFAULT_DEADLINE_MS, - ledger: CleanupLedger = CleanupLedger(), ) raises -> SpawnedMaxLocalStub: return _spawn_max_local( - port, scripts^, "scripted", len(scripts), True, deadline_ms, ledger + port, scripts^, "scripted", len(scripts), True, deadline_ms, guard ) @@ -705,7 +785,7 @@ def _spawn_max_local( requests: Int, scripted: Bool, deadline_ms: Int, - ledger: CleanupLedger, + mut guard: CleanupGuard, ) raises -> SpawnedMaxLocalStub: var pipe = make_pipe() var pid = fork_owned_or_close(pipe.copy()) @@ -719,17 +799,17 @@ def _spawn_max_local( var report = _serve_max_local_for( port, scripts^, mode, requests, scripted ) - write_raw(1, report_line(report) + "\n") + _ = write_raw(1, report_line(report) + "\n") child_exit(0 if report.ok else 125) except: var failed = ServeReport( False, "startup", mode, "serve_failed", 0, 0 ) - write_raw(1, report_line(failed) + "\n") + _ = write_raw(1, report_line(failed) + "\n") child_exit(125) close_fd(pipe.write_fd) var state = piped_child_state( - pid, pipe.read_fd, deadline_ms, requests, ledger + pid, pipe.read_fd, deadline_ms, requests, UnsafePointer(to=guard) ) var ready_line = "" try: diff --git a/tests/parent_lifecycle.mojo b/tests/parent_lifecycle.mojo @@ -55,6 +55,17 @@ comptime SIG_IGN: Int = 1 comptime F_GETFD: Int = 1 comptime CENSUS_MAX_FDS: Int = 1048576 +# ``read_fd``/``_write_fd`` return this when the caller-supplied remaining +# budget expires while a syscall is retried after EINTR. It is a distinct +# cause from a real ``read(2)``/``write(2)`` error, so a bounded retry loop can +# never be reported as plain I/O failure or as success. +comptime IO_DEADLINE_EXPIRED: Int = -2 + +# ``read_fd``/``_write_fd`` test-only retry seam: a negative fault count means +# "retry forever", so the deadline branch is reached deterministically. It is +# only valid together with a non-negative deadline. +comptime IO_FAULT_EINTR_UNBOUNDED: Int = -1 + comptime FIXTURE_DEFAULT_DEADLINE_MS: Int = 20000 comptime TERMINATION_GRACE_MS: Int = 2000 comptime LIFECYCLE_POLL_SLICE_MS: Int = 25 @@ -274,22 +285,79 @@ def poll_three( ) -def read_fd(fd: Int, buf: UnsafePointer[Byte, ...], max_bytes: Int) -> Int: +def read_fd( + fd: Int, + buf: UnsafePointer[Byte, ...], + max_bytes: Int, + deadline_ms: Int = -1, + fault_eintr_count: Int = 0, +) -> Int: + """Read up to ``max_bytes`` with an EINTR retry bounded by the caller. + + Returns the byte count, or the real ``read(2)`` error value. Returns + ``IO_DEADLINE_EXPIRED`` when a retry after ``EINTR`` would outlive the + caller-supplied remaining budget, so the retry loop returns to the deadline + owner instead of spinning indefinitely. ``deadline_ms < 0`` keeps the + legacy unbounded behaviour for short diagnostic writes/reads. + + ``fault_eintr_count`` is a test-only seam: a positive count makes exactly + that many attempts behave like ``EINTR`` and ``IO_FAULT_EINTR_UNBOUNDED`` + retries forever, so the bounded-retry branch is deterministically + executable. The seam changes no host signal state and requires a deadline. + """ + var start = now_ms() + var faults = fault_eintr_count + if faults != 0 and deadline_ms < 0: + return IO_DEADLINE_EXPIRED while True: + var synthetic = faults != 0 + if faults > 0: + faults -= 1 + if synthetic: + if now_ms() - start >= deadline_ms: + return IO_DEADLINE_EXPIRED + continue var n = Int( external_call["read", c_ssize_t](fd, buf, c_size_t(max_bytes)) ) if n >= 0 or get_errno() != ErrNo.EINTR: return n - - -def _write_fd(fd: Int, ptr: UnsafePointer[UInt8, ...], n: Int) -> Int: + if deadline_ms >= 0 and now_ms() - start >= deadline_ms: + return IO_DEADLINE_EXPIRED + + +def _write_fd( + fd: Int, + ptr: UnsafePointer[UInt8, ...], + n: Int, + deadline_ms: Int = -1, + fault_eintr_count: Int = 0, +) -> Int: + """Write ``n`` bytes with an EINTR retry bounded by ``deadline_ms``. + + See ``read_fd``: ``IO_DEADLINE_EXPIRED`` is returned instead of retrying + past the caller's finite budget, and ``fault_eintr_count`` is the bounded + test-only seam for that branch. + """ + var start = now_ms() + var faults = fault_eintr_count + if faults != 0 and deadline_ms < 0: + return IO_DEADLINE_EXPIRED while True: + var synthetic = faults != 0 + if faults > 0: + faults -= 1 + if synthetic: + if now_ms() - start >= deadline_ms: + return IO_DEADLINE_EXPIRED + continue var written = Int( external_call["write", c_ssize_t](fd, ptr, c_size_t(n)) ) if written >= 0 or get_errno() != ErrNo.EINTR: return written + if deadline_ms >= 0 and now_ms() - start >= deadline_ms: + return IO_DEADLINE_EXPIRED def write_raw(fd: Int, text: String) -> Int: @@ -331,7 +399,14 @@ def write_fd_bounded(fd: Int, data: String, deadline_ms: Int) -> String: return "write_pipe_closed" var chunk = min(WRITE_CHUNK_BYTES, total - sent) var slice = data[byte = sent : sent + chunk] - var n = _write_fd(fd, slice.as_bytes().unsafe_ptr(), chunk) + var n = _write_fd( + fd, + slice.as_bytes().unsafe_ptr(), + chunk, + deadline_ms - (now_ms() - start), + ) + if n == IO_DEADLINE_EXPIRED: + return "write_deadline_expired" if n <= 0: return "write_failed" sent += n @@ -344,18 +419,33 @@ struct ChunkWrite(Movable): var written: Int -def write_fd_chunk(fd: Int, data: String, offset: Int) -> ChunkWrite: +def write_fd_chunk( + fd: Int, + data: String, + offset: Int, + deadline_ms: Int = -1, + fault_eintr_count: Int = 0, +) -> ChunkWrite: """Write one ``PIPE_BUF``-safe chunk after a readiness poll. Returns the bounded failure reason and the bytes actually written so a - caller can interleave writing with draining other descriptors. + caller can interleave writing with draining other descriptors. The write's + own EINTR retry is bounded by the caller's remaining ``deadline_ms``. """ var total = data.byte_length() if offset >= total: return ChunkWrite("", 0) var chunk = min(WRITE_CHUNK_BYTES, total - offset) var slice = data[byte = offset : offset + chunk] - var n = _write_fd(fd, slice.as_bytes().unsafe_ptr(), chunk) + var n = _write_fd( + fd, + slice.as_bytes().unsafe_ptr(), + chunk, + deadline_ms, + fault_eintr_count, + ) + if n == IO_DEADLINE_EXPIRED: + return ChunkWrite("write_deadline_expired", 0) if n <= 0: return ChunkWrite("write_failed", 0) return ChunkWrite("", n) @@ -450,7 +540,11 @@ struct BoundedLineReader(Movable): if ev == 0: continue var buf = InlineArray[Byte, 512](fill=0) - var n = read_fd(self.fd, buf.unsafe_ptr(), 512) + var n = read_fd( + self.fd, buf.unsafe_ptr(), 512, deadline_ms - (now_ms() - start) + ) + if n == IO_DEADLINE_EXPIRED: + raise Error("read_deadline_expired") if n < 0: raise Error("read_error") if n == 0: @@ -487,7 +581,11 @@ def read_all_bounded( raise Error("read_error") if ev == 0: continue - var n = read_fd(fd, buf.unsafe_ptr(), 4096) + var n = read_fd( + fd, buf.unsafe_ptr(), 4096, deadline_ms - (now_ms() - start) + ) + if n == IO_DEADLINE_EXPIRED: + raise Error("read_deadline_expired") if n < 0: raise Error("read_error") if n == 0: @@ -513,39 +611,191 @@ def _utf8_line(bytes: List[UInt8]) raises -> String: raise Error("invalid_utf8") -struct CleanupLedger(Copyable, Movable): - """Caller-owned record of owned-child cleanup failures. +@fieldwise_init +struct CleanupFailure(Copyable, Movable): + """One recorded owned-child cleanup problem. + + ``recoverable`` is true when the exact-owned child (``pid``/``report_fd``) + is still retained and a later retry can attempt cleanup again; ``resolved`` + flips once a retry proves cleanup or a subsequent reap collects it. + ``fd_closed`` makes the retained report descriptor close-once: whichever + party recovers first closes it, and no other call may target that number + again after it could have been reused. + """ + + var text: String + var pid: Int + var report_fd: Int + var recoverable: Bool + var resolved: Bool + var fd_closed: Bool + + def __copyinit__(out self, existing: Self): + self.text = existing.text + self.pid = existing.pid + self.report_fd = existing.report_fd + self.recoverable = existing.recoverable + self.resolved = existing.resolved + self.fd_closed = existing.fd_closed + + +def _new_failure( + text: String, pid: Int, report_fd: Int, recoverable: Bool +) -> CleanupFailure: + return CleanupFailure( + text=String(text), + pid=pid, + report_fd=report_fd, + recoverable=recoverable, + resolved=False, + fd_closed=report_fd < 0, + ) + + +struct CleanupGuard(Movable): + """Required, caller-held observation point for owned-child cleanup truth. - 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. + ADR-0017 PC02 and ADR-0018 RA01/RA02 require that an ordinary supported + provider scope cannot silently discard a cleanup failure and that a failed + cleanup keeps a usable ownership handle. This guard is therefore a + *required* parameter of every supported spawner (there is no inert default) + and it is owned by the calling test, so its record survives the ``with`` + scope where the fixture handle is destroyed. + + ``assert_clean`` is the enforcement point: every caller must invoke it, and + it raises the exact retained detail so the owning test fails truthfully. + ``retain`` additionally keeps the exact pid/report descriptor of an unproved + cleanup so ``recover_all`` can retry cleanup instead of leaving only a + message. The guard never signals an identity it was not given. """ - var errors: Optional[UnsafePointer[List[String], MutAnyOrigin]] + var failures: List[CleanupFailure] def __init__(out self): - self.errors = None + self.failures = List[CleanupFailure]() + + def record(mut self, text: String): + """Record a cleanup problem that has no retained recoverable child.""" + self.failures.append(_new_failure(text, 0, -1, False)) - def __init__(out self, errors: UnsafePointer[List[String], MutAnyOrigin]): - self.errors = errors + def retain(mut self, pid: Int, report_fd: Int, text: String): + """Record an unproved cleanup and keep its exact ownership handle.""" + self.failures.append(_new_failure(text, pid, report_fd, True)) - def record(self, text: String): - if self.errors: - self.errors.value()[].append(text) + def resolve_pid(mut self, pid: Int): + """Mark retained ownership resolved once cleanup is proved elsewhere. + + Used when a later retry (``reap``/``cleanup``/``terminate``) collects a + child whose earlier attempt was unproved, so a transient failure does + not fail the owning test after a successful recovery. + """ + if pid <= 0: + return + for index in range(len(self.failures)): + if self.failures[index].pid == pid: + self.failures[index].resolved = True + + def pending(self) -> Int: + """Number of unresolved recorded cleanup problems.""" + var count = 0 + for index in range(len(self.failures)): + if not self.failures[index].resolved: + count += 1 + return count + + def retained(self) -> Int: + """Number of unresolved entries that still hold a recoverable child.""" + var count = 0 + for index in range(len(self.failures)): + if ( + self.failures[index].recoverable + and not self.failures[index].resolved + ): + count += 1 + return count def count(self) -> Int: - if self.errors: - return len(self.errors.value()[]) - return -1 + return len(self.failures) + + def first_pending(self) -> String: + for index in range(len(self.failures)): + if not self.failures[index].resolved: + return String(self.failures[index].text) + return "" + + def first(self) -> String: + if len(self.failures) > 0: + return String(self.failures[0].text) + return "" def last(self) -> String: - if self.errors: - var recorded = self.errors.value()[] - if len(recorded) > 0: - return String(recorded[len(recorded) - 1]) + if len(self.failures) > 0: + return String(self.failures[len(self.failures) - 1].text) return "" + def is_clean(self) -> Bool: + return self.pending() == 0 + + def close_retained_fd(mut self, fd: Int) -> Bool: + """Close a retained report descriptor at most once. + + Returns True when this guard owns an entry for ``fd`` (closed now or + already closed earlier), so the owning handle must not close that + number again and can never target a reused descriptor. Returns False + when the guard holds no retained entry, letting the normal proved path + keep its own descriptor ownership. + """ + if fd < 0: + return False + var owned = False + for index in range(len(self.failures)): + if self.failures[index].report_fd != fd: + continue + owned = True + if not self.failures[index].fd_closed: + close_fd(fd) + self.failures[index].fd_closed = True + return owned + + def recover_all(mut self) -> Int: + """Retry cleanup for every retained exact-owned child. + + Returns the number of children whose cleanup is still unproved. A + child whose retry proves cleanup is marked resolved and its retained + report descriptor is closed exactly once, so recovery is real rather + than a message-only ledger entry and leaves no leaked descriptor. Only + pids this guard was handed are ever signalled. + """ + var unresolved = 0 + for index in range(len(self.failures)): + if not self.failures[index].recoverable: + continue + if self.failures[index].resolved: + continue + var pid = self.failures[index].pid + if pid <= 0: + unresolved += 1 + continue + var st = terminate_owned(pid, TERMINATION_GRACE_MS) + if st.cleanup_proved(): + self.failures[index].resolved = True + _ = self.close_retained_fd(self.failures[index].report_fd) + else: + unresolved += 1 + return unresolved + + def assert_clean(mut self) raises: + """Fail the owning test if any cleanup problem remains unresolved.""" + var remaining = self.pending() + if remaining == 0: + return + raise Error( + "cleanup-unproved: " + + String(remaining) + + " unresolved owned-child cleanup failure(s); first: " + + self.first_pending() + ) + # ── Child lifecycle ───────────────────────────────────────────────────────── @@ -681,10 +931,49 @@ def pid_not_waitable(pid: Int) -> Bool: return wait_nohang(pid).state == "gone" +def pid_running(pid: Int) -> Bool: + """True when ``pid`` is still a live owned child of this process. + + A non-blocking observation only: it never reaps, signals or releases the + exact ownership the guard retains for recovery. + """ + return wait_nohang(pid).state == "running" + + # ── Shared owned-child lifecycle state ────────────────────────────────────── @fieldwise_init +struct LifecycleFaults(Movable): + """Bounded test-only failure seam for real owned child resources. + + The seam never replaces the exact ownership identity: it makes one wait or + one cleanup attempt report an unproved result while the real forked child + keeps running, so the following retry exercises the real recoverable path + instead of an artificial ``pid=0`` handle. Every counter is consumed once. + """ + + var wait_errors: Int + var nonterminal: Int + var cleanup_failures: Int + var wait_delay_ms: Int + + def __init__(out self): + self.wait_errors = 0 + self.nonterminal = 0 + self.cleanup_failures = 0 + self.wait_delay_ms = 0 + + def active(self) -> Bool: + return ( + self.wait_errors > 0 + or self.nonterminal > 0 + or self.cleanup_failures > 0 + or self.wait_delay_ms > 0 + ) + + +@fieldwise_init struct PipedChildState(Movable): """Single mutable lifecycle record shared by every copy of one handle. @@ -692,7 +981,7 @@ struct PipedChildState(Movable): 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. + ``guard`` is the required caller-held cleanup observation point. """ var pid: Int @@ -715,7 +1004,8 @@ struct PipedChildState(Movable): var observed_valid: Bool var last_terminated: Bool var spawn_ms: Int - var ledger: CleanupLedger + var guard: UnsafePointer[CleanupGuard, MutAnyOrigin] + var faults: LifecycleFaults def store( mut self, @@ -734,9 +1024,73 @@ struct PipedChildState(Movable): self.connections = connections def close_reader(mut self): - if not self.closed: + if self.closed: + return + # A descriptor retained by the guard for recovery is closed by the + # guard exactly once, so the handle can never close a reused number. + if not self.guard[].close_retained_fd(self.report_fd): close_fd(self.report_fd) - self.closed = True + self.closed = True + + def work_remaining_ms(mut self) -> Int: + """Remaining part of the one declared work budget for this scope. + + ADR-0018 RA03: startup/body/wait/report/drain share the single finite + budget measured from ``spawn_ms``. A nonpositive result must never be + turned into another successful interval; callers fail and use only the + separate bounded cleanup allowance. + """ + return self.deadline_ms - (now_ms() - self.spawn_ms) + + def record_unproved( + mut self, pid: Int, report_fd: Int, label: String, detail: String + ): + """Record an unproved cleanup and retain its usable ownership.""" + self.cleanup_error = detail + self.guard[].retain( + pid, report_fd, label + " pid=" + String(pid) + " " + detail + ) + + def wait_once(mut self, pid: Int) -> ProcessStatus: + """Observe one child state, applying the bounded test fault seam. + + A ``wait_error`` from this observation is *not* cached as a terminal + result by callers, so a transient wait problem stays retryable. + """ + if self.faults.wait_errors > 0: + self.faults.wait_errors -= 1 + return ProcessStatus( + "wait_error", False, -1, 0, -1, "injected_wait_error" + ) + if self.faults.nonterminal > 0: + self.faults.nonterminal -= 1 + return ProcessStatus("stopped", False, -1, 0, 0x7F, "") + return wait_nohang(pid) + + def wait_until(mut self, pid: Int, deadline_ms: Int) -> ProcessStatus: + """Bounded non-blocking wait that honours the fault seam.""" + var start = now_ms() + while True: + var st = self.wait_once(pid) + if st.state != "running" and st.state != "interrupted": + return st^ + if now_ms() - start >= deadline_ms: + return st^ + sleep_ms(5) + + def terminate_once(mut self, pid: Int, grace_ms: Int) -> ProcessStatus: + """Cleanup attempt that honours the bounded test fault seam. + + A forced cleanup failure leaves the real forked child running and keeps + the exact pid/report descriptor, so recovery is a real retry rather + than a synthetic identity change. + """ + if self.faults.cleanup_failures > 0: + self.faults.cleanup_failures -= 1 + return ProcessStatus( + "wait_error", False, -1, 0, -1, "injected_cleanup_failure" + ) + return terminate_owned(pid, grace_ms) def read_line(mut self, max_bytes: Int, deadline_ms: Int) raises -> String: """Bounded line read that retains surplus as undecoded bytes. @@ -784,7 +1138,14 @@ struct PipedChildState(Movable): if ev == 0: continue var buf = InlineArray[Byte, 512](fill=0) - var n = read_fd(self.report_fd, buf.unsafe_ptr(), 512) + var n = read_fd( + self.report_fd, + buf.unsafe_ptr(), + 512, + deadline_ms - (now_ms() - start), + ) + if n == IO_DEADLINE_EXPIRED: + raise Error("read_deadline_expired") if n < 0: raise Error("read_error") if n == 0: @@ -815,7 +1176,14 @@ struct PipedChildState(Movable): if ev == 0: continue var buf = InlineArray[Byte, 1024](fill=0) - var n = read_fd(self.report_fd, buf.unsafe_ptr(), 1024) + var n = read_fd( + self.report_fd, + buf.unsafe_ptr(), + 1024, + deadline_ms - (now_ms() - start), + ) + if n == IO_DEADLINE_EXPIRED: + raise Error("read_deadline_expired") if n < 0: raise Error("read_error") if n == 0: @@ -832,7 +1200,7 @@ def piped_child_state( report_fd: Int, deadline_ms: Int, expected_requests: Int, - ledger: CleanupLedger, + guard: UnsafePointer[CleanupGuard, MutAnyOrigin], ) -> PipedChildState: """Build one owned-child lifecycle record with explicit ownership truth.""" return PipedChildState( @@ -856,7 +1224,8 @@ def piped_child_state( observed_valid=False, last_terminated=False, spawn_ms=now_ms(), - ledger=ledger.copy(), + guard=guard, + faults=LifecycleFaults(), ) @@ -867,19 +1236,26 @@ def finalize_owned_failure( Uses the same ownership truth as ``cleanup``: ownership is released only when no waitable child can remain. An unproved or uncertain termination - keeps ``reaped=False`` and records the failure in the caller-owned ledger, - so a startup failure can neither claim nor hide an uncollected child. The - report descriptor is invalidated in both cases because the failed startup - will never consume a report. + keeps ``reaped=False``, records the failure in the required guard and + *retains the usable ownership handle* (exact pid and report descriptor) + through ``CleanupGuard.retain``, so a startup failure can neither claim nor + hide an uncollected child and the caller can still recover it. The report + descriptor is invalidated in both cases because the failed startup will + never consume a report. """ - var st = terminate_owned(pid, TERMINATION_GRACE_MS) + var st = state.terminate_once(pid, TERMINATION_GRACE_MS) state.status = st.copy() - state.close_reader() if st.cleanup_proved(): + state.guard[].resolve_pid(pid) state.reaped = True + state.close_reader() else: - state.cleanup_error = "unreaped:" + st.describe() - state.ledger.record(label + " pid=" + String(pid) + " " + st.describe()) + # Retain the usable ownership handle (exact pid and report descriptor) + # so a later retry can recover; the guard closes the descriptor exactly + # once when the recovery succeeds. + state.record_unproved( + pid, state.report_fd, label, "unreaped:" + st.describe() + ) return st^ @@ -933,33 +1309,85 @@ def parse_ready_or_cleanup( ) -def descriptor_census(limit: Int) -> Int: - """Count open descriptors numerically via ``fcntl(F_GETFD)``. +comptime CENSUS_EINTR_RETRY_BOUND: Int = 64 +comptime CENSUS_FAULT_NONE: Int = -1 - Never opens a target path, so device nodes cannot block it and sockets and - high descriptors are counted the same as regular files. Returns -1 when the - census cannot be established (invalid bound or no standard descriptors), - which callers must treat as a failed unavailable census, never a pass. + +def classify_census_errno(errno_value: Int) -> String: + """Classify one ``fcntl(F_GETFD)`` lookup failure (ADR-0018 RA04). + + ``EBADF`` is the only proof that a slot is closed. ``EINTR`` is a bounded + retry. Every other lookup error is ``unavailable``: a partial census must + not silently lower the observed descriptor count. + """ + if errno_value == Int(ErrNo.EBADF.value): + return "closed" + if errno_value == Int(ErrNo.EINTR.value): + return "retry" + return "unavailable" + + +def descriptor_census_with_faults( + limit: Int, fault_fd: Int, fault_errno: Int +) -> Int: + """Numeric descriptor census with a narrow cause-specific fault seam. + + ``fault_fd``/``fault_errno`` are test-only controls that make exactly one + lookup report a chosen errno *without changing any host limit*, so the + EBADF/EINTR/other classification can be proved through the checked caller. + Returns ``-1`` when the census cannot be established: an invalid bound, no + standard descriptor, an exhausted EINTR retry bound, or any non-EBADF, + non-EINTR lookup error that would otherwise hide live descriptors. """ if limit <= 0: return -1 var count = 0 for fd in range(0, limit): + var retries = 0 + var open = False while True: - var rc = Int( - external_call["fcntl", c_int](c_int(fd), c_int(F_GETFD)) - ) + var rc = 0 + var injected = fault_fd >= 0 and fd == fault_fd + if injected: + rc = -1 + else: + rc = Int( + external_call["fcntl", c_int](c_int(fd), c_int(F_GETFD)) + ) if rc >= 0: - count += 1 + open = True + break + var errno_value = Int(get_errno().value) + if injected: + errno_value = fault_errno + var klass = classify_census_errno(errno_value) + if klass == "closed": break - if get_errno() == ErrNo.EINTR: + if klass == "retry": + retries += 1 + if retries > CENSUS_EINTR_RETRY_BOUND: + return -1 continue - break + return -1 + if open: + count += 1 if count == 0: return -1 return count +def descriptor_census(limit: Int) -> Int: + """Count open descriptors numerically via ``fcntl(F_GETFD)``. + + Never opens a target path, so device nodes cannot block it and sockets and + high descriptors are counted the same as regular files. Returns -1 when the + census cannot be established (invalid bound, no standard descriptors, or an + unclassifiable lookup error), which callers must treat as a failed + unavailable census, never a pass. + """ + return descriptor_census_with_faults(limit, CENSUS_FAULT_NONE, 0) + + def fd_scan_limit() -> Int: """Return the OS descriptor-table size, or -1 when unavailable. @@ -980,20 +1408,32 @@ def open_fd_count() -> Int: return descriptor_census(limit) -def open_fd_count_checked(limit: Int = -1) raises -> Int: - """Checked census that fails explicitly when coverage cannot be complete. +def open_fd_count_checked_with_faults( + limit: Int, fault_fd: Int, fault_errno: Int +) raises -> Int: + """Checked census with the narrow cause-specific fault seam (test only). - ``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. + ``fault_fd``/``fault_errno`` inject exactly one classified lookup failure + without changing any host limit, so the checked caller's unavailable-census + propagation can be executed for a non-EBADF/non-EINTR error (ADR-0018 RA04). """ 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) + var count = descriptor_census_with_faults(effective, fault_fd, fault_errno) if count < 0: raise Error("descriptor_census_unavailable") return 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. + """ + return open_fd_count_checked_with_faults(limit, CENSUS_FAULT_NONE, 0) diff --git a/tests/stdio_process_helper.mojo b/tests/stdio_process_helper.mojo @@ -15,6 +15,7 @@ from std.collections import List from std.ffi import CStringSlice, c_int, external_call from parent_lifecycle import ( + IO_DEADLINE_EXPIRED, LIFECYCLE_POLL_SLICE_MS, POLLERR, POLLHUP, @@ -83,13 +84,21 @@ struct DrainOutcome(Movable): def drain_ready( - fd: Int, mut out: List[UInt8], cap: Int, revents: Int + fd: Int, mut out: List[UInt8], cap: Int, revents: Int, deadline_ms: Int ) -> DrainOutcome: - """Read one ready descriptor into a capped buffer; preserve the cause.""" + """Read one ready descriptor into a capped buffer; preserve the cause. + + The read's EINTR retry is bounded by the caller's remaining + ``deadline_ms``, so a retried read returns to the deadline owner instead of + spinning, and an expired retry is a distinct ``read_deadline_expired`` + cause rather than a clean EOF. + """ if (revents & (POLLIN | POLLHUP | POLLERR | POLLNVAL)) == 0: return DrainOutcome(False, "") var buf = InlineArray[Byte, 4096](fill=0) - var n = read_fd(fd, buf.unsafe_ptr(), 4096) + var n = read_fd(fd, buf.unsafe_ptr(), 4096, deadline_ms) + if n == IO_DEADLINE_EXPIRED: + return DrainOutcome(True, "read_deadline_expired") if n < 0: return DrainOutcome(True, "read_error") if n == 0: @@ -115,19 +124,85 @@ def run_stdio_entrypoint( ) +def _stdio_phase_summary( + phase: String, + request_sent: Bool, + stdout_eof: Bool, + stderr_eof: Bool, + output_bytes: Int, + output_valid: Bool, + work_ms: Int, + cleanup_ms: Int, +) -> String: + """Bounded phase/time evidence for a stdio failure (ADR-0018 RA03). + + Work and cleanup time are recorded separately so an expired work budget can + never be confused with the bounded cleanup allowance, and the phase fields + prove which stage the child actually reached. + """ + return ( + "stdio-entrypoint phase=" + + phase + + " request_sent=" + + ("true" if request_sent else "false") + + " stdout_eof=" + + ("true" if stdout_eof else "false") + + " stderr_eof=" + + ("true" if stderr_eof else "false") + + " output_bytes=" + + String(output_bytes) + + " output_valid=" + + ("true" if output_valid else "false") + + " work_ms=" + + String(work_ms) + + " cleanup_ms=" + + String(cleanup_ms) + ) + + +def _output_is_valid_json(bytes: List[UInt8]) -> Bool: + """True when the captured stdout is nonempty and parses as JSON.""" + if len(bytes) == 0: + return False + try: + _ = loads(bytes_to_string(bytes)) + return True + except: + return False + + def _terminate_and_raise( pid: Int, stdin_fd: Int, stdout_fd: Int, stderr_fd: Int, reason: String, + phase: String, + request_sent: Bool, + stdout_eof: Bool, + stderr_eof: Bool, + output_bytes: Int, + output_valid: Bool, + work_ms: Int, ) raises: + var cleanup_start = now_ms() var st = terminate_owned(pid, TERMINATION_GRACE_MS) + var cleanup_ms = now_ms() - cleanup_start close_fd(stdin_fd) close_fd(stdout_fd) close_fd(stderr_fd) raise Error( - "stdio-entrypoint " + _stdio_phase_summary( + phase, + request_sent, + stdout_eof, + stderr_eof, + output_bytes, + output_valid, + work_ms, + cleanup_ms, + ) + + " reason=" + reason + " (child " + st.describe() @@ -145,6 +220,28 @@ def run_stdio_entrypoint_with_2_args( ) +def run_stdio_binary_with_deadline( + binary: String, + request_json: String, + arg0: String, + arg1: String, + deadline_ms: Int, +) raises -> Value: + """Run an already-built test child directly with the same guarantees. + + ADR-0018 RA03 phase isolation: this launch seam performs no compilation, so + a late-exit control can prove the child actually reached valid output, + closed output and the wait phase before its late exit was rejected. The + ordinary compile/run helper keeps its declared total budget. + """ + var args = List[String]() + if arg0 != "": + args.append(arg0) + if arg1 != "": + args.append(arg1) + return _run_stdio_launch(binary, args^, request_json, deadline_ms) + + def run_stdio_entrypoint_with_deadline( entrypoint: String, request_json: String, @@ -152,34 +249,39 @@ def run_stdio_entrypoint_with_deadline( arg1: String, deadline_ms: Int, ) raises -> Value: + var args = List[String]() + args.append("run") + args.append("-I") + args.append("src") + args.append(entrypoint) + if arg0 != "": + args.append(arg0) + if arg1 != "": + args.append(arg1) + return _run_stdio_launch("mojo", args^, request_json, deadline_ms) + + +def _run_stdio_launch( + command: String, + var args: List[String], + request_json: String, + deadline_ms: Int, +) raises -> Value: # Build argv before owning any descriptors so no exception window can leak # pipes between creation and the fork; fork failure alone is handled by the - # rollback helper. - var command = String("mojo") - var include_flag = String("-I") - var include_path = String("src") - var entrypoint_path = String(entrypoint) - var process_arg0 = String(arg0) - var process_arg1 = String(arg1) - var argv = List[Optional[CStringSlice[ImmutAnyOrigin]]](length=8, fill={}) - argv[0] = rebind[CStringSlice[ImmutAnyOrigin]](command.as_c_string_slice()) - argv[1] = rebind[CStringSlice[ImmutAnyOrigin]]("run".as_c_string_slice()) - argv[2] = rebind[CStringSlice[ImmutAnyOrigin]]( - include_flag.as_c_string_slice() - ) - argv[3] = rebind[CStringSlice[ImmutAnyOrigin]]( - include_path.as_c_string_slice() - ) - argv[4] = rebind[CStringSlice[ImmutAnyOrigin]]( - entrypoint_path.as_c_string_slice() + # rollback helper. ``parts`` owns the C-string text for the whole call, so + # every argv entry points at a live buffer until after the fork/exec. + var parts = List[String]() + parts.append(command) + for index in range(len(args)): + parts.append(args[index]) + var argv = List[Optional[CStringSlice[ImmutAnyOrigin]]]( + length=len(parts) + 1, fill={} ) - if process_arg0 != "": - argv[5] = rebind[CStringSlice[ImmutAnyOrigin]]( - process_arg0.as_c_string_slice() - ) - if process_arg1 != "": - argv[6] = rebind[CStringSlice[ImmutAnyOrigin]]( - process_arg1.as_c_string_slice() + var elements = parts.unsafe_ptr() + for index in range(len(parts)): + argv[index] = rebind[CStringSlice[ImmutAnyOrigin]]( + elements[index].as_c_string_slice() ) var pipes = make_three_pipes() @@ -189,7 +291,7 @@ def run_stdio_entrypoint_with_deadline( var stdout_write_fd = pipes.stdout_pipe.write_fd var stderr_read_fd = pipes.stderr_pipe.read_fd var stderr_write_fd = pipes.stderr_pipe.write_fd - var command_ptr = command.as_c_string_slice().unsafe_ptr() + var command_ptr = elements[0].as_c_string_slice().unsafe_ptr() var argv_ptr = argv.unsafe_ptr() var pid = fork_owned_or_close3(pipes) @@ -262,7 +364,12 @@ def run_stdio_entrypoint_with_deadline( write_reason = "write_pipe_closed" stdin_done = True elif (pr.r0 & POLLOUT) != 0: - var cw = write_fd_chunk(stdin_write_fd, request, sent) + var cw = write_fd_chunk( + stdin_write_fd, + request, + sent, + budget - (now_ms() - start), + ) if cw.reason != "": write_reason = cw.reason stdin_done = True @@ -274,7 +381,11 @@ def run_stdio_entrypoint_with_deadline( stdin_write_fd = -1 if not stdout_eof: var d = drain_ready( - stdout_read_fd, stdout, STDIO_MAX_STDOUT_BYTES, pr.r1 + stdout_read_fd, + stdout, + STDIO_MAX_STDOUT_BYTES, + pr.r1, + budget - (now_ms() - start), ) stdout_eof = d.eof if d.reason != "": @@ -282,14 +393,21 @@ def run_stdio_entrypoint_with_deadline( break if not stderr_eof: var d = drain_ready( - stderr_read_fd, stderr_bytes, STDIO_MAX_STDERR_BYTES, pr.r2 + stderr_read_fd, + stderr_bytes, + STDIO_MAX_STDERR_BYTES, + pr.r2, + budget - (now_ms() - start), ) stderr_eof = d.eof if d.reason != "": read_reason = "stderr_" + d.reason break - if write_reason == "" and sent < request.byte_length(): + var request_sent = sent >= request.byte_length() + var work_ms = now_ms() - start + var output_len = len(stdout) + if write_reason == "" and not request_sent: # The loop only ends with stdin finished; guard any path that would # otherwise leave intended request bytes unwritten. write_reason = "write_incomplete" @@ -301,18 +419,50 @@ def run_stdio_entrypoint_with_deadline( stdout_read_fd, stderr_read_fd, "write_" + write_reason, + "write", + request_sent, + stdout_eof, + stderr_eof, + output_len, + _output_is_valid_json(stdout), + work_ms, ) if read_reason != "": _terminate_and_raise( - pid, -1, stdout_read_fd, stderr_read_fd, read_reason + pid, + -1, + stdout_read_fd, + stderr_read_fd, + read_reason, + "read", + request_sent, + stdout_eof, + stderr_eof, + output_len, + _output_is_valid_json(stdout), + work_ms, ) var remaining = budget - (now_ms() - start) if remaining < 1: remaining = 1 var st = wait_bounded(pid, remaining) + work_ms = now_ms() - start if not st.cleanup_proved(): - _terminate_and_raise(pid, -1, stdout_read_fd, stderr_read_fd, "timeout") + _terminate_and_raise( + pid, + -1, + stdout_read_fd, + stderr_read_fd, + "timeout", + "wait", + request_sent, + stdout_eof, + stderr_eof, + output_len, + _output_is_valid_json(stdout), + work_ms, + ) close_fd(stdout_read_fd) close_fd(stderr_read_fd) @@ -321,7 +471,7 @@ def run_stdio_entrypoint_with_deadline( if st.exited and st.exit_code == 127: raise Error( - "stdio-entrypoint exec_failed (stdout=" + "stdio-entrypoint phase=exit exec_failed (stdout=" + output + " stderr=" + diagnostics @@ -329,7 +479,7 @@ def run_stdio_entrypoint_with_deadline( ) if not st.exited or st.exit_code != 0: raise Error( - "stdio-entrypoint child_failed (" + "stdio-entrypoint phase=exit child_failed (" + st.describe() + " stderr=" + diagnostics diff --git a/tests/test_jev.mojo b/tests/test_jev.mojo @@ -274,6 +274,7 @@ def test_circuit_opens_and_recovers() raises: from flare.http import HttpClient +from parent_lifecycle import CleanupGuard from jev_provider_helper import ( reserve_jev_port, spawn_jev_stub_auto, @@ -281,7 +282,8 @@ from jev_provider_helper import ( def test_local_provider_server_serves_scripted_jev() raises: - with spawn_jev_stub_auto("ok", 1) as started: + var guard_1 = CleanupGuard() + with spawn_jev_stub_auto("ok", 1, guard_1) as started: var port = started.port var url = "http://127.0.0.1:" + String(port) + "/v1/systemone" with HttpClient(timeout_ms=5000, max_redirects=0) as client: @@ -293,9 +295,12 @@ def test_local_provider_server_serves_scripted_jev() raises: assert_equal(body["model"].string_value(), "jev-1.13.0") started.stub.wait() + guard_1.assert_clean() + def test_local_provider_server_scripts_transport_failures() raises: - with spawn_jev_stub_auto("rate_limit", 1) as rate_port_started: + var guard_2 = CleanupGuard() + with spawn_jev_stub_auto("rate_limit", 1, guard_2) as rate_port_started: var rate_port = rate_port_started.port with HttpClient(timeout_ms=5000, max_redirects=0) as client: var response = client.post( @@ -304,7 +309,10 @@ def test_local_provider_server_scripts_transport_failures() raises: assert_equal(response.status, 429) rate_port_started.stub.wait() - with spawn_jev_stub_auto("malformed_json", 1) as malformed_port_started: + var guard_3 = CleanupGuard() + with spawn_jev_stub_auto( + "malformed_json", 1, guard_3 + ) as malformed_port_started: var malformed_port = malformed_port_started.port with HttpClient(timeout_ms=5000, max_redirects=0) as client: var response = client.post( @@ -316,7 +324,10 @@ def test_local_provider_server_scripts_transport_failures() raises: assert_equal(response.text(), "not json") malformed_port_started.stub.wait() - with spawn_jev_stub_auto("server_error", 1) as err_port_started: + var guard_4 = CleanupGuard() + with spawn_jev_stub_auto( + "server_error", 1, guard_4 + ) as err_port_started: var err_port = err_port_started.port with HttpClient(timeout_ms=5000, max_redirects=0) as client: var response = client.post( @@ -328,6 +339,10 @@ def test_local_provider_server_scripts_transport_failures() raises: assert_equal(response.status, 500) err_port_started.stub.wait() + guard_4.assert_clean() + guard_3.assert_clean() + guard_2.assert_clean() + from hyf_provider.jev_client import post_jev_systemone, validate_jev_base_url from json import loads as _loads @@ -351,7 +366,8 @@ def test_jev_endpoint_policy_and_loopback_client() raises: with assert_raises(): _ = validate_jev_base_url("ftp://api.typesafe.ai") - with spawn_jev_stub_auto("ok", 1) as started: + var guard_5 = CleanupGuard() + with spawn_jev_stub_auto("ok", 1, guard_5) as started: var port = started.port var outcome = post_jev_systemone( "http://127.0.0.1:" + String(port), @@ -362,6 +378,8 @@ def test_jev_endpoint_policy_and_loopback_client() raises: assert_true(outcome.body_text.find("jev-1.13.0") >= 0) started.stub.wait() + guard_5.assert_clean() + from flare.tls import TlsVerify from flare.tls import TlsConfig @@ -380,7 +398,8 @@ def test_tls_and_redirect_policy() raises: assert_tls_verification_required(TlsConfig.insecure()) assert_true(not redirects_forward_credentials()) - with spawn_jev_stub_auto("redirect", 1) as started: + var guard_6 = CleanupGuard() + with spawn_jev_stub_auto("redirect", 1, guard_6) as started: var port = started.port with assert_raises(): with HttpClient(timeout_ms=5000, max_redirects=0) as client: @@ -389,6 +408,8 @@ def test_tls_and_redirect_policy() raises: ) started.stub.wait() + guard_6.assert_clean() + from hyf_provider.jev_client import failure_kind_for_status, retry_decision @@ -410,7 +431,8 @@ def test_bounded_retry_and_budget_behavior() raises: def test_transport_cleanup_and_local_cancellation() raises: - with spawn_jev_stub_auto("ok", 1) as started: + var guard_7 = CleanupGuard() + with spawn_jev_stub_auto("ok", 1, guard_7) as started: var port = started.port var outcome = post_jev_systemone( "http://127.0.0.1:" + String(port), @@ -428,3 +450,5 @@ def test_transport_cleanup_and_local_cancellation() raises: _loads('{"model":"jev-1.13.0","state":"s","questions":{}}'), 150, ) + + guard_7.assert_clean() diff --git a/tests/test_provider_adapter.mojo b/tests/test_provider_adapter.mojo @@ -26,6 +26,7 @@ from hyf_runtime.config import ( HyfServiceRuntimeConfig, default_loaded_runtime_config, ) +from parent_lifecycle import CleanupGuard from max_local_process_helper import ( reserve_loopback_port, spawn_max_local_stub, @@ -261,8 +262,9 @@ def test_max_local_transport_boundary_rejects_invalid_health_url() raises: def test_max_local_transport_boundary_reports_unknown_chat_transport() raises: + var guard_1 = CleanupGuard() with spawn_max_local_stub( - 0, "query_rewrite_malformed_http", 1 + 0, "query_rewrite_malformed_http", 1, guard_1 ) as provider_stub: var provider_port = provider_stub.port var outcome = post_max_local_chat_completion( @@ -276,9 +278,14 @@ def test_max_local_transport_boundary_reports_unknown_chat_transport() raises: provider_stub.wait() + guard_1.assert_clean() + def test_max_local_transport_boundary_reports_unknown_health_transport() raises: - with spawn_max_local_stub(0, "health_malformed_http", 1) as provider_stub: + var guard_2 = CleanupGuard() + with spawn_max_local_stub( + 0, "health_malformed_http", 1, guard_2 + ) as provider_stub: var provider_port = provider_stub.port var outcome = get_max_local_health( _provider_config_for_port(provider_port) @@ -291,6 +298,8 @@ def test_max_local_transport_boundary_reports_unknown_health_transport() raises: provider_stub.wait() + guard_2.assert_clean() + def test_query_rewrite_request_body_sets_schema_contract() raises: var context = default_request_context() diff --git a/tests/test_provider_helpers.mojo b/tests/test_provider_helpers.mojo @@ -14,13 +14,15 @@ from flare.tcp import TcpListener, TcpStream from parent_lifecycle import ( CENSUS_MAX_FDS, - CleanupLedger, + CleanupGuard, PipedChildState, ProcessStatus, child_exit, + classify_census_errno, classify_wait_errno, close_fd, descriptor_census, + descriptor_census_with_faults, dup2_fd, finalize_owned_failure, fork_owned_or_close, @@ -31,9 +33,11 @@ from parent_lifecycle import ( now_ms, open_fd_count, open_fd_count_checked, + open_fd_count_checked_with_faults, parse_ready_line, parse_ready_or_cleanup, pid_not_waitable, + pid_running, piped_child_state, read_all_bounded, read_line_bounded, @@ -185,7 +189,8 @@ struct FramingFailure(Movable): def test_max_local_stub_reads_fragmented_large_body() raises: - with spawn_max_local_stub(0, "echo_body_bytes", 1) as stub: + var guard_1 = CleanupGuard() + with spawn_max_local_stub(0, "echo_body_bytes", 1, guard_1) as stub: var body = String("") for _ in range(9000): body += "x" @@ -193,10 +198,13 @@ def test_max_local_stub_reads_fragmented_large_body() raises: assert_true(response.find('"received_bytes":9000') >= 0) stub.wait() + guard_1.assert_clean() + def test_max_local_stub_counts_every_wire_attempt() raises: var requests = 3 - with spawn_max_local_stub(0, "count_requests", requests) as stub: + var guard_2 = CleanupGuard() + with spawn_max_local_stub(0, "count_requests", requests, guard_2) as stub: for index in range(requests): var response = _request( stub.port, "POST", "/v1/chat/completions", "{}" @@ -206,11 +214,14 @@ def test_max_local_stub_counts_every_wire_attempt() raises: ) stub.wait() + guard_2.assert_clean() + def test_max_local_stub_rejects_unknown_path() raises: # FX02/FX04: an unexpected route must fail fixture verification, not be # answered 404 and then reported as a successful stub run. - with spawn_max_local_stub(0, "query_rewrite_ok", 1) as stub: + var guard_3 = CleanupGuard() + with spawn_max_local_stub(0, "query_rewrite_ok", 1, guard_3) as stub: var response = _request(stub.port, "POST", "/not-a-route", "{}") assert_true(response.find("404") < 0) stub.reap() @@ -218,17 +229,23 @@ def test_max_local_stub_rejects_unknown_path() raises: assert_equal(stub.phase(), "exchange") assert_equal(stub.reason(), "unexpected_path") + guard_3.assert_clean() + def test_max_local_stub_binds_and_reports_port() raises: - with spawn_max_local_stub(0, "count_requests", 1) as stub: + var guard_4 = CleanupGuard() + with spawn_max_local_stub(0, "count_requests", 1, guard_4) as stub: assert_true(stub.port > 0) var response = _request(stub.port, "POST", "/v1/chat/completions", "{}") assert_true(response.find('"request_index":1') >= 0) stub.wait() + guard_4.assert_clean() + def test_jev_stub_observes_bearer_sentinel_at_intended_origin() raises: - with spawn_jev_stub_auto("echo_authorization", 1) as started: + var guard_5 = CleanupGuard() + with spawn_jev_stub_auto("echo_authorization", 1, guard_5) as started: var response = _request( started.port, "POST", @@ -239,9 +256,12 @@ def test_jev_stub_observes_bearer_sentinel_at_intended_origin() raises: assert_true(response.find("hyf-sentinel-token") >= 0) started.stub.wait() + guard_5.assert_clean() + def test_max_local_stub_stalled_child_is_reaped() raises: - with spawn_max_local_stub(0, "stall", 1) as stub: + var guard_6 = CleanupGuard() + with spawn_max_local_stub(0, "stall", 1, guard_6) as stub: _raw_send_only( stub.port, ( @@ -252,8 +272,9 @@ def test_max_local_stub_stalled_child_is_reaped() raises: stub.terminate() assert_true(pid_not_waitable(stub.pid)) + # ── FX01: explicit ordered scripted exchanges ─────────────────────────────── -# ── FX01: explicit ordered scripted exchanges ─────────────────────────────── + guard_6.assert_clean() def test_max_local_scripted_matches_explicit_exchange() raises: @@ -267,7 +288,8 @@ def test_max_local_scripted_matches_explicit_exchange() raises: script.response_headers = "x-scripted: yes" script.delay_ms = 20 scripts.append(script^) - with spawn_max_local_scripted(0, scripts^) as stub: + var guard_7 = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard_7) as stub: var response = _request( stub.port, "POST", @@ -282,11 +304,14 @@ def test_max_local_scripted_matches_explicit_exchange() raises: assert_equal(stub.request_count(), 1) assert_equal(stub.connection_count(), 1) + guard_7.assert_clean() + def test_max_local_scripted_rejects_wrong_method() raises: var scripts = List[ExchangeScript]() scripts.append(_default_script()) - with spawn_max_local_scripted(0, scripts^) as stub: + var guard_8 = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard_8) as stub: _raw_send_only( stub.port, ( @@ -298,11 +323,14 @@ def test_max_local_scripted_rejects_wrong_method() raises: assert_equal(stub.phase(), "exchange") assert_equal(stub.reason(), "method_mismatch") + guard_8.assert_clean() + def test_max_local_scripted_rejects_wrong_path() raises: var scripts = List[ExchangeScript]() scripts.append(_default_script()) - with spawn_max_local_scripted(0, scripts^) as stub: + var guard_9 = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard_9) as stub: _raw_send_only( stub.port, ( @@ -314,13 +342,16 @@ def test_max_local_scripted_rejects_wrong_path() raises: assert_equal(stub.phase(), "exchange") assert_equal(stub.reason(), "path_mismatch") + guard_9.assert_clean() + def test_max_local_scripted_rejects_wrong_selected_header() raises: var scripts = List[ExchangeScript]() var script = _default_script() script.headers = "x-sentinel:expected" scripts.append(script^) - with spawn_max_local_scripted(0, scripts^) as stub: + var guard_10 = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard_10) as stub: _raw_send_only( stub.port, ( @@ -333,6 +364,8 @@ def test_max_local_scripted_rejects_wrong_selected_header() raises: assert_equal(stub.phase(), "exchange") assert_equal(stub.reason(), "header_mismatch:x-sentinel") + guard_10.assert_clean() + def test_max_local_scripted_rejects_wrong_body() raises: var scripts = List[ExchangeScript]() @@ -340,7 +373,8 @@ def test_max_local_scripted_rejects_wrong_body() raises: script.check_body = True script.body = '{"expected":true}' scripts.append(script^) - with spawn_max_local_scripted(0, scripts^) as stub: + var guard_11 = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard_11) as stub: _raw_send_only( stub.port, ( @@ -352,8 +386,9 @@ def test_max_local_scripted_rejects_wrong_body() raises: assert_equal(stub.phase(), "exchange") assert_equal(stub.reason(), "body_mismatch") + # ── FX02: unexpected/extra/missing/unconsumed accounting ──────────────────── -# ── FX02: unexpected/extra/missing/unconsumed accounting ──────────────────── + guard_11.assert_clean() def test_scripted_rejects_extra_pipelined_exchange() raises: @@ -361,7 +396,8 @@ def test_scripted_rejects_extra_pipelined_exchange() raises: var script = _default_script() script.close_connection = False scripts.append(script^) - with spawn_max_local_scripted(0, scripts^) as stub: + var guard_12 = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard_12) as stub: var first_frame = ( "POST /v1/chat/completions HTTP/1.1\r\nhost: 127.0.0.1\r\n" "content-length: 2\r\nconnection: keep-alive\r\n\r\n{}" @@ -372,11 +408,14 @@ def test_scripted_rejects_extra_pipelined_exchange() raises: assert_equal(stub.phase(), "accounting") assert_equal(stub.reason(), "extra_exchange_after_completion") + guard_12.assert_clean() + def test_scripted_rejects_missing_exchange() raises: var scripts = List[ExchangeScript]() scripts.append(_default_script()) - with spawn_max_local_scripted(0, scripts^) as stub: + var guard_13 = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard_13) as stub: var client = TcpStream.connect(SocketAddr.localhost(UInt16(stub.port))) client.close() stub.reap() @@ -384,6 +423,8 @@ def test_scripted_rejects_missing_exchange() raises: assert_equal(stub.reason(), "missing_exchanges") assert_equal(stub.request_count(), 0) + guard_13.assert_clean() + def test_scripted_reports_unconsumed_remaining_scripts() raises: var scripts = List[ExchangeScript]() @@ -391,7 +432,8 @@ def test_scripted_reports_unconsumed_remaining_scripts() raises: first.close_connection = False scripts.append(first^) scripts.append(_default_script()) - with spawn_max_local_scripted(0, scripts^) as stub: + var guard_14 = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard_14) as stub: var client = TcpStream.connect(SocketAddr.localhost(UInt16(stub.port))) _write_all( client, @@ -407,17 +449,20 @@ def test_scripted_reports_unconsumed_remaining_scripts() raises: assert_equal(stub.reason(), "missing_exchanges") assert_equal(stub.request_count(), 1) + # ── FX03: strict lexical framing ──────────────────────────────────────────── -# ── FX03: strict lexical framing ──────────────────────────────────────────── + guard_14.assert_clean() def _framing_failure(raw: String) raises -> FramingFailure: var scripts = List[ExchangeScript]() scripts.append(_default_script()) - with spawn_max_local_scripted(0, scripts^) as stub: + var guard = CleanupGuard() + var snapshot = FramingFailure(False, "", "-", "", 0, 0) + with spawn_max_local_scripted(0, scripts^, guard) as stub: _raw_send_only(stub.port, raw) stub.reap() - return FramingFailure( + snapshot = FramingFailure( stub.ok(), stub.phase(), stub.failure_case(), @@ -425,6 +470,8 @@ def _framing_failure(raw: String) raises -> FramingFailure: stub.request_count(), stub.connection_count(), ) + guard.assert_clean() + return snapshot^ def test_strict_framing_lexical_content_length() raises: @@ -540,7 +587,8 @@ def test_strict_framing_header_cap_exceeded() raises: def test_jev_echo_authorization_ignores_x_authorization() raises: - with spawn_jev_stub_auto("echo_authorization", 1) as started: + var guard_16 = CleanupGuard() + with spawn_jev_stub_auto("echo_authorization", 1, guard_16) as started: var response = _request( started.port, "POST", @@ -552,6 +600,8 @@ def test_jev_echo_authorization_ignores_x_authorization() raises: assert_true(response.find("spoof") < 0) started.stub.wait() + guard_16.assert_clean() + def test_strict_framing_split_utf8_body() raises: var scripts = List[ExchangeScript]() @@ -564,7 +614,8 @@ def test_strict_framing_split_utf8_body() raises: payload += "é" script.body = payload scripts.append(script^) - with spawn_max_local_scripted(0, scripts^) as stub: + var guard_17 = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard_17) as stub: var client = TcpStream.connect(SocketAddr.localhost(UInt16(stub.port))) var head = ( "POST /v1/chat/completions HTTP/1.1\r\nhost: 127.0.0.1\r\n" @@ -586,6 +637,8 @@ def test_strict_framing_split_utf8_body() raises: assert_true(response.find("200") >= 0) stub.wait() + guard_17.assert_clean() + def test_strict_framing_surplus_retained_for_second_frame() raises: var scripts = List[ExchangeScript]() @@ -598,7 +651,8 @@ def test_strict_framing_surplus_retained_for_second_frame() raises: "second", "POST", "/v1/chat/completions", 200, '{"n":2}' ) scripts.append(second^) - with spawn_max_local_scripted(0, scripts^) as stub: + var guard_18 = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard_18) as stub: var first_frame = ( "POST /v1/chat/completions HTTP/1.1\r\nhost: 127.0.0.1\r\n" "content-length: 2\r\nconnection: keep-alive\r\n\r\n{}" @@ -611,8 +665,9 @@ def test_strict_framing_surplus_retained_for_second_frame() raises: assert_equal(stub.request_count(), 2) assert_equal(stub.connection_count(), 1) + # ── FX04: route/method before auth, exact headers, safe escaping ──────────── -# ── FX04: route/method before auth, exact headers, safe escaping ──────────── + guard_18.assert_clean() def test_jev_scripted_wrong_route_auth_not_bypassed() raises: @@ -622,7 +677,8 @@ def test_jev_scripted_wrong_route_auth_not_bypassed() raises: ) script.require_bearer = True scripts.append(script^) - with spawn_jev_scripted_auto(scripts^) as started: + var guard_19 = CleanupGuard() + with spawn_jev_scripted_auto(scripts^, guard_19) as started: _raw_send_only( started.port, ( @@ -635,6 +691,8 @@ def test_jev_scripted_wrong_route_auth_not_bypassed() raises: assert_equal(started.stub.phase(), "exchange") assert_equal(started.stub.reason(), "path_mismatch") + guard_19.assert_clean() + def test_scripted_rejects_duplicate_authorization() raises: var scripts = List[ExchangeScript]() @@ -643,7 +701,8 @@ def test_scripted_rejects_duplicate_authorization() raises: ) script.require_bearer = True scripts.append(script^) - with spawn_max_local_scripted(0, scripts^) as stub: + var guard_20 = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard_20) as stub: _raw_send_only( stub.port, ( @@ -656,6 +715,8 @@ def test_scripted_rejects_duplicate_authorization() raises: assert_equal(stub.phase(), "exchange") assert_equal(stub.reason(), "auth_duplicate") + guard_20.assert_clean() + def test_json_escape_control_characters() raises: assert_equal(json_escape("a\nb\tc"), '"a\\nb\\tc"') @@ -775,7 +836,8 @@ def test_scripted_persistent_counters_and_close_semantics() raises: ) second.response_headers = "x-step: two" scripts.append(second^) - with spawn_jev_scripted_auto(scripts^) as started: + var guard_21 = CleanupGuard() + with spawn_jev_scripted_auto(scripts^, guard_21) as started: var first_frame = ( "POST /v1/systemone HTTP/1.1\r\nhost: 127.0.0.1\r\n" "authorization: Bearer t\r\ncontent-length: 2\r\n" @@ -791,8 +853,9 @@ def test_scripted_persistent_counters_and_close_semantics() raises: assert_equal(started.stub.request_count(), 2) assert_equal(started.stub.connection_count(), 1) + # ── FX06/FX07/FX08: parent lifecycle and cause-specific failures ──────────── -# ── FX06/FX07/FX08: parent lifecycle and cause-specific failures ──────────── + guard_21.assert_clean() def test_startup_failure_distinct_from_exchange_failure() raises: @@ -803,8 +866,10 @@ def test_startup_failure_distinct_from_exchange_failure() raises: var port = Int(blocker.local_addr().port) var message = "" try: - with spawn_max_local_stub(port, "count_requests", 1) as stub: + var guard_22 = CleanupGuard() + with spawn_max_local_stub(port, "count_requests", 1, guard_22) as stub: stub.terminate() + guard_22.assert_clean() except e: message = String(e) blocker.close() @@ -816,22 +881,28 @@ def test_startup_failure_distinct_from_exchange_failure() raises: def test_provider_stub_parent_deadline_watchdog() raises: # No client connects, so the child blocks in accept until the parent's own # finite deadline fires and the owned child is terminated and reaped. - with spawn_max_local_stub(0, "count_requests", 1, 800) as stub: + var guard_23 = CleanupGuard() + with spawn_max_local_stub(0, "count_requests", 1, guard_23, 800) as stub: stub.reap() assert_true(not stub.ok()) assert_equal(stub.phase(), "watchdog") assert_equal(stub.reason(), "timeout") assert_true(pid_not_waitable(stub.pid)) + guard_23.assert_clean() + def test_jev_stub_parent_deadline_watchdog() raises: - with spawn_jev_stub_auto("ok", 1, 800) as started: + var guard_24 = CleanupGuard() + with spawn_jev_stub_auto("ok", 1, guard_24, 800) as started: started.stub.reap() assert_true(not started.stub.ok()) assert_equal(started.stub.phase(), "watchdog") assert_equal(started.stub.reason(), "timeout") assert_true(pid_not_waitable(started.stub.pid)) + guard_24.assert_clean() + def test_bounded_read_caps_fail_for_intended_cause() raises: var ready_pipe = make_pipe() @@ -880,16 +951,20 @@ def test_write_deadline_and_closed_pipe_causes() raises: def test_owned_child_reaped_after_early_terminate() raises: - with spawn_max_local_stub(0, "count_requests", 1) as stub: + var guard_25 = CleanupGuard() + with spawn_max_local_stub(0, "count_requests", 1, guard_25) as stub: stub.terminate() assert_true(pid_not_waitable(stub.pid)) + guard_25.assert_clean() + def test_repeated_failures_leave_no_owned_child() raises: for _ in range(3): var scripts = List[ExchangeScript]() scripts.append(_default_script()) - with spawn_max_local_scripted(0, scripts^) as stub: + var guard_26 = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard_26) as stub: _raw_send_only( stub.port, ( @@ -902,21 +977,26 @@ def test_repeated_failures_leave_no_owned_child() raises: assert_equal(stub.reason(), "path_mismatch") assert_true(pid_not_waitable(stub.pid)) + guard_26.assert_clean() + def test_repeated_teardown_does_not_leak_descriptors() raises: var before = open_fd_count() assert_true(before > 0) for _ in range(5): - with spawn_max_local_stub(0, "count_requests", 1) as stub: + var guard_27 = CleanupGuard() + with spawn_max_local_stub(0, "count_requests", 1, guard_27) as stub: stub.terminate() assert_true(pid_not_waitable(stub.pid)) + guard_27.assert_clean() var after = open_fd_count() assert_true(after > 0) assert_true(after <= before) def test_timeout_terminates_and_reaps_stalled_child() raises: - with spawn_max_local_stub(0, "stall", 1) as stub: + var guard_28 = CleanupGuard() + with spawn_max_local_stub(0, "stall", 1, guard_28) as stub: _raw_send_only( stub.port, ( @@ -927,14 +1007,17 @@ def test_timeout_terminates_and_reaps_stalled_child() raises: stub.terminate() assert_true(pid_not_waitable(stub.pid)) + guard_28.assert_clean() + def _owned_report_child( - exit_code: Int, report: String + mut guard: CleanupGuard, exit_code: Int, report: String ) raises -> SpawnedMaxLocalStub: """Fork a test-owned child that writes ``report`` to stdout and exits. Lets LC02 prove real forged/empty/mismatched reports fail for their cause, - independent of the fixture serve loop. + independent of the fixture serve loop. The caller holds ``guard`` so the + cleanup outcome stays observable. """ var pipe = make_pipe() var pid = fork_pid() @@ -947,12 +1030,14 @@ def _owned_report_child( _ = 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()) + var state = piped_child_state( + pid, pipe.read_fd, 2000, 1, UnsafePointer(to=guard) + ) return SpawnedMaxLocalStub(pid, 0, state^) def _owned_jev_report_child( - exit_code: Int, report: String + mut guard: CleanupGuard, exit_code: Int, report: String ) raises -> SpawnedJevStub: """Same controlled report child, reaped through the Jev provider path.""" var pipe = make_pipe() @@ -966,7 +1051,9 @@ def _owned_jev_report_child( _ = 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()) + var state = piped_child_state( + pid, pipe.read_fd, 2000, 1, UnsafePointer(to=guard) + ) return SpawnedJevStub(pid, 0, state^) @@ -977,9 +1064,11 @@ def test_scope_cleanup_on_assertion_failure() raises: var held_pid = 0 var caught = False try: - with spawn_max_local_stub(0, "count_requests", 1) as stub: + var guard_29 = CleanupGuard() + with spawn_max_local_stub(0, "count_requests", 1, guard_29) as stub: held_pid = stub.pid assert_true(False) + guard_29.assert_clean() except: caught = True assert_true(caught) @@ -991,34 +1080,42 @@ def test_scope_cleanup_on_generic_error() raises: var held_pid = 0 var message = "" try: - with spawn_max_local_stub(0, "count_requests", 1) as stub: + var guard_30 = CleanupGuard() + with spawn_max_local_stub(0, "count_requests", 1, guard_30) as stub: held_pid = stub.pid raise Error("intentional scope error") + guard_30.assert_clean() except e: message = String(e) assert_equal(message, "intentional scope error") assert_true(pid_not_waitable(held_pid)) -def _early_return_owner() raises -> Int: - with spawn_max_local_stub(0, "count_requests", 1) as stub: +def _early_return_owner(mut guard: CleanupGuard) raises -> Int: + # The guard is caller-held so the cleanup outcome stays observable even + # though this scope exits through a ``return`` before any post-scope line. + with spawn_max_local_stub(0, "count_requests", 1, guard) as stub: return stub.pid return 0 def test_scope_cleanup_on_early_return() raises: - var held_pid = _early_return_owner() + var guard = CleanupGuard() + var held_pid = _early_return_owner(guard) assert_true(held_pid > 0) assert_true(pid_not_waitable(held_pid)) + guard.assert_clean() def test_jev_scope_cleanup_on_assertion_failure() raises: var held_pid = 0 var caught = False try: - with spawn_jev_stub_auto("ok", 1) as started: + var guard_32 = CleanupGuard() + with spawn_jev_stub_auto("ok", 1, guard_32) as started: held_pid = started.stub.pid assert_true(False) + guard_32.assert_clean() except: caught = True assert_true(caught) @@ -1030,39 +1127,71 @@ def test_jev_scope_cleanup_on_generic_error() raises: var held_pid = 0 var message = "" try: - with spawn_jev_stub_auto("ok", 1) as started: + var guard_33 = CleanupGuard() + with spawn_jev_stub_auto("ok", 1, guard_33) as started: held_pid = started.stub.pid raise Error("intentional jev scope error") + guard_33.assert_clean() except e: message = String(e) assert_equal(message, "intentional jev scope error") assert_true(pid_not_waitable(held_pid)) -def _jev_early_return_owner() raises -> Int: - with spawn_jev_stub_auto("ok", 1) as started: +def _jev_early_return_owner(mut guard: CleanupGuard) raises -> Int: + # The guard is caller-held so the cleanup outcome stays observable even + # though this scope exits through a ``return`` before any post-scope line. + with spawn_jev_stub_auto("ok", 1, guard) as started: return started.stub.pid return 0 def test_jev_scope_cleanup_on_early_return() raises: - var held_pid = _jev_early_return_owner() + var guard = CleanupGuard() + var held_pid = _jev_early_return_owner(guard) + assert_true(held_pid > 0) + assert_true(pid_not_waitable(held_pid)) + guard.assert_clean() + + +def test_jev_early_return_cleanup_failure_is_observable() raises: + # RA01: an early return that leaves cleanup unproved must still fail the + # owning test through the caller-held guard, for the Jev provider path. + var guard = CleanupGuard() + var held_pid = 0 + with spawn_jev_stub_auto("ok", 1, guard) as started: + held_pid = started.stub.pid + started.stub.inject_cleanup_failure() + var failure = "" + try: + guard.assert_clean() + except e: + failure = String(e) + assert_true(failure.find("cleanup-unproved") >= 0) assert_true(held_pid > 0) + # The retained exact-owned child is still recoverable, not a message only. + assert_true(guard.retained() >= 1) + assert_equal(guard.recover_all(), 0) + guard.assert_clean() assert_true(pid_not_waitable(held_pid)) def test_jev_startup_failure_is_truthful_and_cause_specific() raises: # PC02/LC01: the Jev startup-readiness failure path executes against a real # owned child, reports its cause and exposes the finalizer's cleanup truth. + # The caller-held guard is the enforcement point: a startup finalization + # that could not prove cleanup fails this test instead of being discarded. var blocker = TcpListener.bind(SocketAddr.localhost(0)) var port = Int(blocker.local_addr().port) var message = "" + var guard = CleanupGuard() try: - var started = spawn_jev_stub(port, "ok", 1) + var started = spawn_jev_stub(port, "ok", 1, guard) started.cleanup() except e: message = String(e) blocker.close() + guard.assert_clean() assert_true(message.find("phase=startup") >= 0) assert_true(message.find("reason=serve_failed") >= 0) assert_true(message.find("cleanup=") >= 0) @@ -1075,16 +1204,19 @@ def test_wait_error_taxonomy_distinguishes_causes() raises: assert_equal(classify_wait_errno(9999), "wait_error") assert_equal(wait_nohang(0).state, "wait_error") var live_pid = 0 - with spawn_max_local_stub(0, "count_requests", 1) as stub: + var guard_35 = CleanupGuard() + with spawn_max_local_stub(0, "count_requests", 1, guard_35) as stub: live_pid = stub.pid assert_equal(wait_nohang(live_pid).state, "running") stub.terminate() assert_equal(wait_nohang(live_pid).state, "gone") + guard_35.assert_clean() assert_true(pid_not_waitable(live_pid)) def test_repeated_reap_and_terminate_are_owned_and_idempotent() raises: - with spawn_max_local_stub(0, "count_requests", 1) as stub: + var guard_36 = CleanupGuard() + with spawn_max_local_stub(0, "count_requests", 1, guard_36) as stub: var owned_pid = stub.pid stub.terminate() stub.terminate() @@ -1094,19 +1226,25 @@ def test_repeated_reap_and_terminate_are_owned_and_idempotent() raises: assert_true(cached.cleanup_proved()) assert_true(not stub.ok()) + # ── LC02: strict result truth ─────────────────────────────────────────────── -# ── LC02: strict result truth ─────────────────────────────────────────────── + guard_36.assert_clean() def test_result_truth_rejects_empty_exit_zero_report() raises: - var stub = _owned_report_child(0, "") + var guard_101 = CleanupGuard() + var stub = _owned_report_child(guard_101, 0, "") stub.reap() assert_true(not stub.ok()) assert_equal(stub.reason(), "missing_report") + guard_101.assert_clean() + def test_result_truth_rejects_forged_success_with_nonzero_exit() raises: + var guard_102 = CleanupGuard() var stub = _owned_report_child( + guard_102, 7, "result ok phase=complete case=- reason=ok requests=1 connections=1\n", ) @@ -1114,9 +1252,13 @@ def test_result_truth_rejects_forged_success_with_nonzero_exit() raises: assert_true(not stub.ok()) assert_true(stub.reason().startswith("report_status_mismatch")) + guard_102.assert_clean() + def test_result_truth_accepts_matching_report_and_exit() raises: + var guard_103 = CleanupGuard() var stub = _owned_report_child( + guard_103, 0, "result ok phase=complete case=- reason=ok requests=1 connections=1\n", ) @@ -1126,6 +1268,8 @@ def test_result_truth_accepts_matching_report_and_exit() raises: assert_equal(stub.request_count(), 1) assert_equal(stub.connection_count(), 1) + guard_103.assert_clean() + def test_parse_report_rejects_malformed_inputs() raises: assert_equal(parse_report("").phase, "parse") @@ -1213,12 +1357,13 @@ def test_coalesced_ready_and_report_lines_retain_surplus() raises: " connections=1\n" ), ) + var guard = CleanupGuard() var state = piped_child_state( pid=0, report_fd=pipe.read_fd, deadline_ms=500, expected_requests=1, - ledger=CleanupLedger(), + guard=UnsafePointer(to=guard), ) var ready = state.read_line(2048, 500) var report = state.read_line(2048, 500) @@ -1226,6 +1371,7 @@ def test_coalesced_ready_and_report_lines_retain_surplus() raises: close_fd(pipe.write_fd) assert_equal(ready, "ready 4242") assert_true(report.startswith("result ok")) + guard.assert_clean() # ── LC05: framing and descriptor census ───────────────────────────────────── @@ -1266,7 +1412,8 @@ def test_header_value_rejects_control_bytes_and_trims_ows() raises: ) script.headers = "x-ows:value" scripts.append(script^) - with spawn_max_local_scripted(0, scripts^) as stub: + var guard_37 = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard_37) as stub: var response = _request( stub.port, "POST", @@ -1277,6 +1424,8 @@ def test_header_value_rejects_control_bytes_and_trims_ows() raises: assert_true(response.find("200") >= 0) stub.wait() + guard_37.assert_clean() + def test_descriptor_census_detects_planted_high_fd() raises: var before = open_fd_count() @@ -1333,7 +1482,9 @@ def test_stdio_fork_failure_closes_all_owned_pipes() raises: def test_result_truth_rejects_duplicate_report_line() raises: + var guard_104 = CleanupGuard() var stub = _owned_report_child( + guard_104, 0, ( "result ok phase=complete case=- reason=ok requests=1" @@ -1345,33 +1496,38 @@ def test_result_truth_rejects_duplicate_report_line() raises: assert_true(not stub.ok()) assert_equal(stub.reason(), "duplicate_report") + guard_104.assert_clean() + def test_cleanup_failure_is_observable() raises: - # 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, - deadline_ms=100, - expected_requests=1, - ledger=ledger, - ) + # RA01/RA02: a cleanup failure must be observable through the required + # caller-held guard, must not claim the child was collected, and must retain + # retryable ownership rather than only a message. + var guard = CleanupGuard() + var state = piped_child_state(0, -1, 100, 1, UnsafePointer(to=guard)) 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) + assert_equal(guard.count(), 1) + assert_true(guard.first().find("unproved") >= 0) + assert_equal(guard.retained(), 1) # 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) + assert_equal(guard.count(), 2) + var failure = "" + try: + guard.assert_clean() + except e: + failure = String(e) + assert_true(failure.find("cleanup-unproved") >= 0) def test_result_truth_rejects_wrong_request_count() raises: + var guard_105 = CleanupGuard() var stub = _owned_report_child( + guard_105, 0, "result ok phase=complete case=- reason=ok requests=2 connections=1\n", ) @@ -1379,9 +1535,13 @@ def test_result_truth_rejects_wrong_request_count() raises: assert_true(not stub.ok()) assert_equal(stub.reason(), "request_count_mismatch") + guard_105.assert_clean() + def test_result_truth_rejects_invalid_connection_count() raises: + var guard_106 = CleanupGuard() var stub = _owned_report_child( + guard_106, 0, "result ok phase=complete case=- reason=ok requests=1 connections=5\n", ) @@ -1389,6 +1549,8 @@ def test_result_truth_rejects_invalid_connection_count() raises: assert_true(not stub.ok()) assert_equal(stub.reason(), "connection_count_invalid") + guard_106.assert_clean() + 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 @@ -1506,7 +1668,9 @@ def test_malformed_ready_line_terminates_owned_child() raises: def test_status_observation_preserves_ownership_and_report() raises: # LC01/LC02: observing an exited child must not consume the report. + var guard_107 = CleanupGuard() var stub = _owned_report_child( + guard_107, 0, "result ok phase=complete case=- reason=ok requests=1 connections=1\n", ) @@ -1524,9 +1688,12 @@ def test_status_observation_preserves_ownership_and_report() raises: assert_equal(stub.request_count(), 1) assert_true(stub.status().cleanup_proved()) + guard_107.assert_clean() + def test_status_observation_then_terminate_is_safe() raises: - var stub = _owned_report_child(0, "") + var guard_108 = CleanupGuard() + var stub = _owned_report_child(guard_108, 0, "") sleep_ms(100) _ = stub.status() var owned_pid = stub.pid @@ -1536,14 +1703,16 @@ def test_status_observation_then_terminate_is_safe() raises: assert_true(pid_not_waitable(owned_pid)) assert_true(not stub.ok()) + # ── LC05: coalesced header cap accounting ─────────────────────────────────── -# ── LC05: coalesced header cap accounting ─────────────────────────────────── + guard_108.assert_clean() def test_coalesced_large_body_does_not_charge_header_cap() raises: var scripts = List[ExchangeScript]() scripts.append(_default_script()) - with spawn_max_local_scripted(0, scripts^) as stub: + var guard_38 = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard_38) as stub: var header_filler = String("") for _ in range(20000): header_filler += "a" @@ -1563,8 +1732,9 @@ def test_coalesced_large_body_does_not_charge_header_cap() raises: assert_true(response.find("200") >= 0) stub.wait() + # ── PC01/PC02/PC04: complete report truth and retained ownership ───────────── -# ── PC01/PC02/PC04: complete report truth and retained ownership ───────────── + guard_38.assert_clean() comptime VALID_REPORT = ( @@ -1577,30 +1747,37 @@ def _report_controls( ) raises: """Every report-framing control must fail for its cause on BOTH providers. """ - var stub = _owned_report_child(exit_code, report) + var guard_109 = CleanupGuard() + var stub = _owned_report_child(guard_109, 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 guard_110 = CleanupGuard() + var jev_stub = _owned_jev_report_child(guard_110, 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)) + guard_109.assert_clean() + guard_110.assert_clean() + 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) + var guard_111 = CleanupGuard() + var stub = _owned_report_child(guard_111, 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) + var guard_112 = CleanupGuard() + var jev_stub = _owned_jev_report_child(guard_112, 0, VALID_REPORT) jev_stub.reap() assert_true(jev_stub.ok()) assert_equal(jev_stub.request_count(), 1) @@ -1619,10 +1796,12 @@ def test_report_stream_controls_both_providers() raises: split_first += tail assert_equal(split_first.byte_length(), 700) _report_controls(split_first + VALID_REPORT, 0, "duplicate_report") - var split_positive = _owned_report_child(0, split_first) + var guard_113 = CleanupGuard() + var split_positive = _owned_report_child(guard_113, 0, split_first) split_positive.reap() assert_true(split_positive.ok()) - var jev_split = _owned_jev_report_child(0, split_first) + var guard_114 = CleanupGuard() + var jev_split = _owned_jev_report_child(guard_114, 0, split_first) jev_split.reap() assert_true(jev_split.ok()) var unterminated = String( @@ -1685,29 +1864,45 @@ def test_report_stream_controls_both_providers() raises: oversize += tail _report_controls(oversize, 0, "ready_output_overflow") + guard_111.assert_clean() + guard_112.assert_clean() + guard_113.assert_clean() + guard_114.assert_clean() + def test_startup_failure_cleanup_ownership_is_truthful() raises: - # PC02/LC01: the shared startup-failure finalizer used by both provider + # PC02/RA02: the shared startup-failure finalizer used by both provider # spawners must not claim an unproved termination as reaped; it records the - # exact pid and reap status in the caller-owned ledger. - var recorded = List[String]() - var ledger = CleanupLedger(UnsafePointer(to=recorded)) - var state = piped_child_state(0, -1, 100, 1, ledger) + # exact pid/reap status in the required guard and retains a usable + # ownership handle for recovery. + var guard = CleanupGuard() + var state = piped_child_state(0, -1, 100, 1, UnsafePointer(to=guard)) var status = finalize_owned_failure(state, 0, "startup cleanup unproved") assert_true(not status.cleanup_proved()) assert_true(not state.reaped) assert_true(state.cleanup_error.startswith("unreaped")) - assert_equal(len(recorded), 1) - assert_true(recorded[0].find("startup cleanup unproved") >= 0) - assert_true(recorded[0].find("pid=0") >= 0) + assert_equal(guard.count(), 1) + assert_equal(guard.retained(), 1) + assert_true(guard.first().find("startup cleanup unproved") >= 0) + assert_true(guard.first().find("pid=0") >= 0) - var stub = _owned_report_child(0, "") + var stub = _owned_report_child(guard, 0, VALID_REPORT) var owned = stub.pid var proved = finalize_owned_failure(stub.state, owned, "startup cleanup") assert_true(proved.cleanup_proved()) assert_true(stub.state.reaped) assert_true(pid_not_waitable(owned)) - assert_equal(len(recorded), 1) + assert_equal(guard.count(), 1) + # The real child was collected, so only the synthetic entry above remains. + assert_true(not guard.is_clean()) + # The required check reports that unresolved entry truthfully rather than + # silently discarding it. + var synthetic = "" + try: + guard.assert_clean() + except e: + synthetic = String(e) + assert_true(synthetic.find("cleanup-unproved") >= 0) def test_descriptor_read_error_is_distinct_from_eof() raises: @@ -1723,7 +1918,10 @@ def test_descriptor_read_error_is_distinct_from_eof() raises: 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 guard = CleanupGuard() + var state = piped_child_state( + closed_fd, closed_fd, 200, 1, UnsafePointer(to=guard) + ) var state_message = "" try: _ = state.read_line(64, 200) @@ -1731,6 +1929,7 @@ def test_descriptor_read_error_is_distinct_from_eof() raises: state_message = String(e) assert_equal(state_message, "read_error") assert_true(not state.last_terminated) + guard.assert_clean() def test_multibyte_surplus_is_not_decoded_prematurely() raises: @@ -1741,12 +1940,13 @@ def test_multibyte_surplus_is_not_decoded_prematurely() raises: var lead = List[UInt8]() lead.append(UInt8(0xC3)) _ = write_raw_bytes(pipe.write_fd, lead) + var guard = CleanupGuard() var state = piped_child_state( pid=0, report_fd=pipe.read_fd, deadline_ms=500, expected_requests=1, - ledger=CleanupLedger(), + guard=UnsafePointer(to=guard), ) var ready = state.read_line(64, 500) assert_equal(ready, "ready 4242") @@ -1760,68 +1960,215 @@ def test_multibyte_surplus_is_not_decoded_prematurely() raises: close_fd(pipe.write_fd) assert_equal(letter, "\u00e9") assert_true(state.last_terminated) + guard.assert_clean() 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) + # RA02: a controlled cleanup failure on a real owned child (the fault seam + # reports an unproved cleanup while the real forked child keeps running) + # must not mark it reaped or discard ownership; the retry collects the exact + # same identity and the earlier entry is resolved, not left as a message. + var guard = CleanupGuard() + var fd_before = open_fd_count_checked() + var stub = spawn_max_local_stub(0, "count_requests", 1, guard, 2000) var actual = stub.pid - stub.pid = 0 + stub.state.faults.cleanup_failures = 1 stub.cleanup() assert_true(stub.cleanup_error().find("unreaped") >= 0) - assert_equal(len(recorded), 1) - stub.pid = actual + assert_true(not stub.status().cleanup_proved()) + assert_equal(guard.count(), 1) + assert_equal(guard.retained(), 1) + assert_true(pid_running(actual)) stub.cleanup() assert_true(stub.status().cleanup_proved()) assert_true(pid_not_waitable(actual)) + guard.assert_clean() + assert_true(open_fd_count_checked() <= fd_before) - var jev_stub = spawn_jev_stub_auto("ok", 1, 2000, ledger) + var jev_guard = CleanupGuard() + var jev_fd_before = open_fd_count_checked() + var jev_stub = spawn_jev_stub_auto("ok", 1, jev_guard, 2000) var jev_actual = jev_stub.stub.pid - jev_stub.stub.pid = 0 + jev_stub.stub.state.faults.cleanup_failures = 1 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 + assert_equal(jev_guard.retained(), 1) + assert_true(pid_running(jev_actual)) jev_stub.stub.cleanup() assert_true(jev_stub.stub.status().cleanup_proved()) assert_true(pid_not_waitable(jev_actual)) + jev_guard.assert_clean() + assert_true(open_fd_count_checked() <= jev_fd_before) + + +def test_cleanup_recovery_through_guard_both_providers() raises: + # RA02: a startup/scope cleanup failure must retain a *usable* ownership + # handle on the required guard, and recovery must collect the exact real + # child rather than only reporting a string. + var guard = CleanupGuard() + var fd_before = open_fd_count_checked() + var stub = spawn_max_local_stub(0, "count_requests", 1, guard, 2000) + var actual = stub.pid + stub.state.faults.cleanup_failures = 1 + stub.cleanup() + assert_equal(guard.retained(), 1) + assert_equal(guard.recover_all(), 0) + assert_true(pid_not_waitable(actual)) + guard.assert_clean() + # Recovery closes the retained report descriptor exactly once. + assert_true(open_fd_count_checked() <= fd_before) + + var jev_guard = CleanupGuard() + var jev_fd_before = open_fd_count_checked() + var jev_stub = spawn_jev_stub_auto("ok", 1, jev_guard, 2000) + var jev_actual = jev_stub.stub.pid + jev_stub.stub.state.faults.cleanup_failures = 1 + jev_stub.stub.cleanup() + assert_equal(jev_guard.retained(), 1) + assert_equal(jev_guard.recover_all(), 0) + assert_true(pid_not_waitable(jev_actual)) + jev_guard.assert_clean() + assert_true(open_fd_count_checked() <= jev_fd_before) def test_reap_wait_error_retains_ownership_both_providers() raises: - # PC02: an unproved/uncertain wait consumed by reap() must not mark the - # child collected or discard retryable ownership; a later retry with the - # restored exact identity still collects it. - var recorded = List[String]() - var ledger = CleanupLedger(UnsafePointer(to=recorded)) - var stub = spawn_max_local_stub(0, "count_requests", 1, 2000, ledger) + # RA02: a transient wait error consumed by reap() must stay retryable and + # must not become a cached terminal result; the retry reports success. + var guard = CleanupGuard() + var stub = _owned_report_child(guard, 0, VALID_REPORT) var actual = stub.pid - stub.pid = 0 + stub.state.faults.wait_errors = 1 stub.reap() assert_true(not stub.ok()) + assert_equal(stub.reason(), "wait_error") assert_true(not stub.status().cleanup_proved()) - assert_equal(len(recorded), 1) - stub.pid = actual - stub.cleanup() + assert_equal(guard.retained(), 1) + # The transient error must not be cached as a terminal observation. + assert_true(stub.status().state != "wait_error") + # A retry must observe the real child and still decode its valid report. + stub.reap() + assert_true(stub.ok()) + assert_equal(stub.phase(), "complete") + assert_equal(stub.request_count(), 1) assert_true(stub.status().cleanup_proved()) assert_true(pid_not_waitable(actual)) + guard.assert_clean() + # A repeated reap is idempotent and does not re-wait or re-read. + stub.reap() + assert_true(stub.ok()) + assert_equal(stub.request_count(), 1) - 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.reap() - assert_true(not jev_stub.stub.ok()) - assert_true(not jev_stub.stub.status().cleanup_proved()) - assert_equal(len(recorded), 2) - jev_stub.stub.pid = jev_actual - jev_stub.stub.cleanup() - assert_true(jev_stub.stub.status().cleanup_proved()) + var jev_guard = CleanupGuard() + var jev_stub = _owned_jev_report_child(jev_guard, 0, VALID_REPORT) + var jev_actual = jev_stub.pid + jev_stub.state.faults.wait_errors = 1 + jev_stub.reap() + assert_true(not jev_stub.ok()) + assert_equal(jev_stub.reason(), "wait_error") + assert_equal(jev_guard.retained(), 1) + jev_stub.reap() + assert_true(jev_stub.ok()) + assert_equal(jev_stub.request_count(), 1) + assert_true(pid_not_waitable(jev_actual)) + jev_guard.assert_clean() + + +def test_unexpected_nonterminal_status_fails_closed_both_providers() raises: + # RA02: an unexpected nonterminal wait status must fail closed and keep the + # exact ownership instead of falling through to a report success. + var guard = CleanupGuard() + var stub = _owned_report_child(guard, 0, VALID_REPORT) + var actual = stub.pid + stub.state.faults.nonterminal = 1 + stub.reap() + assert_true(not stub.ok()) + assert_equal(stub.reason(), "unexpected_status") + assert_true(not stub.status().cleanup_proved()) + assert_equal(guard.retained(), 1) + # Recovery still collects the exact owned child. + guard.recover_all() + guard.assert_clean() + assert_true(pid_not_waitable(actual)) + + var jev_guard = CleanupGuard() + var jev_stub = _owned_jev_report_child(jev_guard, 0, VALID_REPORT) + var jev_actual = jev_stub.pid + jev_stub.state.faults.nonterminal = 1 + jev_stub.reap() + assert_true(not jev_stub.ok()) + assert_equal(jev_stub.reason(), "unexpected_status") + jev_guard.recover_all() + jev_guard.assert_clean() assert_true(pid_not_waitable(jev_actual)) +def test_cleanup_failure_fails_normal_scope_exit_both_providers() raises: + # RA01: an ordinary supported provider scope that exits normally must not + # silently discard a cleanup failure. The exact-owned child is retained for + # recovery, and the required guard check fails the owning test. + var guard = CleanupGuard() + var fd_before = open_fd_count_checked() + var held_pid = 0 + with spawn_max_local_stub(0, "count_requests", 1, guard) as stub: + held_pid = stub.pid + stub.inject_cleanup_failure() + assert_true(pid_running(held_pid)) + var failure = "" + try: + guard.assert_clean() + except e: + failure = String(e) + assert_true(failure.find("cleanup-unproved") >= 0) + assert_equal(guard.recover_all(), 0) + guard.assert_clean() + assert_true(pid_not_waitable(held_pid)) + assert_true(open_fd_count_checked() <= fd_before) + + var jev_guard = CleanupGuard() + var jev_fd_before = open_fd_count_checked() + var jev_pid = 0 + with spawn_jev_stub_auto("ok", 1, jev_guard) as started: + jev_pid = started.stub.pid + started.stub.inject_cleanup_failure() + var jev_failure = "" + try: + jev_guard.assert_clean() + except e: + jev_failure = String(e) + assert_true(jev_failure.find("cleanup-unproved") >= 0) + assert_equal(jev_guard.recover_all(), 0) + jev_guard.assert_clean() + assert_true(pid_not_waitable(jev_pid)) + assert_true(open_fd_count_checked() <= jev_fd_before) + + +def test_cleanup_failure_separately_exposed_with_body_cause() raises: + # RA01: on the exception path the body/assertion cause must be preserved + # exactly while the cleanup failure is separately exposed by the guard. + var guard = CleanupGuard() + var body_message = "" + var held_pid = 0 + try: + with spawn_max_local_stub(0, "count_requests", 1, guard) as stub: + held_pid = stub.pid + stub.inject_cleanup_failure() + raise Error("intentional scope error") + except e: + body_message = String(e) + # Body cause preserved, not replaced by the cleanup failure. + assert_equal(body_message, "intentional scope error") + assert_equal(guard.retained(), 1) + var failure = "" + try: + guard.assert_clean() + except e: + failure = String(e) + assert_true(failure.find("cleanup-unproved") >= 0) + assert_equal(guard.recover_all(), 0) + guard.assert_clean() + assert_true(pid_not_waitable(held_pid)) + + def test_provider_reap_descriptor_read_error_both_paths() raises: # PC01: the descriptor-read control executes through BOTH provider reap # paths, not only the shared reader: a real owned child with an unavailable @@ -1833,12 +2180,16 @@ def test_provider_reap_descriptor_read_error_both_paths() raises: var pid = fork_pid() if pid == 0: child_exit(0) - var state = piped_child_state(pid, closed_fd, 2000, 1, CleanupLedger()) + var guard = CleanupGuard() + var state = piped_child_state( + pid, closed_fd, 2000, 1, UnsafePointer(to=guard) + ) var stub = SpawnedMaxLocalStub(pid, 0, state^) stub.reap() assert_true(not stub.ok()) assert_equal(stub.reason(), "read_error") assert_true(pid_not_waitable(pid)) + guard.assert_clean() var jev_pipe = make_pipe() var jev_closed_fd = jev_pipe.read_fd @@ -1847,31 +2198,201 @@ def test_provider_reap_descriptor_read_error_both_paths() raises: var jev_pid = fork_pid() if jev_pid == 0: child_exit(0) + var jev_guard = CleanupGuard() var jev_state = piped_child_state( - jev_pid, jev_closed_fd, 2000, 1, CleanupLedger() + jev_pid, jev_closed_fd, 2000, 1, UnsafePointer(to=jev_guard) ) var jev_stub = SpawnedJevStub(jev_pid, 0, jev_state^) jev_stub.reap() assert_true(not jev_stub.ok()) assert_equal(jev_stub.reason(), "read_error") assert_true(pid_not_waitable(jev_pid)) + jev_guard.assert_clean() def test_cleanup_failure_survives_scope_exit() raises: - # PC02: cleanup failure must remain observable after the owning handle is - # destroyed at scope exit, for BOTH provider handles. - var recorded = List[String]() - var ledger = CleanupLedger(UnsafePointer(to=recorded)) - var state = piped_child_state(0, -1, 100, 1, ledger) + # RA01/PC02: cleanup failure must remain observable after the owning handle + # is destroyed at scope exit, for BOTH provider handles, because the guard + # is owned by the calling test rather than the handle. + var guard = CleanupGuard() + var state = piped_child_state(0, -1, 100, 1, UnsafePointer(to=guard)) with SpawnedMaxLocalStub(0, 0, state^) as holder: _ = holder - assert_equal(len(recorded), 1) - assert_true(recorded[0].find("unproved") >= 0) - var jev_state = piped_child_state(0, -1, 100, 1, ledger) + assert_equal(guard.count(), 1) + assert_true(guard.first().find("unproved") >= 0) + var jev_state = piped_child_state(0, -1, 100, 1, UnsafePointer(to=guard)) with SpawnedJevStub(0, 0, jev_state^) as jev_holder: _ = jev_holder - assert_equal(len(recorded), 2) - assert_true(recorded[1].find("unproved") >= 0) + assert_equal(guard.count(), 2) + assert_true(guard.first().find("unproved") >= 0) + var failure = "" + try: + guard.assert_clean() + except e: + failure = String(e) + assert_true(failure.find("cleanup-unproved") >= 0) + + +@fieldwise_init +struct ForkedReportChild(Movable): + """A real owned child whose bounded report emission is controlled by delay. + """ + + var pid: Int + var read_fd: Int + + +def _fork_budget_child( + delay_ms: Int, report: String +) raises -> ForkedReportChild: + """Fork a real owned child that writes ``report`` after ``delay_ms``.""" + var pipe = make_pipe() + var pid = fork_pid() + if pid == 0: + close_fd(pipe.read_fd) + sleep_ms(delay_ms) + if report != "": + _ = write_raw(pipe.write_fd, report) + close_fd(pipe.write_fd) + child_exit(0) + close_fd(pipe.write_fd) + return ForkedReportChild(pid, pipe.read_fd) + + +def test_report_after_work_budget_is_rejected_both_providers() raises: + # RA03: the period-9 counterexample. The child emits a *valid* report after + # 200 ms while the declared budget is 100 ms and the parent observes it at + # 300 ms; the late report must be rejected and the exact-owned child + # collected within the bounded cleanup allowance, never accepted with a + # fresh success interval. + for provider in range(2): + var guard = CleanupGuard() + var probe = _fork_budget_child(200, VALID_REPORT) + var start = now_ms() + var state = piped_child_state( + probe.pid, probe.read_fd, 100, 1, UnsafePointer(to=guard) + ) + # The parent consumes its declared budget without observing the child, + # exactly as in the period-9 counterexample. + sleep_ms(300) + if provider == 0: + var stub = SpawnedMaxLocalStub(probe.pid, 0, state^) + stub.reap() + assert_true(not stub.ok()) + assert_equal(stub.reason(), "work_budget_expired") + assert_true(stub.status().cleanup_proved()) + else: + var stub = SpawnedJevStub(probe.pid, 0, state^) + stub.reap() + assert_true(not stub.ok()) + assert_equal(stub.reason(), "work_budget_expired") + assert_true(stub.status().cleanup_proved()) + var elapsed = now_ms() - start + # Only the bounded cleanup allowance is added after the budget expired. + assert_true(elapsed <= 100 + 2000 + 500) + assert_true(pid_not_waitable(probe.pid)) + guard.assert_clean() + + +def test_report_within_work_budget_succeeds_both_providers() raises: + # RA03 positive control: the same budgeted path accepts a report that is + # emitted and observed inside the declared budget. + for provider in range(2): + var guard = CleanupGuard() + var probe = _fork_budget_child(20, VALID_REPORT) + var state = piped_child_state( + probe.pid, probe.read_fd, 4000, 1, UnsafePointer(to=guard) + ) + if provider == 0: + var stub = SpawnedMaxLocalStub(probe.pid, 0, state^) + stub.reap() + assert_true(stub.ok()) + assert_equal(stub.phase(), "complete") + assert_equal(stub.request_count(), 1) + else: + var stub = SpawnedJevStub(probe.pid, 0, state^) + stub.reap() + assert_true(stub.ok()) + assert_equal(stub.phase(), "complete") + assert_equal(stub.request_count(), 1) + assert_true(pid_not_waitable(probe.pid)) + guard.assert_clean() + + +def test_report_drain_deadline_is_bounded() raises: + # RA03: the report EOF-drain stage honours a finite deadline instead of a + # fresh unbounded interval; a report descriptor that never reaches EOF must + # fail with the drain deadline cause, never hang or succeed. + var pipe = make_pipe() + _ = write_raw( + pipe.write_fd, + "result ok phase=complete case=- reason=ok requests=1 connections=1\n", + ) + var guard = CleanupGuard() + var state = piped_child_state( + 0, pipe.read_fd, 500, 1, UnsafePointer(to=guard) + ) + var line = state.read_line(STRICT_MAX_REPORT_BYTES, 500) + assert_true(state.last_terminated) + assert_true(line.startswith("result ok")) + var reason = "" + try: + _ = state.drain_surplus(STRICT_MAX_REPORT_BYTES, 60) + except e: + reason = String(e) + state.close_reader() + close_fd(pipe.write_fd) + assert_equal(reason, "read_deadline_expired") + guard.assert_clean() + + +def test_work_budget_delay_rejects_before_report_both_providers() raises: + # RA03: time consumed inside the wait phase counts against the same finite + # budget, so an exhausted budget rejects even when a valid report is already + # buffered on the owned descriptor. + for provider in range(2): + var guard = CleanupGuard() + var probe = _fork_budget_child(0, VALID_REPORT) + var state = piped_child_state( + probe.pid, probe.read_fd, 400, 1, UnsafePointer(to=guard) + ) + state.faults.wait_delay_ms = 900 + if provider == 0: + var stub = SpawnedMaxLocalStub(probe.pid, 0, state^) + stub.reap() + assert_true(not stub.ok()) + assert_equal(stub.reason(), "work_budget_expired") + assert_true(stub.status().cleanup_proved()) + else: + var stub = SpawnedJevStub(probe.pid, 0, state^) + stub.reap() + assert_true(not stub.ok()) + assert_equal(stub.reason(), "work_budget_expired") + assert_true(stub.status().cleanup_proved()) + assert_true(pid_not_waitable(probe.pid)) + guard.assert_clean() + + +def test_descriptor_census_error_classes() raises: + # RA04: EBADF is a closed slot, EINTR is a bounded retry, and any other + # lookup error makes the census unavailable instead of silently lowering + # the count. The seam changes no host limit. + assert_equal(classify_census_errno(9), "closed") + assert_equal(classify_census_errno(4), "retry") + assert_equal(classify_census_errno(5), "unavailable") + var before = open_fd_count_checked() + assert_true(before > 0) + # EBADF at an open slot is a closed slot: the count is simply lower. + assert_equal(descriptor_census_with_faults(3, 0, 9), 2) + # Any other lookup error makes the whole census unavailable. + assert_equal(descriptor_census_with_faults(3, 0, 5), -1) + var unavailable = "" + try: + _ = open_fd_count_checked_with_faults(3, 0, 5) + except e: + unavailable = String(e) + assert_equal(unavailable, "descriptor_census_unavailable") + assert_equal(open_fd_count_checked(), before) def test_descriptor_census_unavailable_propagates() raises: diff --git a/tests/test_repo_local_process_contract.mojo b/tests/test_repo_local_process_contract.mojo @@ -4,12 +4,23 @@ from safe_tempdir import SafeTempDir from json import Value from fixture_assertions import load_scenario_request_json -from parent_lifecycle import POLLIN, close_fd, make_pipe, now_ms, write_raw +from parent_lifecycle import ( + IO_DEADLINE_EXPIRED, + IO_FAULT_EINTR_UNBOUNDED, + POLLIN, + close_fd, + make_pipe, + now_ms, + read_fd, + write_fd_chunk, + write_raw, +) from stdio_process_helper import ( HYF_PATHS_PROFILE_ENV, HYF_PATHS_REPO_LOCAL_ROOT_ENV, ScopedEnvVar, drain_ready, + run_stdio_binary_with_deadline, run_stdio_entrypoint, run_stdio_entrypoint_with_deadline, ) @@ -127,9 +138,11 @@ def test_run_stdio_entrypoint_rejects_unread_request() raises: def test_run_stdio_entrypoint_rejects_late_success() raises: - # PC03: a child that emits valid JSON, closes stdio and delays exit past the - # declared budget must be rejected within the budget; the bounded cleanup - # grace is not extra successful work. + # PC03: the ordinary compile/run helper keeps its declared total budget, so + # a late exit is rejected inside it and the bounded cleanup grace is not + # extra successful work. The phase-isolated proof of the late-exit child is + # test_run_stdio_binary_rejects_late_exit_after_wait_phase below; this lane + # only characterizes the one-shot compile/run helper's own budget. var start = now_ms() var message = "" try: @@ -150,6 +163,50 @@ def test_run_stdio_entrypoint_rejects_late_success() raises: assert_true(elapsed <= 3700) +comptime LATE_EXIT_SH_SCRIPT = String( + "printf '{\"ok\":true}\\n'; exec 1>&- 2>&-; exec sleep 3" +) + + +def test_run_stdio_binary_rejects_late_exit_after_wait_phase() raises: + # RA03 phase isolation: an already-launched control child (compilation is + # entirely outside this assertion) must reach valid output, closed output + # and the wait phase *before* its late exit is rejected. A generic + # compile/startup/read timeout cannot satisfy these phase assertions. + var start = now_ms() + var message = "" + try: + _ = run_stdio_binary_with_deadline( + "/bin/sh", "{}", "-c", LATE_EXIT_SH_SCRIPT, 1200 + ) + except e: + message = String(e) + var elapsed = now_ms() - start + assert_true(message.find("stdio-entrypoint") >= 0) + assert_true(message.find("phase=wait") >= 0) + assert_true(message.find("request_sent=true") >= 0) + assert_true(message.find("stdout_eof=true") >= 0) + assert_true(message.find("stderr_eof=true") >= 0) + assert_true(message.find("output_valid=true") >= 0) + assert_true(message.find("output_bytes=") >= 0) + assert_true(message.find("reason=timeout") >= 0) + assert_true(message.find("cleanup_error=") >= 0) + assert_true(message.find("child_failed") < 0) + # Work time stays inside the declared budget; only the bounded cleanup + # allowance is added afterwards. + assert_true(elapsed >= 1150) + assert_true(elapsed <= 1200 + 2500) + + +def test_run_stdio_binary_reports_success_within_budget() raises: + # RA03 positive control for the same launch seam: a child that emits valid + # output and exits inside the budget is accepted. + var response = run_stdio_binary_with_deadline( + "/bin/sh", "{}", "-c", "printf '{\"ok\":true}\\n'", 5000 + ) + assert_true(response["ok"].bool_value()) + + def test_run_stdio_entrypoint_classifies_loader_failure() raises: var message = "" try: @@ -197,7 +254,7 @@ def test_diagnostics_overflow_is_bounded_with_cause() raises: var overflow_reason = "" for _ in range(20): _ = write_raw(pipe.write_fd, chunk) - var d = drain_ready(pipe.read_fd, out, 65536, POLLIN) + var d = drain_ready(pipe.read_fd, out, 65536, POLLIN, 2000) if d.reason != "": overflow_reason = d.reason break @@ -210,11 +267,43 @@ def test_diagnostics_overflow_is_bounded_with_cause() raises: 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) + var d = drain_ready(-1, out, 16, POLLIN, 2000) assert_true(d.eof) assert_equal(d.reason, "read_error") +def test_diagnostics_retry_is_bounded_by_deadline() raises: + # RA03: an EINTR-style retry inside the stdio drain path must return to the + # deadline owner instead of retrying without bound. The retry count is a + # bounded test seam; no host signal state is involved. + var pipe = make_pipe() + var buf = InlineArray[Byte, 8](fill=0) + var start = now_ms() + var rc = read_fd( + pipe.read_fd, buf.unsafe_ptr(), 8, 60, IO_FAULT_EINTR_UNBOUNDED + ) + var elapsed = now_ms() - start + close_fd(pipe.read_fd) + close_fd(pipe.write_fd) + assert_equal(rc, IO_DEADLINE_EXPIRED) + assert_true(elapsed >= 50) + + +def test_diagnostics_write_retry_is_bounded_by_deadline() raises: + # RA03: the same bound holds for the stdio request-write retry loop. + var pipe = make_pipe() + var payload = String('{"ok":true}') + var start = now_ms() + var cw = write_fd_chunk( + pipe.write_fd, payload, 0, 60, IO_FAULT_EINTR_UNBOUNDED + ) + var elapsed = now_ms() - start + close_fd(pipe.read_fd) + close_fd(pipe.write_fd) + assert_equal(cw.reason, "write_deadline_expired") + assert_true(elapsed >= 50) + + 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. diff --git a/tests/test_stdio_contract.mojo b/tests/test_stdio_contract.mojo @@ -10,6 +10,7 @@ from fixture_assertions import ( load_scenario_request_json, status_request_with_invalid_version_json, ) +from parent_lifecycle import CleanupGuard from max_local_process_helper import ( reserve_loopback_port, spawn_max_local_stub, @@ -198,7 +199,8 @@ def _assert_query_rewrite_provider_fallback_with_deadline( requests: Int, ) raises: with SafeTempDir() as temp_dir: - with spawn_max_local_stub(0, mode, requests) as provider_stub: + var guard_1 = CleanupGuard() + with spawn_max_local_stub(0, mode, requests, guard_1) as provider_stub: var provider_port = provider_stub.port var startup_config_path = ( Path(temp_dir) / "explicit-hyf-config.toml" @@ -253,6 +255,8 @@ def _assert_query_rewrite_provider_fallback_with_deadline( provider_stub.wait() + guard_1.assert_clean() + def _assert_query_rewrite_provider_fallback( mode: String, expected_reason: String, request_timeout_ms: Int @@ -1003,7 +1007,10 @@ def test_status_reports_unconfigured_assisted_runtime_truthfully() raises: def test_status_reports_non_2xx_max_local_health_truthfully() raises: with SafeTempDir() as temp_dir: - with spawn_max_local_stub(0, "health_non_2xx", 1) as provider_stub: + var guard_2 = CleanupGuard() + with spawn_max_local_stub( + 0, "health_non_2xx", 1, guard_2 + ) as provider_stub: var provider_port = provider_stub.port var startup_config_path = ( Path(temp_dir) / "explicit-hyf-config.toml" @@ -1046,10 +1053,15 @@ def test_status_reports_non_2xx_max_local_health_truthfully() raises: provider_stub.wait() + guard_2.assert_clean() + def test_status_reports_ready_max_local_provider_truthfully() raises: with SafeTempDir() as temp_dir: - with spawn_max_local_stub(0, "query_rewrite_ok", 1) as provider_stub: + var guard_3 = CleanupGuard() + with spawn_max_local_stub( + 0, "query_rewrite_ok", 1, guard_3 + ) as provider_stub: var provider_port = provider_stub.port var startup_config_path = ( Path(temp_dir) / "explicit-hyf-config.toml" @@ -1134,10 +1146,15 @@ def test_status_reports_ready_max_local_provider_truthfully() raises: provider_stub.wait() + guard_3.assert_clean() + def test_status_bounds_max_local_health_probe_timeout() raises: with SafeTempDir() as temp_dir: - with spawn_max_local_stub(0, "health_timeout", 1) as provider_stub: + var guard_4 = CleanupGuard() + with spawn_max_local_stub( + 0, "health_timeout", 1, guard_4 + ) as provider_stub: var provider_port = provider_stub.port var startup_config_path = ( Path(temp_dir) / "explicit-hyf-config.toml" @@ -1180,6 +1197,8 @@ def test_status_bounds_max_local_health_probe_timeout() raises: provider_stub.wait() + guard_4.assert_clean() + def test_status_rejects_invalid_max_local_runtime_config() raises: var prefix = ( @@ -1517,7 +1536,10 @@ def test_capabilities_reports_configured_provider_runtime_truthfully() raises: def test_capabilities_reports_ready_max_local_provider_truthfully() raises: with SafeTempDir() as temp_dir: - with spawn_max_local_stub(0, "query_rewrite_ok", 1) as provider_stub: + var guard_5 = CleanupGuard() + with spawn_max_local_stub( + 0, "query_rewrite_ok", 1, guard_5 + ) as provider_stub: var provider_port = provider_stub.port var startup_config_path = ( Path(temp_dir) / "explicit-hyf-config.toml" @@ -1580,10 +1602,15 @@ def test_capabilities_reports_ready_max_local_provider_truthfully() raises: provider_stub.wait() + guard_5.assert_clean() + def test_capabilities_bounds_max_local_health_probe_timeout() raises: with SafeTempDir() as temp_dir: - with spawn_max_local_stub(0, "health_timeout", 1) as provider_stub: + var guard_6 = CleanupGuard() + with spawn_max_local_stub( + 0, "health_timeout", 1, guard_6 + ) as provider_stub: var provider_port = provider_stub.port var startup_config_path = ( Path(temp_dir) / "explicit-hyf-config.toml" @@ -1634,6 +1661,8 @@ def test_capabilities_bounds_max_local_health_probe_timeout() raises: provider_stub.wait() + guard_6.assert_clean() + def test_query_rewrite_falls_back_deterministically_when_provider_is_unavailable() raises: with SafeTempDir() as temp_dir: @@ -1760,7 +1789,10 @@ def test_assisted_semantic_rank_falls_back_as_unsupported_provider_capability() def test_query_rewrite_uses_max_local_provider_when_ready() raises: with SafeTempDir() as temp_dir: - with spawn_max_local_stub(0, "query_rewrite_ok", 2) as provider_stub: + var guard_7 = CleanupGuard() + with spawn_max_local_stub( + 0, "query_rewrite_ok", 2, guard_7 + ) as provider_stub: var provider_port = provider_stub.port var startup_config_path = ( Path(temp_dir) / "explicit-hyf-config.toml" @@ -1842,6 +1874,8 @@ def test_query_rewrite_uses_max_local_provider_when_ready() raises: provider_stub.wait() + guard_7.assert_clean() + def test_query_rewrite_falls_back_on_provider_non_2xx() raises: _assert_query_rewrite_provider_fallback(