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^