hyf

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

bounded_call_helper.mojo (43405B)


      1 """Minimal test-only parent-bounded in-process call runner (H007 BC01-BC03).
      2 
      3 An elapsed-time assertion after a synchronous provider call is not a parent
      4 deadline: if the call never returns, the owning test hangs and the observed
      5 behavior cannot be characterized. This helper forks one exact-owned child that
      6 performs the risky provider call or raw header/body observation and writes a
      7 single bounded report line; the parent enforces one spawn-relative work budget,
      8 validates the complete report through EOF, verifies a natural child exit zero,
      9 and, when the budget expires, terminates and reaps the exact owned child
     10 through the shared lifecycle primitives.
     11 
     12 Ownership and cleanup reuse the qualified ``PipedChildState``/``CleanupGuard``
     13 primitives rather than a private lifecycle: the exact child and its report
     14 descriptor are scope-owned immediately after spawn, each owned descriptor is
     15 closed at most once, and an unproved cleanup retains the usable ownership
     16 handle so the caller can recover it.
     17 
     18 The report grammar is one line of ``key=value`` fields validated by the parent:
     19 
     20 ``report kind=<kind> correlation=<n> outcome=ok status=<code> [timing fields]``
     21 ``report kind=<kind> correlation=<n> outcome=fail cause=<token> reason=<token>``
     22 
     23 A wrong kind/correlation, an unknown/duplicate/missing field, a malformed or
     24 unterminated/duplicated report, surplus bytes through EOF, a read/poll error or
     25 a non-zero/signaled child exit is a harness failure and can never be reported as
     26 a completed call. A well-formed ``outcome=fail`` report is a characterized
     27 domain/provider failure, which is distinct from a harness failure.
     28 
     29 It is test-only tooling: no product policy, schema, dependency or lock change.
     30 """
     31 
     32 from std.collections import List
     33 
     34 from json import Value, loads
     35 
     36 from flare.net import SocketAddr
     37 from flare.tcp import TcpStream
     38 
     39 from parent_lifecycle import (
     40     SIGKILL,
     41     TERMINATION_GRACE_MS,
     42     CleanupGuard,
     43     ProcessStatus,
     44     child_exit,
     45     close_fd,
     46     fork_owned_or_close,
     47     kill_pid,
     48     make_pipe,
     49     now_ms,
     50     owned_pid,
     51     piped_child_state,
     52     read_fd,
     53     set_alarm,
     54     sleep_ms,
     55     write_raw,
     56     write_raw_bytes,
     57 )
     58 
     59 from hyf_core.request_context import default_request_context
     60 from hyf_provider.client import post_max_local_chat_completion
     61 from hyf_provider.config import MaxLocalProviderConfig
     62 from hyf_provider.jev_client import post_jev_systemone
     63 from hyf_provider.schema import build_query_rewrite_request_body
     64 
     65 
     66 comptime BOUNDED_CALL_MAX_REPORT_BYTES = 4096
     67 comptime BOUNDED_CALL_CHILD_ALARM_SECONDS = 120
     68 
     69 
     70 def _field_allowed(key: String) -> Bool:
     71     if key == "kind" or key == "correlation" or key == "outcome":
     72         return True
     73     if key == "status" or key == "cause" or key == "reason":
     74         return True
     75     if key == "latency_ms" or key == "head_ms" or key == "total_ms":
     76         return True
     77     if key == "body_bytes" or key == "body_match":
     78         return True
     79     if key == "surplus_bytes":
     80         return True
     81     return key == "declared_bytes" or key == "length_match"
     82 
     83 
     84 def _field_numeric(key: String) -> Bool:
     85     if key == "status" or key == "correlation" or key == "latency_ms":
     86         return True
     87     if key == "head_ms" or key == "total_ms":
     88         return True
     89     return (
     90         key == "body_bytes" or key == "declared_bytes" or key == "surplus_bytes"
     91     )
     92 
     93 
     94 def _all_digits(text: String) -> Bool:
     95     if text.byte_length() == 0:
     96         return False
     97     for byte in text.as_bytes():
     98         var b = Int(byte)
     99         if b < 48 or b > 57:
    100             return False
    101     return True
    102 
    103 
    104 @fieldwise_init
    105 struct BoundedCallOutcome(Movable):
    106     """Validated fields of one bounded report line.
    107 
    108     ``ok``/``problem`` describe the validation result; ``outcome`` plus the
    109     call-specific fields describe the declared call outcome.
    110     """
    111 
    112     var ok: Bool
    113     var problem: String
    114     var kind: String
    115     var correlation: Int
    116     var outcome: String
    117     var status: Int
    118     var cause: String
    119     var reason: String
    120     var latency_ms: Int
    121     var head_ms: Int
    122     var total_ms: Int
    123     var body_bytes: Int
    124     var body_match: String
    125     var declared_bytes: Int
    126     var length_match: String
    127     var surplus_bytes: Int
    128 
    129     def domain_failure(self) -> Bool:
    130         return self.ok and self.outcome == "fail"
    131 
    132     def describe(self) -> String:
    133         return (
    134             "kind="
    135             + self.kind
    136             + " correlation="
    137             + String(self.correlation)
    138             + " outcome="
    139             + self.outcome
    140             + " status="
    141             + String(self.status)
    142             + " cause="
    143             + self.cause
    144             + " reason="
    145             + self.reason
    146             + " latency_ms="
    147             + String(self.latency_ms)
    148             + " head_ms="
    149             + String(self.head_ms)
    150             + " total_ms="
    151             + String(self.total_ms)
    152             + " body_bytes="
    153             + String(self.body_bytes)
    154             + " body_match="
    155             + self.body_match
    156             + " declared_bytes="
    157             + String(self.declared_bytes)
    158             + " length_match="
    159             + self.length_match
    160             + " surplus_bytes="
    161             + String(self.surplus_bytes)
    162         )
    163 
    164 
    165 def _empty_outcome(problem: String) -> BoundedCallOutcome:
    166     return BoundedCallOutcome(
    167         ok=False,
    168         problem=problem,
    169         kind="",
    170         correlation=-1,
    171         outcome="",
    172         status=0,
    173         cause="",
    174         reason="",
    175         latency_ms=-1,
    176         head_ms=-1,
    177         total_ms=-1,
    178         body_bytes=-1,
    179         body_match="",
    180         declared_bytes=-1,
    181         length_match="",
    182         surplus_bytes=-1,
    183     )
    184 
    185 
    186 def parse_bounded_report(
    187     text: String, expected_kind: String, correlation: Int
    188 ) raises -> BoundedCallOutcome:
    189     """Validate one complete bounded report line for the declared call.
    190 
    191     Rejects a wrong kind or correlation, an unknown/duplicate/missing field, a
    192     malformed token and an incompatible outcome (a status on a failure, or a
    193     cause/reason on a success). Every rejection is a bounded problem token, so
    194     a harness failure is never mistaken for a characterized call outcome.
    195     """
    196     var tokens = text.strip().split(" ")
    197     if len(tokens) < 1 or String(tokens[0]) != "report":
    198         return _empty_outcome("report_prefix")
    199     var keys = List[String]()
    200     var values = List[String]()
    201     for index in range(1, len(tokens)):
    202         var token = String(tokens[index])
    203         var split = token.find("=")
    204         if split <= 0 or split == token.byte_length() - 1:
    205             return _empty_outcome("report_token_grammar")
    206         var key = String(token[byte=0:split])
    207         if not _field_allowed(key):
    208             return _empty_outcome("report_unknown_field_" + key)
    209         for seen in range(len(keys)):
    210             if keys[seen] == key:
    211                 return _empty_outcome("report_duplicate_field_" + key)
    212         keys.append(key)
    213         values.append(String(token[byte = split + 1 :]))
    214     var kind = ""
    215     var correlation_text = ""
    216     var outcome = ""
    217     var status_text = ""
    218     var cause = ""
    219     var reason = ""
    220     var latency_text = ""
    221     var head_text = ""
    222     var total_text = ""
    223     var bytes_text = ""
    224     var match_text = ""
    225     var declared_text = ""
    226     var length_text = ""
    227     var surplus_text = ""
    228     for index in range(len(keys)):
    229         var key = keys[index]
    230         var value = values[index]
    231         if key == "kind":
    232             kind = value
    233         elif key == "correlation":
    234             correlation_text = value
    235         elif key == "outcome":
    236             outcome = value
    237         elif key == "status":
    238             status_text = value
    239         elif key == "cause":
    240             cause = value
    241         elif key == "reason":
    242             reason = value
    243         elif key == "latency_ms":
    244             latency_text = value
    245         elif key == "head_ms":
    246             head_text = value
    247         elif key == "total_ms":
    248             total_text = value
    249         elif key == "body_bytes":
    250             bytes_text = value
    251         elif key == "body_match":
    252             match_text = value
    253         elif key == "declared_bytes":
    254             declared_text = value
    255         elif key == "length_match":
    256             length_text = value
    257         elif key == "surplus_bytes":
    258             surplus_text = value
    259         if _field_numeric(key) and not _all_digits(value):
    260             return _empty_outcome("report_non_numeric_" + key)
    261     if kind == "" or correlation_text == "" or outcome == "":
    262         return _empty_outcome("report_missing_field")
    263     if kind != expected_kind:
    264         return _empty_outcome("report_kind_mismatch")
    265     if Int(correlation_text) != correlation:
    266         return _empty_outcome("report_correlation_mismatch")
    267     if outcome != "ok" and outcome != "fail":
    268         return _empty_outcome("report_outcome_unknown_" + outcome)
    269     if outcome == "ok":
    270         if status_text == "":
    271             return _empty_outcome("report_status_missing")
    272         var code = Int(status_text)
    273         if code < 100 or code > 599:
    274             return _empty_outcome("report_status_invalid")
    275         if cause != "" or reason != "":
    276             return _empty_outcome("report_incompatible_outcome")
    277     else:
    278         if cause == "" or reason == "":
    279             return _empty_outcome("report_cause_missing")
    280         if status_text != "":
    281             return _empty_outcome("report_incompatible_outcome")
    282     return BoundedCallOutcome(
    283         ok=True,
    284         problem="",
    285         kind=kind,
    286         correlation=Int(correlation_text),
    287         outcome=outcome,
    288         status=Int(status_text) if status_text != "" else 0,
    289         cause=cause,
    290         reason=reason,
    291         latency_ms=Int(latency_text) if latency_text != "" else -1,
    292         head_ms=Int(head_text) if head_text != "" else -1,
    293         total_ms=Int(total_text) if total_text != "" else -1,
    294         body_bytes=Int(bytes_text) if bytes_text != "" else -1,
    295         body_match=match_text,
    296         declared_bytes=Int(declared_text) if declared_text != "" else -1,
    297         length_match=length_text,
    298         surplus_bytes=Int(surplus_text) if surplus_text != "" else -1,
    299     )
    300 
    301 
    302 def _classify_child_raise(text: String) -> String:
    303     """Bounded cause token for a raise observed inside the bounded child.
    304 
    305     The classification is derived from the rendered error text, never from a
    306     caller-supplied string, and is reported as a domain outcome — never as a
    307     peer close or as a harness success.
    308     """
    309     if text.find("refused") >= 0 or text.find("Refused") >= 0:
    310         return "connection_refused"
    311     if text.find("Timeout") >= 0 or text.find("timeout") >= 0:
    312         return "timeout"
    313     if text.find("descriptor") >= 0 or text.find("EBADF") >= 0:
    314         return "invalid_descriptor"
    315     return "raised"
    316 
    317 
    318 def _report_prefix(kind: String, correlation: Int) -> String:
    319     return "report kind=" + kind + " correlation=" + String(correlation) + " "
    320 
    321 
    322 def _long_token(count: Int) -> String:
    323     var text = ""
    324     for _ in range(count):
    325         text += "x"
    326     return text^
    327 
    328 
    329 def _mutant_payload(copies: Int) -> String:
    330     """Deliberately misleading child report for the parent-consumer controls.
    331 
    332     This is a *child-producer* mutation only (ADR-0021 BC02): it changes what
    333     the forked child writes and never the parent consumer under test, which is
    334     the same consumer every real bounded call uses.
    335     """
    336     var line = "report kind=mutant correlation=0 outcome=ok status=200"
    337     var payload = line
    338     for index in range(1, copies):
    339         payload += "\n" + line
    340     return payload^
    341 
    342 
    343 def _raw_request_text(path: String) -> String:
    344     return (
    345         "POST "
    346         + path
    347         + " HTTP/1.1\r\nhost: 127.0.0.1\r\ncontent-length: 2\r\n"
    348         "connection: close\r\n\r\n{}"
    349     )
    350 
    351 
    352 comptime RAW_MAX_HEADER_BYTES: Int = 65536
    353 comptime RAW_MAX_BODY_BYTES: Int = 1048576
    354 comptime RAW_READ_CHUNK_BYTES: Int = 1024
    355 
    356 
    357 def _find_header_terminator(bytes: List[UInt8]) -> Int:
    358     """Index just past the first CRLFCRLF, or -1 while the head is incomplete.
    359     """
    360     var n = len(bytes)
    361     if n < 4:
    362         return -1
    363     for index in range(0, n - 3):
    364         if (
    365             Int(bytes[index]) == 13
    366             and Int(bytes[index + 1]) == 10
    367             and Int(bytes[index + 2]) == 13
    368             and Int(bytes[index + 3]) == 10
    369         ):
    370             return index + 4
    371     return -1
    372 
    373 
    374 def _trim_ows_text(value: String) -> String:
    375     """Trim only legal HTTP optional whitespace (SP / HTAB)."""
    376     var start = 0
    377     var end = value.byte_length()
    378     var bytes = value.as_bytes()
    379     while start < end and (Int(bytes[start]) == 32 or Int(bytes[start]) == 9):
    380         start += 1
    381     while end > start and (
    382         Int(bytes[end - 1]) == 32 or Int(bytes[end - 1]) == 9
    383     ):
    384         end -= 1
    385     if start == 0 and end == value.byte_length():
    386         return String(value)
    387     return String(value[byte=start:end])
    388 
    389 
    390 def _token_name(value: String) -> Bool:
    391     """RFC 7230 token check for an exact header field name."""
    392     if value.byte_length() == 0:
    393         return False
    394     for byte in value.as_bytes():
    395         var b = Int(byte)
    396         if b >= 48 and b <= 57:
    397             continue
    398         if b >= 65 and b <= 90:
    399             continue
    400         if b >= 97 and b <= 122:
    401             continue
    402         if (
    403             b == 33
    404             or b == 35
    405             or b == 36
    406             or b == 37
    407             or b == 38
    408             or b == 39
    409             or b == 42
    410             or b == 43
    411             or b == 45
    412             or b == 46
    413             or b == 94
    414             or b == 95
    415             or b == 96
    416             or b == 124
    417             or b == 126
    418         ):
    419             continue
    420         return False
    421     return True
    422 
    423 
    424 def _value_legal(value: String) -> Bool:
    425     """Reject control bytes in a header value (HTAB is the only legal one)."""
    426     for byte in value.as_bytes():
    427         var b = Int(byte)
    428         if b == 9:
    429             continue
    430         if b < 32 or b == 127:
    431             return False
    432     return True
    433 
    434 
    435 def _ascii_digit(value: Int) -> Bool:
    436     return value >= 48 and value <= 57
    437 
    438 
    439 @fieldwise_init
    440 struct HeadFraming(Movable):
    441     """Exact framing parsed from one response head.
    442 
    443     ``problem`` is an explicit bounded grammar failure token; otherwise
    444     ``status`` is the leading three-digit status and ``declared`` is the exact
    445     Content-Length (``-1`` when the header is absent).
    446     """
    447 
    448     var problem: String
    449     var status: Int
    450     var declared: Int
    451 
    452 
    453 def _parse_response_head(head_text: String) raises -> HeadFraming:
    454     """Parse the leading HTTP/1.1 status line and exact framing headers.
    455 
    456     OB01: the raw observer accepts only an exact leading
    457     ``HTTP/1.1|HTTP/1.0 SP three-digit-status`` line and line-delimited
    458     ``name: value`` headers. A malformed, duplicated, conflicting, negative,
    459     junk or overflowing framing declaration and an unsupported
    460     Transfer-Encoding are explicit bounded problem tokens. A status-looking
    461     string inside a header value is never read as the response status, and a
    462     different field name (``X-Content-Length``) is preserved as an unknown
    463     header rather than interpreted as ``Content-Length``.
    464     """
    465     var lines = head_text.split("\r\n")
    466     if len(lines) < 1:
    467         return HeadFraming("raw_status_missing", 0, -1)
    468     var status_line = String(lines[0])
    469     var version = ""
    470     if status_line.startswith("HTTP/1.1 "):
    471         version = "HTTP/1.1 "
    472     elif status_line.startswith("HTTP/1.0 "):
    473         version = "HTTP/1.0 "
    474     else:
    475         return HeadFraming("raw_status_missing", 0, -1)
    476     var rest = String(status_line[byte = version.byte_length() :])
    477     if rest.byte_length() < 3:
    478         return HeadFraming("raw_status_malformed", 0, -1)
    479     var digits = String(rest[byte=0:3])
    480     for byte in digits.as_bytes():
    481         if not _ascii_digit(Int(byte)):
    482             return HeadFraming("raw_status_malformed", 0, -1)
    483     if rest.byte_length() > 3:
    484         if Int(rest.as_bytes()[3]) != 32:
    485             return HeadFraming("raw_status_malformed", 0, -1)
    486     var status_value = Int(digits)
    487     if status_value < 100 or status_value > 599:
    488         # A three-digit but out-of-range status is malformed, never a
    489         # successfully observed call.
    490         return HeadFraming("raw_status_malformed", 0, -1)
    491     var declared = -1
    492     var have_length = False
    493     var have_transfer = False
    494     for index in range(1, len(lines)):
    495         var line = String(lines[index])
    496         if line.byte_length() == 0:
    497             continue
    498         var first = Int(line.as_bytes()[0])
    499         if first == 32 or first == 9:
    500             # obs-fold / a continuation line is not a legal standalone header.
    501             return HeadFraming("raw_header_malformed", 0, -1)
    502         var colon = line.find(":")
    503         if colon <= 0:
    504             return HeadFraming("raw_header_malformed", 0, -1)
    505         var name = String(line[byte=0:colon])
    506         if not _token_name(name):
    507             return HeadFraming("raw_header_malformed", 0, -1)
    508         var value = _trim_ows_text(String(line[byte = colon + 1 :]))
    509         if not _value_legal(value):
    510             return HeadFraming("raw_header_malformed", 0, -1)
    511         var lower = name.lower()
    512         if lower == "content-length":
    513             if have_length:
    514                 # Includes an identical duplicate: a second framing field is
    515                 # ambiguous even when the value agrees.
    516                 return HeadFraming("raw_content_length_duplicate", 0, -1)
    517             if value.byte_length() == 0:
    518                 return HeadFraming("raw_content_length_malformed", 0, -1)
    519             for byte in value.as_bytes():
    520                 if not _ascii_digit(Int(byte)):
    521                     return HeadFraming("raw_content_length_malformed", 0, -1)
    522             if value.byte_length() > 9:
    523                 return HeadFraming("raw_content_length_overflow", 0, -1)
    524             var parsed = Int(value)
    525             if parsed > RAW_MAX_BODY_BYTES:
    526                 return HeadFraming("raw_body_overflow", 0, -1)
    527             declared = parsed
    528             have_length = True
    529         elif lower == "transfer-encoding":
    530             if have_transfer:
    531                 return HeadFraming("raw_transfer_encoding_duplicate", 0, -1)
    532             have_transfer = True
    533             if value.lower() != "identity":
    534                 return HeadFraming("raw_transfer_encoding_unsupported", 0, -1)
    535     if have_transfer and have_length:
    536         return HeadFraming("raw_framing_conflict", 0, -1)
    537     return HeadFraming("", status_value, declared)
    538 
    539 
    540 def _bytes_contain_from(bytes: List[UInt8], start: Int, needle: String) -> Bool:
    541     """Byte-level substring search over an undecoded body buffer.
    542 
    543     Searching raw bytes (not a decoded String) keeps an invalid UTF-8 body or
    544     a multi-byte character split across reader chunks from corrupting the
    545     observation.
    546     """
    547     var nlen = needle.byte_length()
    548     if nlen == 0:
    549         return True
    550     var limit = len(bytes) - nlen
    551     if limit < start:
    552         return False
    553     var nb = needle.as_bytes()
    554     for base in range(start, limit + 1):
    555         var matched = True
    556         for k in range(nlen):
    557             if Int(bytes[base + k]) != Int(nb[k]):
    558                 matched = False
    559                 break
    560         if matched:
    561             return True
    562     return False
    563 
    564 
    565 @fieldwise_init
    566 struct RawAccounting(Movable):
    567     """Exact raw-response accounting from one bounded byte buffer.
    568 
    569     A non-empty ``problem`` is a bounded raw-observation failure (incomplete
    570     head or undecodable ASCII head); otherwise ``status``/``body_bytes``/
    571     ``body_match``/``declared_bytes``/``length_match`` describe the observation.
    572     """
    573 
    574     var problem: String
    575     var status: Int
    576     var body_bytes: Int
    577     var body_match: String
    578     var declared_bytes: Int
    579     var length_match: String
    580     var surplus_bytes: Int
    581 
    582 
    583 def account_raw_bytes(raw: List[UInt8], expect: String) raises -> RawAccounting:
    584     """Account header end, body byte count and length/match from raw bytes.
    585 
    586     The buffer is treated as one byte stream: packet/read boundaries are never
    587     protocol boundaries, and the body begins exactly after CRLFCRLF. A valid
    588     declared Content-Length is compared with the observed body byte count and
    589     any already-buffered bytes beyond it are reported explicitly as surplus.
    590     Malformed, duplicated, conflicting, negative, junk or overflowing framing
    591     and an unsupported Transfer-Encoding are explicit problem tokens rather
    592     than silently absent fields.
    593     """
    594     var head_end = _find_header_terminator(raw)
    595     if head_end < 0:
    596         return RawAccounting(
    597             "raw_head_incomplete", 0, 0, "unknown", -1, "unknown", 0
    598         )
    599     var head_bytes = List[UInt8]()
    600     for index in range(head_end):
    601         head_bytes.append(raw[index])
    602     var head_text = ""
    603     var decoded = True
    604     try:
    605         head_text = String(
    606             from_utf8=Span(ptr=head_bytes.unsafe_ptr(), length=len(head_bytes))
    607         )
    608     except:
    609         decoded = False
    610     if not decoded:
    611         return RawAccounting(
    612             "raw_head_decode", 0, 0, "unknown", -1, "unknown", 0
    613         )
    614     var framing = _parse_response_head(head_text)
    615     if framing.problem != "":
    616         return RawAccounting(framing.problem, 0, 0, "unknown", -1, "unknown", 0)
    617     var body_bytes = len(raw) - head_end
    618     var match_text = "unknown"
    619     if expect != "":
    620         match_text = "yes" if _bytes_contain_from(
    621             raw, head_end, expect
    622         ) else "no"
    623     var length_match = "unknown"
    624     var surplus = 0
    625     if framing.declared >= 0:
    626         length_match = "yes" if body_bytes == framing.declared else "no"
    627         if body_bytes > framing.declared:
    628             surplus = body_bytes - framing.declared
    629     return RawAccounting(
    630         "",
    631         framing.status,
    632         body_bytes,
    633         match_text,
    634         framing.declared,
    635         length_match,
    636         surplus,
    637     )
    638 
    639 
    640 def _plan_read_size(plan: List[Int], index: Int, default: Int) -> Int:
    641     """Deterministic bounded read size from an optional chunk plan (OB02).
    642 
    643     The plan element at ``index`` caps one actual read; once the plan is
    644     exhausted the last element stays in force. An absent plan uses ``default``.
    645     This is a disclosed synthetic acquisition seam, used only to force a header
    646     terminator or a multi-byte body character to split across real incremental
    647     reads; it is not a claim about packet boundaries.
    648     """
    649     if len(plan) == 0:
    650         return default
    651     var i = index if index < len(plan) else len(plan) - 1
    652     var want = plan[i]
    653     if want <= 0:
    654         want = 1
    655     return want if want < default else default
    656 
    657 
    658 def _child_raw_head_body(
    659     head: String,
    660     port: Int,
    661     path: String,
    662     expect: String,
    663     chunk_plan: List[Int],
    664 ) raises -> String:
    665     """Owned raw client: record head/body arrival timing for a scripted peer.
    666 
    667     The observation is independent of the product caller, which exposes only a
    668     parsed response, so a headers-before-body-stall claim can be proved from the
    669     wire while still running under the parent's finite deadline.
    670 
    671     OB01/OB02: the head grammar is parsed exactly, the header cap counts bytes
    672     through the terminator only (a coalesced body is not charged to it), a valid
    673     declared Content-Length is read as an exact message body with any surplus
    674     reported explicitly, a short read to EOF is an explicit incomplete
    675     observation, and a no-length response completes only at EOF inside the
    676     body cap. ``chunk_plan`` is a disclosed deterministic read seam that runs
    677     the same incremental loop.
    678     """
    679     var client = TcpStream.connect(SocketAddr.localhost(UInt16(port)))
    680     var start = now_ms()
    681     client.write_all(Span[UInt8, _](_raw_request_text(path).as_bytes()))
    682     var raw = List[UInt8]()
    683     var buffer = InlineArray[Byte, RAW_READ_CHUNK_BYTES](fill=0)
    684     var plan_index = 0
    685     var head_end = -1
    686     while head_end < 0:
    687         var want = _plan_read_size(chunk_plan, plan_index, RAW_READ_CHUNK_BYTES)
    688         plan_index += 1
    689         var n = client.read(buffer.unsafe_ptr(), want)
    690         if n <= 0:
    691             client.close()
    692             return head + "outcome=fail cause=raw_eof reason=head_eof"
    693         for index in range(n):
    694             raw.append(UInt8(Int(buffer[index])))
    695         head_end = _find_header_terminator(raw)
    696         if head_end < 0 and len(raw) > RAW_MAX_HEADER_BYTES:
    697             client.close()
    698             return (
    699                 head
    700                 + "outcome=fail cause=raw_head_overflow reason=head_overflow"
    701             )
    702     if head_end > RAW_MAX_HEADER_BYTES:
    703         client.close()
    704         return (
    705             head + "outcome=fail cause=raw_head_overflow reason=head_overflow"
    706         )
    707     var head_ms = now_ms() - start
    708     var head_bytes = List[UInt8]()
    709     for index in range(head_end):
    710         head_bytes.append(raw[index])
    711     var head_text = ""
    712     var decoded = True
    713     try:
    714         head_text = String(
    715             from_utf8=Span(ptr=head_bytes.unsafe_ptr(), length=len(head_bytes))
    716         )
    717     except:
    718         decoded = False
    719     if not decoded:
    720         client.close()
    721         return head + "outcome=fail cause=raw_head_decode reason=head_decode"
    722     var framing = _parse_response_head(head_text)
    723     if framing.problem != "":
    724         client.close()
    725         return (
    726             head
    727             + "outcome=fail cause="
    728             + framing.problem
    729             + " reason="
    730             + framing.problem
    731         )
    732     var declared = framing.declared
    733     var status = framing.status
    734     var body_bytes = len(raw) - head_end
    735     var incomplete = False
    736     var overflow = False
    737     while True:
    738         if declared >= 0 and body_bytes >= declared:
    739             break
    740         if declared < 0 and body_bytes >= RAW_MAX_BODY_BYTES:
    741             # D44/IL01: the 1,048,576-byte body cap is inclusive and a no-length
    742             # (EOF-delimited) response is complete at the exact cap only when
    743             # the peer has actually closed. Distinguish a real EOF from a real
    744             # extra byte with at most one one-byte sentinel read under the same
    745             # parent deadline; the sentinel is never appended to the capped
    746             # body, so the reported body_bytes stays at the inclusive maximum.
    747             # A peer that instead stalls leaves this read blocking, and the
    748             # existing parent deadline/cleanup reports a stopped (timeout)
    749             # call rather than an invented overflow.
    750             var sentinel = InlineArray[Byte, 1](fill=0)
    751             var extra = client.read(sentinel.unsafe_ptr(), 1)
    752             if extra > 0:
    753                 overflow = True
    754             break
    755         var want2 = _plan_read_size(
    756             chunk_plan, plan_index, RAW_READ_CHUNK_BYTES
    757         )
    758         plan_index += 1
    759         var remaining = RAW_MAX_BODY_BYTES - body_bytes
    760         if want2 > remaining:
    761             want2 = remaining
    762         if want2 <= 0:
    763             overflow = True
    764             break
    765         var n2 = client.read(buffer.unsafe_ptr(), want2)
    766         if n2 <= 0:
    767             # A length-delimited response that ends before its declared body is
    768             # explicitly incomplete; a no-length response completes at EOF.
    769             if declared >= 0 and body_bytes < declared:
    770                 incomplete = True
    771             break
    772         for index in range(n2):
    773             raw.append(UInt8(Int(buffer[index])))
    774         body_bytes += n2
    775     var total_ms = now_ms() - start
    776     client.close()
    777     var match_text = "unknown"
    778     if expect != "":
    779         match_text = "yes" if _bytes_contain_from(
    780             raw, head_end, expect
    781         ) else "no"
    782     var length_match = "unknown"
    783     var surplus = 0
    784     if declared >= 0:
    785         length_match = "yes" if body_bytes == declared else "no"
    786         if body_bytes > declared:
    787             surplus = body_bytes - declared
    788     var declared_field = ""
    789     if declared >= 0:
    790         declared_field = " declared_bytes=" + String(declared)
    791     var tail = (
    792         " head_ms="
    793         + String(head_ms)
    794         + " total_ms="
    795         + String(total_ms)
    796         + " body_bytes="
    797         + String(body_bytes)
    798         + " body_match="
    799         + match_text
    800         + " length_match="
    801         + length_match
    802         + " surplus_bytes="
    803         + String(surplus)
    804         + declared_field
    805     )
    806     if incomplete:
    807         return (
    808             head
    809             + "outcome=fail cause=raw_body_incomplete reason=body_eof"
    810             + tail
    811         )
    812     if overflow:
    813         return (
    814             head
    815             + "outcome=fail cause=raw_body_overflow reason=body_overflow"
    816             + tail
    817         )
    818     return head + "outcome=ok status=" + String(status) + tail
    819 
    820 
    821 def _child_report(
    822     kind: String,
    823     port: Int,
    824     timeout_ms: Int,
    825     correlation: Int,
    826     raw_path: String,
    827     raw_expect: String,
    828     raw_chunk_plan: List[Int],
    829 ) raises -> String:
    830     """One bounded report line from inside the forked provider-call child."""
    831     var head = _report_prefix(kind, correlation)
    832     if kind == "never_return":
    833         # A deliberately non-returning control: it cannot complete within any
    834         # parent deadline, so the parent must stop and reap it.
    835         while True:
    836             sleep_ms(1000)
    837         return head + "outcome=fail cause=never reason=never_return"
    838     if kind == "max_local":
    839         var config = MaxLocalProviderConfig(
    840             base_url="http://127.0.0.1:" + String(port) + "/v1/",
    841             health_url="http://127.0.0.1:" + String(port) + "/health",
    842             model="max-local-query-rewrite",
    843             request_timeout_ms=timeout_ms,
    844         )
    845         var context = default_request_context()
    846         var body = build_query_rewrite_request_body(
    847             config, "eggs near me", context
    848         )
    849         var outcome = post_max_local_chat_completion(config, body)
    850         if outcome.failure:
    851             return (
    852                 head
    853                 + "outcome=fail cause="
    854                 + outcome.failure.value().kind
    855                 + " reason="
    856                 + outcome.failure.value().reason
    857             )
    858         return (
    859             head
    860             + "outcome=ok status="
    861             + String(outcome.response.value().status)
    862             + " latency_ms="
    863             + String(outcome.response.value().latency_ms)
    864         )
    865     if kind == "jev":
    866         var payload = loads('{"model":"jev-1.13.0","state":"s","questions":{}}')
    867         var response = post_jev_systemone(
    868             "http://127.0.0.1:" + String(port), payload, timeout_ms
    869         )
    870         return head + "outcome=ok status=" + String(response.status)
    871     if kind == "raw_head_body" or kind == "raw_jev_head_body":
    872         return _child_raw_head_body(
    873             head, port, raw_path, raw_expect, raw_chunk_plan
    874         )
    875     return head + "outcome=fail cause=unknown_kind reason=unknown_kind"
    876 
    877 
    878 @fieldwise_init
    879 struct BoundedCallReport(Movable):
    880     """Bounded parent observation of one risky provider call.
    881 
    882     ``completed`` means exactly one complete bounded report was validated, no
    883     surplus byte followed it through EOF, and the child exited naturally with
    884     status zero inside one spawn-relative work budget. ``outcome`` is the
    885     declared call outcome (``ok`` or a characterized ``fail``); ``problem`` is
    886     non-empty only for a harness failure (malformed/duplicate/unterminated
    887     report, wrong correlation, early EOF, read/poll error, cap overflow,
    888     non-zero or signaled exit) and can never be reported as a completed call.
    889     ``stopped`` means the parent work budget expired first and the exact owned
    890     child was terminated and reaped under the separate bounded cleanup
    891     allowance.
    892     """
    893 
    894     var completed: Bool
    895     var stopped: Bool
    896     var outcome: String
    897     var status: Int
    898     var cause: String
    899     var reason: String
    900     var latency_ms: Int
    901     var head_ms: Int
    902     var total_ms: Int
    903     var body_bytes: Int
    904     var body_match: String
    905     var declared_bytes: Int
    906     var length_match: String
    907     var surplus_bytes: Int
    908     var problem: String
    909     var report: String
    910     var elapsed_ms: Int
    911     var child_status: String
    912     var cleanup_proved: Bool
    913 
    914     def ok(self) -> Bool:
    915         return self.completed and self.outcome == "ok"
    916 
    917     def domain_failure(self) -> Bool:
    918         return self.completed and self.outcome == "fail"
    919 
    920     def describe(self) -> String:
    921         return (
    922             "completed="
    923             + String(self.completed)
    924             + " stopped="
    925             + String(self.stopped)
    926             + " outcome="
    927             + self.outcome
    928             + " status="
    929             + String(self.status)
    930             + " cause="
    931             + self.cause
    932             + " reason="
    933             + self.reason
    934             + " latency_ms="
    935             + String(self.latency_ms)
    936             + " head_ms="
    937             + String(self.head_ms)
    938             + " total_ms="
    939             + String(self.total_ms)
    940             + " body_bytes="
    941             + String(self.body_bytes)
    942             + " body_match="
    943             + self.body_match
    944             + " declared_bytes="
    945             + String(self.declared_bytes)
    946             + " length_match="
    947             + self.length_match
    948             + " surplus_bytes="
    949             + String(self.surplus_bytes)
    950             + " problem="
    951             + (self.problem if self.problem != "" else "-")
    952             + " elapsed_ms="
    953             + String(self.elapsed_ms)
    954             + " child="
    955             + self.child_status
    956             + " cleanup="
    957             + ("proved" if self.cleanup_proved else "unproved")
    958             + " report="
    959             + self.report
    960         )
    961 
    962 
    963 def run_bounded_call(
    964     kind: String,
    965     port: Int,
    966     timeout_ms: Int,
    967     deadline_ms: Int,
    968     mut guard: CleanupGuard,
    969     correlation: Int,
    970     raw_path: String = "/v1/chat/completions",
    971     raw_expect: String = "",
    972     fault_cleanup_failures: Int = 0,
    973     fault_wait_errors: Int = 0,
    974     fault_poll_eintrs: Int = 0,
    975     fault_poll_errors: Int = 0,
    976     raw_chunk_plan: List[Int] = List[Int](),
    977 ) raises -> BoundedCallReport:
    978     """Run one risky provider call under a parent-enforced finite deadline.
    979 
    980     One spawn-relative work budget covers the report read, the surplus drain and
    981     the natural child exit. Success requires all three plus a validated report;
    982     a report alone never justifies success, and a child the parent had to stop
    983     is never reported as completed. Cleanup keeps a bounded allowance separate
    984     from the work budget and never turns an expired or failed result into
    985     success.
    986     """
    987     if deadline_ms <= 0:
    988         raise Error("bounded call: invalid deadline")
    989     var pipe = make_pipe()
    990     var pid = fork_owned_or_close(pipe.copy())
    991     var start = now_ms()
    992     if pid == 0:
    993         close_fd(pipe.read_fd)
    994         _ = set_alarm(BOUNDED_CALL_CHILD_ALARM_SECONDS)
    995         var payload = ""
    996         var terminator = "\n"
    997         var exit_code = 0
    998         var self_signal = 0
    999         # Child-producer-only controls for the parent consumer: each produces a
   1000         # deliberately incomplete, duplicated, non-zero-exit or over-running
   1001         # child result, and the unchanged parent consumer must reject it.
   1002         if kind == "mutate_exit7":
   1003             payload = _mutant_payload(1)
   1004             exit_code = 7
   1005         elif kind == "mutate_unterminated":
   1006             payload = _mutant_payload(1)
   1007             terminator = ""
   1008         elif kind == "mutate_duplicate":
   1009             payload = _mutant_payload(2)
   1010         elif kind == "mutate_huge":
   1011             # Child-producer-only control: a report line past the bounded cap.
   1012             payload = (
   1013                 _report_prefix("mutant", correlation)
   1014                 + "outcome=ok status=200 pad="
   1015                 + _long_token(5000)
   1016             )
   1017         elif kind == "mutate_valid":
   1018             # A well-formed child report, used to reach the child-exit wait
   1019             # phase with the bounded wait-fault seam.
   1020             payload = (
   1021                 _report_prefix(kind, correlation) + "outcome=ok status=200"
   1022             )
   1023         elif kind == "mutate_silent":
   1024             # RP02 child-producer-only control: the child closes its report pipe
   1025             # without writing a byte, so the parent must report an early EOF
   1026             # rather than an empty completed call.
   1027             payload = ""
   1028         elif kind == "mutate_signaled":
   1029             # RP02 child-producer-only control: a valid report followed by a
   1030             # signaled termination of the exact owned child; the parent must
   1031             # report the signal, never a completed call.
   1032             payload = (
   1033                 _report_prefix(kind, correlation) + "outcome=ok status=200"
   1034             )
   1035             self_signal = SIGKILL
   1036         elif kind == "mutate_invalid_utf8":
   1037             # RP02 child-producer-only control: a terminated report line whose
   1038             # body contains an invalid UTF-8 byte, so the parent's bounded
   1039             # decode rejects it instead of accepting a corrupted line.
   1040             var bytes = List[UInt8]()
   1041             var prefix = (
   1042                 _report_prefix(kind, correlation)
   1043                 + "outcome=ok status=200 mark="
   1044             )
   1045             for byte in prefix.as_bytes():
   1046                 bytes.append(UInt8(Int(byte)))
   1047             bytes.append(UInt8(0xFF))
   1048             bytes.append(UInt8(10))
   1049             _ = write_raw_bytes(pipe.write_fd, bytes^)
   1050             close_fd(pipe.write_fd)
   1051             child_exit(0)
   1052         elif kind == "mutate_late_exit":
   1053             # RP02 child-producer-only control: the report pipe is closed after
   1054             # one complete report while the child stays alive past the budget,
   1055             # isolating the late-exit wait phase from any report drain.
   1056             _ = write_raw(
   1057                 pipe.write_fd,
   1058                 _report_prefix(kind, correlation) + "outcome=ok status=200\n",
   1059             )
   1060             close_fd(pipe.write_fd)
   1061             sleep_ms(2000)
   1062             child_exit(0)
   1063         elif kind == "mutate_delayed":
   1064             _ = write_raw(pipe.write_fd, _mutant_payload(1) + "\n")
   1065             sleep_ms(2000)
   1066             close_fd(pipe.write_fd)
   1067             child_exit(0)
   1068         else:
   1069             try:
   1070                 payload = _child_report(
   1071                     kind,
   1072                     port,
   1073                     timeout_ms,
   1074                     correlation,
   1075                     raw_path,
   1076                     raw_expect,
   1077                     raw_chunk_plan,
   1078                 )
   1079             except e:
   1080                 payload = (
   1081                     _report_prefix(kind, correlation)
   1082                     + "outcome=fail cause=raised reason="
   1083                     + _classify_child_raise(String(e))
   1084                 )
   1085         if payload != "":
   1086             _ = write_raw(pipe.write_fd, payload + terminator)
   1087         close_fd(pipe.write_fd)
   1088         if self_signal != 0:
   1089             _ = kill_pid(owned_pid(), self_signal)
   1090         child_exit(exit_code)
   1091 
   1092     close_fd(pipe.write_fd)
   1093     var state = piped_child_state(
   1094         pid, pipe.read_fd, deadline_ms, 1, UnsafePointer(to=guard)
   1095     )
   1096     # Bounded test-only fault seam: one forced unproved cleanup or transient
   1097     # wait error against the real exact-owned child, so the retained-ownership
   1098     # and recovery path is exercised rather than a synthetic identity.
   1099     state.faults.cleanup_failures = fault_cleanup_failures
   1100     state.faults.wait_errors = fault_wait_errors
   1101     state.faults.poll_eintrs = fault_poll_eintrs
   1102     state.faults.poll_errors = fault_poll_errors
   1103     var deadline_hit = False
   1104     var problem = ""
   1105     var report_text = ""
   1106     # 1. Exactly one complete, newline-terminated bounded report line.
   1107     if state.work_remaining_ms() <= 0:
   1108         deadline_hit = True
   1109     if not deadline_hit:
   1110         try:
   1111             report_text = state.read_line(
   1112                 BOUNDED_CALL_MAX_REPORT_BYTES, state.work_remaining_ms()
   1113             )
   1114         except e:
   1115             var cause = String(e)
   1116             if cause == "read_deadline_expired":
   1117                 deadline_hit = True
   1118             else:
   1119                 problem = cause
   1120     if not deadline_hit and problem == "":
   1121         if not state.last_terminated:
   1122             if report_text.byte_length() == 0:
   1123                 problem = "report_early_eof"
   1124             else:
   1125                 problem = "report_unterminated"
   1126     # 2. No surplus bytes through EOF.
   1127     if not deadline_hit and problem == "":
   1128         try:
   1129             var surplus = state.drain_surplus(
   1130                 BOUNDED_CALL_MAX_REPORT_BYTES, state.work_remaining_ms()
   1131             )
   1132             if surplus > 0:
   1133                 problem = "report_duplicate_report"
   1134         except e:
   1135             var cause = String(e)
   1136             if cause == "read_deadline_expired":
   1137                 deadline_hit = True
   1138             else:
   1139                 problem = cause
   1140     # 3. Natural child exit zero inside the same work budget.
   1141     var status = ProcessStatus("pending", False, -1, 0, 0, "")
   1142     if not deadline_hit and problem == "":
   1143         if state.work_remaining_ms() <= 0:
   1144             deadline_hit = True
   1145         else:
   1146             status = state.wait_until(pid, state.work_remaining_ms())
   1147             if status.state == "running" or status.state == "interrupted":
   1148                 deadline_hit = True
   1149             elif status.state == "wait_error":
   1150                 problem = "child_wait_error"
   1151             elif not status.exited:
   1152                 problem = "child_signal_" + String(status.signal)
   1153             elif status.exit_code != 0:
   1154                 problem = "child_exit_" + String(status.exit_code)
   1155     # Cleanup: bounded allowance, separate from the work budget. A terminal
   1156     # status already observed for this exact child is cached and never re-waited
   1157     # (a second wait could only report "gone"), and an unproved cleanup retains
   1158     # the usable ownership handle instead of closing its descriptor.
   1159     if not status.cleanup_proved():
   1160         var terminal = state.terminate_once(pid, TERMINATION_GRACE_MS)
   1161         state.status = terminal.copy()
   1162         status = terminal.copy()
   1163     if status.cleanup_proved():
   1164         state.reaped = True
   1165         state.guard[].resolve_pid(pid)
   1166         state.close_reader()
   1167     else:
   1168         state.record_unproved(
   1169             pid,
   1170             state.report_fd,
   1171             "bounded call cleanup unproved",
   1172             "unreaped:" + status.describe(),
   1173         )
   1174     var cleanup_proved = status.cleanup_proved()
   1175     # 4. Report validation for the declared call.
   1176     var parsed = _empty_outcome("")
   1177     var completed = (not deadline_hit) and problem == ""
   1178     if completed:
   1179         try:
   1180             parsed = parse_bounded_report(report_text, kind, correlation)
   1181         except:
   1182             parsed = _empty_outcome("report_parse_raised")
   1183         if not parsed.ok:
   1184             completed = False
   1185             problem = parsed.problem
   1186     var elapsed_ms = now_ms() - start
   1187     return BoundedCallReport(
   1188         completed=completed,
   1189         stopped=deadline_hit,
   1190         outcome=parsed.outcome,
   1191         status=parsed.status,
   1192         cause=parsed.cause,
   1193         reason=parsed.reason,
   1194         latency_ms=parsed.latency_ms,
   1195         head_ms=parsed.head_ms,
   1196         total_ms=parsed.total_ms,
   1197         body_bytes=parsed.body_bytes,
   1198         body_match=parsed.body_match,
   1199         declared_bytes=parsed.declared_bytes,
   1200         length_match=parsed.length_match,
   1201         surplus_bytes=parsed.surplus_bytes,
   1202         problem=problem,
   1203         report=String(report_text.strip()),
   1204         elapsed_ms=elapsed_ms,
   1205         child_status=status.describe(),
   1206         cleanup_proved=cleanup_proved,
   1207     )