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 )