hyf

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

stdio_process_helper.mojo (15238B)


      1 """Bounded, parent-owned stdio entrypoint runner for HYF process tests.
      2 
      3 ADR-0014 D33 FX06/FX08 and ADR-0015 D35 LC04: the *parent* enforces one finite
      4 startup/write/read/wait budget, bounds stdout/stderr/ready output, drains
      5 stdout and stderr concurrently with writing the request so a chatty child
      6 cannot deadlock, ignores SIGPIPE so a peer-close race cannot kill the parent,
      7 and reaps/closes every owned descriptor on success, assertion failure, timeout
      8 or early return. The child ``alarm`` is defense in depth only.
      9 
     10 The request/response API is unchanged so existing call sites keep working.
     11 """
     12 
     13 import std.os
     14 from std.collections import List
     15 from std.ffi import CStringSlice, c_int, external_call
     16 
     17 from parent_lifecycle import (
     18     IO_DEADLINE_EXPIRED,
     19     LIFECYCLE_POLL_SLICE_MS,
     20     POLLERR,
     21     POLLHUP,
     22     POLLIN,
     23     POLLNVAL,
     24     POLLOUT,
     25     TERMINATION_GRACE_MS,
     26     bytes_to_string,
     27     child_exit,
     28     close_fd,
     29     dup2_fd,
     30     fork_owned_or_close3,
     31     make_three_pipes,
     32     now_ms,
     33     poll_three,
     34     read_fd,
     35     set_alarm,
     36     terminate_owned,
     37     wait_bounded,
     38     write_fd_chunk,
     39 )
     40 from safe_tempdir import SafeTempDir
     41 
     42 from json import Value, loads
     43 
     44 
     45 comptime HYF_PATHS_PROFILE_ENV = "HYF_PATHS_PROFILE"
     46 comptime HYF_PATHS_REPO_LOCAL_ROOT_ENV = "HYF_PATHS_REPO_LOCAL_ROOT"
     47 
     48 # The entrypoint compiles and runs a real Mojo program; this is a finite
     49 # build/run lane budget, not a fixture lifetime and not a product SLO. One
     50 # budget covers compile, write, drain and wait; only the bounded cleanup grace
     51 # is added for reaping.
     52 comptime STDIO_ENTRYPOINT_DEADLINE_MS = 120000
     53 comptime STDIO_CHILD_ALARM_SECONDS = 180
     54 comptime STDIO_MAX_STDOUT_BYTES = 2097152
     55 comptime STDIO_MAX_STDERR_BYTES = 65536
     56 
     57 
     58 struct ScopedEnvVar:
     59     var name: String
     60     var value: String
     61     var previous: String
     62     var had_previous: Bool
     63 
     64     def __init__(out self, name: String, value: String):
     65         self.name = String(name)
     66         self.value = String(value)
     67         self.previous = std.os.getenv(name)
     68         self.had_previous = self.previous != ""
     69 
     70     def __enter__(mut self) raises:
     71         _ = std.os.setenv(self.name, self.value, overwrite=True)
     72 
     73     def __exit__(mut self):
     74         if self.had_previous:
     75             _ = std.os.setenv(self.name, self.previous, overwrite=True)
     76         else:
     77             _ = std.os.unsetenv(self.name)
     78 
     79 
     80 @fieldwise_init
     81 struct DrainOutcome(Movable):
     82     var eof: Bool
     83     var reason: String
     84 
     85 
     86 def drain_ready(
     87     fd: Int, mut out: List[UInt8], cap: Int, revents: Int, deadline_ms: Int
     88 ) -> DrainOutcome:
     89     """Read one ready descriptor into a capped buffer; preserve the cause.
     90 
     91     The read's EINTR retry is bounded by the caller's remaining
     92     ``deadline_ms``, so a retried read returns to the deadline owner instead of
     93     spinning, and an expired retry is a distinct ``read_deadline_expired``
     94     cause rather than a clean EOF.
     95     """
     96     if (revents & (POLLIN | POLLHUP | POLLERR | POLLNVAL)) == 0:
     97         return DrainOutcome(False, "")
     98     var buf = InlineArray[Byte, 4096](fill=0)
     99     var n = read_fd(fd, buf.unsafe_ptr(), 4096, deadline_ms)
    100     if n == IO_DEADLINE_EXPIRED:
    101         return DrainOutcome(True, "read_deadline_expired")
    102     if n < 0:
    103         return DrainOutcome(True, "read_error")
    104     if n == 0:
    105         return DrainOutcome(True, "")
    106     if len(out) + n > cap:
    107         return DrainOutcome(True, "stream_overflow")
    108     for index in range(n):
    109         out.append(UInt8(Int(buf[index])))
    110     return DrainOutcome(False, "")
    111 
    112 
    113 def run_stdio_entrypoint(
    114     entrypoint: String, request_json: String
    115 ) raises -> Value:
    116     return run_stdio_entrypoint_with_2_args(entrypoint, request_json, "", "")
    117 
    118 
    119 def run_stdio_entrypoint(
    120     entrypoint: String, request_json: String, arg0: String, arg1: String
    121 ) raises -> Value:
    122     return run_stdio_entrypoint_with_2_args(
    123         entrypoint, request_json, arg0, arg1
    124     )
    125 
    126 
    127 def _stdio_phase_summary(
    128     phase: String,
    129     request_sent: Bool,
    130     stdout_eof: Bool,
    131     stderr_eof: Bool,
    132     output_bytes: Int,
    133     output_valid: Bool,
    134     work_ms: Int,
    135     cleanup_ms: Int,
    136 ) -> String:
    137     """Bounded phase/time evidence for a stdio failure (ADR-0018 RA03).
    138 
    139     Work and cleanup time are recorded separately so an expired work budget can
    140     never be confused with the bounded cleanup allowance, and the phase fields
    141     prove which stage the child actually reached.
    142     """
    143     return (
    144         "stdio-entrypoint phase="
    145         + phase
    146         + " request_sent="
    147         + ("true" if request_sent else "false")
    148         + " stdout_eof="
    149         + ("true" if stdout_eof else "false")
    150         + " stderr_eof="
    151         + ("true" if stderr_eof else "false")
    152         + " output_bytes="
    153         + String(output_bytes)
    154         + " output_valid="
    155         + ("true" if output_valid else "false")
    156         + " work_ms="
    157         + String(work_ms)
    158         + " cleanup_ms="
    159         + String(cleanup_ms)
    160     )
    161 
    162 
    163 def _output_is_valid_json(bytes: List[UInt8]) -> Bool:
    164     """True when the captured stdout is nonempty and parses as JSON."""
    165     if len(bytes) == 0:
    166         return False
    167     try:
    168         _ = loads(bytes_to_string(bytes))
    169         return True
    170     except:
    171         return False
    172 
    173 
    174 def _terminate_and_raise(
    175     pid: Int,
    176     stdin_fd: Int,
    177     stdout_fd: Int,
    178     stderr_fd: Int,
    179     reason: String,
    180     phase: String,
    181     request_sent: Bool,
    182     stdout_eof: Bool,
    183     stderr_eof: Bool,
    184     output_bytes: Int,
    185     output_valid: Bool,
    186     work_ms: Int,
    187 ) raises:
    188     var cleanup_start = now_ms()
    189     var st = terminate_owned(pid, TERMINATION_GRACE_MS)
    190     var cleanup_ms = now_ms() - cleanup_start
    191     close_fd(stdin_fd)
    192     close_fd(stdout_fd)
    193     close_fd(stderr_fd)
    194     raise Error(
    195         _stdio_phase_summary(
    196             phase,
    197             request_sent,
    198             stdout_eof,
    199             stderr_eof,
    200             output_bytes,
    201             output_valid,
    202             work_ms,
    203             cleanup_ms,
    204         )
    205         + " reason="
    206         + reason
    207         + " (child "
    208         + st.describe()
    209         + " cleanup_error="
    210         + String("" if st.cleanup_proved() else "unreaped")
    211         + ")"
    212     )
    213 
    214 
    215 def run_stdio_entrypoint_with_2_args(
    216     entrypoint: String, request_json: String, arg0: String, arg1: String
    217 ) raises -> Value:
    218     return run_stdio_entrypoint_with_deadline(
    219         entrypoint, request_json, arg0, arg1, STDIO_ENTRYPOINT_DEADLINE_MS
    220     )
    221 
    222 
    223 def run_stdio_binary_with_deadline(
    224     binary: String,
    225     request_json: String,
    226     arg0: String,
    227     arg1: String,
    228     deadline_ms: Int,
    229 ) raises -> Value:
    230     """Run an already-built test child directly with the same guarantees.
    231 
    232     ADR-0018 RA03 phase isolation: this launch seam performs no compilation, so
    233     a late-exit control can prove the child actually reached valid output,
    234     closed output and the wait phase before its late exit was rejected. The
    235     ordinary compile/run helper keeps its declared total budget.
    236     """
    237     var args = List[String]()
    238     if arg0 != "":
    239         args.append(arg0)
    240     if arg1 != "":
    241         args.append(arg1)
    242     return _run_stdio_launch(binary, args^, request_json, deadline_ms)
    243 
    244 
    245 def run_stdio_entrypoint_with_deadline(
    246     entrypoint: String,
    247     request_json: String,
    248     arg0: String,
    249     arg1: String,
    250     deadline_ms: Int,
    251 ) raises -> Value:
    252     var args = List[String]()
    253     args.append("run")
    254     args.append("-I")
    255     args.append("src")
    256     args.append(entrypoint)
    257     if arg0 != "":
    258         args.append(arg0)
    259     if arg1 != "":
    260         args.append(arg1)
    261     return _run_stdio_launch("mojo", args^, request_json, deadline_ms)
    262 
    263 
    264 def _run_stdio_launch(
    265     command: String,
    266     var args: List[String],
    267     request_json: String,
    268     deadline_ms: Int,
    269 ) raises -> Value:
    270     # Build argv before owning any descriptors so no exception window can leak
    271     # pipes between creation and the fork; fork failure alone is handled by the
    272     # rollback helper. ``parts`` owns the C-string text for the whole call, so
    273     # every argv entry points at a live buffer until after the fork/exec.
    274     var parts = List[String]()
    275     parts.append(command)
    276     for index in range(len(args)):
    277         parts.append(args[index])
    278     var argv = List[Optional[CStringSlice[ImmutAnyOrigin]]](
    279         length=len(parts) + 1, fill={}
    280     )
    281     var elements = parts.unsafe_ptr()
    282     for index in range(len(parts)):
    283         argv[index] = rebind[CStringSlice[ImmutAnyOrigin]](
    284             elements[index].as_c_string_slice()
    285         )
    286 
    287     var pipes = make_three_pipes()
    288     var stdin_read_fd = pipes.stdin_pipe.read_fd
    289     var stdin_write_fd = pipes.stdin_pipe.write_fd
    290     var stdout_read_fd = pipes.stdout_pipe.read_fd
    291     var stdout_write_fd = pipes.stdout_pipe.write_fd
    292     var stderr_read_fd = pipes.stderr_pipe.read_fd
    293     var stderr_write_fd = pipes.stderr_pipe.write_fd
    294     var command_ptr = elements[0].as_c_string_slice().unsafe_ptr()
    295     var argv_ptr = argv.unsafe_ptr()
    296 
    297     var pid = fork_owned_or_close3(pipes)
    298     if pid == 0:
    299         if dup2_fd(stdin_read_fd, 0) < 0:
    300             child_exit(126)
    301         if dup2_fd(stdout_write_fd, 1) < 0:
    302             child_exit(126)
    303         if dup2_fd(stderr_write_fd, 2) < 0:
    304             child_exit(126)
    305         close_fd(stdin_read_fd)
    306         close_fd(stdin_write_fd)
    307         close_fd(stdout_read_fd)
    308         close_fd(stdout_write_fd)
    309         close_fd(stderr_read_fd)
    310         close_fd(stderr_write_fd)
    311         _ = set_alarm(STDIO_CHILD_ALARM_SECONDS)
    312         _ = external_call["execvp", c_int](command_ptr, argv_ptr)
    313         child_exit(127)
    314 
    315     close_fd(stdin_read_fd)
    316     close_fd(stdout_write_fd)
    317     close_fd(stderr_write_fd)
    318 
    319     var request = request_json + "\n"
    320     var sent = 0
    321     var stdout = List[UInt8]()
    322     var stderr_bytes = List[UInt8]()
    323     var stdout_eof = False
    324     var stderr_eof = False
    325     var stdin_done = False
    326     var write_reason = ""
    327     var read_reason = ""
    328     var start = now_ms()
    329     var budget = deadline_ms
    330     if budget <= 0:
    331         budget = 1
    332 
    333     while not (stdin_done and stdout_eof and stderr_eof):
    334         var elapsed = now_ms() - start
    335         if elapsed >= budget:
    336             read_reason = "read_deadline_expired"
    337             break
    338         var remaining = budget - elapsed
    339         var slice_ms = min(LIFECYCLE_POLL_SLICE_MS, remaining)
    340         if slice_ms < 1:
    341             slice_ms = 1
    342         var ev_stdin = 0 if stdin_done else POLLOUT
    343         var ev_out = 0 if stdout_eof else POLLIN
    344         var ev_err = 0 if stderr_eof else POLLIN
    345         var pr = poll_three(
    346             stdin_write_fd,
    347             ev_stdin,
    348             stdout_read_fd,
    349             ev_out,
    350             stderr_read_fd,
    351             ev_err,
    352             slice_ms,
    353         )
    354         if pr.count < 0:
    355             read_reason = "poll_failed"
    356             break
    357         if pr.count == 0:
    358             continue
    359         if not stdin_done:
    360             if (pr.r0 & (POLLERR | POLLHUP | POLLNVAL)) != 0:
    361                 # An early peer close or error is a cause-specific failure even
    362                 # when the child later exits 0 with valid stdout: the intended
    363                 # request bytes were not delivered.
    364                 write_reason = "write_pipe_closed"
    365                 stdin_done = True
    366             elif (pr.r0 & POLLOUT) != 0:
    367                 var cw = write_fd_chunk(
    368                     stdin_write_fd,
    369                     request,
    370                     sent,
    371                     budget - (now_ms() - start),
    372                 )
    373                 if cw.reason != "":
    374                     write_reason = cw.reason
    375                     stdin_done = True
    376                 else:
    377                     sent += cw.written
    378                     if sent >= request.byte_length():
    379                         stdin_done = True
    380                         close_fd(stdin_write_fd)
    381                         stdin_write_fd = -1
    382         if not stdout_eof:
    383             var d = drain_ready(
    384                 stdout_read_fd,
    385                 stdout,
    386                 STDIO_MAX_STDOUT_BYTES,
    387                 pr.r1,
    388                 budget - (now_ms() - start),
    389             )
    390             stdout_eof = d.eof
    391             if d.reason != "":
    392                 read_reason = "stdout_" + d.reason
    393                 break
    394         if not stderr_eof:
    395             var d = drain_ready(
    396                 stderr_read_fd,
    397                 stderr_bytes,
    398                 STDIO_MAX_STDERR_BYTES,
    399                 pr.r2,
    400                 budget - (now_ms() - start),
    401             )
    402             stderr_eof = d.eof
    403             if d.reason != "":
    404                 read_reason = "stderr_" + d.reason
    405                 break
    406 
    407     var request_sent = sent >= request.byte_length()
    408     var work_ms = now_ms() - start
    409     var output_len = len(stdout)
    410     if write_reason == "" and not request_sent:
    411         # The loop only ends with stdin finished; guard any path that would
    412         # otherwise leave intended request bytes unwritten.
    413         write_reason = "write_incomplete"
    414     close_fd(stdin_write_fd)
    415     if write_reason != "":
    416         _terminate_and_raise(
    417             pid,
    418             -1,
    419             stdout_read_fd,
    420             stderr_read_fd,
    421             "write_" + write_reason,
    422             "write",
    423             request_sent,
    424             stdout_eof,
    425             stderr_eof,
    426             output_len,
    427             _output_is_valid_json(stdout),
    428             work_ms,
    429         )
    430     if read_reason != "":
    431         _terminate_and_raise(
    432             pid,
    433             -1,
    434             stdout_read_fd,
    435             stderr_read_fd,
    436             read_reason,
    437             "read",
    438             request_sent,
    439             stdout_eof,
    440             stderr_eof,
    441             output_len,
    442             _output_is_valid_json(stdout),
    443             work_ms,
    444         )
    445 
    446     var remaining = budget - (now_ms() - start)
    447     if remaining < 1:
    448         remaining = 1
    449     var st = wait_bounded(pid, remaining)
    450     work_ms = now_ms() - start
    451     if not st.cleanup_proved():
    452         _terminate_and_raise(
    453             pid,
    454             -1,
    455             stdout_read_fd,
    456             stderr_read_fd,
    457             "timeout",
    458             "wait",
    459             request_sent,
    460             stdout_eof,
    461             stderr_eof,
    462             output_len,
    463             _output_is_valid_json(stdout),
    464             work_ms,
    465         )
    466     close_fd(stdout_read_fd)
    467     close_fd(stderr_read_fd)
    468 
    469     var output = bytes_to_string(stdout)
    470     var diagnostics = bytes_to_string(stderr_bytes)
    471 
    472     if st.exited and st.exit_code == 127:
    473         raise Error(
    474             "stdio-entrypoint phase=exit exec_failed (stdout="
    475             + output
    476             + " stderr="
    477             + diagnostics
    478             + ")"
    479         )
    480     if not st.exited or st.exit_code != 0:
    481         raise Error(
    482             "stdio-entrypoint phase=exit child_failed ("
    483             + st.describe()
    484             + " stderr="
    485             + diagnostics
    486             + ")"
    487         )
    488     if output == "":
    489         raise Error("hyf process returned no stdout payload")
    490     return loads(output)
    491 
    492 
    493 def run_hyf_stdio(request_json: String) raises -> Value:
    494     var response = Value(None)
    495     with SafeTempDir() as temp_dir:
    496         with ScopedEnvVar(HYF_PATHS_PROFILE_ENV, "repo_local"):
    497             with ScopedEnvVar(HYF_PATHS_REPO_LOCAL_ROOT_ENV, temp_dir):
    498                 response = run_stdio_entrypoint("src/main.mojo", request_json)
    499     return response^