commit b0f82196be5f87db02102408f8bb47d397648cb4
parent 60abdb01a4ec3bd2df11b27bfc1110984d514619
Author: triesap <tyson@radroots.org>
Date: Thu, 24 Sep 2026 00:04:13 +0000
H007: complete raw byte accounting, bounded poll retry and real write-error proof
- retain body bytes coalesced with CRLFCRLF and account declared/observed length
- bound poll EINTR retry without refreshing the budget and execute consumer error controls
- exercise a real descriptor/write failure and label synthetic errno seams
- rerun affected lanes after the shared-helper changes; test-only
Diffstat:
6 files changed, 1291 insertions(+), 48 deletions(-)
diff --git a/tests/bounded_call_helper.mojo b/tests/bounded_call_helper.mojo
@@ -37,19 +37,23 @@ from flare.net import SocketAddr
from flare.tcp import TcpStream
from parent_lifecycle import (
+ SIGKILL,
TERMINATION_GRACE_MS,
CleanupGuard,
ProcessStatus,
child_exit,
close_fd,
fork_owned_or_close,
+ kill_pid,
make_pipe,
now_ms,
+ owned_pid,
piped_child_state,
read_fd,
set_alarm,
sleep_ms,
write_raw,
+ write_raw_bytes,
)
from hyf_core.request_context import default_request_context
@@ -70,7 +74,9 @@ def _field_allowed(key: String) -> Bool:
return True
if key == "latency_ms" or key == "head_ms" or key == "total_ms":
return True
- return key == "body_bytes" or key == "body_match"
+ if key == "body_bytes" or key == "body_match":
+ return True
+ return key == "declared_bytes" or key == "length_match"
def _field_numeric(key: String) -> Bool:
@@ -78,7 +84,7 @@ def _field_numeric(key: String) -> Bool:
return True
if key == "head_ms" or key == "total_ms":
return True
- return key == "body_bytes"
+ return key == "body_bytes" or key == "declared_bytes"
def _all_digits(text: String) -> Bool:
@@ -112,6 +118,8 @@ struct BoundedCallOutcome(Movable):
var total_ms: Int
var body_bytes: Int
var body_match: String
+ var declared_bytes: Int
+ var length_match: String
def domain_failure(self) -> Bool:
return self.ok and self.outcome == "fail"
@@ -140,6 +148,10 @@ struct BoundedCallOutcome(Movable):
+ String(self.body_bytes)
+ " body_match="
+ self.body_match
+ + " declared_bytes="
+ + String(self.declared_bytes)
+ + " length_match="
+ + self.length_match
)
@@ -158,6 +170,8 @@ def _empty_outcome(problem: String) -> BoundedCallOutcome:
total_ms=-1,
body_bytes=-1,
body_match="",
+ declared_bytes=-1,
+ length_match="",
)
@@ -200,6 +214,8 @@ def parse_bounded_report(
var total_text = ""
var bytes_text = ""
var match_text = ""
+ var declared_text = ""
+ var length_text = ""
for index in range(len(keys)):
var key = keys[index]
var value = values[index]
@@ -225,6 +241,10 @@ def parse_bounded_report(
bytes_text = value
elif key == "body_match":
match_text = value
+ elif key == "declared_bytes":
+ declared_text = value
+ elif key == "length_match":
+ length_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 == "":
@@ -262,6 +282,8 @@ def parse_bounded_report(
total_ms=Int(total_text) if total_text != "" else -1,
body_bytes=Int(bytes_text) if bytes_text != "" else -1,
body_match=match_text,
+ declared_bytes=Int(declared_text) if declared_text != "" else -1,
+ length_match=length_text,
)
@@ -333,6 +355,166 @@ def _raw_status_code(head: String) raises -> Int:
return Int(String(head[byte = base : base + length]))
+comptime RAW_MAX_HEADER_BYTES: Int = 65536
+comptime RAW_MAX_BODY_BYTES: Int = 1048576
+comptime RAW_READ_CHUNK_BYTES: Int = 1024
+
+
+def _find_header_terminator(bytes: List[UInt8]) -> Int:
+ """Index just past the first CRLFCRLF, or -1 while the head is incomplete.
+ """
+ var n = len(bytes)
+ if n < 4:
+ return -1
+ for index in range(0, n - 3):
+ if (
+ Int(bytes[index]) == 13
+ and Int(bytes[index + 1]) == 10
+ and Int(bytes[index + 2]) == 13
+ and Int(bytes[index + 3]) == 10
+ ):
+ return index + 4
+ return -1
+
+
+def _declared_content_length(head_text: String) raises -> Int:
+ """Lexical Content-Length from a decoded head, or -1 when absent.
+
+ Case-insensitive header name, ASCII digits only, with an explicit length
+ guard. A malformed or oversized value is reported as absent rather than
+ silently truncated; the caller treats that as an unknown declared size.
+ """
+ var marker = "content-length"
+ var m = marker.byte_length()
+ var bytes = head_text.as_bytes()
+ var n = len(bytes)
+ var base = 0
+ while base + m <= n:
+ var matched = True
+ for k in range(m):
+ var hb = Int(bytes[base + k])
+ if hb >= 65 and hb <= 90:
+ hb += 32
+ if hb != Int(marker.as_bytes()[k]):
+ matched = False
+ break
+ if matched:
+ var i = base + m
+ while i < n and (Int(bytes[i]) == 32 or Int(bytes[i]) == 9):
+ i += 1
+ if i < n and Int(bytes[i]) == 58:
+ i += 1
+ while i < n and (Int(bytes[i]) == 32 or Int(bytes[i]) == 9):
+ i += 1
+ var j = i
+ while j < n and Int(bytes[j]) >= 48 and Int(bytes[j]) <= 57:
+ j += 1
+ if j == i or j - i > 9:
+ return -1
+ return Int(String(head_text[byte=i:j]))
+ base += 1
+ return -1
+
+
+def _bytes_contain_from(bytes: List[UInt8], start: Int, needle: String) -> Bool:
+ """Byte-level substring search over an undecoded body buffer.
+
+ Searching raw bytes (not a decoded String) keeps an invalid UTF-8 body or
+ a multi-byte character split across reader chunks from corrupting the
+ observation.
+ """
+ var nlen = needle.byte_length()
+ if nlen == 0:
+ return True
+ var limit = len(bytes) - nlen
+ if limit < start:
+ return False
+ var nb = needle.as_bytes()
+ for base in range(start, limit + 1):
+ var matched = True
+ for k in range(nlen):
+ if Int(bytes[base + k]) != Int(nb[k]):
+ matched = False
+ break
+ if matched:
+ return True
+ return False
+
+
+@fieldwise_init
+struct RawAccounting(Movable):
+ """Exact raw-response accounting from one bounded byte buffer.
+
+ A non-empty ``problem`` is a bounded raw-observation failure (incomplete
+ head or undecodable ASCII head); otherwise ``status``/``body_bytes``/
+ ``body_match``/``declared_bytes``/``length_match`` describe the observation.
+ """
+
+ var problem: String
+ var status: Int
+ var body_bytes: Int
+ var body_match: String
+ var declared_bytes: Int
+ var length_match: String
+
+
+def account_raw_bytes(raw: List[UInt8], expect: String) raises -> RawAccounting:
+ """Account header end, body byte count and length/match from raw bytes.
+
+ The buffer is treated as one byte stream: packet/read boundaries are never
+ protocol boundaries, and the body begins exactly after CRLFCRLF. A valid
+ declared Content-Length is compared with the observed body byte count.
+ """
+ var head_end = _find_header_terminator(raw)
+ if head_end < 0:
+ return RawAccounting(
+ "raw_head_incomplete", 0, 0, "unknown", -1, "unknown"
+ )
+ var head_bytes = List[UInt8]()
+ for index in range(head_end):
+ head_bytes.append(raw[index])
+ var head_text = ""
+ var decoded = True
+ try:
+ head_text = String(
+ from_utf8=Span(ptr=head_bytes.unsafe_ptr(), length=len(head_bytes))
+ )
+ except:
+ decoded = False
+ if not decoded:
+ return RawAccounting("raw_head_decode", 0, 0, "unknown", -1, "unknown")
+ var body_bytes = len(raw) - head_end
+ var declared = _declared_content_length(head_text)
+ var match_text = "unknown"
+ if expect != "":
+ match_text = "yes" if _bytes_contain_from(
+ raw, head_end, expect
+ ) else "no"
+ var length_match = "unknown"
+ if declared >= 0:
+ length_match = "yes" if body_bytes == declared else "no"
+ var status = _raw_status_code(head_text)
+ if status == 0:
+ # A header block with no parseable HTTP/1.1 status line is a malformed
+ # observation, never a successful call with an invented status.
+ return RawAccounting(
+ "raw_status_missing",
+ 0,
+ body_bytes,
+ match_text,
+ declared,
+ length_match,
+ )
+ return RawAccounting(
+ "",
+ status,
+ body_bytes,
+ match_text,
+ declared,
+ length_match,
+ )
+
+
def _child_raw_head_body(
head: String, port: Int, path: String, expect: String
) raises -> String:
@@ -341,46 +523,96 @@ def _child_raw_head_body(
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.
+
+ RP01: bytes coalesced with the header terminator stay with the body, the
+ accumulation is byte-oriented and bounded, and decode happens once over the
+ complete buffer so a split multi-byte character is never misread as a
+ packet boundary. A valid declared Content-Length bounds the read; the
+ observed body byte count is compared with it and reported truthfully.
"""
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)
+ var raw = List[UInt8]()
+ var buffer = InlineArray[Byte, RAW_READ_CHUNK_BYTES](fill=0)
+ var head_end = -1
+ while head_end < 0:
+ var n = client.read(buffer.unsafe_ptr(), RAW_READ_CHUNK_BYTES)
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))
- )
+ for index in range(n):
+ raw.append(UInt8(Int(buffer[index])))
+ if len(raw) > RAW_MAX_HEADER_BYTES:
+ client.close()
+ return (
+ head
+ + "outcome=fail cause=raw_head_overflow reason=head_overflow"
+ )
+ head_end = _find_header_terminator(raw)
var head_ms = now_ms() - start
- var body = String("")
+ var declared = -1
+ var head_bytes = List[UInt8]()
+ for index in range(head_end):
+ head_bytes.append(raw[index])
+ var decoded = True
+ try:
+ declared = _declared_content_length(
+ String(
+ from_utf8=Span(
+ ptr=head_bytes.unsafe_ptr(), length=len(head_bytes)
+ )
+ )
+ )
+ except:
+ decoded = False
+ if not decoded:
+ client.close()
+ return head + "outcome=fail cause=raw_head_decode reason=head_decode"
+ var body_bytes = len(raw) - head_end
while True:
- var n2 = client.read(buffer.unsafe_ptr(), 1024)
+ if declared >= 0 and body_bytes >= declared:
+ break
+ if body_bytes >= RAW_MAX_BODY_BYTES:
+ break
+ var want = min(RAW_READ_CHUNK_BYTES, RAW_MAX_BODY_BYTES - body_bytes)
+ var n2 = client.read(buffer.unsafe_ptr(), want)
if n2 <= 0:
break
- body += String(
- unsafe_from_utf8=Span(ptr=buffer.unsafe_ptr(), length=Int(n2))
- )
+ for index in range(n2):
+ raw.append(UInt8(Int(buffer[index])))
+ body_bytes += 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"
+ var accounting = account_raw_bytes(raw^, expect)
+ if accounting.problem != "":
+ return (
+ head
+ + "outcome=fail cause="
+ + accounting.problem
+ + " reason="
+ + accounting.problem
+ )
+ var declared_field = ""
+ if accounting.declared_bytes >= 0:
+ # A negative declared length is omitted rather than emitted, because the
+ # report grammar accepts only non-negative numeric fields.
+ declared_field = " declared_bytes=" + String(accounting.declared_bytes)
return (
head
+ "outcome=ok status="
- + String(_raw_status_code(head_text))
+ + String(accounting.status)
+ " head_ms="
+ String(head_ms)
+ " total_ms="
+ String(total_ms)
+ " body_bytes="
- + String(body.byte_length())
+ + String(accounting.body_bytes)
+ " body_match="
- + body_match
+ + accounting.body_match
+ + " length_match="
+ + accounting.length_match
+ + declared_field
)
@@ -465,6 +697,8 @@ struct BoundedCallReport(Movable):
var total_ms: Int
var body_bytes: Int
var body_match: String
+ var declared_bytes: Int
+ var length_match: String
var problem: String
var report: String
var elapsed_ms: Int
@@ -501,6 +735,10 @@ struct BoundedCallReport(Movable):
+ String(self.body_bytes)
+ " body_match="
+ self.body_match
+ + " declared_bytes="
+ + String(self.declared_bytes)
+ + " length_match="
+ + self.length_match
+ " problem="
+ (self.problem if self.problem != "" else "-")
+ " elapsed_ms="
@@ -525,6 +763,8 @@ def run_bounded_call(
raw_expect: String = "",
fault_cleanup_failures: Int = 0,
fault_wait_errors: Int = 0,
+ fault_poll_eintrs: Int = 0,
+ fault_poll_errors: Int = 0,
) raises -> BoundedCallReport:
"""Run one risky provider call under a parent-enforced finite deadline.
@@ -546,6 +786,7 @@ def run_bounded_call(
var payload = ""
var terminator = "\n"
var exit_code = 0
+ var self_signal = 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.
@@ -570,6 +811,46 @@ def run_bounded_call(
payload = (
_report_prefix(kind, correlation) + "outcome=ok status=200"
)
+ elif kind == "mutate_silent":
+ # RP02 child-producer-only control: the child closes its report pipe
+ # without writing a byte, so the parent must report an early EOF
+ # rather than an empty completed call.
+ payload = ""
+ elif kind == "mutate_signaled":
+ # RP02 child-producer-only control: a valid report followed by a
+ # signaled termination of the exact owned child; the parent must
+ # report the signal, never a completed call.
+ payload = (
+ _report_prefix(kind, correlation) + "outcome=ok status=200"
+ )
+ self_signal = SIGKILL
+ elif kind == "mutate_invalid_utf8":
+ # RP02 child-producer-only control: a terminated report line whose
+ # body contains an invalid UTF-8 byte, so the parent's bounded
+ # decode rejects it instead of accepting a corrupted line.
+ var bytes = List[UInt8]()
+ var prefix = (
+ _report_prefix(kind, correlation)
+ + "outcome=ok status=200 mark="
+ )
+ for byte in prefix.as_bytes():
+ bytes.append(UInt8(Int(byte)))
+ bytes.append(UInt8(0xFF))
+ bytes.append(UInt8(10))
+ _ = write_raw_bytes(pipe.write_fd, bytes^)
+ close_fd(pipe.write_fd)
+ child_exit(0)
+ elif kind == "mutate_late_exit":
+ # RP02 child-producer-only control: the report pipe is closed after
+ # one complete report while the child stays alive past the budget,
+ # isolating the late-exit wait phase from any report drain.
+ _ = write_raw(
+ pipe.write_fd,
+ _report_prefix(kind, correlation) + "outcome=ok status=200\n",
+ )
+ close_fd(pipe.write_fd)
+ sleep_ms(2000)
+ child_exit(0)
elif kind == "mutate_delayed":
_ = write_raw(pipe.write_fd, _mutant_payload(1) + "\n")
sleep_ms(2000)
@@ -589,6 +870,8 @@ def run_bounded_call(
if payload != "":
_ = write_raw(pipe.write_fd, payload + terminator)
close_fd(pipe.write_fd)
+ if self_signal != 0:
+ _ = kill_pid(owned_pid(), self_signal)
child_exit(exit_code)
close_fd(pipe.write_fd)
@@ -600,6 +883,8 @@ def run_bounded_call(
# and recovery path is exercised rather than a synthetic identity.
state.faults.cleanup_failures = fault_cleanup_failures
state.faults.wait_errors = fault_wait_errors
+ state.faults.poll_eintrs = fault_poll_eintrs
+ state.faults.poll_errors = fault_poll_errors
var deadline_hit = False
var problem = ""
var report_text = ""
@@ -696,6 +981,8 @@ def run_bounded_call(
total_ms=parsed.total_ms,
body_bytes=parsed.body_bytes,
body_match=parsed.body_match,
+ declared_bytes=parsed.declared_bytes,
+ length_match=parsed.length_match,
problem=problem,
report=String(report_text.strip()),
elapsed_ms=elapsed_ms,
diff --git a/tests/parent_lifecycle.mojo b/tests/parent_lifecycle.mojo
@@ -69,6 +69,9 @@ 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
+# Upper bound on the EINTR retry loop when no finite deadline is supplied, so a
+# repeated signal can never spin unbounded; a finite deadline always wins.
+comptime POLL_EINTR_RETRY_BOUND: Int = 64
def now_ms() -> Int:
@@ -219,26 +222,69 @@ def sleep_ms(ms: Int):
_ = external_call["usleep", c_int](c_int(ms * 1000))
-def poll_fd(fd: Int, events: Int, timeout_ms: Int) -> Int:
- """Poll one descriptor.
-
- Returns the ``revents`` mask, ``0`` on timeout and ``-1`` on a real
- ``poll(2)`` error so callers can distinguish a read-phase error from an
- ordinary timeout instead of collapsing both to ``0``.
+def poll_fd(
+ fd: Int,
+ events: Int,
+ timeout_ms: Int,
+ deadline_ms: Int = -1,
+ fault_eintr_count: Int = 0,
+ fault_error_count: Int = 0,
+) -> Int:
+ """Poll one descriptor with a bounded EINTR retry.
+
+ Returns the ``revents`` mask, ``0`` on timeout, ``-1`` on a real
+ ``poll(2)`` error, and ``IO_DEADLINE_EXPIRED`` when a retry after ``EINTR``
+ would outlive the caller-supplied absolute remaining budget. ``deadline_ms``
+ is the remaining part of the caller's budget and is never refreshed by the
+ retry loop, so a signal storm can neither spin forever nor be reported as a
+ real read/poll failure.
+
+ ``fault_eintr_count``/``fault_error_count`` are bounded test-only seams that
+ exercise the actual consumer's retry/error branches; they change no host
+ signal state. A negative ``fault_eintr_count`` retries until the deadline so
+ the bounded-retry branch is deterministically executable.
"""
- var cell = InlineArray[Int32, 2](fill=0)
- cell[0] = Int32(fd)
- cell[1] = Int32(events)
- var n = Int(
- external_call["poll", c_int](
- cell.unsafe_ptr(), c_uint(1), c_int(timeout_ms)
+ var start = now_ms()
+ var eintrs = fault_eintr_count
+ var errors = fault_error_count
+ var attempts = 0
+ while True:
+ var slice = timeout_ms
+ if deadline_ms >= 0:
+ var remaining = deadline_ms - (now_ms() - start)
+ if remaining <= 0:
+ return IO_DEADLINE_EXPIRED
+ if remaining < slice:
+ slice = remaining
+ if errors != 0:
+ if errors > 0:
+ errors -= 1
+ return -1
+ if eintrs != 0:
+ if eintrs > 0:
+ eintrs -= 1
+ attempts += 1
+ if deadline_ms < 0 and attempts > POLL_EINTR_RETRY_BOUND:
+ return -1
+ continue
+ var cell = InlineArray[Int32, 2](fill=0)
+ cell[0] = Int32(fd)
+ cell[1] = Int32(events)
+ var n = Int(
+ external_call["poll", c_int](
+ cell.unsafe_ptr(), c_uint(1), c_int(slice)
+ )
)
- )
- if n < 0:
- return -1
- if n == 0:
- return 0
- return (Int(cell[1]) >> 16) & 0xFFFF
+ if n < 0:
+ if get_errno() != ErrNo.EINTR:
+ return -1
+ attempts += 1
+ if deadline_ms < 0 and attempts > POLL_EINTR_RETRY_BOUND:
+ return -1
+ continue
+ if n == 0:
+ return 0
+ return (Int(cell[1]) >> 16) & 0xFFFF
@fieldwise_init
@@ -390,7 +436,14 @@ def write_fd_bounded(fd: Int, data: String, deadline_ms: Int) -> String:
while sent < total:
if now_ms() - start >= deadline_ms:
return "write_deadline_expired"
- var ev = poll_fd(fd, POLLOUT, LIFECYCLE_POLL_SLICE_MS)
+ var ev = poll_fd(
+ fd,
+ POLLOUT,
+ LIFECYCLE_POLL_SLICE_MS,
+ deadline_ms - (now_ms() - start),
+ )
+ if ev == IO_DEADLINE_EXPIRED:
+ return "write_deadline_expired"
if ev < 0:
return "write_poll_error"
if ev == 0:
@@ -534,7 +587,14 @@ struct BoundedLineReader(Movable):
return bytes_to_string(out)
if now_ms() - start >= deadline_ms:
raise Error("read_deadline_expired")
- var ev = poll_fd(self.fd, POLLIN, LIFECYCLE_POLL_SLICE_MS)
+ var ev = poll_fd(
+ self.fd,
+ POLLIN,
+ LIFECYCLE_POLL_SLICE_MS,
+ deadline_ms - (now_ms() - start),
+ )
+ if ev == IO_DEADLINE_EXPIRED:
+ raise Error("read_deadline_expired")
if ev < 0:
raise Error("read_error")
if ev == 0:
@@ -576,7 +636,14 @@ def read_all_bounded(
while True:
if now_ms() - start >= deadline_ms:
raise Error("read_deadline_expired")
- var ev = poll_fd(fd, POLLIN, LIFECYCLE_POLL_SLICE_MS)
+ var ev = poll_fd(
+ fd,
+ POLLIN,
+ LIFECYCLE_POLL_SLICE_MS,
+ deadline_ms - (now_ms() - start),
+ )
+ if ev == IO_DEADLINE_EXPIRED:
+ raise Error("read_deadline_expired")
if ev < 0:
raise Error("read_error")
if ev == 0:
@@ -967,12 +1034,16 @@ struct LifecycleFaults(Movable):
var nonterminal: Int
var cleanup_failures: Int
var wait_delay_ms: Int
+ var poll_eintrs: Int
+ var poll_errors: Int
def __init__(out self):
self.wait_errors = 0
self.nonterminal = 0
self.cleanup_failures = 0
self.wait_delay_ms = 0
+ self.poll_eintrs = 0
+ self.poll_errors = 0
def active(self) -> Bool:
return (
@@ -980,6 +1051,8 @@ struct LifecycleFaults(Movable):
or self.nonterminal > 0
or self.cleanup_failures > 0
or self.wait_delay_ms > 0
+ or self.poll_eintrs != 0
+ or self.poll_errors != 0
)
@@ -1102,6 +1175,25 @@ struct PipedChildState(Movable):
)
return terminate_owned(pid, grace_ms)
+ def _poll_eintr_fault(mut self) -> Int:
+ """Consume one bounded test-only EINTR fault, if one is armed.
+
+ A positive count is consumed once so a following poll is real; a
+ negative count stays armed so the bounded-retry deadline branch is
+ reached deterministically. This changes no host signal state.
+ """
+ var value = self.faults.poll_eintrs
+ if value > 0:
+ self.faults.poll_eintrs -= 1
+ return value
+
+ def _poll_error_fault(mut self) -> Int:
+ """Consume one bounded test-only real poll-error fault, if armed."""
+ var value = self.faults.poll_errors
+ if value > 0:
+ self.faults.poll_errors -= 1
+ return value
+
def read_line(mut self, max_bytes: Int, deadline_ms: Int) raises -> String:
"""Bounded line read that retains surplus as undecoded bytes.
@@ -1142,7 +1234,16 @@ struct PipedChildState(Movable):
return _utf8_line(line^)
if now_ms() - start >= deadline_ms:
raise Error("read_deadline_expired")
- var ev = poll_fd(self.report_fd, POLLIN, LIFECYCLE_POLL_SLICE_MS)
+ var ev = poll_fd(
+ self.report_fd,
+ POLLIN,
+ LIFECYCLE_POLL_SLICE_MS,
+ deadline_ms - (now_ms() - start),
+ self._poll_eintr_fault(),
+ self._poll_error_fault(),
+ )
+ if ev == IO_DEADLINE_EXPIRED:
+ raise Error("read_deadline_expired")
if ev < 0:
raise Error("read_error")
if ev == 0:
@@ -1180,7 +1281,16 @@ struct PipedChildState(Movable):
while not self.eof:
if now_ms() - start >= deadline_ms:
raise Error("read_deadline_expired")
- var ev = poll_fd(self.report_fd, POLLIN, LIFECYCLE_POLL_SLICE_MS)
+ var ev = poll_fd(
+ self.report_fd,
+ POLLIN,
+ LIFECYCLE_POLL_SLICE_MS,
+ deadline_ms - (now_ms() - start),
+ self._poll_eintr_fault(),
+ self._poll_error_fault(),
+ )
+ if ev == IO_DEADLINE_EXPIRED:
+ raise Error("read_deadline_expired")
if ev < 0:
raise Error("read_error")
if ev == 0:
diff --git a/tests/strict_fixture.mojo b/tests/strict_fixture.mojo
@@ -24,6 +24,7 @@ distinct and verified.
"""
from std.collections import List
+from std.ffi import ErrNo, get_errno
from flare.net import Timeout
from flare.tcp import TcpListener
@@ -474,6 +475,8 @@ struct ExchangeScript(Copyable, Movable):
var expected_close_cause: String
var expected_close_phase: String
var inject_write_error: String
+ var inject_write_errno: Int
+ var close_before_response: Bool
def __copyinit__(out self, existing: Self):
self.case_label = existing.case_label
@@ -495,6 +498,8 @@ struct ExchangeScript(Copyable, Movable):
self.expected_close_cause = existing.expected_close_cause
self.expected_close_phase = existing.expected_close_phase
self.inject_write_error = existing.inject_write_error
+ self.inject_write_errno = existing.inject_write_errno
+ self.close_before_response = existing.close_before_response
def exchange_script(
@@ -524,6 +529,8 @@ def exchange_script(
expected_close_cause="",
expected_close_phase="",
inject_write_error="",
+ inject_write_errno=-1,
+ close_before_response=False,
)
@@ -635,6 +642,8 @@ def classify_write_error_cause(text: String) -> String:
# an injected error can never be accepted by declaring a peer-close
# cause (ADR-0020 TC01).
return "unrelated_error"
+ if text.find("Bad file descriptor") >= 0 or text.find("(errno 9)") >= 0:
+ return "invalid_descriptor"
if text.find("ConnectionReset") >= 0:
return "peer_reset"
if text.find("BrokenPipe") >= 0:
@@ -644,6 +653,47 @@ def classify_write_error_cause(text: String) -> String:
return "unrelated_error"
+def write_errno_class(errno_value: Int) -> String:
+ """Bounded class for a real write(2)/send(2) errno.
+
+ Mirrors the flare ``TcpStream.write`` mapping, so synthetic seams and real
+ failures are classified by the same vocabulary and a timeout or descriptor
+ error can be compared directly against a declared peer-close cause.
+ """
+ if errno_value == Int(ErrNo.EPIPE.value):
+ return "broken_pipe"
+ if errno_value == Int(ErrNo.ECONNRESET.value):
+ return "peer_reset"
+ if errno_value == Int(ErrNo.EAGAIN.value) or errno_value == Int(
+ ErrNo.EWOULDBLOCK.value
+ ):
+ return "write_timeout"
+ if errno_value == Int(ErrNo.EBADF.value):
+ return "invalid_descriptor"
+ return "unrelated_error"
+
+
+def render_write_api_error(errno_value: Int) -> String:
+ """Render an errno exactly as the flare write API would render it.
+
+ Used by the narrowly scoped syscall-result seam so the synthetic failure is
+ classified by the *same* path as a real write failure, and by the real
+ descriptor control to render the observed errno. It is a rendering helper
+ only and changes no product or fork source.
+ """
+ if errno_value == Int(ErrNo.EAGAIN.value) or errno_value == Int(
+ ErrNo.EWOULDBLOCK.value
+ ):
+ return "Timeout: send"
+ if errno_value == Int(ErrNo.EPIPE.value):
+ return "BrokenPipe"
+ if errno_value == Int(ErrNo.ECONNRESET.value):
+ return "ConnectionReset"
+ if errno_value == Int(ErrNo.EBADF.value):
+ return "NetworkError(errno 9): Bad file descriptor (send)"
+ return "NetworkError(errno " + String(errno_value) + "): send error"
+
+
def is_peer_close_cause(cause: String) -> Bool:
"""True only for the bounded peer-close classes a script may expect.
@@ -761,6 +811,13 @@ def serve_scripts(
connection_count,
)
request_count = next_index
+ if script.close_before_response:
+ # RP01 raw-observation control: a peer that closes the
+ # connection without sending a response head at all. The
+ # raw observer must report a bounded raw_eof rather than
+ # hang or invent a status.
+ close_connection = True
+ break
if script.delay_ms > 0:
usleep(script.delay_ms * 1000)
var phase = "delayed_write"
@@ -793,6 +850,13 @@ def serve_scripts(
"injected_error: "
+ script.inject_write_error
)
+ if script.inject_write_errno >= 0:
+ raise Error(
+ "synthetic_errno: "
+ + render_write_api_error(
+ script.inject_write_errno
+ )
+ )
reader.write_all(String(rendered[byte=head_end:]))
else:
phase = "head_write"
@@ -803,6 +867,13 @@ def serve_scripts(
raise Error(
"injected_error: " + script.inject_write_error
)
+ if script.inject_write_errno >= 0:
+ raise Error(
+ "synthetic_errno: "
+ + render_write_api_error(
+ script.inject_write_errno
+ )
+ )
reader.write_all(
render_response(
script,
@@ -821,7 +892,23 @@ def serve_scripts(
# is itself restricted to a real peer-close class, so a
# timeout or unrelated error cannot be waived by declaring
# it as the expected cause.
- var observed = classify_write_error_cause(String(e))
+ # Record the raw errno observed here as well as the
+ # classification. A synthetic seam is labelled explicitly
+ # so it is never presented as a real OS timeout/EBADF.
+ var raw_errno = Int(get_errno().value)
+ var observed_text = String(e)
+ var synthetic = (
+ observed_text.find("synthetic_errno") >= 0
+ or observed_text.find("injected_error") >= 0
+ )
+ var observed = classify_write_error_cause(observed_text)
+ var evidence = (
+ "synthdecl"
+ + String(
+ script.inject_write_errno
+ ) if synthetic else "realerrno"
+ + String(raw_errno)
+ )
if (
script.expect_peer_close
and is_peer_close_cause(script.expected_close_cause)
@@ -838,6 +925,8 @@ def serve_scripts(
+ observed
+ "_"
+ phase
+ + "_"
+ + evidence
)
close_connection = True
else:
@@ -845,7 +934,12 @@ def serve_scripts(
False,
"peer_close",
script.case_label,
- "unexpected_write_" + observed + "_" + phase,
+ "unexpected_write_"
+ + observed
+ + "_"
+ + phase
+ + "_"
+ + evidence,
request_count,
connection_count,
)
diff --git a/tests/test_jev.mojo b/tests/test_jev.mojo
@@ -1,4 +1,5 @@
from std.collections import List
+from std.ffi import ErrNo
from std.testing import TestSuite, assert_equal, assert_raises, assert_true
from hyf_assist.questions import (
@@ -502,13 +503,25 @@ from jev_provider_helper import (
require_bearer_for,
spawn_jev_scripted_auto,
)
-from strict_fixture import ExchangeScript, exchange_script
+from strict_fixture import (
+ ExchangeScript,
+ exchange_script,
+ is_peer_close_cause,
+ write_errno_class,
+)
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
+# H007 RP01: distinct correlations for the repaired Jev raw byte-accounting
+# controls, so the Jev raw caller is executed and not just a shared branch.
+comptime JEV_CORRELATION_RAW_COALESCED = 203
+comptime JEV_CORRELATION_RAW_TRUNCATED = 204
+# H007 RP03: distinct correlations for the Jev write-error controls.
+comptime JEV_CORRELATION_SYNTHETIC_DESCRIPTOR = 205
+comptime JEV_CORRELATION_REAL_ERRNO = 206
def _raw_jev_request_text(path: String) -> String:
@@ -601,6 +614,73 @@ def test_jev_strict_delayed_success_under_bounded_harness() raises:
guard.assert_clean()
+def test_jev_raw_head_body_accounts_coalesced_body() raises:
+ # RP01: the repaired raw observer is executed through the Jev raw caller.
+ # A body coalesced with the header terminator must be counted exactly.
+ var guard = CleanupGuard()
+ var scripts = List[ExchangeScript]()
+ scripts.append(
+ exchange_script(
+ "jev_coalesced_body", "POST", "/v1/systemone", 200, "JEVBODY"
+ )
+ )
+ with spawn_jev_scripted_auto(scripts^, guard) as started:
+ var report = run_bounded_call(
+ "raw_jev_head_body",
+ started.port,
+ 5000,
+ 5000,
+ guard,
+ JEV_CORRELATION_RAW_COALESCED,
+ "/v1/systemone",
+ "JEVBODY",
+ )
+ assert_true(report.ok())
+ assert_equal(report.status, 200)
+ assert_equal(report.body_bytes, 7)
+ assert_equal(report.body_match, "yes")
+ assert_equal(report.declared_bytes, 7)
+ assert_equal(report.length_match, "yes")
+ started.stub.wait()
+ assert_true(started.stub.ok())
+ guard.assert_clean()
+
+
+def test_jev_raw_head_body_reports_truncated_length_mismatch() raises:
+ # RP01: the Jev raw caller reports an explicit truncation truthfully: five
+ # observed bytes against a declared twenty, with a length mismatch.
+ var guard = CleanupGuard()
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "jev_truncated_body", "POST", "/v1/systemone", 200, ""
+ )
+ script.raw_response = (
+ "HTTP/1.1 200 OK\r\ncontent-type: application/json\r\n"
+ "content-length: 20\r\nconnection: close\r\n\r\nSHORT"
+ )
+ scripts.append(script^)
+ with spawn_jev_scripted_auto(scripts^, guard) as started:
+ var report = run_bounded_call(
+ "raw_jev_head_body",
+ started.port,
+ 5000,
+ 5000,
+ guard,
+ JEV_CORRELATION_RAW_TRUNCATED,
+ "/v1/systemone",
+ "SHORT",
+ )
+ assert_true(report.ok())
+ assert_equal(report.status, 200)
+ assert_equal(report.body_bytes, 5)
+ assert_equal(report.body_match, "yes")
+ assert_equal(report.declared_bytes, 20)
+ assert_equal(report.length_match, "no")
+ 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.
@@ -623,6 +703,77 @@ def test_jev_scripted_permitted_peer_close_is_declared() raises:
guard.assert_clean()
+def _jev_realerrno_from_reason(reason: String) raises -> Int:
+ """Raw errno recorded in a bounded Jev write-failure reason, or -1."""
+ var marker = "realerrno"
+ var at = reason.find(marker)
+ if at < 0:
+ return -1
+ var digits = String(reason[byte = at + marker.byte_length() :])
+ var end = digits.find("_")
+ if end >= 0:
+ digits = String(digits[byte=0:end])
+ if digits.byte_length() == 0:
+ return -1
+ return Int(digits)
+
+
+def test_jev_scripted_synthetic_invalid_descriptor_is_not_peer_close() raises:
+ # RP03: the Jev serve path exercises the same narrowly scoped syscall-result
+ # seam, mapped to the flare write API's EBADF rendering. The invalid
+ # descriptor is classified, labelled synthetic and rejected even though a
+ # peer close was declared.
+ var guard = CleanupGuard()
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "jev_synthetic_descriptor", "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_errno = Int(ErrNo.EBADF.value)
+ 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_true(
+ started.stub.reason().find(
+ "unexpected_write_invalid_descriptor_body_stall_synthdecl"
+ )
+ >= 0
+ )
+ guard.assert_clean()
+
+
+def test_jev_scripted_real_peer_close_records_raw_errno() raises:
+ # RP03: a real Jev peer close makes the fixture's real send(2) fail; the
+ # recorded raw errno's class must equal the observed classification.
+ var guard = CleanupGuard()
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "jev_real_errno", "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())
+ var token = started.stub.failure_case()
+ assert_true(token.find("realerrno") >= 0)
+ var errno = _jev_realerrno_from_reason(token)
+ assert_true(errno > 0)
+ var observed = write_errno_class(errno)
+ assert_true(observed == "broken_pipe" or observed == "peer_reset")
+ assert_true(token.find("_peer_close_" + observed + "_") >= 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.
diff --git a/tests/test_provider_adapter.mojo b/tests/test_provider_adapter.mojo
@@ -1,4 +1,5 @@
from std.collections import List
+from std.ffi import ErrNo
from std.testing import TestSuite, assert_equal, assert_raises, assert_true
from json import Value, loads
@@ -35,10 +36,17 @@ from max_local_process_helper import (
)
from bounded_call_helper import (
BoundedCallReport,
+ account_raw_bytes,
parse_bounded_report,
run_bounded_call,
)
-from strict_fixture import ExchangeScript, exchange_script, json_escape
+from strict_fixture import (
+ ExchangeScript,
+ exchange_script,
+ is_peer_close_cause,
+ json_escape,
+ write_errno_class,
+)
# H007 BC02: each bounded-call invocation carries its own correlation value, so
# a report produced for one call can never be accepted for another.
@@ -57,6 +65,27 @@ comptime BOUNDED_CORRELATION_MUTANT_DUPLICATE = 113
comptime BOUNDED_CORRELATION_MUTANT_DELAYED = 114
comptime BOUNDED_CORRELATION_MUTANT_HUGE = 115
comptime BOUNDED_CORRELATION_WAIT_ERROR = 116
+# H007 RP01: distinct correlations for the repaired raw byte-accounting controls.
+comptime BOUNDED_CORRELATION_RAW_COALESCED = 117
+comptime BOUNDED_CORRELATION_RAW_EMPTY = 118
+comptime BOUNDED_CORRELATION_RAW_TRUNCATED = 119
+comptime BOUNDED_CORRELATION_RAW_SPLIT_UTF8 = 120
+# H007 RP02: distinct correlations for the real-consumer error/retry controls.
+comptime BOUNDED_CORRELATION_EARLY_EOF = 121
+comptime BOUNDED_CORRELATION_SIGNALED = 122
+comptime BOUNDED_CORRELATION_INVALID_UTF8 = 123
+comptime BOUNDED_CORRELATION_LATE_EXIT = 124
+comptime BOUNDED_CORRELATION_EINTR_RETRY = 125
+comptime BOUNDED_CORRELATION_EINTR_DEADLINE = 126
+comptime BOUNDED_CORRELATION_POLL_ERROR = 127
+# H007 RP03: distinct correlations for the real/synthetic write-error controls.
+comptime BOUNDED_CORRELATION_SYNTHETIC_TIMEOUT = 128
+comptime BOUNDED_CORRELATION_SYNTHETIC_DESCRIPTOR = 129
+comptime BOUNDED_CORRELATION_REAL_ERRNO = 130
+comptime BOUNDED_CORRELATION_DECLARED_VS_SYNTHETIC = 131
+# H007 RP01: malformed/error raw-observation controls.
+comptime BOUNDED_CORRELATION_RAW_EOF = 132
+comptime BOUNDED_CORRELATION_RAW_MALFORMED = 133
from flare.net import SocketAddr
from flare.tcp import TcpStream
@@ -618,6 +647,250 @@ def test_provider_stall_sends_headers_before_body() raises:
guard.assert_clean()
+def test_provider_raw_head_body_accounts_coalesced_body() raises:
+ # RP01: a body that arrives coalesced with the header terminator must be
+ # accounted as body bytes, never discarded. The scripted exchange writes
+ # head and body together, so the unchanged raw observer must report the
+ # exact eight-byte BODYMARK body with a matching declared length.
+ var scripts = List[ExchangeScript]()
+ scripts.append(
+ exchange_script(
+ "coalesced_body", "POST", "/v1/chat/completions", 200, "BODYMARK"
+ )
+ )
+ var guard = CleanupGuard()
+ with spawn_max_local_scripted(0, scripts^, guard) as provider_stub:
+ var report = run_bounded_call(
+ "raw_head_body",
+ provider_stub.port,
+ 5000,
+ 5000,
+ guard,
+ BOUNDED_CORRELATION_RAW_COALESCED,
+ "/v1/chat/completions",
+ "BODYMARK",
+ )
+ assert_true(report.ok())
+ assert_equal(report.status, 200)
+ assert_equal(report.body_bytes, 8)
+ assert_equal(report.body_match, "yes")
+ assert_equal(report.declared_bytes, 8)
+ assert_equal(report.length_match, "yes")
+ provider_stub.wait()
+ assert_true(provider_stub.ok())
+ guard.assert_clean()
+
+
+def test_provider_raw_head_body_accounts_empty_body() raises:
+ # RP01: a declared zero-length body is an exact observation, not a missing
+ # one: body_bytes is 0 while the declared length still matches.
+ var scripts = List[ExchangeScript]()
+ scripts.append(
+ exchange_script("empty_body", "POST", "/v1/chat/completions", 200, "")
+ )
+ var guard = CleanupGuard()
+ with spawn_max_local_scripted(0, scripts^, guard) as provider_stub:
+ var report = run_bounded_call(
+ "raw_head_body",
+ provider_stub.port,
+ 5000,
+ 5000,
+ guard,
+ BOUNDED_CORRELATION_RAW_EMPTY,
+ "/v1/chat/completions",
+ "",
+ )
+ assert_true(report.ok())
+ assert_equal(report.status, 200)
+ assert_equal(report.body_bytes, 0)
+ assert_equal(report.body_match, "unknown")
+ assert_equal(report.declared_bytes, 0)
+ assert_equal(report.length_match, "yes")
+ provider_stub.wait()
+ assert_true(provider_stub.ok())
+ guard.assert_clean()
+
+
+def test_provider_raw_head_body_reports_truncated_length_mismatch() raises:
+ # RP01: an explicit truncation control. The peer declares twenty body bytes
+ # but sends five and closes; the observer must report the exact five
+ # observed bytes, the twenty declared and a length mismatch, never the
+ # declared length as if it had arrived.
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "truncated_body", "POST", "/v1/chat/completions", 200, ""
+ )
+ script.raw_response = (
+ "HTTP/1.1 200 OK\r\ncontent-type: application/json\r\n"
+ "content-length: 20\r\nconnection: close\r\n\r\nSHORT"
+ )
+ scripts.append(script^)
+ var guard = CleanupGuard()
+ with spawn_max_local_scripted(0, scripts^, guard) as provider_stub:
+ var report = run_bounded_call(
+ "raw_head_body",
+ provider_stub.port,
+ 5000,
+ 5000,
+ guard,
+ BOUNDED_CORRELATION_RAW_TRUNCATED,
+ "/v1/chat/completions",
+ "SHORT",
+ )
+ assert_true(report.ok())
+ assert_equal(report.status, 200)
+ assert_equal(report.body_bytes, 5)
+ assert_equal(report.body_match, "yes")
+ assert_equal(report.declared_bytes, 20)
+ assert_equal(report.length_match, "no")
+ provider_stub.wait()
+ assert_true(provider_stub.ok())
+ guard.assert_clean()
+
+
+def test_provider_raw_head_body_split_utf8_body_is_preserved() raises:
+ # RP01: a multi-byte body is preserved whole. The observer accumulates
+ # bytes and searches them directly, so a split character can never be
+ # misread as a read/packet boundary.
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "split_utf8_body",
+ "POST",
+ "/v1/chat/completions",
+ 200,
+ '{"mark":"a☃b"}',
+ )
+ script.stall_after_head_ms = 250
+ scripts.append(script^)
+ var guard = CleanupGuard()
+ with spawn_max_local_scripted(0, scripts^, guard) as provider_stub:
+ var report = run_bounded_call(
+ "raw_head_body",
+ provider_stub.port,
+ 5000,
+ 5000,
+ guard,
+ BOUNDED_CORRELATION_RAW_SPLIT_UTF8,
+ "/v1/chat/completions",
+ "a☃b",
+ )
+ assert_true(report.ok())
+ assert_equal(report.status, 200)
+ assert_equal(report.body_bytes, '{"mark":"a☃b"}'.byte_length())
+ assert_equal(report.body_match, "yes")
+ assert_equal(report.length_match, "yes")
+ provider_stub.wait()
+ assert_true(provider_stub.ok())
+ guard.assert_clean()
+
+
+def test_raw_accounting_preserves_split_utf8_bytes() raises:
+ # RP01: the byte accounting itself is chunk-independent. The same bytes are
+ # supplied as two pieces whose boundary falls inside the three-byte
+ # character, and the observation is still exact.
+ var head = "HTTP/1.1 200 OK\r\ncontent-length: 5\r\n\r\n"
+ var body = "a☃b"
+ var raw = List[UInt8]()
+ for byte in head.as_bytes():
+ raw.append(UInt8(Int(byte)))
+ for byte in body.as_bytes():
+ raw.append(UInt8(Int(byte)))
+ assert_equal(len(raw), head.byte_length() + body.byte_length())
+ var boundary = head.byte_length() + 3
+ var split = List[UInt8]()
+ for index in range(boundary):
+ split.append(raw[index])
+ for index in range(boundary, len(raw)):
+ split.append(raw[index])
+ var accounting = account_raw_bytes(split^, "☃")
+ assert_equal(accounting.problem, "")
+ assert_equal(accounting.status, 200)
+ assert_equal(accounting.body_bytes, 5)
+ assert_equal(accounting.body_match, "yes")
+ assert_equal(accounting.declared_bytes, 5)
+ assert_equal(accounting.length_match, "yes")
+
+
+def test_raw_accounting_reports_incomplete_head() raises:
+ # RP01: a buffer with no header terminator is an explicit incomplete-head
+ # observation, never a completed call with invented fields.
+ var raw = List[UInt8]()
+ for byte in "HTTP/1.1 200 OK\r\ncontent-length: 3".as_bytes():
+ raw.append(UInt8(Int(byte)))
+ var accounting = account_raw_bytes(raw^, "abc")
+ assert_equal(accounting.problem, "raw_head_incomplete")
+ assert_equal(accounting.body_bytes, 0)
+ assert_equal(accounting.status, 0)
+
+
+def test_raw_accounting_reports_undecodable_head() raises:
+ # RP01: an undecodable header block is a bounded decode failure, not a
+ # silently accepted status/body observation.
+ var raw = List[UInt8]()
+ for byte in "HTTP/1.1 200 OK\r\nx: ".as_bytes():
+ raw.append(UInt8(Int(byte)))
+ raw.append(UInt8(0xFF))
+ for byte in "\r\n\r\n".as_bytes():
+ raw.append(UInt8(Int(byte)))
+ var accounting = account_raw_bytes(raw^, "")
+ assert_equal(accounting.problem, "raw_head_decode")
+ assert_equal(accounting.status, 0)
+
+
+def test_provider_raw_head_body_reports_peer_close_without_head() raises:
+ # RP01: a peer that closes before sending any response head is a bounded
+ # raw_eof domain failure, not a hang and not an invented status.
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "close_without_head", "POST", "/v1/chat/completions", 200, "{}"
+ )
+ script.close_before_response = True
+ scripts.append(script^)
+ var guard = CleanupGuard()
+ with spawn_max_local_scripted(0, scripts^, guard) as provider_stub:
+ var report = run_bounded_call(
+ "raw_head_body",
+ provider_stub.port,
+ 5000,
+ 5000,
+ guard,
+ BOUNDED_CORRELATION_RAW_EOF,
+ "/v1/chat/completions",
+ "x",
+ )
+ assert_true(report.completed)
+ assert_true(report.domain_failure())
+ assert_equal(report.cause, "raw_eof")
+ provider_stub.wait()
+ assert_true(provider_stub.ok())
+ guard.assert_clean()
+
+
+def test_provider_raw_head_body_reports_malformed_status() raises:
+ # RP01: a malformed response with a terminator but no HTTP status line must
+ # be reported as a bounded malformed observation (status_missing), never
+ # read as a successful response with an invented status.
+ var guard = CleanupGuard()
+ with spawn_max_local_stub(
+ 0, "query_rewrite_malformed_http", 1, guard
+ ) as provider_stub:
+ var report = run_bounded_call(
+ "raw_head_body",
+ provider_stub.port,
+ 5000,
+ 5000,
+ guard,
+ BOUNDED_CORRELATION_RAW_MALFORMED,
+ "/v1/chat/completions",
+ "x",
+ )
+ assert_true(report.completed)
+ assert_true(report.domain_failure())
+ assert_equal(report.cause, "raw_status_missing")
+ provider_stub.wait()
+ 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
@@ -747,6 +1020,141 @@ def test_provider_scripted_injected_invalid_descriptor_fails() raises:
guard.assert_clean()
+def _realerrno_from_reason(reason: String) raises -> Int:
+ """Raw errno recorded in a bounded write-failure reason, or -1."""
+ var marker = "realerrno"
+ var at = reason.find(marker)
+ if at < 0:
+ return -1
+ var digits = String(reason[byte = at + marker.byte_length() :])
+ var end = digits.find("_")
+ if end >= 0:
+ digits = String(digits[byte=0:end])
+ if digits.byte_length() == 0:
+ return -1
+ return Int(digits)
+
+
+def test_provider_scripted_synthetic_write_timeout_is_not_peer_close() raises:
+ # RP03: a narrowly scoped syscall-result seam mapped exactly to the flare
+ # write API's EAGAIN/EWOULDBLOCK rendering. The synthetic failure exercises
+ # the actual classification and serve rejection path, is labelled synthdecl
+ # and can never be accepted as a peer close.
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "synthetic_timeout",
+ "POST",
+ "/v1/chat/completions",
+ 200,
+ '{"choices":[]}',
+ )
+ script.stall_after_head_ms = 200
+ script.inject_write_errno = Int(ErrNo.EAGAIN.value)
+ 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(not is_peer_close_cause("write_timeout"))
+ assert_true(
+ provider_stub.reason().find(
+ "unexpected_write_write_timeout_body_stall_synthdecl"
+ )
+ >= 0
+ )
+ guard.assert_clean()
+
+
+def test_provider_scripted_synthetic_invalid_descriptor_is_not_peer_close() raises:
+ # RP03: the same seam mapped to the flare write API's EBADF rendering. The
+ # invalid descriptor is classified as invalid_descriptor, labelled synthetic
+ # and rejected even when a peer close was declared.
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "synthetic_descriptor",
+ "POST",
+ "/v1/chat/completions",
+ 200,
+ '{"choices":[]}',
+ )
+ script.stall_after_head_ms = 200
+ script.inject_write_errno = Int(ErrNo.EBADF.value)
+ 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_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_invalid_descriptor_body_stall_synthdecl"
+ )
+ >= 0
+ )
+ guard.assert_clean()
+
+
+def test_provider_scripted_declared_peer_close_rejects_synthetic_timeout() raises:
+ # RP03: declaring a real peer close does not waive a synthetic timeout. The
+ # declared cause is compared against the observed class, so a non-peer-close
+ # failure stays a failure.
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "declared_vs_synthetic",
+ "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_errno = Int(ErrNo.EAGAIN.value)
+ 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_true(provider_stub.reason().find("write_timeout") >= 0)
+ guard.assert_clean()
+
+
+def test_provider_scripted_real_peer_close_records_raw_errno() raises:
+ # RP03: the real (non-synthetic) write failure. A real peer close makes the
+ # fixture's real send(2) fail; the fixture records the raw errno alongside
+ # the classification, and the errno's class must equal the observed class.
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "real_errno", "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())
+ var token = provider_stub.failure_case()
+ assert_true(token.find("realerrno") >= 0)
+ var errno = _realerrno_from_reason(token)
+ assert_true(errno > 0)
+ var observed = write_errno_class(errno)
+ assert_true(observed == "broken_pipe" or observed == "peer_reset")
+ assert_true(token.find("_peer_close_" + observed + "_") >= 0)
+ guard.assert_clean()
+
+
def test_provider_scripted_declared_non_peer_cause_is_rejected() raises:
# TC01: the declared expected cause is restricted to a real peer-close class,
# so a write timeout can never be waived by declaring it as expected.
@@ -1149,6 +1557,175 @@ def test_bounded_call_transient_wait_error_is_not_success() raises:
guard.assert_clean()
+def test_bounded_call_rejects_early_eof_report() raises:
+ # RP02: an early EOF with no report byte is a harness failure, never an
+ # empty completed call, through the actual bounded consumer.
+ var guard = CleanupGuard()
+ var fd_before = open_fd_count_checked()
+ var report = run_bounded_call(
+ "mutate_silent",
+ 0,
+ 300,
+ 500,
+ guard,
+ BOUNDED_CORRELATION_EARLY_EOF,
+ )
+ assert_true(not report.completed)
+ assert_equal(report.problem, "report_early_eof")
+ assert_true(report.cleanup_proved)
+ guard.assert_clean()
+ assert_equal(open_fd_count_checked(), fd_before)
+
+
+def test_bounded_call_rejects_signaled_child_after_valid_report() raises:
+ # RP02: a valid report followed by a signaled child is a harness failure.
+ # The parent must observe the signal, not accept the report or kill the
+ # child and call it complete.
+ var guard = CleanupGuard()
+ var fd_before = open_fd_count_checked()
+ var report = run_bounded_call(
+ "mutate_signaled",
+ 0,
+ 300,
+ 500,
+ guard,
+ BOUNDED_CORRELATION_SIGNALED,
+ )
+ assert_true(not report.completed)
+ assert_true(not report.stopped)
+ assert_equal(report.problem, "child_signal_9")
+ assert_true(report.child_status.find("signal=9") >= 0)
+ assert_true(report.cleanup_proved)
+ guard.assert_clean()
+ assert_equal(open_fd_count_checked(), fd_before)
+
+
+def test_bounded_call_rejects_invalid_utf8_report() raises:
+ # RP02: an invalid UTF-8 byte in the report line is rejected by the bounded
+ # decode, so a corrupted report can never be accepted as a completed call.
+ var guard = CleanupGuard()
+ var fd_before = open_fd_count_checked()
+ var report = run_bounded_call(
+ "mutate_invalid_utf8",
+ 0,
+ 300,
+ 500,
+ guard,
+ BOUNDED_CORRELATION_INVALID_UTF8,
+ )
+ assert_true(not report.completed)
+ assert_equal(report.problem, "invalid_utf8")
+ assert_true(report.cleanup_proved)
+ guard.assert_clean()
+ assert_equal(open_fd_count_checked(), fd_before)
+
+
+def test_bounded_call_rejects_phase_isolated_late_exit() raises:
+ # RP02: after the report pipe is closed with one complete report, a child
+ # that outlives the budget must be isolated in the wait phase and stopped
+ # and reaped, never reported as completed.
+ var guard = CleanupGuard()
+ var fd_before = open_fd_count_checked()
+ var report = run_bounded_call(
+ "mutate_late_exit",
+ 0,
+ 300,
+ 500,
+ guard,
+ BOUNDED_CORRELATION_LATE_EXIT,
+ )
+ 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()
+ assert_equal(open_fd_count_checked(), fd_before)
+
+
+def test_bounded_call_retries_synthetic_eintr_without_error() raises:
+ # RP02: a bounded test-only EINTR seam exercises the actual poll consumer's
+ # retry branch. One interrupted poll must be retried, not reported as a read
+ # error, and the complete report must still be accepted. The seam is
+ # synthetic and disclosed as such; it changes no host signal state.
+ var guard = CleanupGuard()
+ var report = run_bounded_call(
+ "mutate_valid",
+ 0,
+ 300,
+ 500,
+ guard,
+ BOUNDED_CORRELATION_EINTR_RETRY,
+ "/v1/chat/completions",
+ "",
+ 0,
+ 0,
+ 1,
+ 0,
+ )
+ assert_true(report.completed)
+ assert_equal(report.outcome, "ok")
+ assert_equal(report.status, 200)
+ assert_equal(report.problem, "")
+ assert_true(report.cleanup_proved)
+ guard.assert_clean()
+
+
+def test_bounded_call_eintr_retry_is_bounded_by_deadline() raises:
+ # RP02: an unbounded synthetic EINTR storm must terminate at the caller's
+ # absolute deadline, never refresh the budget and never spin. The actual
+ # consumer reports the expiry, so the parent stops and reaps the child.
+ var guard = CleanupGuard()
+ var report = run_bounded_call(
+ "mutate_valid",
+ 0,
+ 300,
+ 500,
+ guard,
+ BOUNDED_CORRELATION_EINTR_DEADLINE,
+ "/v1/chat/completions",
+ "",
+ 0,
+ 0,
+ -1,
+ 0,
+ )
+ 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_call_poll_error_is_bounded_failure() raises:
+ # RP02: a real poll error (bounded test-only seam) is surfaced as
+ # read_error by the actual consumer, and the exception path still proves
+ # owned cleanup with no child or descriptor growth.
+ var guard = CleanupGuard()
+ var fd_before = open_fd_count_checked()
+ var report = run_bounded_call(
+ "mutate_valid",
+ 0,
+ 300,
+ 500,
+ guard,
+ BOUNDED_CORRELATION_POLL_ERROR,
+ "/v1/chat/completions",
+ "",
+ 0,
+ 0,
+ 0,
+ 1,
+ )
+ assert_true(not report.completed)
+ assert_equal(report.problem, "read_error")
+ assert_true(report.cleanup_proved)
+ guard.assert_clean()
+ assert_equal(guard.pending(), 0)
+ assert_equal(open_fd_count_checked(), fd_before)
+
+
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).
diff --git a/tests/test_provider_helpers.mojo b/tests/test_provider_helpers.mojo
@@ -6,7 +6,7 @@ exception, signal or unrelated startup failure.
"""
from std.testing import TestSuite, assert_true, assert_equal
-from std.ffi import ErrNo, c_int, external_call
+from std.ffi import ErrNo, c_int, external_call, get_errno
from flare.net import SocketAddr
from flare.net.socket import RawSocket
@@ -53,11 +53,15 @@ from strict_fixture import (
ExchangeScript,
FramedRequest,
authorization_reason,
+ classify_write_error_cause,
exchange_script,
+ is_peer_close_cause,
json_escape,
parse_report,
+ render_write_api_error,
report_status_matches_exit,
verify_exchange,
+ write_errno_class,
)
from max_local_process_helper import (
SpawnedMaxLocalStub,
@@ -950,6 +954,26 @@ def test_write_deadline_and_closed_pipe_causes() raises:
assert_equal(deadline_reason, "write_deadline_expired")
+def test_real_write_to_closed_descriptor_is_ebadf() raises:
+ # RP03: a real (non-synthetic) descriptor/write failure. An exact-owned pipe
+ # write end is closed and then written to, so the OS returns a real EBADF.
+ # The raw errno is recorded and the actual classification function maps the
+ # rendered write-API cause to invalid_descriptor, which can never be a peer
+ # close. No product or fork source is involved.
+ var pipe = make_pipe()
+ close_fd(pipe.write_fd)
+ var written = write_raw(pipe.write_fd, "x")
+ assert_true(written < 0)
+ var errno_value = Int(get_errno().value)
+ assert_equal(errno_value, Int(ErrNo.EBADF.value))
+ assert_equal(write_errno_class(errno_value), "invalid_descriptor")
+ var rendered = render_write_api_error(errno_value)
+ assert_true(rendered.find("Bad file descriptor") >= 0)
+ assert_equal(classify_write_error_cause(rendered), "invalid_descriptor")
+ assert_true(not is_peer_close_cause("invalid_descriptor"))
+ close_fd(pipe.read_fd)
+
+
def test_owned_child_reaped_after_early_terminate() raises:
var guard_25 = CleanupGuard()
with spawn_max_local_stub(0, "count_requests", 1, guard_25) as stub: