parent_lifecycle.mojo (54261B)
1 """Governed test-only POSIX process/pipe lifecycle and deadline helpers. 2 3 ADR-0012 D29, ADR-0014 D33 FX06/FX08 and ADR-0015 D35 LC01-LC06 require 4 *parent-enforced*, finite startup/read/write/wait deadlines, exception-safe 5 cleanup/reaping for provider fixture children and ``run_stdio_entrypoint``, a 6 truthful wait-error taxonomy, byte caps enforced inside each read chunk and a 7 non-opening numeric descriptor census. A child ``alarm(2)`` watchdog is 8 defense in depth, never the parent's lifecycle proof. 9 10 Mojo's ``std.os.Pipe`` wrapper does not expose stable raw descriptors for 11 ``poll(2)``-based bounded I/O, so this module owns the required POSIX surface 12 directly. It is test-only tooling: it changes no HYF product policy. 13 14 Ownership rules: 15 16 * ``make_pipe``/``fork_pid`` create resources owned by the calling test. 17 * ``terminate_owned`` only signals a pid this process forked and has not yet 18 reaped, so a reused PID can never be targeted by a repeated teardown. 19 * ``PipedChildState`` is the single shared mutable lifecycle record for one 20 owned fixture child; every copy of a provider handle shares it, so cleanup is 21 idempotent across the ``with`` manager, the body handle and repeated calls. 22 * No broad ``pkill``/name matching is performed anywhere. 23 24 This toolchain exposes ``std.ffi.get_errno``/``ErrNo``, so the wait-error 25 taxonomy classifies a negative ``waitpid`` by real ``EINTR`` (retain and 26 retry ownership), ``ECHILD`` (no waitable owned child remains) and every other 27 errno (``wait_error``, never completed cleanup). 28 """ 29 30 from std.collections import List 31 from std.ffi import ( 32 ErrNo, 33 c_int, 34 c_uint, 35 c_ssize_t, 36 c_size_t, 37 external_call, 38 get_errno, 39 ) 40 from std.sys._libc import close 41 from std.time import perf_counter_ns 42 43 44 comptime WNOHANG: Int = 1 45 comptime POLLIN: Int = 1 46 comptime POLLOUT: Int = 4 47 comptime POLLERR: Int = 8 48 comptime POLLHUP: Int = 16 49 comptime POLLNVAL: Int = 32 50 comptime SIGALRM: Int = 14 51 comptime SIGPIPE: Int = 13 52 comptime SIGKILL: Int = 9 53 comptime SIGTERM: Int = 15 54 comptime SIG_IGN: Int = 1 55 comptime F_GETFD: Int = 1 56 comptime CENSUS_MAX_FDS: Int = 1048576 57 58 # ``read_fd``/``_write_fd`` return this when the caller-supplied remaining 59 # budget expires while a syscall is retried after EINTR. It is a distinct 60 # cause from a real ``read(2)``/``write(2)`` error, so a bounded retry loop can 61 # never be reported as plain I/O failure or as success. 62 comptime IO_DEADLINE_EXPIRED: Int = -2 63 64 # ``read_fd``/``_write_fd`` test-only retry seam: a negative fault count means 65 # "retry forever", so the deadline branch is reached deterministically. It is 66 # only valid together with a non-negative deadline. 67 comptime IO_FAULT_EINTR_UNBOUNDED: Int = -1 68 69 comptime FIXTURE_DEFAULT_DEADLINE_MS: Int = 20000 70 comptime TERMINATION_GRACE_MS: Int = 2000 71 comptime LIFECYCLE_POLL_SLICE_MS: Int = 25 72 # Upper bound on the EINTR retry loop when no finite deadline is supplied, so a 73 # repeated signal can never spin unbounded; a finite deadline always wins. 74 comptime POLL_EINTR_RETRY_BOUND: Int = 64 75 76 77 def now_ms() -> Int: 78 return Int(perf_counter_ns() // 1_000_000) 79 80 81 # ── Raw descriptor helpers ────────────────────────────────────────────────── 82 83 84 @fieldwise_init 85 struct PipeFds(Copyable, Movable): 86 var read_fd: Int 87 var write_fd: Int 88 89 def __copyinit__(out self, existing: Self): 90 self.read_fd = existing.read_fd 91 self.write_fd = existing.write_fd 92 93 94 @fieldwise_init 95 struct PipeTriple(Copyable, Movable): 96 var stdin_pipe: PipeFds 97 var stdout_pipe: PipeFds 98 var stderr_pipe: PipeFds 99 100 def __copyinit__(out self, existing: Self): 101 self.stdin_pipe = existing.stdin_pipe.copy() 102 self.stdout_pipe = existing.stdout_pipe.copy() 103 self.stderr_pipe = existing.stderr_pipe.copy() 104 105 106 def close_pipe(pipe: PipeFds): 107 close_fd(pipe.read_fd) 108 close_fd(pipe.write_fd) 109 110 111 def make_three_pipes(inject_fail_after: Int = -1) raises -> PipeTriple: 112 """Create three owned pipes, closing earlier ones if any creation fails. 113 114 ``inject_fail_after`` is a test-only control: when >= 0 a failure is 115 raised after that many successful pipes. The injected failure and a real 116 ``make_pipe`` failure share the same rollback handler. 117 """ 118 var fds = InlineArray[Int, 6](fill=-1) 119 try: 120 for index in range(3): 121 if inject_fail_after >= 0 and index == inject_fail_after: 122 raise Error("lifecycle: injected pipe creation failure") 123 var pipe = make_pipe() 124 fds[index * 2] = pipe.read_fd 125 fds[index * 2 + 1] = pipe.write_fd 126 except: 127 for slot in range(6): 128 close_fd(fds[slot]) 129 raise 130 return PipeTriple( 131 PipeFds(fds[0], fds[1]), 132 PipeFds(fds[2], fds[3]), 133 PipeFds(fds[4], fds[5]), 134 ) 135 136 137 def fork_owned_or_close( 138 pipe: PipeFds, inject_failure: Bool = False 139 ) raises -> Int: 140 """Fork the owned child, closing both pipe ends if the fork fails. 141 142 A real ``fork`` failure and the test-only injected failure share the same 143 rollback handler, so partial-startup cleanup is execution-proven. 144 """ 145 try: 146 if inject_failure: 147 raise Error("lifecycle: injected fork failure") 148 return fork_pid() 149 except: 150 close_pipe(pipe.copy()) 151 raise 152 return -1 153 154 155 def fork_owned_or_close3( 156 pipes: PipeTriple, inject_failure: Bool = False 157 ) raises -> Int: 158 """Fork the stdio child, closing all three pipe pairs if the fork fails.""" 159 try: 160 if inject_failure: 161 raise Error("lifecycle: injected fork failure") 162 return fork_pid() 163 except: 164 close_pipe(pipes.stdin_pipe.copy()) 165 close_pipe(pipes.stdout_pipe.copy()) 166 close_pipe(pipes.stderr_pipe.copy()) 167 raise 168 return -1 169 170 171 def ignore_sigpipe(): 172 """Ignore SIGPIPE so a peer-close race surfaces as EPIPE, not parent death. 173 174 Inherited across ``fork``, so owned fixture children get the same bounded 175 write-failure behaviour instead of dying on an interrupted report write. 176 """ 177 _ = external_call["signal", Int](c_int(SIGPIPE), c_int(SIG_IGN)) 178 179 180 def make_pipe() raises -> PipeFds: 181 ignore_sigpipe() 182 var fds = InlineArray[c_int, 2](fill=0) 183 if Int(external_call["pipe", c_int](fds.unsafe_ptr())) != 0: 184 raise Error("lifecycle: pipe failed") 185 return PipeFds(Int(fds[0]), Int(fds[1])) 186 187 188 def close_fd(fd: Int): 189 if fd >= 0: 190 _ = close(c_int(fd)) 191 192 193 def fork_pid() raises -> Int: 194 var pid = Int(external_call["fork", c_int]()) 195 if pid < 0: 196 raise Error("lifecycle: fork failed") 197 return pid 198 199 200 def child_exit(code: Int): 201 _ = external_call["_exit", c_int](c_int(code)) 202 203 204 def dup2_fd(oldfd: Int, newfd: Int) -> Int: 205 return Int(external_call["dup2", c_int](c_int(oldfd), c_int(newfd))) 206 207 208 def kill_pid(pid: Int, sig: Int) -> Int: 209 return Int(external_call["kill", c_int](c_int(pid), c_int(sig))) 210 211 212 def owned_pid() -> Int: 213 return Int(external_call["getpid", c_int]()) 214 215 216 def set_alarm(seconds: Int) -> Int: 217 return Int(external_call["alarm", c_uint](c_uint(seconds))) 218 219 220 def sleep_ms(ms: Int): 221 if ms > 0: 222 _ = external_call["usleep", c_int](c_int(ms * 1000)) 223 224 225 def poll_fd( 226 fd: Int, 227 events: Int, 228 timeout_ms: Int, 229 deadline_ms: Int = -1, 230 fault_eintr_count: Int = 0, 231 fault_error_count: Int = 0, 232 ) -> Int: 233 """Poll one descriptor with a bounded EINTR retry. 234 235 Returns the ``revents`` mask, ``0`` on timeout, ``-1`` on a real 236 ``poll(2)`` error, and ``IO_DEADLINE_EXPIRED`` when a retry after ``EINTR`` 237 would outlive the caller-supplied absolute remaining budget. ``deadline_ms`` 238 is the remaining part of the caller's budget and is never refreshed by the 239 retry loop, so a signal storm can neither spin forever nor be reported as a 240 real read/poll failure. 241 242 ``fault_eintr_count``/``fault_error_count`` are bounded test-only seams that 243 exercise the actual consumer's retry/error branches; they change no host 244 signal state. A negative ``fault_eintr_count`` retries until the deadline so 245 the bounded-retry branch is deterministically executable. 246 """ 247 var start = now_ms() 248 var eintrs = fault_eintr_count 249 var errors = fault_error_count 250 var attempts = 0 251 while True: 252 var slice = timeout_ms 253 if deadline_ms >= 0: 254 var remaining = deadline_ms - (now_ms() - start) 255 if remaining <= 0: 256 return IO_DEADLINE_EXPIRED 257 if remaining < slice: 258 slice = remaining 259 if errors != 0: 260 if errors > 0: 261 errors -= 1 262 return -1 263 if eintrs != 0: 264 if eintrs > 0: 265 eintrs -= 1 266 attempts += 1 267 if deadline_ms < 0 and attempts > POLL_EINTR_RETRY_BOUND: 268 return -1 269 continue 270 var cell = InlineArray[Int32, 2](fill=0) 271 cell[0] = Int32(fd) 272 cell[1] = Int32(events) 273 var n = Int( 274 external_call["poll", c_int]( 275 cell.unsafe_ptr(), c_uint(1), c_int(slice) 276 ) 277 ) 278 if n < 0: 279 if get_errno() != ErrNo.EINTR: 280 return -1 281 attempts += 1 282 if deadline_ms < 0 and attempts > POLL_EINTR_RETRY_BOUND: 283 return -1 284 continue 285 if n == 0: 286 return 0 287 return (Int(cell[1]) >> 16) & 0xFFFF 288 289 290 @fieldwise_init 291 struct PollThree(Movable): 292 """Bounded three-descriptor ``poll(2)`` result; ``count`` is -1 on error.""" 293 294 var count: Int 295 var r0: Int 296 var r1: Int 297 var r2: Int 298 299 300 def poll_three( 301 fd0: Int, 302 events0: Int, 303 fd1: Int, 304 events1: Int, 305 fd2: Int, 306 events2: Int, 307 timeout_ms: Int, 308 ) -> PollThree: 309 """Poll exactly three descriptors with one finite timeout.""" 310 var cell = InlineArray[Int32, 12](fill=0) 311 cell[0] = Int32(fd0) 312 cell[1] = Int32(events0) 313 cell[2] = Int32(fd1) 314 cell[3] = Int32(events1) 315 cell[4] = Int32(fd2) 316 cell[5] = Int32(events2) 317 var n = Int( 318 external_call["poll", c_int]( 319 cell.unsafe_ptr(), c_uint(3), c_int(timeout_ms) 320 ) 321 ) 322 if n < 0: 323 return PollThree(-1, 0, 0, 0) 324 if n == 0: 325 return PollThree(0, 0, 0, 0) 326 return PollThree( 327 n, 328 (Int(cell[1]) >> 16) & 0xFFFF, 329 (Int(cell[3]) >> 16) & 0xFFFF, 330 (Int(cell[5]) >> 16) & 0xFFFF, 331 ) 332 333 334 def read_fd( 335 fd: Int, 336 buf: UnsafePointer[Byte, ...], 337 max_bytes: Int, 338 deadline_ms: Int = -1, 339 fault_eintr_count: Int = 0, 340 ) -> Int: 341 """Read up to ``max_bytes`` with an EINTR retry bounded by the caller. 342 343 Returns the byte count, or the real ``read(2)`` error value. Returns 344 ``IO_DEADLINE_EXPIRED`` when a retry after ``EINTR`` would outlive the 345 caller-supplied remaining budget, so the retry loop returns to the deadline 346 owner instead of spinning indefinitely. ``deadline_ms < 0`` keeps the 347 legacy unbounded behaviour for short diagnostic writes/reads. 348 349 ``fault_eintr_count`` is a test-only seam: a positive count makes exactly 350 that many attempts behave like ``EINTR`` and ``IO_FAULT_EINTR_UNBOUNDED`` 351 retries forever, so the bounded-retry branch is deterministically 352 executable. The seam changes no host signal state and requires a deadline. 353 """ 354 var start = now_ms() 355 var faults = fault_eintr_count 356 if faults != 0 and deadline_ms < 0: 357 return IO_DEADLINE_EXPIRED 358 while True: 359 var synthetic = faults != 0 360 if faults > 0: 361 faults -= 1 362 if synthetic: 363 if now_ms() - start >= deadline_ms: 364 return IO_DEADLINE_EXPIRED 365 continue 366 var n = Int( 367 external_call["read", c_ssize_t](fd, buf, c_size_t(max_bytes)) 368 ) 369 if n >= 0 or get_errno() != ErrNo.EINTR: 370 return n 371 if deadline_ms >= 0 and now_ms() - start >= deadline_ms: 372 return IO_DEADLINE_EXPIRED 373 374 375 def _write_fd( 376 fd: Int, 377 ptr: UnsafePointer[UInt8, ...], 378 n: Int, 379 deadline_ms: Int = -1, 380 fault_eintr_count: Int = 0, 381 ) -> Int: 382 """Write ``n`` bytes with an EINTR retry bounded by ``deadline_ms``. 383 384 See ``read_fd``: ``IO_DEADLINE_EXPIRED`` is returned instead of retrying 385 past the caller's finite budget, and ``fault_eintr_count`` is the bounded 386 test-only seam for that branch. 387 """ 388 var start = now_ms() 389 var faults = fault_eintr_count 390 if faults != 0 and deadline_ms < 0: 391 return IO_DEADLINE_EXPIRED 392 while True: 393 var synthetic = faults != 0 394 if faults > 0: 395 faults -= 1 396 if synthetic: 397 if now_ms() - start >= deadline_ms: 398 return IO_DEADLINE_EXPIRED 399 continue 400 var written = Int( 401 external_call["write", c_ssize_t](fd, ptr, c_size_t(n)) 402 ) 403 if written >= 0 or get_errno() != ErrNo.EINTR: 404 return written 405 if deadline_ms >= 0 and now_ms() - start >= deadline_ms: 406 return IO_DEADLINE_EXPIRED 407 408 409 def write_raw(fd: Int, text: String) -> Int: 410 """Best-effort blocking write of a small bounded string (no deadline).""" 411 var n = _write_fd(fd, text.as_bytes().unsafe_ptr(), text.byte_length()) 412 return n 413 414 415 def write_raw_bytes(fd: Int, bytes: List[UInt8]) -> Int: 416 """Best-effort blocking write of raw bytes (test-only split controls).""" 417 if len(bytes) == 0: 418 return 0 419 return _write_fd(fd, bytes.unsafe_ptr(), len(bytes)) 420 421 422 comptime WRITE_CHUNK_BYTES: Int = 512 423 424 425 def write_fd_bounded(fd: Int, data: String, deadline_ms: Int) -> String: 426 """Write ``data`` with a parent-enforced deadline. 427 428 Chunks are capped at ``PIPE_BUF``-safe size so a ``POLLOUT`` readiness 429 never lets a blocking write stall past the deadline. Returns ``""`` on 430 success or a bounded reason such as ``write_deadline_expired`` / 431 ``write_pipe_closed``. 432 """ 433 var total = data.byte_length() 434 var sent = 0 435 var start = now_ms() 436 while sent < total: 437 if now_ms() - start >= deadline_ms: 438 return "write_deadline_expired" 439 var ev = poll_fd( 440 fd, 441 POLLOUT, 442 LIFECYCLE_POLL_SLICE_MS, 443 deadline_ms - (now_ms() - start), 444 ) 445 if ev == IO_DEADLINE_EXPIRED: 446 return "write_deadline_expired" 447 if ev < 0: 448 return "write_poll_error" 449 if ev == 0: 450 continue 451 if (ev & (POLLERR | POLLHUP | POLLNVAL)) != 0: 452 return "write_pipe_closed" 453 var chunk = min(WRITE_CHUNK_BYTES, total - sent) 454 var slice = data[byte = sent : sent + chunk] 455 var n = _write_fd( 456 fd, 457 slice.as_bytes().unsafe_ptr(), 458 chunk, 459 deadline_ms - (now_ms() - start), 460 ) 461 if n == IO_DEADLINE_EXPIRED: 462 return "write_deadline_expired" 463 if n <= 0: 464 return "write_failed" 465 sent += n 466 return "" 467 468 469 @fieldwise_init 470 struct ChunkWrite(Movable): 471 var reason: String 472 var written: Int 473 474 475 def write_fd_chunk( 476 fd: Int, 477 data: String, 478 offset: Int, 479 deadline_ms: Int = -1, 480 fault_eintr_count: Int = 0, 481 ) -> ChunkWrite: 482 """Write one ``PIPE_BUF``-safe chunk after a readiness poll. 483 484 Returns the bounded failure reason and the bytes actually written so a 485 caller can interleave writing with draining other descriptors. The write's 486 own EINTR retry is bounded by the caller's remaining ``deadline_ms``. 487 """ 488 var total = data.byte_length() 489 if offset >= total: 490 return ChunkWrite("", 0) 491 var chunk = min(WRITE_CHUNK_BYTES, total - offset) 492 var slice = data[byte = offset : offset + chunk] 493 var n = _write_fd( 494 fd, 495 slice.as_bytes().unsafe_ptr(), 496 chunk, 497 deadline_ms, 498 fault_eintr_count, 499 ) 500 if n == IO_DEADLINE_EXPIRED: 501 return ChunkWrite("write_deadline_expired", 0) 502 if n <= 0: 503 return ChunkWrite("write_failed", 0) 504 return ChunkWrite("", n) 505 506 507 # ── Bounded line reader with surplus retention ────────────────────────────── 508 509 510 @fieldwise_init 511 struct BoundedLineReader(Movable): 512 """Byte-bounded, deadline-bounded line reader that never discards surplus. 513 514 The byte cap is enforced *inside* every read chunk, so a coalesced chunk of 515 ``ready`` + ``report`` lines cannot smuggle an oversized line past the cap 516 and bytes after a returned newline stay available to the next reader call. 517 """ 518 519 var fd: Int 520 var max_bytes: Int 521 var _pending: List[UInt8] 522 var _pos: Int 523 var _eof: Bool 524 var _closed: Bool 525 526 def __init__(out self, fd: Int, max_bytes: Int): 527 self.fd = fd 528 self.max_bytes = max_bytes 529 self._pending = List[UInt8]() 530 self._pos = 0 531 self._eof = False 532 self._closed = False 533 534 def close(mut self): 535 if not self._closed: 536 close_fd(self.fd) 537 self._closed = True 538 539 def has_pending(self) -> Bool: 540 return self._pos < len(self._pending) 541 542 def _compact(mut self): 543 if self._pos == 0: 544 return 545 if self._pos >= len(self._pending): 546 self._pending = List[UInt8]() 547 self._pos = 0 548 return 549 var rest = List[UInt8]() 550 for index in range(self._pos, len(self._pending)): 551 rest.append(self._pending[index]) 552 self._pending = rest^ 553 self._pos = 0 554 555 def _take_available(mut self, mut out: List[UInt8], stop: Int): 556 for index in range(self._pos, stop): 557 out.append(self._pending[index]) 558 559 def read_line(mut self, deadline_ms: Int) raises -> String: 560 """Read one newline-terminated line with bounded size and deadline. 561 562 Raises ``ready_output_overflow`` once the line exceeds ``max_bytes`` 563 and ``read_deadline_expired`` when the deadline elapses first. EOF 564 before a newline returns the bytes read so far (possibly empty). 565 """ 566 var out = List[UInt8]() 567 var start = now_ms() 568 while True: 569 var found = -1 570 for index in range(self._pos, len(self._pending)): 571 if Int(self._pending[index]) == 10: 572 found = index 573 break 574 if found >= 0: 575 self._take_available(out, found) 576 self._pos = found + 1 577 self._compact() 578 if len(out) > self.max_bytes: 579 raise Error("ready_output_overflow") 580 return bytes_to_string(out) 581 self._take_available(out, len(self._pending)) 582 self._pos = len(self._pending) 583 self._compact() 584 if len(out) > self.max_bytes: 585 raise Error("ready_output_overflow") 586 if self._eof: 587 return bytes_to_string(out) 588 if now_ms() - start >= deadline_ms: 589 raise Error("read_deadline_expired") 590 var ev = poll_fd( 591 self.fd, 592 POLLIN, 593 LIFECYCLE_POLL_SLICE_MS, 594 deadline_ms - (now_ms() - start), 595 ) 596 if ev == IO_DEADLINE_EXPIRED: 597 raise Error("read_deadline_expired") 598 if ev < 0: 599 raise Error("read_error") 600 if ev == 0: 601 continue 602 var buf = InlineArray[Byte, 512](fill=0) 603 var n = read_fd( 604 self.fd, buf.unsafe_ptr(), 512, deadline_ms - (now_ms() - start) 605 ) 606 if n == IO_DEADLINE_EXPIRED: 607 raise Error("read_deadline_expired") 608 if n < 0: 609 raise Error("read_error") 610 if n == 0: 611 self._eof = True 612 continue 613 for index in range(n): 614 self._pending.append(UInt8(Int(buf[index]))) 615 616 617 def read_line_bounded( 618 fd: Int, max_bytes: Int, deadline_ms: Int 619 ) raises -> String: 620 """One-shot bounded line read for callers without surplus to preserve.""" 621 var reader = BoundedLineReader(fd, max_bytes) 622 return reader.read_line(deadline_ms) 623 624 625 def read_all_bounded( 626 fd: Int, max_bytes: Int, deadline_ms: Int 627 ) raises -> String: 628 """Read until EOF, with a bounded size cap and parent deadline. 629 630 Raises ``stdout_overflow`` when the cap is exceeded and 631 ``read_deadline_expired`` when the deadline elapses first. 632 """ 633 var out = List[UInt8]() 634 var buf = InlineArray[Byte, 4096](fill=0) 635 var start = now_ms() 636 while True: 637 if now_ms() - start >= deadline_ms: 638 raise Error("read_deadline_expired") 639 var ev = poll_fd( 640 fd, 641 POLLIN, 642 LIFECYCLE_POLL_SLICE_MS, 643 deadline_ms - (now_ms() - start), 644 ) 645 if ev == IO_DEADLINE_EXPIRED: 646 raise Error("read_deadline_expired") 647 if ev < 0: 648 raise Error("read_error") 649 if ev == 0: 650 continue 651 var n = read_fd( 652 fd, buf.unsafe_ptr(), 4096, deadline_ms - (now_ms() - start) 653 ) 654 if n == IO_DEADLINE_EXPIRED: 655 raise Error("read_deadline_expired") 656 if n < 0: 657 raise Error("read_error") 658 if n == 0: 659 break 660 if len(out) + n > max_bytes: 661 raise Error("stdout_overflow") 662 for index in range(n): 663 out.append(UInt8(Int(buf[index]))) 664 return bytes_to_string(out) 665 666 667 def bytes_to_string(bytes: List[UInt8]) raises -> String: 668 if len(bytes) == 0: 669 return "" 670 return String(from_utf8=Span(ptr=bytes.unsafe_ptr(), length=len(bytes))) 671 672 673 def _utf8_line(bytes: List[UInt8]) raises -> String: 674 """Decode a complete line as UTF-8, or fail with a bounded cause.""" 675 try: 676 return bytes_to_string(bytes) 677 except: 678 raise Error("invalid_utf8") 679 680 681 @fieldwise_init 682 struct CleanupFailure(Copyable, Movable): 683 """One recorded owned-child cleanup problem. 684 685 ``recoverable`` is true when the exact-owned child (``pid``/``report_fd``) 686 is still retained and a later retry can attempt cleanup again; ``resolved`` 687 flips once a retry proves cleanup or a subsequent reap collects it. 688 ``fd_closed`` makes the retained report descriptor close-once: whichever 689 party recovers first closes it, and no other call may target that number 690 again after it could have been reused. 691 """ 692 693 var text: String 694 var pid: Int 695 var report_fd: Int 696 var recoverable: Bool 697 var resolved: Bool 698 var fd_closed: Bool 699 700 def __copyinit__(out self, existing: Self): 701 self.text = existing.text 702 self.pid = existing.pid 703 self.report_fd = existing.report_fd 704 self.recoverable = existing.recoverable 705 self.resolved = existing.resolved 706 self.fd_closed = existing.fd_closed 707 708 709 def _new_failure( 710 text: String, pid: Int, report_fd: Int, recoverable: Bool 711 ) -> CleanupFailure: 712 return CleanupFailure( 713 text=String(text), 714 pid=pid, 715 report_fd=report_fd, 716 recoverable=recoverable, 717 resolved=False, 718 fd_closed=report_fd < 0, 719 ) 720 721 722 struct CleanupGuard(Movable): 723 """Required, caller-held observation point for owned-child cleanup truth. 724 725 ADR-0017 PC02 and ADR-0018 RA01/RA02 require that an ordinary supported 726 provider scope cannot silently discard a cleanup failure and that a failed 727 cleanup keeps a usable ownership handle. This guard is therefore a 728 *required* parameter of every supported spawner (there is no inert default) 729 and it is owned by the calling test, so its record survives the ``with`` 730 scope where the fixture handle is destroyed. 731 732 ``assert_clean`` is the enforcement point: every caller must invoke it, and 733 it raises the exact retained detail so the owning test fails truthfully. 734 ``retain`` additionally keeps the exact pid/report descriptor of an unproved 735 cleanup so ``recover_all`` can retry cleanup instead of leaving only a 736 message. The guard never signals an identity it was not given. 737 """ 738 739 var failures: List[CleanupFailure] 740 741 def __init__(out self): 742 self.failures = List[CleanupFailure]() 743 744 def record(mut self, text: String): 745 """Record a cleanup problem that has no retained recoverable child.""" 746 self.failures.append(_new_failure(text, 0, -1, False)) 747 748 def retain(mut self, pid: Int, report_fd: Int, text: String): 749 """Record an unproved cleanup and keep its exact ownership handle.""" 750 self.failures.append(_new_failure(text, pid, report_fd, True)) 751 752 def resolve_pid(mut self, pid: Int): 753 """Mark retained ownership resolved once cleanup is proved elsewhere. 754 755 Used when a later retry (``reap``/``cleanup``/``terminate``) collects a 756 child whose earlier attempt was unproved, so a transient failure does 757 not fail the owning test after a successful recovery. 758 """ 759 if pid <= 0: 760 return 761 for index in range(len(self.failures)): 762 if self.failures[index].pid == pid: 763 self.failures[index].resolved = True 764 765 def pending(self) -> Int: 766 """Number of unresolved recorded cleanup problems.""" 767 var count = 0 768 for index in range(len(self.failures)): 769 if not self.failures[index].resolved: 770 count += 1 771 return count 772 773 def retained(self) -> Int: 774 """Number of unresolved entries that still hold a recoverable child.""" 775 var count = 0 776 for index in range(len(self.failures)): 777 if ( 778 self.failures[index].recoverable 779 and not self.failures[index].resolved 780 ): 781 count += 1 782 return count 783 784 def count(self) -> Int: 785 return len(self.failures) 786 787 def first_pending(self) -> String: 788 for index in range(len(self.failures)): 789 if not self.failures[index].resolved: 790 return String(self.failures[index].text) 791 return "" 792 793 def first(self) -> String: 794 if len(self.failures) > 0: 795 return String(self.failures[0].text) 796 return "" 797 798 def last(self) -> String: 799 if len(self.failures) > 0: 800 return String(self.failures[len(self.failures) - 1].text) 801 return "" 802 803 def is_clean(self) -> Bool: 804 return self.pending() == 0 805 806 def close_retained_fd(mut self, fd: Int) -> Bool: 807 """Close a retained report descriptor at most once. 808 809 Returns True only when this guard still owns an *open* retained 810 descriptor with number ``fd`` and closes it now. An entry whose 811 descriptor was already closed has released that number, so it must not 812 claim or close a later descriptor that reused the same number (the 813 period-11 R73/MC03 defect: a recovered entry matched a new measurement's 814 descriptor and made it skip its own close, leaking one descriptor). 815 Returns False when the guard holds no open retained entry, letting the 816 normal proved path keep its own descriptor ownership. 817 """ 818 if fd < 0: 819 return False 820 var claimed = -1 821 for index in range(len(self.failures)): 822 if self.failures[index].report_fd != fd: 823 continue 824 if self.failures[index].fd_closed: 825 # Already released: this number now belongs to a new owner. 826 continue 827 claimed = index 828 break 829 if claimed < 0: 830 return False 831 close_fd(fd) 832 for index in range(len(self.failures)): 833 if self.failures[index].report_fd == fd: 834 self.failures[index].fd_closed = True 835 return True 836 837 def recover_all(mut self) -> Int: 838 """Retry cleanup for every retained exact-owned child. 839 840 Returns the number of children whose cleanup is still unproved. A 841 child whose retry proves cleanup is marked resolved and its retained 842 report descriptor is closed exactly once, so recovery is real rather 843 than a message-only ledger entry and leaves no leaked descriptor. Only 844 pids this guard was handed are ever signalled. 845 """ 846 var unresolved = 0 847 for index in range(len(self.failures)): 848 if not self.failures[index].recoverable: 849 continue 850 if self.failures[index].resolved: 851 continue 852 var pid = self.failures[index].pid 853 if pid <= 0: 854 unresolved += 1 855 continue 856 var st = terminate_owned(pid, TERMINATION_GRACE_MS) 857 if st.cleanup_proved(): 858 self.failures[index].resolved = True 859 _ = self.close_retained_fd(self.failures[index].report_fd) 860 else: 861 unresolved += 1 862 return unresolved 863 864 def assert_clean(mut self) raises: 865 """Fail the owning test if any cleanup problem remains unresolved.""" 866 var remaining = self.pending() 867 if remaining == 0: 868 return 869 raise Error( 870 "cleanup-unproved: " 871 + String(remaining) 872 + " unresolved owned-child cleanup failure(s); first: " 873 + self.first_pending() 874 ) 875 876 877 # ── Child lifecycle ───────────────────────────────────────────────────────── 878 879 880 @fieldwise_init 881 struct ProcessStatus(Copyable, Movable): 882 var state: String 883 var exited: Bool 884 var exit_code: Int 885 var signal: Int 886 var raw: Int 887 var error: String 888 889 def __copyinit__(out self, existing: Self): 890 self.state = existing.state 891 self.exited = existing.exited 892 self.exit_code = existing.exit_code 893 self.signal = existing.signal 894 self.raw = existing.raw 895 self.error = existing.error 896 897 def reaped(self) -> Bool: 898 return self.state == "reaped" 899 900 def cleanup_proved(self) -> Bool: 901 """True only when no waitable owned child can remain for this pid.""" 902 return self.state == "reaped" or self.state == "gone" 903 904 def describe(self) -> String: 905 if self.state != "reaped": 906 if self.error != "": 907 return self.state + ":" + self.error 908 return self.state 909 if self.exited: 910 return "exited=" + String(self.exit_code) 911 return "signal=" + String(self.signal) 912 913 914 def _decode_status(raw: Int) -> ProcessStatus: 915 var low = raw & 0x7F 916 if low == 0: 917 return ProcessStatus("reaped", True, (raw >> 8) & 0xFF, 0, raw, "") 918 if low == 0x7F: 919 return ProcessStatus("stopped", False, -1, 0, raw, "") 920 return ProcessStatus("reaped", False, -1, low, raw, "") 921 922 923 def classify_wait_errno(errno_value: Int) -> String: 924 """Map a real ``waitpid`` errno to the ownership taxonomy. 925 926 ``interrupted`` (EINTR) retains ownership and is retried; ``gone`` 927 (ECHILD) proves no waitable child remains; anything else is 928 ``wait_error`` and must never be reported as completed cleanup. 929 """ 930 if errno_value == Int(ErrNo.EINTR.value): 931 return "interrupted" 932 if errno_value == Int(ErrNo.ECHILD.value): 933 return "gone" 934 return "wait_error" 935 936 937 def wait_nohang(pid: Int) -> ProcessStatus: 938 """Non-blocking wait with an exact errno ownership taxonomy. 939 940 ``reaped``/``running`` are exact. A negative ``waitpid`` is classified by 941 the real errno: ``EINTR`` retains ownership and is retried by callers; 942 ``ECHILD`` proves no waitable child remains; any other error is a 943 ``wait_error`` that must not be reported as completed cleanup. A 944 non-positive pid is never waited on. 945 """ 946 if pid <= 0: 947 return ProcessStatus("wait_error", False, -1, 0, -1, "invalid_pid") 948 var status = InlineArray[c_int, 1](fill=0) 949 var r = Int( 950 external_call["waitpid", c_int]( 951 c_int(pid), status.unsafe_ptr(), c_int(WNOHANG) 952 ) 953 ) 954 if r == pid: 955 return _decode_status(Int(status[0])) 956 if r == 0: 957 return ProcessStatus("running", False, -1, 0, 0, "") 958 var errno_value = Int(get_errno().value) 959 var klass = classify_wait_errno(errno_value) 960 if klass == "gone": 961 return ProcessStatus("gone", False, -1, -1, -1, "") 962 return ProcessStatus( 963 klass, False, -1, 0, -1, "errno_" + String(errno_value) 964 ) 965 966 967 def wait_bounded(pid: Int, deadline_ms: Int) -> ProcessStatus: 968 var start = now_ms() 969 while True: 970 var st = wait_nohang(pid) 971 if st.state != "running" and st.state != "interrupted": 972 return st^ 973 if now_ms() - start >= deadline_ms: 974 return st^ 975 sleep_ms(5) 976 977 978 def terminate_owned(pid: Int, grace_ms: Int) -> ProcessStatus: 979 """Reap a child this test owns, escalating SIGTERM -> SIGKILL. 980 981 In ``terminate_owned`` a pid already reaped or not waitable (``gone``) is 982 never signaled, so a reused PID from an unrelated process can never be 983 targeted. An ``interrupted`` status still owns a live child and is 984 retried/signaled; a ``wait_error`` leaves identity/ownership unproved and is 985 returned without signalling and never as completed cleanup. 986 """ 987 var st = wait_nohang(pid) 988 if st.cleanup_proved(): 989 return st^ 990 if st.state == "wait_error": 991 # Identity/ownership is unproved; never signal and never claim cleanup. 992 return st^ 993 _ = kill_pid(pid, SIGTERM) 994 st = wait_bounded(pid, grace_ms) 995 if st.state == "running" or st.state == "interrupted": 996 _ = kill_pid(pid, SIGKILL) 997 st = wait_bounded(pid, grace_ms) 998 return st^ 999 1000 1001 def pid_not_waitable(pid: Int) -> Bool: 1002 """True when ``pid`` is neither running nor an unreaped zombie of ours. 1003 1004 Used only to evidence, after ``terminate_owned`` reaped a specific owned 1005 child, that the same pid is no longer waitable. It never reaps a pid this 1006 test did not fork and never scans by process name. 1007 """ 1008 return wait_nohang(pid).state == "gone" 1009 1010 1011 def pid_running(pid: Int) -> Bool: 1012 """True when ``pid`` is still a live owned child of this process. 1013 1014 A non-blocking observation only: it never reaps, signals or releases the 1015 exact ownership the guard retains for recovery. 1016 """ 1017 return wait_nohang(pid).state == "running" 1018 1019 1020 # ── Shared owned-child lifecycle state ────────────────────────────────────── 1021 1022 1023 @fieldwise_init 1024 struct LifecycleFaults(Movable): 1025 """Bounded test-only failure seam for real owned child resources. 1026 1027 The seam never replaces the exact ownership identity: it makes one wait or 1028 one cleanup attempt report an unproved result while the real forked child 1029 keeps running, so the following retry exercises the real recoverable path 1030 instead of an artificial ``pid=0`` handle. Every counter is consumed once. 1031 """ 1032 1033 var wait_errors: Int 1034 var nonterminal: Int 1035 var cleanup_failures: Int 1036 var wait_delay_ms: Int 1037 var poll_eintrs: Int 1038 var poll_errors: Int 1039 1040 def __init__(out self): 1041 self.wait_errors = 0 1042 self.nonterminal = 0 1043 self.cleanup_failures = 0 1044 self.wait_delay_ms = 0 1045 self.poll_eintrs = 0 1046 self.poll_errors = 0 1047 1048 def active(self) -> Bool: 1049 return ( 1050 self.wait_errors > 0 1051 or self.nonterminal > 0 1052 or self.cleanup_failures > 0 1053 or self.wait_delay_ms > 0 1054 or self.poll_eintrs != 0 1055 or self.poll_errors != 0 1056 ) 1057 1058 1059 @fieldwise_init 1060 struct PipedChildState(Movable): 1061 """Single mutable lifecycle record shared by every copy of one handle. 1062 1063 Retained read surplus is an undecoded byte buffer, so a chunk that splits a 1064 multi-byte character is preserved without raising on an incomplete UTF-8 1065 fragment. ``last_terminated`` records whether the most recent line ended 1066 with a newline, so a truncated report can never be read as complete. 1067 ``guard`` is the required caller-held cleanup observation point. 1068 """ 1069 1070 var pid: Int 1071 var report_fd: Int 1072 var pending: List[UInt8] 1073 var eof: Bool 1074 var closed: Bool 1075 var deadline_ms: Int 1076 var expected_requests: Int 1077 var reaped: Bool 1078 var ok: Bool 1079 var phase: String 1080 var case_label: String 1081 var reason: String 1082 var requests: Int 1083 var connections: Int 1084 var cleanup_error: String 1085 var status: ProcessStatus 1086 var observed: ProcessStatus 1087 var observed_valid: Bool 1088 var last_terminated: Bool 1089 var spawn_ms: Int 1090 var guard: UnsafePointer[CleanupGuard, MutAnyOrigin] 1091 var faults: LifecycleFaults 1092 1093 def store( 1094 mut self, 1095 ok: Bool, 1096 phase: String, 1097 case_label: String, 1098 reason: String, 1099 requests: Int, 1100 connections: Int, 1101 ): 1102 self.ok = ok 1103 self.phase = String(phase) 1104 self.case_label = String(case_label) 1105 self.reason = String(reason) 1106 self.requests = requests 1107 self.connections = connections 1108 1109 def close_reader(mut self): 1110 if self.closed: 1111 return 1112 # A descriptor retained by the guard for recovery is closed by the 1113 # guard exactly once, so the handle can never close a reused number. 1114 if not self.guard[].close_retained_fd(self.report_fd): 1115 close_fd(self.report_fd) 1116 self.closed = True 1117 1118 def work_remaining_ms(mut self) -> Int: 1119 """Remaining part of the one declared work budget for this scope. 1120 1121 ADR-0018 RA03: startup/body/wait/report/drain share the single finite 1122 budget measured from ``spawn_ms``. A nonpositive result must never be 1123 turned into another successful interval; callers fail and use only the 1124 separate bounded cleanup allowance. 1125 """ 1126 return self.deadline_ms - (now_ms() - self.spawn_ms) 1127 1128 def record_unproved( 1129 mut self, pid: Int, report_fd: Int, label: String, detail: String 1130 ): 1131 """Record an unproved cleanup and retain its usable ownership.""" 1132 self.cleanup_error = detail 1133 self.guard[].retain( 1134 pid, report_fd, label + " pid=" + String(pid) + " " + detail 1135 ) 1136 1137 def wait_once(mut self, pid: Int) -> ProcessStatus: 1138 """Observe one child state, applying the bounded test fault seam. 1139 1140 A ``wait_error`` from this observation is *not* cached as a terminal 1141 result by callers, so a transient wait problem stays retryable. 1142 """ 1143 if self.faults.wait_errors > 0: 1144 self.faults.wait_errors -= 1 1145 return ProcessStatus( 1146 "wait_error", False, -1, 0, -1, "injected_wait_error" 1147 ) 1148 if self.faults.nonterminal > 0: 1149 self.faults.nonterminal -= 1 1150 return ProcessStatus("stopped", False, -1, 0, 0x7F, "") 1151 return wait_nohang(pid) 1152 1153 def wait_until(mut self, pid: Int, deadline_ms: Int) -> ProcessStatus: 1154 """Bounded non-blocking wait that honours the fault seam.""" 1155 var start = now_ms() 1156 while True: 1157 var st = self.wait_once(pid) 1158 if st.state != "running" and st.state != "interrupted": 1159 return st^ 1160 if now_ms() - start >= deadline_ms: 1161 return st^ 1162 sleep_ms(5) 1163 1164 def terminate_once(mut self, pid: Int, grace_ms: Int) -> ProcessStatus: 1165 """Cleanup attempt that honours the bounded test fault seam. 1166 1167 A forced cleanup failure leaves the real forked child running and keeps 1168 the exact pid/report descriptor, so recovery is a real retry rather 1169 than a synthetic identity change. 1170 """ 1171 if self.faults.cleanup_failures > 0: 1172 self.faults.cleanup_failures -= 1 1173 return ProcessStatus( 1174 "wait_error", False, -1, 0, -1, "injected_cleanup_failure" 1175 ) 1176 return terminate_owned(pid, grace_ms) 1177 1178 def _poll_eintr_fault(mut self) -> Int: 1179 """Consume one bounded test-only EINTR fault, if one is armed. 1180 1181 A positive count is consumed once so a following poll is real; a 1182 negative count stays armed so the bounded-retry deadline branch is 1183 reached deterministically. This changes no host signal state. 1184 """ 1185 var value = self.faults.poll_eintrs 1186 if value > 0: 1187 self.faults.poll_eintrs -= 1 1188 return value 1189 1190 def _poll_error_fault(mut self) -> Int: 1191 """Consume one bounded test-only real poll-error fault, if armed.""" 1192 var value = self.faults.poll_errors 1193 if value > 0: 1194 self.faults.poll_errors -= 1 1195 return value 1196 1197 def read_line(mut self, max_bytes: Int, deadline_ms: Int) raises -> String: 1198 """Bounded line read that retains surplus as undecoded bytes. 1199 1200 The byte cap is enforced inside every read chunk (including a newline 1201 in the same chunk) and bytes after the returned newline stay buffered 1202 for the next consumer. ``last_terminated`` reports whether the returned 1203 line ended with a newline. Raises ``ready_output_overflow`` past the cap, 1204 ``read_deadline_expired`` on the read deadline, ``read_error`` for a 1205 real read/poll failure (never conflated with EOF) and ``invalid_utf8`` 1206 for a line that is not valid UTF-8. 1207 """ 1208 var line = List[UInt8]() 1209 var start = now_ms() 1210 self.last_terminated = False 1211 while True: 1212 var found = -1 1213 for index in range(len(self.pending)): 1214 if Int(self.pending[index]) == 10: 1215 found = index 1216 break 1217 if found >= 0: 1218 for index in range(found): 1219 line.append(self.pending[index]) 1220 var rest = List[UInt8]() 1221 for index in range(found + 1, len(self.pending)): 1222 rest.append(self.pending[index]) 1223 self.pending = rest^ 1224 if len(line) > max_bytes or len(self.pending) > max_bytes: 1225 raise Error("ready_output_overflow") 1226 self.last_terminated = True 1227 return _utf8_line(line^) 1228 for index in range(len(self.pending)): 1229 line.append(self.pending[index]) 1230 self.pending = List[UInt8]() 1231 if len(line) > max_bytes: 1232 raise Error("ready_output_overflow") 1233 if self.eof: 1234 return _utf8_line(line^) 1235 if now_ms() - start >= deadline_ms: 1236 raise Error("read_deadline_expired") 1237 var ev = poll_fd( 1238 self.report_fd, 1239 POLLIN, 1240 LIFECYCLE_POLL_SLICE_MS, 1241 deadline_ms - (now_ms() - start), 1242 self._poll_eintr_fault(), 1243 self._poll_error_fault(), 1244 ) 1245 if ev == IO_DEADLINE_EXPIRED: 1246 raise Error("read_deadline_expired") 1247 if ev < 0: 1248 raise Error("read_error") 1249 if ev == 0: 1250 continue 1251 var buf = InlineArray[Byte, 512](fill=0) 1252 var n = read_fd( 1253 self.report_fd, 1254 buf.unsafe_ptr(), 1255 512, 1256 deadline_ms - (now_ms() - start), 1257 ) 1258 if n == IO_DEADLINE_EXPIRED: 1259 raise Error("read_deadline_expired") 1260 if n < 0: 1261 raise Error("read_error") 1262 if n == 0: 1263 self.eof = True 1264 continue 1265 for index in range(n): 1266 self.pending.append(UInt8(Int(buf[index]))) 1267 1268 def drain_surplus(mut self, max_bytes: Int, deadline_ms: Int) raises -> Int: 1269 """Consume every byte after the last returned line, through EOF. 1270 1271 A real report is exactly one newline-terminated line, so any surplus 1272 byte is a duplicate or trailing report regardless of chunk alignment. 1273 A read/poll failure is a distinct cause and never a clean end of 1274 stream. 1275 """ 1276 var total = len(self.pending) 1277 self.pending = List[UInt8]() 1278 if total > max_bytes: 1279 raise Error("ready_output_overflow") 1280 var start = now_ms() 1281 while not self.eof: 1282 if now_ms() - start >= deadline_ms: 1283 raise Error("read_deadline_expired") 1284 var ev = poll_fd( 1285 self.report_fd, 1286 POLLIN, 1287 LIFECYCLE_POLL_SLICE_MS, 1288 deadline_ms - (now_ms() - start), 1289 self._poll_eintr_fault(), 1290 self._poll_error_fault(), 1291 ) 1292 if ev == IO_DEADLINE_EXPIRED: 1293 raise Error("read_deadline_expired") 1294 if ev < 0: 1295 raise Error("read_error") 1296 if ev == 0: 1297 continue 1298 var buf = InlineArray[Byte, 1024](fill=0) 1299 var n = read_fd( 1300 self.report_fd, 1301 buf.unsafe_ptr(), 1302 1024, 1303 deadline_ms - (now_ms() - start), 1304 ) 1305 if n == IO_DEADLINE_EXPIRED: 1306 raise Error("read_deadline_expired") 1307 if n < 0: 1308 raise Error("read_error") 1309 if n == 0: 1310 self.eof = True 1311 continue 1312 total += n 1313 if total > max_bytes: 1314 raise Error("ready_output_overflow") 1315 return total 1316 1317 1318 def piped_child_state( 1319 pid: Int, 1320 report_fd: Int, 1321 deadline_ms: Int, 1322 expected_requests: Int, 1323 guard: UnsafePointer[CleanupGuard, MutAnyOrigin], 1324 ) -> PipedChildState: 1325 """Build one owned-child lifecycle record with explicit ownership truth.""" 1326 return PipedChildState( 1327 pid=pid, 1328 report_fd=report_fd, 1329 pending=List[UInt8](), 1330 eof=False, 1331 closed=False, 1332 deadline_ms=deadline_ms, 1333 expected_requests=expected_requests, 1334 reaped=False, 1335 ok=False, 1336 phase="pending", 1337 case_label="-", 1338 reason="not_reaped", 1339 requests=0, 1340 connections=0, 1341 cleanup_error="", 1342 status=ProcessStatus("pending", False, -1, 0, 0, ""), 1343 observed=ProcessStatus("pending", False, -1, 0, 0, ""), 1344 observed_valid=False, 1345 last_terminated=False, 1346 spawn_ms=now_ms(), 1347 guard=guard, 1348 faults=LifecycleFaults(), 1349 ) 1350 1351 1352 def finalize_owned_failure( 1353 mut state: PipedChildState, pid: Int, label: String 1354 ) -> ProcessStatus: 1355 """Reap-or-retain an owned child after a startup/readiness failure. 1356 1357 Uses the same ownership truth as ``cleanup``: ownership is released only 1358 when no waitable child can remain. An unproved or uncertain termination 1359 keeps ``reaped=False``, records the failure in the required guard and 1360 *retains the usable ownership handle* (exact pid and report descriptor) 1361 through ``CleanupGuard.retain``, so a startup failure can neither claim nor 1362 hide an uncollected child and the caller can still recover it. The report 1363 descriptor is invalidated in both cases because the failed startup will 1364 never consume a report. 1365 """ 1366 var st = state.terminate_once(pid, TERMINATION_GRACE_MS) 1367 state.status = st.copy() 1368 if st.cleanup_proved(): 1369 state.guard[].resolve_pid(pid) 1370 state.reaped = True 1371 state.close_reader() 1372 else: 1373 # Retain the usable ownership handle (exact pid and report descriptor) 1374 # so a later retry can recover; the guard closes the descriptor exactly 1375 # once when the recovery succeeds. 1376 state.record_unproved( 1377 pid, state.report_fd, label, "unreaped:" + st.describe() 1378 ) 1379 return st^ 1380 1381 1382 def parse_ready_line(line: String, max_bytes: Int) raises -> Int: 1383 """Parse the exact ``ready <port>`` grammar with a valid TCP port range.""" 1384 if line.byte_length() == 0: 1385 raise Error("ready_empty") 1386 if line.byte_length() > max_bytes: 1387 raise Error("ready_too_large") 1388 if not line.startswith("ready "): 1389 raise Error("ready_grammar") 1390 var digits = String(line[byte=6:]) 1391 if digits.byte_length() == 0: 1392 raise Error("ready_missing_port") 1393 if digits.byte_length() > 5: 1394 raise Error("ready_port_range") 1395 for byte in digits.as_bytes(): 1396 var b = Int(byte) 1397 if b < 48 or b > 57: 1398 raise Error("ready_non_digit") 1399 var port = Int(digits) 1400 if port < 1 or port > 65535: 1401 raise Error("ready_port_range") 1402 return port 1403 1404 1405 # ── Non-opening descriptor census ─────────────────────────────────────────── 1406 1407 1408 def parse_ready_or_cleanup( 1409 pid: Int, line: String, max_bytes: Int 1410 ) raises -> Int: 1411 """Parse exact readiness or terminate and prove the owned child is gone. 1412 1413 Gives malformed readiness a real cleanup path with a bounded cause carrying 1414 the offending line and the owned child's reap result. 1415 """ 1416 try: 1417 return parse_ready_line(line, max_bytes) 1418 except e: 1419 var status = terminate_owned(pid, TERMINATION_GRACE_MS) 1420 raise Error( 1421 "ready_invalid:" 1422 + String(e) 1423 + " report=" 1424 + line 1425 + " child=" 1426 + status.describe() 1427 + " cleanup=" 1428 + ("proved" if status.cleanup_proved() else "unreaped") 1429 ) 1430 1431 1432 comptime CENSUS_EINTR_RETRY_BOUND: Int = 64 1433 comptime CENSUS_FAULT_NONE: Int = -1 1434 1435 1436 def classify_census_errno(errno_value: Int) -> String: 1437 """Classify one ``fcntl(F_GETFD)`` lookup failure (ADR-0018 RA04). 1438 1439 ``EBADF`` is the only proof that a slot is closed. ``EINTR`` is a bounded 1440 retry. Every other lookup error is ``unavailable``: a partial census must 1441 not silently lower the observed descriptor count. 1442 """ 1443 if errno_value == Int(ErrNo.EBADF.value): 1444 return "closed" 1445 if errno_value == Int(ErrNo.EINTR.value): 1446 return "retry" 1447 return "unavailable" 1448 1449 1450 def descriptor_census_with_faults( 1451 limit: Int, fault_fd: Int, fault_errno: Int 1452 ) -> Int: 1453 """Numeric descriptor census with a narrow cause-specific fault seam. 1454 1455 ``fault_fd``/``fault_errno`` are test-only controls that make exactly one 1456 lookup report a chosen errno *without changing any host limit*, so the 1457 EBADF/EINTR/other classification can be proved through the checked caller. 1458 Returns ``-1`` when the census cannot be established: an invalid bound, no 1459 standard descriptor, an exhausted EINTR retry bound, or any non-EBADF, 1460 non-EINTR lookup error that would otherwise hide live descriptors. 1461 """ 1462 if limit <= 0: 1463 return -1 1464 var count = 0 1465 for fd in range(0, limit): 1466 var retries = 0 1467 var open = False 1468 while True: 1469 var rc = 0 1470 var injected = fault_fd >= 0 and fd == fault_fd 1471 if injected: 1472 rc = -1 1473 else: 1474 rc = Int( 1475 external_call["fcntl", c_int](c_int(fd), c_int(F_GETFD)) 1476 ) 1477 if rc >= 0: 1478 open = True 1479 break 1480 var errno_value = Int(get_errno().value) 1481 if injected: 1482 errno_value = fault_errno 1483 var klass = classify_census_errno(errno_value) 1484 if klass == "closed": 1485 break 1486 if klass == "retry": 1487 retries += 1 1488 if retries > CENSUS_EINTR_RETRY_BOUND: 1489 return -1 1490 continue 1491 return -1 1492 if open: 1493 count += 1 1494 if count == 0: 1495 return -1 1496 return count 1497 1498 1499 def descriptor_census(limit: Int) -> Int: 1500 """Count open descriptors numerically via ``fcntl(F_GETFD)``. 1501 1502 Never opens a target path, so device nodes cannot block it and sockets and 1503 high descriptors are counted the same as regular files. Returns -1 when the 1504 census cannot be established (invalid bound, no standard descriptors, or an 1505 unclassifiable lookup error), which callers must treat as a failed 1506 unavailable census, never a pass. 1507 """ 1508 return descriptor_census_with_faults(limit, CENSUS_FAULT_NONE, 0) 1509 1510 1511 def fd_scan_limit() -> Int: 1512 """Return the OS descriptor-table size, or -1 when unavailable. 1513 1514 The raw size is returned; callers decide whether complete coverage inside 1515 the admitted range is possible rather than silently truncating the scan. 1516 """ 1517 var n = Int(external_call["getdtablesize", c_int]()) 1518 if n <= 0: 1519 return -1 1520 return n 1521 1522 1523 def open_fd_count() -> Int: 1524 """Numeric open-descriptor census for this process (-1 if unavailable).""" 1525 var limit = fd_scan_limit() 1526 if limit <= 0 or limit > CENSUS_MAX_FDS: 1527 return -1 1528 return descriptor_census(limit) 1529 1530 1531 def open_fd_count_checked_with_faults( 1532 limit: Int, fault_fd: Int, fault_errno: Int 1533 ) raises -> Int: 1534 """Checked census with the narrow cause-specific fault seam (test only). 1535 1536 ``fault_fd``/``fault_errno`` inject exactly one classified lookup failure 1537 without changing any host limit, so the checked caller's unavailable-census 1538 propagation can be executed for a non-EBADF/non-EINTR error (ADR-0018 RA04). 1539 """ 1540 var effective = limit 1541 if effective < 0: 1542 effective = fd_scan_limit() 1543 if effective <= 0 or effective > CENSUS_MAX_FDS: 1544 raise Error("descriptor_census_unavailable") 1545 var count = descriptor_census_with_faults(effective, fault_fd, fault_errno) 1546 if count < 0: 1547 raise Error("descriptor_census_unavailable") 1548 return count 1549 1550 1551 def open_fd_count_checked(limit: Int = -1) raises -> Int: 1552 """Checked census that fails explicitly when coverage cannot be complete. 1553 1554 ``limit`` defaults to the OS descriptor-table size. A nonpositive, 1555 oversized or otherwise unavailable census raises 1556 ``descriptor_census_unavailable`` rather than returning a partial count 1557 that a caller could read as a pass. 1558 """ 1559 return open_fd_count_checked_with_faults(limit, CENSUS_FAULT_NONE, 0)