max_local_process_helper.mojo (29339B)
1 """Strict scripted MaxLocal HTTP fixture (ADR-0012 D29 / ADR-0014 D33). 2 3 Owned by the parent test process: startup, read, write and wait all have 4 parent-enforced deadlines (FX06) and cleanup is exception-safe (FX08). Fixture 5 verification failures are reported with a bounded phase/case/reason so tests 6 can distinguish intended rejections from compiler/loader/startup/signal 7 failures (FX07). 8 """ 9 10 from std.collections import List 11 from std.memory import ArcPointer 12 13 from flare.net import SocketAddr 14 from flare.tcp import TcpListener 15 from flare.utils import usleep 16 17 from parent_lifecycle import ( 18 FIXTURE_DEFAULT_DEADLINE_MS, 19 TERMINATION_GRACE_MS, 20 CleanupGuard, 21 PipedChildState, 22 ProcessStatus, 23 child_exit, 24 close_fd, 25 dup2_fd, 26 finalize_owned_failure, 27 fork_owned_or_close, 28 make_pipe, 29 parse_ready_or_cleanup, 30 piped_child_state, 31 set_alarm, 32 sleep_ms, 33 write_raw, 34 ) 35 from strict_fixture import ( 36 STRICT_MAX_REPORT_BYTES, 37 json_escape, 38 STRICT_COMPLETION_GRACE_MS, 39 ConnectionReader, 40 ExchangeScript, 41 FramedRequest, 42 ServeReport, 43 authorization_reason, 44 exchange_script, 45 parse_report, 46 render_response, 47 report_line, 48 report_status_matches_exit, 49 serve_scripts, 50 validate_convenience, 51 verify_exchange, 52 ) 53 54 55 comptime STUB_ALARM_SECONDS = 20 56 57 58 def allowed_max_local_paths() -> List[String]: 59 var paths = List[String]() 60 paths.append("/health") 61 paths.append("/v1/chat/completions") 62 return paths^ 63 64 65 def allowed_max_local_methods() -> List[String]: 66 var methods = List[String]() 67 methods.append("GET") 68 methods.append("POST") 69 return methods^ 70 71 72 def require_bearer_for(mode: String) -> Bool: 73 # Convenience modes never bypass route/method validation; the 74 # echo_authorization mode validates the exact Authorization header inside 75 # its own handler and answers 401 when it is missing/duplicated/spoofed. 76 return False 77 78 79 # ── Scripted response bodies ──────────────────────────────────────────────── 80 81 82 def _chat_completion(body: String) -> String: 83 return '{"choices":[{"message":{"content":' + json_escape(body) + "}}]}" 84 85 86 def query_rewrite_analysis() -> String: 87 return ( 88 '{"original_text":"local apples pickup weekend",' 89 '"normalized_text":"local apples pickup weekend",' 90 '"rewritten_text":"apples pickup weekend",' 91 '"query_terms":["apples","pickup","weekend"],' 92 '"normalization_signals":["lowercase","local_intent_detected"],' 93 '"ranking_hints":["prefer_local_results","prefer_pickup"],' 94 '"extracted_filters":{' 95 '"local_intent":true,' 96 '"fulfillment":"pickup",' 97 '"time_window":"weekend"' 98 "}}" 99 ) 100 101 102 def _schema_invalid_analysis() -> String: 103 return ( 104 '{"original_text":"local apples pickup weekend",' 105 '"normalized_text":"local apples pickup weekend",' 106 '"query_terms":["apples","pickup","weekend"],' 107 '"normalization_signals":["lowercase","local_intent_detected"],' 108 '"ranking_hints":["prefer_local_results","prefer_pickup"],' 109 '"extracted_filters":{' 110 '"local_intent":true,' 111 '"fulfillment":"pickup",' 112 '"time_window":"weekend"' 113 "}}" 114 ) 115 116 117 def _health_status(mode: String) -> Int: 118 if mode == "health_non_2xx": 119 return 503 120 return 200 121 122 123 def _chat_status(mode: String) -> Int: 124 if mode == "query_rewrite_non_2xx": 125 return 503 126 return 200 127 128 129 def _raw_response(mode: String, path: String) -> String: 130 if path == "/health" and mode == "health_malformed_http": 131 return "not an http response\r\n\r\n" 132 if ( 133 path == "/v1/chat/completions" 134 and mode == "query_rewrite_malformed_http" 135 ): 136 return "not an http response\r\n\r\n" 137 return "" 138 139 140 def _delay_ms(mode: String, path: String) -> Int: 141 if path == "/health" and mode == "health_timeout": 142 return 1000 143 if path == "/health" and mode == "query_rewrite_remaining_deadline_timeout": 144 return 200 145 if path == "/v1/chat/completions": 146 if mode == "query_rewrite_timeout": 147 return 2000 148 if mode == "query_rewrite_remaining_deadline_timeout": 149 return 400 150 if mode == "stall": 151 return 30000 152 return 0 153 154 155 def _chat_body( 156 mode: String, request_index: Int, connection_index: Int 157 ) -> String: 158 if mode == "count_requests": 159 return ( 160 '{"request_index":' 161 + String(request_index) 162 + ',"connection_index":' 163 + String(connection_index) 164 + "}" 165 ) 166 if mode == "query_rewrite_ok": 167 return _chat_completion(query_rewrite_analysis()) 168 if mode == "query_rewrite_non_2xx": 169 return '{"error":{"message":"provider unavailable"}}' 170 if mode == "query_rewrite_invalid_json": 171 return '{"choices":[{"message":{"content":"not json"}}]}' 172 if mode == "query_rewrite_schema_invalid": 173 return _chat_completion(_schema_invalid_analysis()) 174 if mode == "query_rewrite_top_level_string": 175 return '"not object"' 176 if mode == "query_rewrite_top_level_array": 177 return "[]" 178 if mode == "query_rewrite_top_level_null": 179 return "null" 180 if mode == "query_rewrite_empty_choices": 181 return '{"choices":[]}' 182 if mode == "query_rewrite_missing_content": 183 return '{"choices":[{"message":{}}]}' 184 if mode == "query_rewrite_error_payload": 185 return '{"error":{"message":"provider refusal"}}' 186 if mode in ( 187 "query_rewrite_timeout", 188 "query_rewrite_remaining_deadline_timeout", 189 ): 190 return _chat_completion(query_rewrite_analysis()) 191 return '{"error":"unsupported_mode"}' 192 193 194 def _health_body(mode: String) -> String: 195 if mode == "health_non_2xx": 196 return '{"status":"unavailable"}' 197 return '{"status":"ok"}' 198 199 200 def _build_script( 201 mode: String, 202 framed: FramedRequest, 203 request_index: Int, 204 connection_index: Int, 205 ) -> ExchangeScript: 206 var path = framed.path 207 var status = _chat_status(mode) 208 if path == "/health": 209 status = _health_status(mode) 210 var script = exchange_script(mode, framed.method, path, status, "") 211 script.delay_ms = _delay_ms(mode, path) 212 script.close_connection = True 213 script.require_bearer = require_bearer_for(mode) 214 var raw = _raw_response(mode, path) 215 if raw != "": 216 script.raw_response = raw 217 elif mode == "echo_authorization" and path == "/v1/chat/completions": 218 var auth = authorization_reason(framed.headers_raw, True) 219 if auth != "": 220 script.status = 401 221 script.response_body = '{"error":"' + auth + '"}' 222 else: 223 script.echo_authorization = True 224 elif mode == "echo_body_bytes" and path == "/v1/chat/completions": 225 script.response_body = ( 226 '{"received_bytes":' + String(framed.body.byte_length()) + "}" 227 ) 228 elif path == "/health": 229 script.response_body = _health_body(mode) 230 else: 231 script.response_body = _chat_body(mode, request_index, connection_index) 232 return script^ 233 234 235 # ── Child serve loop ──────────────────────────────────────────────────────── 236 237 238 def serve_max_local( 239 port: Int, mode: String, requests: Int 240 ) raises -> ServeReport: 241 var allowed = allowed_max_local_paths() 242 var methods = allowed_max_local_methods() 243 var listener = TcpListener.bind(SocketAddr.localhost(UInt16(port))) 244 var actual_port = Int(listener.local_addr().port) 245 write_raw(1, "ready " + String(actual_port) + "\n") 246 var request_count = 0 247 var connection_count = 0 248 try: 249 while request_count < requests: 250 var stream = listener.accept() 251 connection_count += 1 252 var reader = ConnectionReader(stream^) 253 while request_count < requests: 254 var framed = reader.read() 255 if not framed.ok: 256 if framed.error == "empty": 257 if request_count < requests: 258 return ServeReport( 259 False, 260 "accounting", 261 mode, 262 "missing_exchanges", 263 request_count, 264 connection_count, 265 ) 266 break 267 return ServeReport( 268 False, 269 "read", 270 mode, 271 framed.error, 272 request_count, 273 connection_count, 274 ) 275 var reason = validate_convenience( 276 framed, allowed, methods, require_bearer_for(mode) 277 ) 278 if reason != "": 279 return ServeReport( 280 False, 281 "exchange", 282 mode, 283 reason, 284 request_count, 285 connection_count, 286 ) 287 var next_index = request_count + 1 288 var script = _build_script( 289 mode, framed, next_index, connection_count 290 ) 291 var verify = verify_exchange(script, framed) 292 if verify != "": 293 return ServeReport( 294 False, 295 "exchange", 296 mode, 297 verify, 298 request_count, 299 connection_count, 300 ) 301 if next_index == requests: 302 var extra = reader.probe_completion( 303 STRICT_COMPLETION_GRACE_MS 304 ) 305 if extra != "": 306 return ServeReport( 307 False, 308 "accounting", 309 mode, 310 extra, 311 request_count, 312 connection_count, 313 ) 314 request_count = next_index 315 if script.delay_ms > 0: 316 usleep(script.delay_ms * 1000) 317 reader.write_all( 318 render_response( 319 script, 320 framed.headers_raw, 321 request_count, 322 connection_count, 323 ) 324 ) 325 if script.close_connection: 326 break 327 if request_count < requests: 328 return ServeReport( 329 False, 330 "accounting", 331 mode, 332 "missing_exchanges", 333 request_count, 334 connection_count, 335 ) 336 return ServeReport( 337 True, "complete", mode, "ok", request_count, connection_count 338 ) 339 except: 340 return ServeReport( 341 False, "read", mode, "io_error", request_count, connection_count 342 ) 343 344 345 # ── Parent side ───────────────────────────────────────────────────────────── 346 347 348 struct SpawnedMaxLocalStub(Movable): 349 """Single-owner MaxLocal fixture handle. 350 351 The body receives a :class:`SpawnedMaxLocalView` from ``__enter__`` while 352 ``__exit__`` on this manager performs owned cleanup, so parent assertion, 353 exception and early-return paths all reap and close the child without any 354 shared heap state. Use ``with spawn_max_local_stub(...) as stub:``. 355 """ 356 357 var pid: Int 358 var port: Int 359 var state: PipedChildState 360 361 def __init__(out self, pid: Int, port: Int, var state: PipedChildState): 362 self.pid = pid 363 self.port = port 364 self.state = state^ 365 366 def __enter__(mut self) -> SpawnedMaxLocalView: 367 return SpawnedMaxLocalView(self.pid, self.port, UnsafePointer(to=self)) 368 369 def __exit__(mut self): 370 self.cleanup() 371 372 def cleanup(mut self): 373 """Fast, non-raising owned cleanup for assertion/error/early return. 374 375 Ownership is released only once the child is provably collected. An 376 uncertain wait keeps the handle retryable, records the failure in the 377 required caller-held guard and retains the exact pid/report descriptor 378 for recovery, instead of being silently marked complete. Never raises, 379 so a body/assertion cause is preserved at scope exit. 380 """ 381 if self.state.reaped: 382 return 383 var report_fd = self.state.report_fd 384 var status = self.state.terminate_once(self.pid, TERMINATION_GRACE_MS) 385 self.state.status = status.copy() 386 if status.cleanup_proved(): 387 self.state.reaped = True 388 self.state.close_reader() 389 self.state.guard[].resolve_pid(self.pid) 390 return 391 self.state.record_unproved( 392 self.pid, 393 report_fd, 394 "owned-child cleanup unproved", 395 "unreaped:" + status.describe(), 396 ) 397 398 def ok(self) -> Bool: 399 return self.state.ok 400 401 def phase(self) -> String: 402 return String(self.state.phase) 403 404 def failure_case(self) -> String: 405 return String(self.state.case_label) 406 407 def reason(self) -> String: 408 return String(self.state.reason) 409 410 def request_count(self) -> Int: 411 return self.state.requests 412 413 def connection_count(self) -> Int: 414 return self.state.connections 415 416 def cleanup_error(self) -> String: 417 return String(self.state.cleanup_error) 418 419 def describe(self) -> String: 420 return ( 421 "phase=" 422 + self.state.phase 423 + " case=" 424 + self.state.case_label 425 + " reason=" 426 + self.state.reason 427 + " requests=" 428 + String(self.state.requests) 429 + " connections=" 430 + String(self.state.connections) 431 ) 432 433 def status(mut self) -> ProcessStatus: 434 """Observe child status without losing ownership or report truth. 435 436 A terminal observation is cached so a later ``reap`` never performs a 437 new wait on a stale/reused identity. A transient ``wait_error`` or any 438 other nonterminal result is returned but deliberately *not* cached, so 439 it stays retryable and cannot overwrite a valid terminal result. 440 """ 441 if self.state.reaped: 442 return self.state.status.copy() 443 if self.state.observed_valid: 444 return self.state.observed.copy() 445 var st = self.state.wait_once(self.pid) 446 if st.state == "reaped" or st.state == "gone": 447 self.state.observed = st.copy() 448 self.state.observed_valid = True 449 return st^ 450 451 def reap(mut self): 452 """Reap the owned child within one declared finite work budget. 453 454 Never raises (so ``__exit__`` cannot mask a body error). Startup, wait, 455 report line and EOF drain share ``deadline_ms`` measured from spawn: an 456 expired budget fails with ``work_budget_expired`` and only the bounded 457 cleanup allowance, and no fresh successful interval is granted. 458 """ 459 if self.state.reaped: 460 return 461 if self.state.faults.wait_delay_ms > 0: 462 sleep_ms(self.state.faults.wait_delay_ms) 463 self.state.faults.wait_delay_ms = 0 464 var remaining = self.state.work_remaining_ms() 465 if remaining <= 0: 466 self.state.store( 467 False, "watchdog", "-", "work_budget_expired", 0, 0 468 ) 469 var expired = self.state.terminate_once( 470 self.pid, TERMINATION_GRACE_MS 471 ) 472 self.state.status = expired.copy() 473 if expired.cleanup_proved(): 474 self.state.reaped = True 475 self.state.close_reader() 476 self.state.guard[].resolve_pid(self.pid) 477 else: 478 self.state.reason = "work_budget_expired_unreaped" 479 self.state.record_unproved( 480 self.pid, 481 self.state.report_fd, 482 "owned-child cleanup unproved", 483 "unreaped:" + expired.describe(), 484 ) 485 return 486 var status = ProcessStatus("pending", False, -1, 0, 0, "") 487 if self.state.observed_valid: 488 status = self.state.observed.copy() 489 else: 490 status = self.state.wait_until(self.pid, remaining) 491 self.state.status = status.copy() 492 if status.state == "running" or status.state == "interrupted": 493 var term = self.state.terminate_once(self.pid, TERMINATION_GRACE_MS) 494 self.state.status = term.copy() 495 self.state.store(False, "watchdog", "-", "timeout", 0, 0) 496 if term.cleanup_proved(): 497 self.state.reaped = True 498 self.state.close_reader() 499 self.state.guard[].resolve_pid(self.pid) 500 else: 501 self.state.reason = "timeout_unreaped" 502 self.state.record_unproved( 503 self.pid, 504 self.state.report_fd, 505 "owned-child cleanup unproved", 506 "unreaped:" + term.describe(), 507 ) 508 return 509 if status.state == "gone": 510 # No waitable owned child remains: cleanup is proved without a 511 # signal and the report cannot be trusted, but ownership is done. 512 self.state.store(False, "watchdog", "-", "gone", 0, 0) 513 self.state.reaped = True 514 self.state.close_reader() 515 self.state.guard[].resolve_pid(self.pid) 516 return 517 if status.state == "wait_error": 518 # Identity/ownership is unproved: retain it for a retry and record 519 # the uncertainty instead of claiming the child was collected. 520 self.state.store(False, "watchdog", "-", "wait_error", 0, 0) 521 self.state.record_unproved( 522 self.pid, 523 self.state.report_fd, 524 "owned-child wait unproved", 525 "wait_error:" + status.error, 526 ) 527 return 528 if status.state != "reaped": 529 # Unexpected nonterminal/unknown taxonomy: fail closed and keep the 530 # exact ownership so a later retry can still collect it. 531 self.state.store(False, "watchdog", "-", "unexpected_status", 0, 0) 532 self.state.record_unproved( 533 self.pid, 534 self.state.report_fd, 535 "owned-child unexpected status", 536 "unexpected:" + status.describe(), 537 ) 538 return 539 var report_text = "" 540 var report_error = "" 541 var report_remaining = self.state.work_remaining_ms() 542 if report_remaining <= 0: 543 report_error = "report_budget_expired" 544 else: 545 try: 546 report_text = self.state.read_line( 547 STRICT_MAX_REPORT_BYTES, report_remaining 548 ) 549 if self.state.last_terminated: 550 var drain_remaining = self.state.work_remaining_ms() 551 if drain_remaining <= 0: 552 report_error = "drain_deadline_expired" 553 else: 554 var surplus = self.state.drain_surplus( 555 STRICT_MAX_REPORT_BYTES, drain_remaining 556 ) 557 if surplus > 0: 558 report_error = "duplicate_report" 559 elif report_text.byte_length() > 0: 560 # Bytes at EOF without a terminating newline are a 561 # truncated report, never a complete one. 562 report_error = "unterminated_report" 563 except e: 564 report_error = String(e) 565 self.state.close_reader() 566 if report_text == "" and report_error == "": 567 if status.exited and status.exit_code == 0: 568 self.state.store(False, "startup", "-", "missing_report", 0, 0) 569 elif status.exited: 570 self.state.store( 571 False, 572 "startup", 573 "-", 574 "exit_" + String(status.exit_code), 575 0, 576 0, 577 ) 578 else: 579 self.state.store( 580 False, 581 "watchdog", 582 "-", 583 "signal_" + String(status.signal), 584 0, 585 0, 586 ) 587 self.state.reaped = True 588 self.state.guard[].resolve_pid(self.pid) 589 return 590 if report_error != "": 591 self.state.store(False, "parse", "-", report_error, -1, -1) 592 self.state.reaped = True 593 self.state.guard[].resolve_pid(self.pid) 594 return 595 var parsed = parse_report(report_text) 596 if parsed.phase == "parse": 597 self.state.store(False, "parse", "-", parsed.reason, -1, -1) 598 self.state.reaped = True 599 self.state.guard[].resolve_pid(self.pid) 600 return 601 self.state.store( 602 parsed.ok, 603 parsed.phase, 604 parsed.case_label, 605 parsed.reason, 606 parsed.requests, 607 parsed.connections, 608 ) 609 if not report_status_matches_exit( 610 status.exited, status.exit_code, parsed.ok 611 ): 612 self.state.ok = False 613 self.state.phase = "startup" 614 self.state.reason = "report_status_mismatch_" + status.describe() 615 elif parsed.ok and parsed.requests != self.state.expected_requests: 616 self.state.ok = False 617 self.state.phase = "accounting" 618 self.state.reason = "request_count_mismatch" 619 elif parsed.ok and ( 620 parsed.connections < 1 621 or parsed.connections > self.state.expected_requests 622 ): 623 self.state.ok = False 624 self.state.phase = "accounting" 625 self.state.reason = "connection_count_invalid" 626 self.state.reaped = True 627 self.state.guard[].resolve_pid(self.pid) 628 629 def wait(mut self) raises: 630 self.reap() 631 if not self.state.ok: 632 raise Error("fixture-failure " + self.describe()) 633 634 def terminate(mut self) raises: 635 if self.state.reaped: 636 return 637 var report_fd = self.state.report_fd 638 var status = self.state.terminate_once(self.pid, TERMINATION_GRACE_MS) 639 self.state.status = status.copy() 640 if status.cleanup_proved(): 641 self.state.reaped = True 642 self.state.close_reader() 643 self.state.guard[].resolve_pid(self.pid) 644 return 645 self.state.record_unproved( 646 self.pid, 647 report_fd, 648 "owned-child cleanup unproved", 649 "unreaped:" + status.describe(), 650 ) 651 raise Error("lifecycle: owned child not reaped: " + status.describe()) 652 653 654 struct SpawnedMaxLocalView(Movable): 655 """Body-scope view of an owned MaxLocal fixture handle.""" 656 657 var pid: Int 658 var port: Int 659 var target: UnsafePointer[SpawnedMaxLocalStub, MutAnyOrigin] 660 661 def __init__( 662 out self, 663 pid: Int, 664 port: Int, 665 target: UnsafePointer[SpawnedMaxLocalStub, MutAnyOrigin], 666 ): 667 self.pid = pid 668 self.port = port 669 self.target = target 670 671 def ok(self) -> Bool: 672 return self.target[].ok() 673 674 def phase(self) -> String: 675 return self.target[].phase() 676 677 def failure_case(self) -> String: 678 return self.target[].failure_case() 679 680 def reason(self) -> String: 681 return self.target[].reason() 682 683 def request_count(self) -> Int: 684 return self.target[].request_count() 685 686 def connection_count(self) -> Int: 687 return self.target[].connection_count() 688 689 def cleanup_error(self) -> String: 690 return self.target[].cleanup_error() 691 692 def describe(self) -> String: 693 return self.target[].describe() 694 695 def status(mut self) -> ProcessStatus: 696 return self.target[].status() 697 698 def reap(mut self): 699 self.target[].reap() 700 701 def wait(mut self) raises: 702 self.target[].wait() 703 704 def terminate(mut self) raises: 705 self.target[].terminate() 706 707 def inject_cleanup_failure(mut self): 708 """Test-only bounded fault: force one unproved cleanup attempt. 709 710 The exact-owned child keeps running, so the follow-up retry exercises 711 the real recoverable ownership path rather than a synthetic identity. 712 """ 713 self.target[].state.faults.cleanup_failures += 1 714 715 def inject_wait_error(mut self): 716 """Test-only bounded fault: force one transient wait error.""" 717 self.target[].state.faults.wait_errors += 1 718 719 def inject_nonterminal_status(mut self): 720 """Test-only bounded fault: force one unexpected nonterminal status.""" 721 self.target[].state.faults.nonterminal += 1 722 723 def inject_wait_delay_ms(mut self, ms: Int): 724 """Test-only bounded fault: consume work-budget time in the wait phase. 725 """ 726 self.target[].state.faults.wait_delay_ms = ms 727 728 729 def reserve_loopback_port() raises -> Int: 730 var listener = TcpListener.bind(SocketAddr.localhost(0)) 731 var port = Int(listener.local_addr().port) 732 listener.close() 733 return port 734 735 736 def _serve_max_local_for( 737 port: Int, 738 var scripts: List[ExchangeScript], 739 mode: String, 740 requests: Int, 741 scripted: Bool, 742 ) raises -> ServeReport: 743 if scripted: 744 return serve_max_local_scripted(port, scripts^) 745 return serve_max_local(port, mode, requests) 746 747 748 def serve_max_local_scripted( 749 port: Int, var scripts: List[ExchangeScript] 750 ) raises -> ServeReport: 751 var listener = TcpListener.bind(SocketAddr.localhost(UInt16(port))) 752 var actual_port = Int(listener.local_addr().port) 753 write_raw(1, "ready " + String(actual_port) + "\n") 754 return serve_scripts(listener, scripts^, "scripted") 755 756 757 def spawn_max_local_stub( 758 port: Int, 759 mode: String, 760 requests: Int, 761 mut guard: CleanupGuard, 762 deadline_ms: Int = FIXTURE_DEFAULT_DEADLINE_MS, 763 ) raises -> SpawnedMaxLocalStub: 764 var scripts = List[ExchangeScript]() 765 return _spawn_max_local( 766 port, scripts^, mode, requests, False, deadline_ms, guard 767 ) 768 769 770 def spawn_max_local_scripted( 771 port: Int, 772 var scripts: List[ExchangeScript], 773 mut guard: CleanupGuard, 774 deadline_ms: Int = FIXTURE_DEFAULT_DEADLINE_MS, 775 ) raises -> SpawnedMaxLocalStub: 776 return _spawn_max_local( 777 port, scripts^, "scripted", len(scripts), True, deadline_ms, guard 778 ) 779 780 781 def _spawn_max_local( 782 port: Int, 783 var scripts: List[ExchangeScript], 784 mode: String, 785 requests: Int, 786 scripted: Bool, 787 deadline_ms: Int, 788 mut guard: CleanupGuard, 789 ) raises -> SpawnedMaxLocalStub: 790 var pipe = make_pipe() 791 var pid = fork_owned_or_close(pipe.copy()) 792 if pid == 0: 793 if dup2_fd(pipe.write_fd, 1) < 0: 794 child_exit(126) 795 close_fd(pipe.read_fd) 796 close_fd(pipe.write_fd) 797 _ = set_alarm(STUB_ALARM_SECONDS) 798 try: 799 var report = _serve_max_local_for( 800 port, scripts^, mode, requests, scripted 801 ) 802 _ = write_raw(1, report_line(report) + "\n") 803 child_exit(0 if report.ok else 125) 804 except: 805 var failed = ServeReport( 806 False, "startup", mode, "serve_failed", 0, 0 807 ) 808 _ = write_raw(1, report_line(failed) + "\n") 809 child_exit(125) 810 close_fd(pipe.write_fd) 811 var state = piped_child_state( 812 pid, pipe.read_fd, deadline_ms, requests, UnsafePointer(to=guard) 813 ) 814 var ready_line = "" 815 try: 816 ready_line = state.read_line(STRICT_MAX_REPORT_BYTES, deadline_ms) 817 except e: 818 var st = finalize_owned_failure( 819 state, pid, "max_local startup readiness cleanup unproved" 820 ) 821 raise Error( 822 "max_local stub readiness failed (" 823 + String(e) 824 + " pid=" 825 + String(pid) 826 + " / " 827 + st.describe() 828 + " cleanup=" 829 + ("proved" if st.cleanup_proved() else "unreaped") 830 + ")" 831 ) 832 var reported_port = 0 833 try: 834 reported_port = parse_ready_or_cleanup(pid, ready_line, 256) 835 except e: 836 var st = finalize_owned_failure( 837 state, pid, "max_local startup malformed-readiness cleanup unproved" 838 ) 839 raise Error( 840 "max_local stub malformed readiness (" 841 + String(e) 842 + " pid=" 843 + String(pid) 844 + " cleanup=" 845 + ("proved" if st.cleanup_proved() else "unreaped") 846 + ")" 847 ) 848 return SpawnedMaxLocalStub(pid, reported_port, state^)