hyf

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

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)