commit 8b1c6f9a8dd75115ea2e2033d95b9c45b3c6cc39
parent 35d21546630f1cda8776a3a68044a101e865b15d
Author: triesap <tyson@radroots.org>
Date: Wed, 23 Sep 2026 20:24:00 +0000
H007: bind risky characterization to truthful bounded calls
- Run refusal, delayed-head, body-stall and raw head/body observations
through one parent-enforced finite deadline.
- Require a complete validated report, clean EOF and natural exit zero;
reject nonzero exit, truncation, duplicates, overflow and wait errors.
- Scope-own the exact child and descriptors with cached status, retained
recovery and a bounded cleanup allowance.
- Require an exact declared peer close, reject invalid or absent close
declarations and continue the remaining sequence on both providers.
Diffstat:
4 files changed, 1207 insertions(+), 209 deletions(-)
diff --git a/tests/bounded_call_helper.mojo b/tests/bounded_call_helper.mojo
@@ -1,13 +1,30 @@
-"""Minimal test-only parent-bounded in-process call runner (H007 TC02).
+"""Minimal test-only parent-bounded in-process call runner (H007 BC01-BC03).
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.
+performs the risky provider call or raw header/body observation and writes a
+single bounded report line; the parent enforces one spawn-relative work budget,
+validates the complete report through EOF, verifies a natural child exit zero,
+and, when the budget expires, terminates and reaps the exact owned child
+through the shared lifecycle primitives.
+
+Ownership and cleanup reuse the qualified ``PipedChildState``/``CleanupGuard``
+primitives rather than a private lifecycle: the exact child and its report
+descriptor are scope-owned immediately after spawn, each owned descriptor is
+closed at most once, and an unproved cleanup retains the usable ownership
+handle so the caller can recover it.
+
+The report grammar is one line of ``key=value`` fields validated by the parent:
+
+``report kind=<kind> correlation=<n> outcome=ok status=<code> [timing fields]``
+``report kind=<kind> correlation=<n> outcome=fail cause=<token> reason=<token>``
+
+A wrong kind/correlation, an unknown/duplicate/missing field, a malformed or
+unterminated/duplicated report, surplus bytes through EOF, a read/poll error or
+a non-zero/signaled child exit is a harness failure and can never be reported as
+a completed call. A well-formed ``outcome=fail`` report is a characterized
+domain/provider failure, which is distinct from a harness failure.
It is test-only tooling: no product policy, schema, dependency or lock change.
"""
@@ -16,22 +33,22 @@ from std.collections import List
from json import Value, loads
+from flare.net import SocketAddr
+from flare.tcp import TcpStream
+
from parent_lifecycle import (
- IO_DEADLINE_EXPIRED,
- POLLIN,
TERMINATION_GRACE_MS,
CleanupGuard,
- bytes_to_string,
+ ProcessStatus,
child_exit,
close_fd,
fork_owned_or_close,
make_pipe,
now_ms,
- poll_fd,
+ piped_child_state,
read_fd,
set_alarm,
sleep_ms,
- terminate_owned,
write_raw,
)
@@ -44,52 +61,345 @@ 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
+
+
+def _field_allowed(key: String) -> Bool:
+ if key == "kind" or key == "correlation" or key == "outcome":
+ return True
+ if key == "status" or key == "cause" or key == "reason":
+ return True
+ if key == "latency_ms" or key == "head_ms" or key == "total_ms":
+ return True
+ return key == "body_bytes" or key == "body_match"
+
+
+def _field_numeric(key: String) -> Bool:
+ if key == "status" or key == "correlation" or key == "latency_ms":
+ return True
+ if key == "head_ms" or key == "total_ms":
+ return True
+ return key == "body_bytes"
+
+
+def _all_digits(text: String) -> Bool:
+ if text.byte_length() == 0:
+ return False
+ for byte in text.as_bytes():
+ var b = Int(byte)
+ if b < 48 or b > 57:
+ return False
+ return True
@fieldwise_init
-struct BoundedCallReport(Movable):
- """Bounded parent observation of one risky provider call.
+struct BoundedCallOutcome(Movable):
+ """Validated fields of one bounded report line.
- ``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.
+ ``ok``/``problem`` describe the validation result; ``outcome`` plus the
+ call-specific fields describe the declared call outcome.
"""
- var completed: Bool
- var stopped: Bool
- var report: String
- var elapsed_ms: Int
- var child_status: String
- var cleanup_proved: Bool
+ var ok: Bool
+ var problem: String
+ var kind: String
+ var correlation: Int
+ var outcome: String
+ var status: Int
+ var cause: String
+ var reason: String
+ var latency_ms: Int
+ var head_ms: Int
+ var total_ms: Int
+ var body_bytes: Int
+ var body_match: String
+
+ def domain_failure(self) -> Bool:
+ return self.ok and self.outcome == "fail"
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
+ "kind="
+ + self.kind
+ + " correlation="
+ + String(self.correlation)
+ + " outcome="
+ + self.outcome
+ + " status="
+ + String(self.status)
+ + " cause="
+ + self.cause
+ + " reason="
+ + self.reason
+ + " latency_ms="
+ + String(self.latency_ms)
+ + " head_ms="
+ + String(self.head_ms)
+ + " total_ms="
+ + String(self.total_ms)
+ + " body_bytes="
+ + String(self.body_bytes)
+ + " body_match="
+ + self.body_match
+ )
+
+
+def _empty_outcome(problem: String) -> BoundedCallOutcome:
+ return BoundedCallOutcome(
+ ok=False,
+ problem=problem,
+ kind="",
+ correlation=-1,
+ outcome="",
+ status=0,
+ cause="",
+ reason="",
+ latency_ms=-1,
+ head_ms=-1,
+ total_ms=-1,
+ body_bytes=-1,
+ body_match="",
+ )
+
+
+def parse_bounded_report(
+ text: String, expected_kind: String, correlation: Int
+) raises -> BoundedCallOutcome:
+ """Validate one complete bounded report line for the declared call.
+
+ Rejects a wrong kind or correlation, an unknown/duplicate/missing field, a
+ malformed token and an incompatible outcome (a status on a failure, or a
+ cause/reason on a success). Every rejection is a bounded problem token, so
+ a harness failure is never mistaken for a characterized call outcome.
+ """
+ var tokens = text.strip().split(" ")
+ if len(tokens) < 1 or String(tokens[0]) != "report":
+ return _empty_outcome("report_prefix")
+ var keys = List[String]()
+ var values = List[String]()
+ for index in range(1, len(tokens)):
+ var token = String(tokens[index])
+ var split = token.find("=")
+ if split <= 0 or split == token.byte_length() - 1:
+ return _empty_outcome("report_token_grammar")
+ var key = String(token[byte=0:split])
+ if not _field_allowed(key):
+ return _empty_outcome("report_unknown_field_" + key)
+ for seen in range(len(keys)):
+ if keys[seen] == key:
+ return _empty_outcome("report_duplicate_field_" + key)
+ keys.append(key)
+ values.append(String(token[byte = split + 1 :]))
+ var kind = ""
+ var correlation_text = ""
+ var outcome = ""
+ var status_text = ""
+ var cause = ""
+ var reason = ""
+ var latency_text = ""
+ var head_text = ""
+ var total_text = ""
+ var bytes_text = ""
+ var match_text = ""
+ for index in range(len(keys)):
+ var key = keys[index]
+ var value = values[index]
+ if key == "kind":
+ kind = value
+ elif key == "correlation":
+ correlation_text = value
+ elif key == "outcome":
+ outcome = value
+ elif key == "status":
+ status_text = value
+ elif key == "cause":
+ cause = value
+ elif key == "reason":
+ reason = value
+ elif key == "latency_ms":
+ latency_text = value
+ elif key == "head_ms":
+ head_text = value
+ elif key == "total_ms":
+ total_text = value
+ elif key == "body_bytes":
+ bytes_text = value
+ elif key == "body_match":
+ match_text = value
+ if _field_numeric(key) and not _all_digits(value):
+ return _empty_outcome("report_non_numeric_" + key)
+ if kind == "" or correlation_text == "" or outcome == "":
+ return _empty_outcome("report_missing_field")
+ if kind != expected_kind:
+ return _empty_outcome("report_kind_mismatch")
+ if Int(correlation_text) != correlation:
+ return _empty_outcome("report_correlation_mismatch")
+ if outcome != "ok" and outcome != "fail":
+ return _empty_outcome("report_outcome_unknown_" + outcome)
+ if outcome == "ok":
+ if status_text == "":
+ return _empty_outcome("report_status_missing")
+ var code = Int(status_text)
+ if code < 100 or code > 599:
+ return _empty_outcome("report_status_invalid")
+ if cause != "" or reason != "":
+ return _empty_outcome("report_incompatible_outcome")
+ else:
+ if cause == "" or reason == "":
+ return _empty_outcome("report_cause_missing")
+ if status_text != "":
+ return _empty_outcome("report_incompatible_outcome")
+ return BoundedCallOutcome(
+ ok=True,
+ problem="",
+ kind=kind,
+ correlation=Int(correlation_text),
+ outcome=outcome,
+ status=Int(status_text) if status_text != "" else 0,
+ cause=cause,
+ reason=reason,
+ latency_ms=Int(latency_text) if latency_text != "" else -1,
+ head_ms=Int(head_text) if head_text != "" else -1,
+ total_ms=Int(total_text) if total_text != "" else -1,
+ body_bytes=Int(bytes_text) if bytes_text != "" else -1,
+ body_match=match_text,
+ )
+
+
+def _classify_child_raise(text: String) -> String:
+ """Bounded cause token for a raise observed inside the bounded child.
+
+ The classification is derived from the rendered error text, never from a
+ caller-supplied string, and is reported as a domain outcome — never as a
+ peer close or as a harness success.
+ """
+ if text.find("refused") >= 0 or text.find("Refused") >= 0:
+ return "connection_refused"
+ if text.find("Timeout") >= 0 or text.find("timeout") >= 0:
+ return "timeout"
+ if text.find("descriptor") >= 0 or text.find("EBADF") >= 0:
+ return "invalid_descriptor"
+ return "raised"
+
+
+def _report_prefix(kind: String, correlation: Int) -> String:
+ return "report kind=" + kind + " correlation=" + String(correlation) + " "
+
+
+def _long_token(count: Int) -> String:
+ var text = ""
+ for _ in range(count):
+ text += "x"
+ return text^
+
+
+def _mutant_payload(copies: Int) -> String:
+ """Deliberately misleading child report for the parent-consumer controls.
+
+ This is a *child-producer* mutation only (ADR-0021 BC02): it changes what
+ the forked child writes and never the parent consumer under test, which is
+ the same consumer every real bounded call uses.
+ """
+ var line = "report kind=mutant correlation=0 outcome=ok status=200"
+ var payload = line
+ for index in range(1, copies):
+ payload += "\n" + line
+ return payload^
+
+
+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_status_code(head: String) raises -> Int:
+ var prefix = "HTTP/1.1 "
+ var start = head.find(prefix)
+ if start < 0:
+ return 0
+ var base = start + prefix.byte_length()
+ var length = 0
+ var bytes = head.as_bytes()
+ while base + length < len(bytes):
+ var b = Int(bytes[base + length])
+ if b < 48 or b > 57:
+ break
+ length += 1
+ if length == 0:
+ return 0
+ return Int(String(head[byte = base : base + length]))
+
+
+def _child_raw_head_body(
+ head: String, port: Int, path: String, expect: String
+) raises -> String:
+ """Owned raw client: record head/body arrival timing for a scripted peer.
+
+ The observation is independent of the product caller, which exposes only a
+ parsed response, so a headers-before-body-stall claim can be proved from the
+ wire while still running under the parent's finite deadline.
+ """
+ var client = TcpStream.connect(SocketAddr.localhost(UInt16(port)))
+ var start = now_ms()
+ client.write_all(Span[UInt8, _](_raw_request_text(path).as_bytes()))
+ var head_text = String("")
+ var buffer = InlineArray[Byte, 1024](fill=0)
+ while head_text.find("\r\n\r\n") < 0:
+ var n = client.read(buffer.unsafe_ptr(), 1024)
+ if n <= 0:
+ client.close()
+ return head + "outcome=fail cause=raw_eof reason=head_eof"
+ head_text += 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()
+ var body_match = "unknown"
+ if expect != "":
+ body_match = "yes" if body.find(expect) >= 0 else "no"
+ return (
+ head
+ + "outcome=ok status="
+ + String(_raw_status_code(head_text))
+ + " head_ms="
+ + String(head_ms)
+ + " total_ms="
+ + String(total_ms)
+ + " body_bytes="
+ + String(body.byte_length())
+ + " body_match="
+ + body_match
+ )
-def _child_report(kind: String, port: Int, timeout_ms: Int) raises -> String:
+def _child_report(
+ kind: String,
+ port: Int,
+ timeout_ms: Int,
+ correlation: Int,
+ raw_path: String,
+ raw_expect: String,
+) raises -> String:
"""One bounded report line from inside the forked provider-call child."""
+ var head = _report_prefix(kind, correlation)
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"
+ return head + "outcome=fail cause=never reason=never_return"
if kind == "max_local":
var config = MaxLocalProviderConfig(
base_url="http://127.0.0.1:" + String(port) + "/v1/",
@@ -104,26 +414,104 @@ def _child_report(kind: String, port: Int, timeout_ms: Int) raises -> String:
var outcome = post_max_local_chat_completion(config, body)
if outcome.failure:
return (
- "fail max_local "
+ head
+ + "outcome=fail cause="
+ outcome.failure.value().kind
- + " "
+ + " reason="
+ outcome.failure.value().reason
)
- return "ok max_local " + String(outcome.response.value().status)
+ return (
+ head
+ + "outcome=ok status="
+ + String(outcome.response.value().status)
+ + " latency_ms="
+ + String(outcome.response.value().latency_ms)
+ )
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"
+ return head + "outcome=ok status=" + String(response.status)
+ if kind == "raw_head_body" or kind == "raw_jev_head_body":
+ return _child_raw_head_body(head, port, raw_path, raw_expect)
+ return head + "outcome=fail cause=unknown_kind reason=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
+@fieldwise_init
+struct BoundedCallReport(Movable):
+ """Bounded parent observation of one risky provider call.
+
+ ``completed`` means exactly one complete bounded report was validated, no
+ surplus byte followed it through EOF, and the child exited naturally with
+ status zero inside one spawn-relative work budget. ``outcome`` is the
+ declared call outcome (``ok`` or a characterized ``fail``); ``problem`` is
+ non-empty only for a harness failure (malformed/duplicate/unterminated
+ report, wrong correlation, early EOF, read/poll error, cap overflow,
+ non-zero or signaled exit) and can never be reported as a completed call.
+ ``stopped`` means the parent work budget expired first and the exact owned
+ child was terminated and reaped under the separate bounded cleanup
+ allowance.
+ """
+
+ var completed: Bool
+ var stopped: Bool
+ var outcome: String
+ var status: Int
+ var cause: String
+ var reason: String
+ var latency_ms: Int
+ var head_ms: Int
+ var total_ms: Int
+ var body_bytes: Int
+ var body_match: String
+ var problem: String
+ var report: String
+ var elapsed_ms: Int
+ var child_status: String
+ var cleanup_proved: Bool
+
+ def ok(self) -> Bool:
+ return self.completed and self.outcome == "ok"
+
+ def domain_failure(self) -> Bool:
+ return self.completed and self.outcome == "fail"
+
+ def describe(self) -> String:
+ return (
+ "completed="
+ + String(self.completed)
+ + " stopped="
+ + String(self.stopped)
+ + " outcome="
+ + self.outcome
+ + " status="
+ + String(self.status)
+ + " cause="
+ + self.cause
+ + " reason="
+ + self.reason
+ + " latency_ms="
+ + String(self.latency_ms)
+ + " head_ms="
+ + String(self.head_ms)
+ + " total_ms="
+ + String(self.total_ms)
+ + " body_bytes="
+ + String(self.body_bytes)
+ + " body_match="
+ + self.body_match
+ + " problem="
+ + (self.problem if self.problem != "" else "-")
+ + " elapsed_ms="
+ + String(self.elapsed_ms)
+ + " child="
+ + self.child_status
+ + " cleanup="
+ + ("proved" if self.cleanup_proved else "unproved")
+ + " report="
+ + self.report
+ )
def run_bounded_call(
@@ -132,13 +520,20 @@ def run_bounded_call(
timeout_ms: Int,
deadline_ms: Int,
mut guard: CleanupGuard,
+ correlation: Int,
+ raw_path: String = "/v1/chat/completions",
+ raw_expect: String = "",
+ fault_cleanup_failures: Int = 0,
+ fault_wait_errors: Int = 0,
) 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.
+ One spawn-relative work budget covers the report read, the surplus drain and
+ the natural child exit. Success requires all three plus a validated report;
+ a report alone never justifies success, and a child the parent had to stop
+ is never reported as completed. Cleanup keeps a bounded allowance separate
+ from the work budget and never turns an expired or failed result into
+ success.
"""
if deadline_ms <= 0:
raise Error("bounded call: invalid deadline")
@@ -148,68 +543,161 @@ def run_bounded_call(
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")
+ var payload = ""
+ var terminator = "\n"
+ var exit_code = 0
+ # Child-producer-only controls for the parent consumer: each produces a
+ # deliberately incomplete, duplicated, non-zero-exit or over-running
+ # child result, and the unchanged parent consumer must reject it.
+ if kind == "mutate_exit7":
+ payload = _mutant_payload(1)
+ exit_code = 7
+ elif kind == "mutate_unterminated":
+ payload = _mutant_payload(1)
+ terminator = ""
+ elif kind == "mutate_duplicate":
+ payload = _mutant_payload(2)
+ elif kind == "mutate_huge":
+ # Child-producer-only control: a report line past the bounded cap.
+ payload = (
+ _report_prefix("mutant", correlation)
+ + "outcome=ok status=200 pad="
+ + _long_token(5000)
+ )
+ elif kind == "mutate_valid":
+ # A well-formed child report, used to reach the child-exit wait
+ # phase with the bounded wait-fault seam.
+ payload = (
+ _report_prefix(kind, correlation) + "outcome=ok status=200"
+ )
+ elif kind == "mutate_delayed":
+ _ = write_raw(pipe.write_fd, _mutant_payload(1) + "\n")
+ sleep_ms(2000)
+ close_fd(pipe.write_fd)
+ child_exit(0)
+ else:
+ try:
+ payload = _child_report(
+ kind, port, timeout_ms, correlation, raw_path, raw_expect
+ )
+ except e:
+ payload = (
+ _report_prefix(kind, correlation)
+ + "outcome=fail cause=raised reason="
+ + _classify_child_raise(String(e))
+ )
+ if payload != "":
+ _ = write_raw(pipe.write_fd, payload + terminator)
close_fd(pipe.write_fd)
- child_exit(0)
+ child_exit(exit_code)
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)
+ var state = piped_child_state(
+ pid, pipe.read_fd, deadline_ms, 1, UnsafePointer(to=guard)
+ )
+ # Bounded test-only fault seam: one forced unproved cleanup or transient
+ # wait error against the real exact-owned child, so the retained-ownership
+ # and recovery path is exercised rather than a synthetic identity.
+ state.faults.cleanup_failures = fault_cleanup_failures
+ state.faults.wait_errors = fault_wait_errors
+ var deadline_hit = False
+ var problem = ""
+ var report_text = ""
+ # 1. Exactly one complete, newline-terminated bounded report line.
+ if state.work_remaining_ms() <= 0:
+ deadline_hit = True
+ if not deadline_hit:
+ try:
+ report_text = state.read_line(
+ BOUNDED_CALL_MAX_REPORT_BYTES, state.work_remaining_ms()
+ )
+ except e:
+ var cause = String(e)
+ if cause == "read_deadline_expired":
+ deadline_hit = True
+ else:
+ problem = cause
+ if not deadline_hit and problem == "":
+ if not state.last_terminated:
+ if report_text.byte_length() == 0:
+ problem = "report_early_eof"
+ else:
+ problem = "report_unterminated"
+ # 2. No surplus bytes through EOF.
+ if not deadline_hit and problem == "":
+ try:
+ var surplus = state.drain_surplus(
+ BOUNDED_CALL_MAX_REPORT_BYTES, state.work_remaining_ms()
+ )
+ if surplus > 0:
+ problem = "report_duplicate_report"
+ except e:
+ var cause = String(e)
+ if cause == "read_deadline_expired":
+ deadline_hit = True
+ else:
+ problem = cause
+ # 3. Natural child exit zero inside the same work budget.
+ var status = ProcessStatus("pending", False, -1, 0, 0, "")
+ if not deadline_hit and problem == "":
+ if state.work_remaining_ms() <= 0:
+ deadline_hit = True
+ else:
+ status = state.wait_until(pid, state.work_remaining_ms())
+ if status.state == "running" or status.state == "interrupted":
+ deadline_hit = True
+ elif status.state == "wait_error":
+ problem = "child_wait_error"
+ elif not status.exited:
+ problem = "child_signal_" + String(status.signal)
+ elif status.exit_code != 0:
+ problem = "child_exit_" + String(status.exit_code)
+ # Cleanup: bounded allowance, separate from the work budget. A terminal
+ # status already observed for this exact child is cached and never re-waited
+ # (a second wait could only report "gone"), and an unproved cleanup retains
+ # the usable ownership handle instead of closing its descriptor.
+ if not status.cleanup_proved():
+ var terminal = state.terminate_once(pid, TERMINATION_GRACE_MS)
+ state.status = terminal.copy()
+ status = terminal.copy()
+ if status.cleanup_proved():
+ state.reaped = True
+ state.guard[].resolve_pid(pid)
+ state.close_reader()
else:
- guard.retain(
+ state.record_unproved(
pid,
- pipe.read_fd,
- "bounded call cleanup unproved " + status.describe(),
+ state.report_fd,
+ "bounded call cleanup unproved",
+ "unreaped:" + status.describe(),
)
- close_fd(pipe.read_fd)
+ var cleanup_proved = status.cleanup_proved()
+ # 4. Report validation for the declared call.
+ var parsed = _empty_outcome("")
+ var completed = (not deadline_hit) and problem == ""
+ if completed:
+ try:
+ parsed = parse_bounded_report(report_text, kind, correlation)
+ except:
+ parsed = _empty_outcome("report_parse_raised")
+ if not parsed.ok:
+ completed = False
+ problem = parsed.problem
var elapsed_ms = now_ms() - start
return BoundedCallReport(
- completed=(not timed_out) and report.byte_length() > 0,
- stopped=timed_out,
- report=String(report.strip()),
+ completed=completed,
+ stopped=deadline_hit,
+ outcome=parsed.outcome,
+ status=parsed.status,
+ cause=parsed.cause,
+ reason=parsed.reason,
+ latency_ms=parsed.latency_ms,
+ head_ms=parsed.head_ms,
+ total_ms=parsed.total_ms,
+ body_bytes=parsed.body_bytes,
+ body_match=parsed.body_match,
+ problem=problem,
+ report=String(report_text.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
@@ -654,6 +654,32 @@ def is_peer_close_cause(cause: String) -> Bool:
return cause == "peer_reset" or cause == "broken_pipe"
+def is_expected_close_phase(phase: String) -> Bool:
+ """True only for a write step the fixture can actually be performing.
+
+ EC01: an expected-close declaration must name an observed phase, so an
+ unknown or empty phase is an invalid declaration and cannot be satisfied by
+ any real write step.
+ """
+ return (
+ phase == "delayed_write"
+ or phase == "head_write"
+ or phase == "body_stall"
+ )
+
+
+def is_valid_close_declaration(script: ExchangeScript) -> Bool:
+ """True when an expected-close declaration names a real cause and phase.
+
+ EC01: an invalid cause/phase declaration is rejected before any response
+ work, so a script can never succeed merely because response writing
+ happened to raise an unrelated error.
+ """
+ return is_peer_close_cause(
+ script.expected_close_cause
+ ) and is_expected_close_phase(script.expected_close_phase)
+
+
def serve_scripts(
listener: TcpListener, var scripts: List[ExchangeScript], label: String
) raises -> ServeReport:
@@ -667,11 +693,13 @@ def serve_scripts(
var request_count = 0
var connection_count = 0
var total = len(scripts)
+ var accepted_closes = List[String]()
try:
while request_count < total:
var stream = listener.accept()
connection_count += 1
var reader = ConnectionReader(stream^)
+ var close_connection = False
while request_count < total:
var framed = reader.read()
if not framed.ok:
@@ -705,6 +733,19 @@ def serve_scripts(
request_count,
connection_count,
)
+ # EC01: reject an invalid expected-close declaration before any
+ # response work, even when the write would have succeeded.
+ if script.expect_peer_close and not is_valid_close_declaration(
+ script
+ ):
+ return ServeReport(
+ False,
+ "declaration",
+ script.case_label,
+ "invalid_expected_close_declaration",
+ request_count,
+ connection_count,
+ )
var next_index = request_count + 1
if next_index == total:
var extra = reader.probe_completion(
@@ -787,27 +828,43 @@ def serve_scripts(
and observed == script.expected_close_cause
and phase == script.expected_close_phase
):
- return ServeReport(
- True,
- "complete",
+ # EC01: a permitted close consumes this exchange only.
+ # The connection is finished, so the sequence continues
+ # on a fresh connection and the exact request-count
+ # accounting still requires every remaining script.
+ accepted_closes.append(
script.case_label
+ "_peer_close_"
+ observed
+ "_"
- + phase,
- "ok",
+ + phase
+ )
+ close_connection = True
+ else:
+ return ServeReport(
+ False,
+ "peer_close",
+ script.case_label,
+ "unexpected_write_" + observed + "_" + phase,
request_count,
connection_count,
)
- return ServeReport(
- False,
- "peer_close",
- script.case_label,
- "unexpected_write_" + observed + "_" + phase,
- request_count,
- connection_count,
- )
- if script.close_connection:
+ else:
+ # EC01: the script declared an expected peer close, but the
+ # response write succeeded, so the declared event never
+ # happened. A successful write is not permission to accept
+ # a declared close that was never observed.
+ if script.expect_peer_close:
+ return ServeReport(
+ False,
+ "peer_close",
+ script.case_label,
+ "missing_expected_close_"
+ + script.expected_close_phase,
+ request_count,
+ connection_count,
+ )
+ if script.close_connection or close_connection:
break
if request_count < total:
return ServeReport(
@@ -818,8 +875,15 @@ def serve_scripts(
request_count,
connection_count,
)
+ var case_label = label
+ if len(accepted_closes) > 0:
+ case_label = ""
+ for index in range(len(accepted_closes)):
+ if index > 0:
+ case_label += ";"
+ case_label += accepted_closes[index]
return ServeReport(
- True, "complete", label, "ok", request_count, connection_count
+ True, "complete", case_label, "ok", request_count, connection_count
)
except:
return ServeReport(
diff --git a/tests/test_jev.mojo b/tests/test_jev.mojo
@@ -465,6 +465,10 @@ from strict_fixture import ExchangeScript, exchange_script
from bounded_call_helper import run_bounded_call
from parent_lifecycle import now_ms
+# H007 BC02: a distinct correlation value per bounded-call invocation.
+comptime JEV_CORRELATION_BODY_STALL = 201
+comptime JEV_CORRELATION_DELAYED_SUCCESS = 202
+
def _raw_jev_request_text(path: String) -> String:
return (
@@ -483,6 +487,23 @@ def _raw_jev_send_then_close(port: Int, path: String) raises:
client.close()
+def _raw_jev_send_and_read(port: Int, path: String) raises -> String:
+ """Owned raw client that sends one scripted request and reads the reply."""
+ var client = TcpStream.connect(SocketAddr.localhost(UInt16(port)))
+ client.write_all(Span[UInt8, _](_raw_jev_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 test_jev_headers_then_stall_is_bounded() raises:
# H007/TC01-TC02: the Jev client call against a headers-then-stall peer runs
# under the exact-owned parent-bounded mechanism. Characterized current gap:
@@ -497,13 +518,16 @@ 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 report = run_bounded_call("jev", started.port, 300, 5000, guard)
- assert_true(report.completed)
+ var report = run_bounded_call(
+ "jev", started.port, 300, 5000, guard, JEV_CORRELATION_BODY_STALL
+ )
+ assert_true(report.ok())
assert_true(not report.stopped)
assert_true(report.cleanup_proved)
+ assert_equal(report.status, 200)
+ assert_true(report.latency_ms < 0)
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()
@@ -520,10 +544,17 @@ def test_jev_strict_delayed_success_under_bounded_harness() raises:
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)
+ var report = run_bounded_call(
+ "jev",
+ started.port,
+ 5000,
+ 5000,
+ guard,
+ JEV_CORRELATION_DELAYED_SUCCESS,
+ )
+ assert_true(report.ok())
assert_true(not report.stopped)
- assert_true(report.report.find("ok jev 200") >= 0)
+ assert_equal(report.status, 200)
started.stub.wait()
assert_true(started.stub.ok())
guard.assert_clean()
@@ -571,6 +602,68 @@ def test_jev_scripted_unexpected_peer_close_fails() raises:
guard.assert_clean()
+def test_jev_scripted_missing_expected_close_fails() raises:
+ # EC01: the Jev serve path must reject a declared expected close that never
+ # happened, even though the response write succeeded.
+ var guard = CleanupGuard()
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "jev_missing_expected_close",
+ "POST",
+ "/v1/systemone",
+ 200,
+ '{"ok":true}',
+ )
+ script.expect_peer_close = True
+ script.expected_close_cause = "broken_pipe"
+ script.expected_close_phase = "delayed_write"
+ scripts.append(script^)
+ with spawn_jev_scripted_auto(scripts^, guard) as started:
+ _ = _raw_jev_send_and_read(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("missing_expected_close_delayed_write")
+ >= 0
+ )
+ guard.assert_clean()
+
+
+def test_jev_permitted_close_continues_script_sequence() raises:
+ # EC01: on the Jev serve path a permitted peer close consumes that exchange
+ # only; the remaining scripted exchange is still served and counted.
+ var guard = CleanupGuard()
+ var scripts = List[ExchangeScript]()
+ var closing = exchange_script(
+ "jev_permitted_then_next", "POST", "/v1/systemone", 200, '{"ok":true}'
+ )
+ closing.stall_after_head_ms = 300
+ closing.expect_peer_close = True
+ closing.expected_close_cause = "broken_pipe"
+ closing.expected_close_phase = "body_stall"
+ scripts.append(closing^)
+ scripts.append(
+ exchange_script(
+ "jev_after_permitted_close",
+ "POST",
+ "/v1/systemone",
+ 200,
+ '{"ok":true}',
+ )
+ )
+ with spawn_jev_scripted_auto(scripts^, guard) as started:
+ _raw_jev_send_then_close(started.port, "/v1/systemone")
+ var second = _raw_jev_send_and_read(started.port, "/v1/systemone")
+ started.stub.wait()
+ assert_true(started.stub.ok())
+ assert_equal(started.stub.request_count(), 2)
+ assert_equal(started.stub.connection_count(), 2)
+ assert_true(started.stub.failure_case().find("_peer_close_") >= 0)
+ assert_true(second.find("ok") >= 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.
diff --git a/tests/test_provider_adapter.mojo b/tests/test_provider_adapter.mojo
@@ -27,15 +27,37 @@ from hyf_runtime.config import (
HyfServiceRuntimeConfig,
default_loaded_runtime_config,
)
-from parent_lifecycle import CleanupGuard, now_ms
+from parent_lifecycle import CleanupGuard, now_ms, open_fd_count_checked
from max_local_process_helper import (
reserve_loopback_port,
spawn_max_local_scripted,
spawn_max_local_stub,
)
-from bounded_call_helper import BoundedCallReport, run_bounded_call
+from bounded_call_helper import (
+ BoundedCallReport,
+ parse_bounded_report,
+ run_bounded_call,
+)
from strict_fixture import ExchangeScript, exchange_script
+# H007 BC02: each bounded-call invocation carries its own correlation value, so
+# a report produced for one call can never be accepted for another.
+comptime BOUNDED_CORRELATION_REFUSAL = 101
+comptime BOUNDED_CORRELATION_BODY_STALL = 102
+comptime BOUNDED_CORRELATION_DELAYED_HEAD = 103
+comptime BOUNDED_CORRELATION_DELAYED_SUCCESS = 104
+comptime BOUNDED_CORRELATION_RAW_HEAD_BODY = 105
+comptime BOUNDED_CORRELATION_BOUNDED_STALL = 106
+comptime BOUNDED_CORRELATION_NEVER_RETURN = 107
+comptime BOUNDED_CORRELATION_RETAINED_CLEANUP = 108
+comptime BOUNDED_CORRELATION_REUSED_DESCRIPTOR = 109
+comptime BOUNDED_CORRELATION_MUTANT_EXIT7 = 111
+comptime BOUNDED_CORRELATION_MUTANT_UNTERMINATED = 112
+comptime BOUNDED_CORRELATION_MUTANT_DUPLICATE = 113
+comptime BOUNDED_CORRELATION_MUTANT_DELAYED = 114
+comptime BOUNDED_CORRELATION_MUTANT_HUGE = 115
+comptime BOUNDED_CORRELATION_WAIT_ERROR = 116
+
from flare.net import SocketAddr
from flare.tcp import TcpStream
@@ -394,25 +416,28 @@ def _bounded_timeout_provider_config(
def test_provider_refused_connection_is_bounded_and_specific() raises:
- # H007: a refused connection is a bounded, cause-specific transport failure,
- # not a hang. It is explicitly characterized as a refusal, not as a real
- # connect-timeout scenario, and no provider client policy is changed here.
+ # H007/BC01: a refused connection is a bounded, cause-specific transport
+ # failure, not a hang. The risky product call runs under the parent-enforced
+ # finite deadline through the shared bounded-call consumer, and it is
+ # explicitly characterized as a refusal, not as a real connect-timeout
+ # scenario. No provider client policy is changed here.
var guard = CleanupGuard()
var dead_port = reserve_loopback_port()
- var config = _bounded_timeout_provider_config(dead_port, 300)
- var context = default_request_context()
- var body = build_query_rewrite_request_body(config, "eggs near me", context)
- var start = now_ms()
- var outcome = post_max_local_chat_completion(config, body)
- var elapsed = now_ms() - start
- assert_true(outcome.failure)
- assert_true(not outcome.response)
- assert_equal(outcome.failure.value().kind, "transport")
+ var report = run_bounded_call(
+ "max_local", dead_port, 300, 5000, guard, BOUNDED_CORRELATION_REFUSAL
+ )
+ assert_true(report.completed)
+ assert_true(not report.stopped)
+ assert_true(report.problem == "")
+ assert_true(report.domain_failure())
+ assert_equal(report.outcome, "fail")
+ assert_equal(report.cause, "transport")
# Current characterized gap: a fast refused connection is not distinguished
# from an unknown transport error, because the elapsed time is below the
# declared request budget.
- assert_equal(outcome.failure.value().reason, "unknown_transport")
- assert_true(elapsed < 5000)
+ assert_equal(report.reason, "unknown_transport")
+ assert_true(report.cleanup_proved)
+ assert_true(report.elapsed_ms < 5000)
guard.assert_clean()
@@ -434,21 +459,18 @@ def test_provider_headers_then_stall_is_bounded_and_specific() raises:
var guard_2 = CleanupGuard()
var timeout_ms = 300
with spawn_max_local_scripted(0, scripts^, guard_2) as provider_stub:
- var config = _bounded_timeout_provider_config(
- provider_stub.port, timeout_ms
- )
- var context = default_request_context()
- var body = build_query_rewrite_request_body(
- config, "eggs near me", context
+ var report = run_bounded_call(
+ "max_local",
+ provider_stub.port,
+ timeout_ms,
+ 5000,
+ guard_2,
+ BOUNDED_CORRELATION_BODY_STALL,
)
- var start = now_ms()
- var outcome = post_max_local_chat_completion(config, body)
- var elapsed = now_ms() - start
- assert_true(not outcome.failure)
- assert_true(outcome.response)
- assert_equal(outcome.response.value().status, 200)
- assert_true(elapsed >= 1200)
- assert_true(elapsed < 5000)
+ assert_true(report.ok())
+ assert_equal(report.status, 200)
+ assert_true(report.latency_ms >= 1200)
+ assert_true(report.elapsed_ms < 5000)
provider_stub.wait()
guard_2.assert_clean()
@@ -471,21 +493,18 @@ def test_provider_delayed_head_is_bounded_and_specific() raises:
var guard_3 = CleanupGuard()
var timeout_ms = 300
with spawn_max_local_scripted(0, scripts^, guard_3) as provider_stub:
- var config = _bounded_timeout_provider_config(
- provider_stub.port, timeout_ms
- )
- var context = default_request_context()
- var body = build_query_rewrite_request_body(
- config, "eggs near me", context
+ var report = run_bounded_call(
+ "max_local",
+ provider_stub.port,
+ timeout_ms,
+ 5000,
+ guard_3,
+ BOUNDED_CORRELATION_DELAYED_HEAD,
)
- var start = now_ms()
- var outcome = post_max_local_chat_completion(config, body)
- var elapsed = now_ms() - start
- assert_true(not outcome.failure)
- assert_true(outcome.response)
- assert_equal(outcome.response.value().status, 200)
- assert_true(elapsed >= 1200)
- assert_true(elapsed < 5000)
+ assert_true(report.ok())
+ assert_equal(report.status, 200)
+ assert_true(report.latency_ms >= 1200)
+ assert_true(report.elapsed_ms < 5000)
provider_stub.wait()
guard_3.assert_clean()
@@ -544,12 +563,19 @@ def test_provider_strict_delayed_success_under_bounded_harness() raises:
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
+ "max_local",
+ provider_stub.port,
+ 5000,
+ 5000,
+ guard,
+ BOUNDED_CORRELATION_DELAYED_SUCCESS,
)
- assert_true(report.completed)
+ assert_true(report.ok())
assert_true(not report.stopped)
+ assert_equal(report.status, 200)
+ assert_equal(report.problem, "")
+ assert_true(report.latency_ms >= 300)
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)
@@ -558,9 +584,11 @@ def test_provider_strict_delayed_success_under_bounded_harness() raises:
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.
+ # TC02/BC01: 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. The raw header/body
+ # observation runs through the same parent-bounded consumer as the product
+ # call, so a hanging peer cannot hang the owning test.
var scripts = List[ExchangeScript]()
var script = exchange_script(
"head_then_stall", "POST", "/v1/chat/completions", 200, '{"choices":[]}'
@@ -569,35 +597,22 @@ def test_provider_stall_sends_headers_before_body() raises:
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 report = run_bounded_call(
+ "raw_head_body",
+ provider_stub.port,
+ 5000,
+ 5000,
+ guard,
+ BOUNDED_CORRELATION_RAW_HEAD_BODY,
+ "/v1/chat/completions",
+ "choices",
)
- 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)
+ assert_true(report.ok())
+ assert_equal(report.status, 200)
+ assert_true(report.head_ms < 400)
+ assert_true(report.total_ms >= 700)
+ assert_equal(report.body_match, "yes")
+ assert_true(report.body_bytes > 0)
provider_stub.wait()
assert_true(provider_stub.ok())
guard.assert_clean()
@@ -754,8 +769,111 @@ def test_provider_scripted_declared_non_peer_cause_is_rejected() raises:
_ = _raw_send_and_read(provider_stub.port, "/v1/chat/completions")
provider_stub.reap()
assert_true(not provider_stub.ok())
+ assert_equal(provider_stub.phase(), "declaration")
+ assert_true(
+ provider_stub.reason().find("invalid_expected_close_declaration")
+ >= 0
+ )
+ guard.assert_clean()
+
+
+def test_provider_scripted_invalid_phase_declaration_is_rejected() raises:
+ # EC01: an unknown/empty expected-close phase is an invalid declaration and
+ # cannot be satisfied by any real write step. It is rejected before any
+ # response work even though the response write would have succeeded.
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "invalid_phase",
+ "POST",
+ "/v1/chat/completions",
+ 200,
+ '{"choices":[]}',
+ )
+ script.expect_peer_close = True
+ script.expected_close_cause = "broken_pipe"
+ script.expected_close_phase = "unknown_phase"
+ 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(), "declaration")
+ assert_true(
+ provider_stub.reason().find("invalid_expected_close_declaration")
+ >= 0
+ )
+ guard.assert_clean()
+
+
+def test_provider_scripted_missing_expected_close_fails() raises:
+ # EC01: the script declares an expected peer close but the response write
+ # succeeds, so the declared event never happened. A successful write is not
+ # permission to accept a declared close that was not observed.
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "missing_expected_close",
+ "POST",
+ "/v1/chat/completions",
+ 200,
+ '{"choices":[]}',
+ )
+ 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_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_") >= 0)
+ assert_true(
+ provider_stub.reason().find("missing_expected_close_delayed_write")
+ >= 0
+ )
+ guard.assert_clean()
+
+
+def test_provider_permitted_close_continues_script_sequence() raises:
+ # EC01: a permitted peer close consumes that exchange only. The remaining
+ # scripted exchange must still be served and counted, so a permitted close
+ # can never terminate the whole sequence as successful with unused
+ # exchanges.
+ var scripts = List[ExchangeScript]()
+ var closing = exchange_script(
+ "permitted_then_next",
+ "POST",
+ "/v1/chat/completions",
+ 200,
+ '{"choices":[]}',
+ )
+ closing.stall_after_head_ms = 300
+ closing.expect_peer_close = True
+ closing.expected_close_cause = "broken_pipe"
+ closing.expected_close_phase = "body_stall"
+ scripts.append(closing^)
+ scripts.append(
+ exchange_script(
+ "after_permitted_close",
+ "POST",
+ "/v1/chat/completions",
+ 200,
+ '{"choices":[]}',
+ )
+ )
+ var guard = CleanupGuard()
+ with spawn_max_local_scripted(0, scripts^, guard) as provider_stub:
+ _raw_send_then_close(provider_stub.port, "/v1/chat/completions")
+ var second = _raw_send_and_read(
+ provider_stub.port, "/v1/chat/completions"
+ )
+ provider_stub.wait()
+ assert_true(provider_stub.ok())
+ assert_equal(provider_stub.request_count(), 2)
+ assert_equal(provider_stub.connection_count(), 2)
+ assert_true(provider_stub.failure_case().find("_peer_close_") >= 0)
+ assert_true(second.find("choices") >= 0)
guard.assert_clean()
@@ -771,7 +889,12 @@ def test_provider_stalled_call_is_parent_bounded() raises:
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
+ "max_local",
+ provider_stub.port,
+ 5000,
+ 600,
+ guard,
+ BOUNDED_CORRELATION_BOUNDED_STALL,
)
assert_true(report.stopped)
assert_true(not report.completed)
@@ -787,7 +910,9 @@ def test_provider_never_returning_call_is_stopped_and_reaped() raises:
# 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)
+ var report = run_bounded_call(
+ "never_return", 0, 300, 400, guard, BOUNDED_CORRELATION_NEVER_RETURN
+ )
assert_true(report.stopped)
assert_true(not report.completed)
assert_true(report.cleanup_proved)
@@ -796,6 +921,234 @@ def test_provider_never_returning_call_is_stopped_and_reaped() raises:
guard.assert_clean()
+def test_bounded_call_retained_cleanup_recovers_and_leaks_nothing() raises:
+ # BC03: an unproved bounded-call cleanup must retain usable ownership in the
+ # caller-held guard, recover the exact owned child on retry, and leave no
+ # descriptor or child behind — including a following call whose pipe reuses
+ # the descriptor number the recovery released.
+ var guard = CleanupGuard()
+ var fd_before = open_fd_count_checked()
+ var report = run_bounded_call(
+ "never_return",
+ 0,
+ 300,
+ 400,
+ guard,
+ BOUNDED_CORRELATION_RETAINED_CLEANUP,
+ "/v1/chat/completions",
+ "",
+ 1,
+ 0,
+ )
+ assert_true(report.stopped)
+ assert_true(not report.completed)
+ assert_true(not report.cleanup_proved)
+ assert_true(guard.retained() >= 1)
+ assert_equal(guard.recover_all(), 0)
+ guard.assert_clean()
+ var follow = run_bounded_call(
+ "never_return",
+ 0,
+ 300,
+ 400,
+ guard,
+ BOUNDED_CORRELATION_REUSED_DESCRIPTOR,
+ )
+ assert_true(follow.stopped)
+ assert_true(follow.cleanup_proved)
+ guard.assert_clean()
+ assert_equal(open_fd_count_checked(), fd_before)
+ assert_equal(guard.pending(), 0)
+
+
+def test_bounded_call_rejects_nonzero_exit_after_valid_report() raises:
+ # BC02 before/after control: the period-12 child-only mutation produced a
+ # valid-looking report followed by exit 7 and the old consumer called it
+ # completed. The unchanged parent consumer must reject the non-zero exit
+ # after observing the natural exit, never kill the child and call it done.
+ var guard = CleanupGuard()
+ var report = run_bounded_call(
+ "mutate_exit7",
+ 0,
+ 300,
+ 500,
+ guard,
+ BOUNDED_CORRELATION_MUTANT_EXIT7,
+ )
+ assert_true(not report.completed)
+ assert_true(not report.stopped)
+ assert_equal(report.problem, "child_exit_7")
+ assert_true(report.child_status.find("exited=7") >= 0)
+ assert_true(report.cleanup_proved)
+ guard.assert_clean()
+
+
+def test_bounded_call_rejects_unterminated_report() raises:
+ # BC02 before/after control: an unterminated report is a harness failure,
+ # never a completed call.
+ var guard = CleanupGuard()
+ var report = run_bounded_call(
+ "mutate_unterminated",
+ 0,
+ 300,
+ 500,
+ guard,
+ BOUNDED_CORRELATION_MUTANT_UNTERMINATED,
+ )
+ assert_true(not report.completed)
+ assert_equal(report.problem, "report_unterminated")
+ assert_true(report.cleanup_proved)
+ guard.assert_clean()
+
+
+def test_bounded_call_rejects_duplicate_report() raises:
+ # BC02 before/after control: a second report line is surplus through EOF and
+ # can never be accepted, whatever its chunk alignment.
+ var guard = CleanupGuard()
+ var report = run_bounded_call(
+ "mutate_duplicate",
+ 0,
+ 300,
+ 500,
+ guard,
+ BOUNDED_CORRELATION_MUTANT_DUPLICATE,
+ )
+ assert_true(not report.completed)
+ assert_equal(report.problem, "report_duplicate_report")
+ assert_true(report.cleanup_proved)
+ guard.assert_clean()
+
+
+def test_bounded_call_rejects_overrunning_child_after_report() raises:
+ # BC02/BC03 before/after control: the old consumer reported a valid-looking
+ # completed call after killing a child that stayed alive past the budget.
+ # The repaired consumer must stop and reap it and report an expiry, not
+ # success.
+ var guard = CleanupGuard()
+ var report = run_bounded_call(
+ "mutate_delayed",
+ 0,
+ 300,
+ 500,
+ guard,
+ BOUNDED_CORRELATION_MUTANT_DELAYED,
+ )
+ assert_true(not report.completed)
+ assert_true(report.stopped)
+ assert_true(report.cleanup_proved)
+ assert_true(report.elapsed_ms >= 400)
+ assert_true(report.elapsed_ms < 3000)
+ guard.assert_clean()
+
+
+def test_bounded_report_grammar_rejects_invalid_fields() raises:
+ # BC02: the shared parent consumer validates the declared call grammar. A
+ # wrong correlation, an unknown/duplicate/missing field, a malformed token
+ # or an incompatible outcome is rejected as a bounded problem, never as a
+ # completed call for the declared kind.
+ var good = "report kind=max_local correlation=7 outcome=ok status=200"
+ var accepted = parse_bounded_report(good, "max_local", 7)
+ assert_true(accepted.ok)
+ assert_equal(accepted.status, 200)
+ var failing_kind = parse_bounded_report(
+ (
+ "report kind=jev correlation=7 outcome=fail cause=transport"
+ " reason=unknown_transport"
+ ),
+ "jev",
+ 7,
+ )
+ assert_true(failing_kind.domain_failure())
+ assert_equal(failing_kind.cause, "transport")
+ var cases = List[String]()
+ cases.append(
+ "not_a_report kind=max_local correlation=7 outcome=ok status=200"
+ )
+ cases.append("report kind=max_local correlation=7 outcome=ok")
+ cases.append("report kind=max_local outcome=ok status=200")
+ cases.append(
+ "report kind=max_local correlation=7 outcome=ok status=200 bogus=1"
+ )
+ cases.append(
+ "report kind=max_local correlation=7 outcome=ok status=200 status=200"
+ )
+ cases.append("report kind=jev correlation=7 outcome=ok status=200")
+ cases.append("report kind=max_local correlation=8 outcome=ok status=200")
+ cases.append("report kind=max_local correlation=7 outcome=maybe status=200")
+ cases.append(
+ "report kind=max_local correlation=7 outcome=fail cause=transport "
+ "reason=unknown_transport status=200"
+ )
+ cases.append("report kind=max_local correlation=7 outcome=ok status=abc")
+ cases.append("report kind=max_local correlation=7 outcome=ok status=999")
+ cases.append(
+ "report kind=max_local correlation=7 outcome=fail cause=transport"
+ )
+ cases.append("report x")
+ var expected = List[String]()
+ expected.append("report_prefix")
+ expected.append("report_status_missing")
+ expected.append("report_missing_field")
+ expected.append("report_unknown_field_bogus")
+ expected.append("report_duplicate_field_status")
+ expected.append("report_kind_mismatch")
+ expected.append("report_correlation_mismatch")
+ expected.append("report_outcome_unknown_maybe")
+ expected.append("report_incompatible_outcome")
+ expected.append("report_non_numeric_status")
+ expected.append("report_status_invalid")
+ expected.append("report_cause_missing")
+ expected.append("report_token_grammar")
+ for index in range(len(cases)):
+ var parsed = parse_bounded_report(cases[index], "max_local", 7)
+ assert_true(not parsed.ok)
+ assert_equal(parsed.problem, expected[index])
+
+
+def test_bounded_call_rejects_report_past_the_byte_cap() raises:
+ # BC02: a report line past the bounded cap is a cap-overflow harness
+ # failure, never a completed call.
+ var guard = CleanupGuard()
+ var report = run_bounded_call(
+ "mutate_huge",
+ 0,
+ 300,
+ 500,
+ guard,
+ BOUNDED_CORRELATION_MUTANT_HUGE,
+ )
+ assert_true(not report.completed)
+ assert_equal(report.problem, "ready_output_overflow")
+ assert_true(report.cleanup_proved)
+ guard.assert_clean()
+
+
+def test_bounded_call_transient_wait_error_is_not_success() raises:
+ # BC02/BC03: a transient child-exit wait error is a bounded harness failure
+ # that stays retryable; it can never be reported as a completed call, and
+ # the following cleanup still proves the exact owned child is collected.
+ # The wait fault is a labelled test-only seam (never a real owned child
+ # replaced by a synthetic identity).
+ var guard = CleanupGuard()
+ var report = run_bounded_call(
+ "mutate_valid",
+ 0,
+ 300,
+ 500,
+ guard,
+ BOUNDED_CORRELATION_WAIT_ERROR,
+ "/v1/chat/completions",
+ "",
+ 0,
+ 1,
+ )
+ assert_true(not report.completed)
+ assert_equal(report.problem, "child_wait_error")
+ assert_true(report.cleanup_proved)
+ assert_true(report.child_status.find("exited=0") >= 0)
+ guard.assert_clean()
+
+
def test_maxlocal_wire_attempt_counts_are_exact() raises:
# H008: a retryable non-2xx does not trigger a hidden retry. Each explicit
# call makes exactly one counted wire attempt (request and connection).