hyf

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

commit 986fa89e553578f3686d83a5fda9f502a7fd1eff
parent b4fae0bcd78143075dacd4f4ae441b86406f3c1a
Author: triesap <tyson@radroots.org>
Date:   Wed, 23 Sep 2026 18:17:55 +0000

H007: complete strict transport characterization with bounded controls

- Replace the delayed-fixture catch-all with explicit close cause and phase.
- Add a parent-bounded exact-owned call runner with a never-returning control.
- Prove headers-before-stall and staged no-reset budget behavior.
- Characterize MaxLocal and Jev delayed success, permitted and unexpected errors.

Diffstat:
Atests/bounded_call_helper.mojo | 216+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mtests/strict_fixture.mojo | 78++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------
Mtests/test_jev.mojo | 142++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------
Mtests/test_provider_adapter.mojo | 252+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mtests/test_runtime_paths.mojo | 16++++++++++++----
5 files changed, 673 insertions(+), 31 deletions(-)

diff --git a/tests/bounded_call_helper.mojo b/tests/bounded_call_helper.mojo @@ -0,0 +1,216 @@ +"""Minimal test-only parent-bounded in-process call runner (H007 TC02). + +An elapsed-time assertion after a synchronous provider call is not a parent +deadline: if the call never returns, the owning test hangs and the observed +behavior cannot be characterized. This helper forks one exact-owned child that +performs the risky provider call and writes a single bounded report line; the +parent enforces one finite deadline, drains that report and, on expiry, +terminates and reaps the exact owned child through the shared lifecycle +primitives. A deliberate never-returning control therefore becomes a bounded, +cause-specific parent observation ("stopped") instead of a hang. + +It is test-only tooling: no product policy, schema, dependency or lock change. +""" + +from std.collections import List + +from json import Value, loads + +from parent_lifecycle import ( + IO_DEADLINE_EXPIRED, + POLLIN, + TERMINATION_GRACE_MS, + CleanupGuard, + bytes_to_string, + child_exit, + close_fd, + fork_owned_or_close, + make_pipe, + now_ms, + poll_fd, + read_fd, + set_alarm, + sleep_ms, + terminate_owned, + write_raw, +) + +from hyf_core.request_context import default_request_context +from hyf_provider.client import post_max_local_chat_completion +from hyf_provider.config import MaxLocalProviderConfig +from hyf_provider.jev_client import post_jev_systemone +from hyf_provider.schema import build_query_rewrite_request_body + + +comptime BOUNDED_CALL_MAX_REPORT_BYTES = 4096 +comptime BOUNDED_CALL_CHILD_ALARM_SECONDS = 120 +comptime BOUNDED_CALL_POLL_SLICE_MS = 25 + + +@fieldwise_init +struct BoundedCallReport(Movable): + """Bounded parent observation of one risky provider call. + + ``completed`` means the child returned and reported inside the parent + deadline; ``stopped`` means the parent deadline expired first and the parent + terminated and reaped the exact owned child. ``report`` is the child's + bounded single-line outcome, and ``cleanup_proved`` records that no waitable + owned child remains. + """ + + var completed: Bool + var stopped: Bool + var report: String + var elapsed_ms: Int + var child_status: String + var cleanup_proved: Bool + + def describe(self) -> String: + return ( + "completed=" + + String(self.completed) + + " stopped=" + + String(self.stopped) + + " elapsed_ms=" + + String(self.elapsed_ms) + + " child=" + + self.child_status + + " cleanup=" + + ("proved" if self.cleanup_proved else "unproved") + + " report=" + + self.report + ) + + +def _child_report(kind: String, port: Int, timeout_ms: Int) raises -> String: + """One bounded report line from inside the forked provider-call child.""" + if kind == "never_return": + # A deliberately non-returning control: it cannot complete within any + # parent deadline, so the parent must stop and reap it. + while True: + sleep_ms(1000) + return "never_returned" + if kind == "max_local": + var config = MaxLocalProviderConfig( + base_url="http://127.0.0.1:" + String(port) + "/v1/", + health_url="http://127.0.0.1:" + String(port) + "/health", + model="max-local-query-rewrite", + request_timeout_ms=timeout_ms, + ) + var context = default_request_context() + var body = build_query_rewrite_request_body( + config, "eggs near me", context + ) + var outcome = post_max_local_chat_completion(config, body) + if outcome.failure: + return ( + "fail max_local " + + outcome.failure.value().kind + + " " + + outcome.failure.value().reason + ) + return "ok max_local " + String(outcome.response.value().status) + if kind == "jev": + var payload = loads('{"model":"jev-1.13.0","state":"s","questions":{}}') + var response = post_jev_systemone( + "http://127.0.0.1:" + String(port), payload, timeout_ms + ) + return "ok jev " + String(response.status) + return "fail unknown_kind" + + +def _has_report_line(bytes: List[UInt8]) -> Bool: + for index in range(len(bytes)): + if Int(bytes[index]) == 10: + return True + return False + + +def run_bounded_call( + kind: String, + port: Int, + timeout_ms: Int, + deadline_ms: Int, + mut guard: CleanupGuard, +) raises -> BoundedCallReport: + """Run one risky provider call under a parent-enforced finite deadline. + + The child is exact-owned: the parent drains its bounded report and then + always collects it (reaping a returned child, terminating a stopping one). + Cleanup is reported through the caller-held guard, so an unproved collection + cannot be silently discarded. + """ + if deadline_ms <= 0: + raise Error("bounded call: invalid deadline") + var pipe = make_pipe() + var pid = fork_owned_or_close(pipe.copy()) + var start = now_ms() + if pid == 0: + close_fd(pipe.read_fd) + _ = set_alarm(BOUNDED_CALL_CHILD_ALARM_SECONDS) + try: + var report = _child_report(kind, port, timeout_ms) + _ = write_raw(pipe.write_fd, report + "\n") + except e: + _ = write_raw(pipe.write_fd, "fail raised\n") + close_fd(pipe.write_fd) + child_exit(0) + + close_fd(pipe.write_fd) + var pending = List[UInt8]() + var eof = False + var timed_out = False + var buffer = InlineArray[Byte, 1024](fill=0) + while not eof: + var remaining = deadline_ms - (now_ms() - start) + if remaining <= 0: + timed_out = True + break + var slice = min(BOUNDED_CALL_POLL_SLICE_MS, remaining) + if slice < 1: + slice = 1 + var events = poll_fd(pipe.read_fd, POLLIN, slice) + if events < 0: + timed_out = True + break + if events == 0: + if _has_report_line(pending): + break + continue + var n = read_fd(pipe.read_fd, buffer.unsafe_ptr(), 1024, remaining) + if n == IO_DEADLINE_EXPIRED or n < 0: + timed_out = True + break + if n == 0: + eof = True + continue + if len(pending) + n > BOUNDED_CALL_MAX_REPORT_BYTES: + timed_out = True + break + for index in range(n): + pending.append(UInt8(Int(buffer[index]))) + if _has_report_line(pending): + break + var report = bytes_to_string(pending) + # Always collect the exact owned child: a returned child is reaped and a + # never-returning child is stopped here by the parent. + var status = terminate_owned(pid, TERMINATION_GRACE_MS) + var cleanup_proved = status.cleanup_proved() + if cleanup_proved: + guard.resolve_pid(pid) + else: + guard.retain( + pid, + pipe.read_fd, + "bounded call cleanup unproved " + status.describe(), + ) + close_fd(pipe.read_fd) + var elapsed_ms = now_ms() - start + return BoundedCallReport( + completed=(not timed_out) and report.byte_length() > 0, + stopped=timed_out, + report=String(report.strip()), + elapsed_ms=elapsed_ms, + child_status=status.describe(), + cleanup_proved=cleanup_proved, + ) diff --git a/tests/strict_fixture.mojo b/tests/strict_fixture.mojo @@ -470,6 +470,10 @@ struct ExchangeScript(Copyable, Movable): var stall_after_head_ms: Int var close_connection: Bool var echo_authorization: Bool + var expect_peer_close: Bool + var expected_close_cause: String + var expected_close_phase: String + var inject_write_error: String def __copyinit__(out self, existing: Self): self.case_label = existing.case_label @@ -487,6 +491,10 @@ struct ExchangeScript(Copyable, Movable): self.stall_after_head_ms = existing.stall_after_head_ms self.close_connection = existing.close_connection self.echo_authorization = existing.echo_authorization + self.expect_peer_close = existing.expect_peer_close + self.expected_close_cause = existing.expected_close_cause + self.expected_close_phase = existing.expected_close_phase + self.inject_write_error = existing.inject_write_error def exchange_script( @@ -512,6 +520,10 @@ def exchange_script( stall_after_head_ms=0, close_connection=True, echo_authorization=False, + expect_peer_close=False, + expected_close_cause="", + expected_close_phase="", + inject_write_error="", ) @@ -609,6 +621,24 @@ def render_response( # ── Convenience-mode validation + reports ──────────────────────────────────── +def classify_write_error_cause(text: String) -> String: + """Bounded cause class for a response-write error observed by the fixture. + + Mojo's error model cannot discriminate handler types, so the exact rendered + cause is classified. A script that declares an expected peer close must + match one of these bounded classes *and* the exact phase; every other write + error (an injected handler failure, a write timeout, an invalid descriptor) + is surfaced as a bounded fixture failure rather than tolerated. + """ + if text.find("ConnectionReset") >= 0: + return "peer_reset" + if text.find("BrokenPipe") >= 0: + return "broken_pipe" + if text.find("Timeout") >= 0: + return "write_timeout" + return "unrelated_error" + + def serve_scripts( listener: TcpListener, var scripts: List[ExchangeScript], label: String ) raises -> ServeReport: @@ -677,8 +707,9 @@ def serve_scripts( request_count = next_index if script.delay_ms > 0: usleep(script.delay_ms * 1000) - var tolerant = ( - script.stall_after_head_ms > 0 or script.delay_ms > 0 + var phase = ( + "body_stall" if script.stall_after_head_ms + > 0 else "delayed_write" ) try: if script.stall_after_head_ms > 0: @@ -699,10 +730,14 @@ def serve_scripts( var head_end = separator + 4 reader.write_all(String(rendered[byte=0:head_end])) usleep(script.stall_after_head_ms * 1000) + if script.inject_write_error != "": + raise Error(script.inject_write_error) reader.write_all(String(rendered[byte=head_end:])) else: reader.write_all(rendered) else: + if script.inject_write_error != "": + raise Error(script.inject_write_error) reader.write_all( render_response( script, @@ -711,12 +746,39 @@ def serve_scripts( connection_count, ) ) - except: - # A deliberate delay or stall lets the peer time out and - # close; that peer close must not turn the bounded stall - # control into a fixture failure. - if not tolerant: - raise + except e: + # A delay/stall is not permission to swallow every write + # error. Only an explicitly declared expected peer close + # whose exact bounded cause and phase match is accepted; + # unexpected/wrong-phase closes, write timeouts, invalid + # descriptors and unrelated handler errors fail the fixture + # with a bounded, cause-specific reason. + var observed = classify_write_error_cause(String(e)) + if ( + script.expect_peer_close + and observed == script.expected_close_cause + and phase == script.expected_close_phase + ): + return ServeReport( + True, + "complete", + script.case_label + + "_peer_close_" + + observed + + "_" + + phase, + "ok", + request_count, + connection_count, + ) + return ServeReport( + False, + "peer_close", + script.case_label, + "unexpected_write_" + observed + "_" + phase, + request_count, + connection_count, + ) if script.close_connection: break if request_count < total: diff --git a/tests/test_jev.mojo b/tests/test_jev.mojo @@ -462,14 +462,33 @@ from jev_provider_helper import ( spawn_jev_scripted_auto, ) from strict_fixture import ExchangeScript, exchange_script +from bounded_call_helper import run_bounded_call from parent_lifecycle import now_ms +def _raw_jev_request_text(path: String) -> String: + return ( + "POST " + + path + + " HTTP/1.1\r\nhost: 127.0.0.1\r\ncontent-length: 2\r\n" + "connection: close\r\n\r\n{}" + ) + + +def _raw_jev_send_then_close(port: Int, path: String) raises: + """Owned raw client that closes right after its request (real peer close). + """ + var client = TcpStream.connect(SocketAddr.localhost(UInt16(port))) + client.write_all(Span[UInt8, _](_raw_jev_request_text(path).as_bytes())) + client.close() + + def test_jev_headers_then_stall_is_bounded() raises: - # H007: current Jev client behavior for a response that sends headers and - # then stalls the body is a bounded transport failure inside the declared - # timeout, not a hang. No provider client policy is changed here; this - # characterizes the pre-migration gap. + # H007/TC01-TC02: the Jev client call against a headers-then-stall peer runs + # under the exact-owned parent-bounded mechanism. Characterized current gap: + # the declared timeout does not bound a body-read stall, so the call returns + # only after the stall completes (elapsed >= 1200 ms) instead of failing + # inside the declared 300 ms timeout. No provider client policy is changed. var guard = CleanupGuard() var scripts = List[ExchangeScript]() var script = exchange_script( @@ -478,22 +497,107 @@ def test_jev_headers_then_stall_is_bounded() raises: script.stall_after_head_ms = 1200 scripts.append(script^) with spawn_jev_scripted_auto(scripts^, guard) as started: - var timeout_ms = 300 - var start = now_ms() - var raised = False - try: - _ = post_jev_systemone( - "http://127.0.0.1:" + String(started.port), - _loads('{"model":"jev-1.13.0","state":"s","questions":{}}'), - timeout_ms, - ) - except: - raised = True - var elapsed = now_ms() - start - assert_true(not raised) - assert_true(elapsed >= 1200) - assert_true(elapsed < 5000) + var report = run_bounded_call("jev", started.port, 300, 5000, guard) + assert_true(report.completed) + assert_true(not report.stopped) + assert_true(report.cleanup_proved) + assert_true(report.elapsed_ms >= 1200) + assert_true(report.elapsed_ms < 5000) + assert_true(report.report.find("ok jev") >= 0) + started.stub.wait() + guard.assert_clean() + + +def test_jev_strict_delayed_success_under_bounded_harness() raises: + # TC01/TC02: a delayed response is not permission to swallow errors; with a + # client budget above the delay the strict scripted Jev exchange succeeds + # and the risky call is parent-bounded. + var guard = CleanupGuard() + var scripts = List[ExchangeScript]() + var script = exchange_script( + "jev_delayed_success", "POST", "/v1/systemone", 200, '{"ok":true}' + ) + script.delay_ms = 300 + scripts.append(script^) + with spawn_jev_scripted_auto(scripts^, guard) as started: + var report = run_bounded_call("jev", started.port, 5000, 5000, guard) + assert_true(report.completed) + assert_true(not report.stopped) + assert_true(report.report.find("ok jev 200") >= 0) + started.stub.wait() + assert_true(started.stub.ok()) + guard.assert_clean() + + +def test_jev_scripted_permitted_peer_close_is_declared() raises: + # TC01: the Jev serve path verifies an explicitly declared expected peer + # close (exact bounded cause and phase) and records the observed outcome. + var guard = CleanupGuard() + var scripts = List[ExchangeScript]() + var script = exchange_script( + "jev_permitted_close", "POST", "/v1/systemone", 200, '{"ok":true}' + ) + script.stall_after_head_ms = 300 + script.expect_peer_close = True + script.expected_close_cause = "broken_pipe" + script.expected_close_phase = "body_stall" + scripts.append(script^) + with spawn_jev_scripted_auto(scripts^, guard) as started: + _raw_jev_send_then_close(started.port, "/v1/systemone") started.stub.wait() + assert_true(started.stub.ok()) + assert_true(started.stub.failure_case().find("_peer_close_") >= 0) + assert_true(started.stub.failure_case().find("_body_stall") >= 0) + guard.assert_clean() + + +def test_jev_scripted_unexpected_peer_close_fails() raises: + # TC01: an undeclared peer close through the Jev serve path fails with a + # bounded, cause-specific reason instead of being swallowed. + var guard = CleanupGuard() + var scripts = List[ExchangeScript]() + var script = exchange_script( + "jev_unexpected_close", "POST", "/v1/systemone", 200, '{"ok":true}' + ) + script.stall_after_head_ms = 300 + scripts.append(script^) + with spawn_jev_scripted_auto(scripts^, guard) as started: + _raw_jev_send_then_close(started.port, "/v1/systemone") + started.stub.reap() + assert_true(not started.stub.ok()) + assert_equal(started.stub.phase(), "peer_close") + assert_true(started.stub.reason().find("unexpected_write_") >= 0) + assert_true(started.stub.reason().find("_body_stall") >= 0) + guard.assert_clean() + + +def test_jev_scripted_injected_write_error_fails() raises: + # TC01: an unrelated injected handler error (here an invalid descriptor) + # must fail the Jev serve path even when a peer close was declared. + var guard = CleanupGuard() + var scripts = List[ExchangeScript]() + var script = exchange_script( + "jev_injected_error", "POST", "/v1/systemone", 200, '{"ok":true}' + ) + script.stall_after_head_ms = 200 + script.expect_peer_close = True + script.expected_close_cause = "broken_pipe" + script.expected_close_phase = "body_stall" + script.inject_write_error = "Bad file descriptor" + scripts.append(script^) + with spawn_jev_scripted_auto(scripts^, guard) as started: + _ = _raw_jev_exchange( + started.port, _raw_jev_request_text("/v1/systemone") + ) + started.stub.reap() + assert_true(not started.stub.ok()) + assert_equal(started.stub.phase(), "peer_close") + assert_true( + started.stub.reason().find( + "unexpected_write_unrelated_error_body_stall" + ) + >= 0 + ) guard.assert_clean() diff --git a/tests/test_provider_adapter.mojo b/tests/test_provider_adapter.mojo @@ -33,8 +33,12 @@ from max_local_process_helper import ( spawn_max_local_scripted, spawn_max_local_stub, ) +from bounded_call_helper import BoundedCallReport, run_bounded_call from strict_fixture import ExchangeScript, exchange_script +from flare.net import SocketAddr +from flare.tcp import TcpStream + def _provider_runtime_config() -> HyfLoadedRuntimeConfig: var runtime = HyfExecutionRuntimeConfig() @@ -485,5 +489,253 @@ def test_provider_delayed_head_is_bounded_and_specific() raises: guard_3.assert_clean() +def _raw_request_text(path: String) -> String: + return ( + "POST " + + path + + " HTTP/1.1\r\nhost: 127.0.0.1\r\ncontent-length: 2\r\n" + "connection: close\r\n\r\n{}" + ) + + +def _raw_send_then_close(port: Int, path: String) raises: + """A deliberately owned raw client that closes right after its request. + + Used to produce a real peer close during the scripted delayed/body-stall + write without any host or policy change. + """ + var client = TcpStream.connect(SocketAddr.localhost(UInt16(port))) + client.write_all(Span[UInt8, _](_raw_request_text(path).as_bytes())) + client.close() + + +def _raw_send_and_read(port: Int, path: String) raises -> String: + var client = TcpStream.connect(SocketAddr.localhost(UInt16(port))) + client.write_all(Span[UInt8, _](_raw_request_text(path).as_bytes())) + var response = String("") + var buffer = InlineArray[Byte, 1024](fill=0) + while True: + var n = client.read(buffer.unsafe_ptr(), 1024) + if n <= 0: + break + response += String( + unsafe_from_utf8=Span(ptr=buffer.unsafe_ptr(), length=Int(n)) + ) + client.close() + return response^ + + +def _delayed_script(label: String, delay_ms: Int) -> ExchangeScript: + var script = exchange_script( + label, "POST", "/v1/chat/completions", 200, '{"choices":[]}' + ) + script.delay_ms = delay_ms + return script^ + + +def test_provider_strict_delayed_success_under_bounded_harness() raises: + # TC01/TC02: a delayed response is not permission to swallow errors. With a + # client budget above the delay the strict scripted exchange must succeed, + # and the risky call runs under the exact-owned parent-bounded mechanism. + var scripts = List[ExchangeScript]() + var script = _delayed_script("delayed_success", 300) + scripts.append(script^) + var guard = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard) as provider_stub: + var report = run_bounded_call( + "max_local", provider_stub.port, 5000, 5000, guard + ) + assert_true(report.completed) + assert_true(not report.stopped) + assert_true(report.cleanup_proved) + assert_true(report.report.find("ok max_local 200") >= 0) + provider_stub.wait() + assert_true(provider_stub.ok()) + assert_equal(provider_stub.request_count(), 1) + assert_equal(provider_stub.connection_count(), 1) + guard.assert_clean() + + +def test_provider_stall_sends_headers_before_body() raises: + # TC02: prove the fixture delivered the response head before the bounded + # body stall, so the stall control characterizes a body-read stall rather + # than an unrelated connect/refusal condition. + var scripts = List[ExchangeScript]() + var script = exchange_script( + "head_then_stall", "POST", "/v1/chat/completions", 200, '{"choices":[]}' + ) + script.stall_after_head_ms = 800 + scripts.append(script^) + var guard = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard) as provider_stub: + var client = TcpStream.connect( + SocketAddr.localhost(UInt16(provider_stub.port)) + ) + var start = now_ms() + client.write_all( + Span[UInt8, _](_raw_request_text("/v1/chat/completions").as_bytes()) + ) + var head = String("") + var buffer = InlineArray[Byte, 1024](fill=0) + while head.find("\r\n\r\n") < 0: + var n = client.read(buffer.unsafe_ptr(), 1024) + assert_true(n > 0) + head += String( + unsafe_from_utf8=Span(ptr=buffer.unsafe_ptr(), length=Int(n)) + ) + var head_ms = now_ms() - start + var body = String("") + while True: + var n2 = client.read(buffer.unsafe_ptr(), 1024) + if n2 <= 0: + break + body += String( + unsafe_from_utf8=Span(ptr=buffer.unsafe_ptr(), length=Int(n2)) + ) + var total_ms = now_ms() - start + client.close() + assert_true(head_ms < 400) + assert_true(total_ms >= 700) + assert_true(body.find("choices") >= 0) + provider_stub.wait() + assert_true(provider_stub.ok()) + guard.assert_clean() + + +def test_provider_scripted_permitted_peer_close_is_declared() raises: + # TC01: a script may declare an expected peer close with an exact bounded + # cause and phase. The fixture verifies that declaration and records the + # observed outcome; it no longer tolerates arbitrary write errors. + var scripts = List[ExchangeScript]() + var script = exchange_script( + "permitted_close", "POST", "/v1/chat/completions", 200, '{"choices":[]}' + ) + script.stall_after_head_ms = 300 + script.expect_peer_close = True + script.expected_close_cause = "broken_pipe" + script.expected_close_phase = "body_stall" + scripts.append(script^) + var guard = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard) as provider_stub: + _raw_send_then_close(provider_stub.port, "/v1/chat/completions") + provider_stub.wait() + assert_true(provider_stub.ok()) + assert_true(provider_stub.failure_case().find("_peer_close_") >= 0) + assert_true(provider_stub.failure_case().find("_body_stall") >= 0) + guard.assert_clean() + + +def test_provider_scripted_unexpected_peer_close_fails() raises: + # TC01: an ordinary delayed/stalled script that never declared an expected + # close must fail with a bounded, cause-specific reason instead of + # swallowing the write error. + var scripts = List[ExchangeScript]() + var script = exchange_script( + "unexpected_close", + "POST", + "/v1/chat/completions", + 200, + '{"choices":[]}', + ) + script.stall_after_head_ms = 300 + scripts.append(script^) + var guard = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard) as provider_stub: + _raw_send_then_close(provider_stub.port, "/v1/chat/completions") + provider_stub.reap() + assert_true(not provider_stub.ok()) + assert_equal(provider_stub.phase(), "peer_close") + assert_true(provider_stub.reason().find("unexpected_write_") >= 0) + assert_true(provider_stub.reason().find("_body_stall") >= 0) + guard.assert_clean() + + +def test_provider_scripted_wrong_phase_close_fails() raises: + # TC01: a declaration whose phase does not match the observed phase is + # rejected, so a close in the wrong phase can never pass as expected. + var scripts = List[ExchangeScript]() + var script = exchange_script( + "wrong_phase", "POST", "/v1/chat/completions", 200, '{"choices":[]}' + ) + script.stall_after_head_ms = 300 + script.expect_peer_close = True + script.expected_close_cause = "broken_pipe" + script.expected_close_phase = "delayed_write" + scripts.append(script^) + var guard = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard) as provider_stub: + _raw_send_then_close(provider_stub.port, "/v1/chat/completions") + provider_stub.reap() + assert_true(not provider_stub.ok()) + assert_true(provider_stub.reason().find("_body_stall") >= 0) + guard.assert_clean() + + +def test_provider_scripted_injected_write_error_fails() raises: + # TC01: an unrelated injected handler error (here a write timeout) must not + # be accepted even when a peer close was declared, because the bounded + # observed cause differs from the declared cause. + var scripts = List[ExchangeScript]() + var script = exchange_script( + "injected_error", "POST", "/v1/chat/completions", 200, '{"choices":[]}' + ) + script.stall_after_head_ms = 200 + script.expect_peer_close = True + script.expected_close_cause = "broken_pipe" + script.expected_close_phase = "body_stall" + script.inject_write_error = "Timeout" + scripts.append(script^) + var guard = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard) as provider_stub: + _ = _raw_send_and_read(provider_stub.port, "/v1/chat/completions") + provider_stub.reap() + assert_true(not provider_stub.ok()) + assert_equal(provider_stub.phase(), "peer_close") + assert_true( + provider_stub.reason().find( + "unexpected_write_write_timeout_body_stall" + ) + >= 0 + ) + guard.assert_clean() + + +def test_provider_stalled_call_is_parent_bounded() raises: + # TC02: a real risky client call against a stalling peer runs under the + # parent deadline and is stopped/reaped by the parent instead of hanging. + var scripts = List[ExchangeScript]() + var script = exchange_script( + "bounded_stall", "POST", "/v1/chat/completions", 200, '{"choices":[]}' + ) + script.stall_after_head_ms = 3000 + scripts.append(script^) + var guard = CleanupGuard() + with spawn_max_local_scripted(0, scripts^, guard) as provider_stub: + var report = run_bounded_call( + "max_local", provider_stub.port, 5000, 600, guard + ) + assert_true(report.stopped) + assert_true(not report.completed) + assert_true(report.cleanup_proved) + assert_true(report.elapsed_ms >= 500) + assert_true(report.elapsed_ms < 3000) + provider_stub.wait() + guard.assert_clean() + + +def test_provider_never_returning_call_is_stopped_and_reaped() raises: + # TC02: a deliberate never-returning control proves the parent can stop and + # reap the call; this is the bounded-harness proof an elapsed assertion + # after a synchronous call cannot provide. + var guard = CleanupGuard() + var report = run_bounded_call("never_return", 0, 300, 400, guard) + assert_true(report.stopped) + assert_true(not report.completed) + assert_true(report.cleanup_proved) + assert_true(report.elapsed_ms >= 350) + assert_true(report.elapsed_ms < 5000) + guard.assert_clean() + + def main() raises: TestSuite.discover_tests[__functions_in_module()]().run() diff --git a/tests/test_runtime_paths.mojo b/tests/test_runtime_paths.mojo @@ -320,10 +320,18 @@ def test_budget_boundaries_are_deterministic() raises: assert_true(not budget_exhausted(exact, 249_999_999)) assert_true(budget_exhausted(exact, 250_000_000)) - # A later stage cannot extend the budget: the same cap is reused. - var stage_two = budget_from_clock(250, 400, 249_000_000) - assert_equal(stage_two.cap_ms, 250) - assert_equal(budget_remaining_ms(stage_two, 249_000_000), 250) + # The same budget object across staged clock advancement must not reset: a + # freshly constructed object would report the full cap again, which cannot + # prove the original budget was preserved (ADR-0020 TC03). + var staged = budget_from_clock(250, 400, 0) + assert_equal(staged.cap_ms, 250) + assert_equal(budget_remaining_ms(staged, 0), 250) + assert_equal(budget_remaining_ms(staged, 200_000_000), 50) + assert_equal(budget_remaining_ms(staged, 249_000_000), 1) + assert_equal(budget_remaining_ms(staged, 250_000_000), 0) + assert_equal(budget_remaining_ms(staged, 400_000_000), 0) + assert_true(not budget_exhausted(staged, 249_999_999)) + assert_true(budget_exhausted(staged, 250_000_000)) from hyf_runtime.jev_composition import compose_jev