hyf

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

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^)