hyf

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

commit de333d775a6eeb233cd3c8b3e080dc558d4030bf
parent 2c421b37ae844fddc711727d898e27a0131d4c43
Author: triesap <tyson@radroots.org>
Date:   Wed, 23 Sep 2026 15:48:05 +0000

test: requalify H005A measurement accounting, budget and ownership

ADR-0019 D39 MR01-MR05 repair of the test-only measurement tooling and its
callers. No product source, schema, dependency, lock, feature or environment
change.

- MR01: account the whole stdout stream through EOF. Exactly the expected
  warmup/measured frames may succeed; a retained surplus byte, a coalesced,
  split or unterminated trailing frame, early EOF and a nonzero exit fail with
  a bounded cause. stdout/stderr are drained concurrently with request writes
  and exit observation.
- MR02: one spawn-relative work budget covers sampling, request writes,
  response reads, the EOF drain and the child exit; an expired budget is never
  clamped into a fresh successful interval and a late valid response/exit is
  never accepted. Byte caps distinguish timeout, read error, EOF and overflow;
  stderr read errors and overflow fail explicitly; poll EINTR is retried within
  the deadline.
- MR03: stdin/stdout/stderr each close exactly once on success, invalid output,
  sampling failure, early EOF, nonzero exit and timeout; EOF is never treated
  as descriptor closure; terminal state is cached before any further wait or
  signal and an unproved cleanup retains exact child/descriptor ownership.
- MR04: the product binary is bound to a clean verified source revision/tree and
  a deterministic manifest of the tracked src build inputs, with capture/build/
  run drift rejected; startup is measured from spawn to the first validated
  response with per-request timing and sampling instrumentation separated; the
  standalone runner records numeric direct-provider elapsed/requests/connections
  and client/schema construction.

25/25 test-measurement; measure-h005a 1100/1100 frames; check-build,
check-format, test-architecture, test-provider-helpers, test-repo-local-process
and test-jev pass; test-stdio keeps the exact 24 D16 signatures.

Diffstat:
Mtests/measurement_process_helper.mojo | 1158++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------------
Mtests/measurement_runner.mojo | 123++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---
Mtests/test_measurement_contract.mojo | 599++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------------
3 files changed, 1561 insertions(+), 319 deletions(-)

diff --git a/tests/measurement_process_helper.mojo b/tests/measurement_process_helper.mojo @@ -10,25 +10,59 @@ It exists only for tests. It changes no HYF product policy, no schema and no dependency: the process is launched from a build of the existing product entry point and the samples are read back with the bounded lifecycle primitives. -Explicit failure policy (R56/R57): +ADR-0019 D39 MR01-MR05 requalification: + +* MR01 — the *whole* stdout stream is accounted for through EOF. Exactly the + expected warmup/measured frames may succeed; a retained surplus byte, a + trailing frame in the same or a later chunk, an unterminated tail, early EOF + and a nonzero exit all fail with a bounded cause. stdout/stderr are drained + concurrently with request writes and exit observation. +* MR02 — one spawn-relative measurement work deadline covers sampling, request + writes, response reads, the EOF drain and the child exit. An expired budget + is never clamped into a fresh successful interval and a late valid response + is never accepted. Build and pre-spawn identity capture are separate, + explicitly bounded phases. Byte/output caps distinguish timeout, read error, + EOF and overflow; stderr read errors and overflow fail explicitly. +* MR03 — every owned descriptor (stdin, stdout, stderr) is closed exactly once + on success, invalid output, sampling failure, early EOF, nonzero exit and + timeout. EOF is not descriptor closure. The terminal child state is cached + before any further wait/signal, and an unproved cleanup retains the exact + child/descriptor ownership in the caller-held guard. +* MR04 — the product binary is bound to a clean, verified source/tree identity + plus a deterministic content manifest of the tracked build inputs; capture, + build and run drift is rejected. Startup is measured from spawn to the first + validated response, sampling is recorded separately as instrumentation + overhead, and the secret-free HYF_PATHS profile is read back from the live + environment. + +Explicit failure policy (R56/R57, R69/R70/R71): * a frame that is not JSON, has the wrong correlation or outcome, or a child that exits nonzero, fails the measurement instead of reporting success; * unavailable ``ps``/``lsof`` sampling raises a sampling error rather than recording a placeholder value; -* every read/write/wait is parent-bounded by a finite deadline and the owned - child is always terminated and reaped through the shared ownership guard. +* every read/write/wait is parent-bounded by the one finite work deadline and + the owned child is always terminated and reaped through the shared ownership + guard. """ import std.os from std.collections import List -from std.ffi import CStringSlice, c_int, external_call +from std.ffi import ( + ErrNo, + CStringSlice, + c_int, + c_ssize_t, + c_size_t, + c_uint, + external_call, + get_errno, +) from json import Value, loads from parent_lifecycle import ( IO_DEADLINE_EXPIRED, - LIFECYCLE_POLL_SLICE_MS, POLLERR, POLLHUP, POLLIN, @@ -36,7 +70,6 @@ from parent_lifecycle import ( POLLOUT, TERMINATION_GRACE_MS, CleanupGuard, - PipedChildState, ProcessStatus, child_exit, close_fd, @@ -44,10 +77,10 @@ from parent_lifecycle import ( fork_owned_or_close3, make_three_pipes, now_ms, - piped_child_state, poll_three, read_fd, set_alarm, + sleep_ms, terminate_owned, wait_bounded, write_fd_chunk, @@ -55,11 +88,16 @@ from parent_lifecycle import ( comptime MEASUREMENT_FRAME_BYTES = 1048576 -comptime MEASUREMENT_FRAME_DEADLINE_MS = 5000 -comptime MEASUREMENT_SAMPLE_DEADLINE_MS = 10000 comptime MEASUREMENT_CHILD_ALARM_SECONDS = 900 -comptime MEASUREMENT_MAX_CAPTURE_BYTES = 65536 +comptime MEASUREMENT_MAX_STDERR_BYTES = 65536 +comptime MEASUREMENT_SAMPLE_DEADLINE_MS = 10000 +comptime MEASUREMENT_BUILD_DEADLINE_MS = 600000 comptime MEASUREMENT_SAMPLE_INTERVAL = 50 +# Bounded polling granularity. Every wait is sliced by this value so an +# idle-but-live child is observed incrementally; the recorded request latency +# therefore carries at most one polling slice of tolerance rather than a fixed +# idle delay. Instrumentation (sampling subprocesses) is timed separately. +comptime MEASUREMENT_POLL_SLICE_MS = 10 # ── Bounded captured commands (identity and sampling) ─────────────────────── @@ -157,7 +195,7 @@ def run_capture( if elapsed >= deadline_ms: read_reason = "sample_deadline_expired" break - var slice_ms = min(LIFECYCLE_POLL_SLICE_MS, deadline_ms - elapsed) + var slice_ms = min(MEASUREMENT_POLL_SLICE_MS * 2, deadline_ms - elapsed) if slice_ms < 1: slice_ms = 1 var pr = poll_three( @@ -204,8 +242,20 @@ def run_capture( ) var remaining = deadline_ms - (now_ms() - start) - if remaining < 1: - remaining = 1 + if remaining <= 0: + # A bounded sampling phase never converts an expired budget into a new + # successful interval: the exact owned child is reaped or retained. + var expired = terminate_owned(pid, TERMINATION_GRACE_MS) + guard.retain( + pid, -1, "measurement sample command " + command + " deadline" + ) + if expired.cleanup_proved(): + guard.resolve_pid(pid) + raise Error( + "measurement: sample command " + + command + + " sample_deadline_expired" + ) var st = wait_bounded(pid, remaining) if not st.cleanup_proved(): var term = terminate_owned(pid, TERMINATION_GRACE_MS) @@ -238,14 +288,17 @@ def drain_capture( if (revents & (POLLIN | POLLHUP | POLLERR | POLLNVAL)) == 0: return DrainResult(False, "") var buf = InlineArray[Byte, 4096](fill=0) - var n = read_fd(fd, buf.unsafe_ptr(), 4096, deadline_ms) + var budget = deadline_ms + if budget < 1: + budget = 1 + var n = read_fd(fd, buf.unsafe_ptr(), 4096, budget) if n == IO_DEADLINE_EXPIRED: return DrainResult(True, "read_deadline_expired") if n < 0: return DrainResult(True, "read_error") if n == 0: return DrainResult(True, "") - if len(out) + n > MEASUREMENT_MAX_CAPTURE_BYTES: + if len(out) + n > MEASUREMENT_MAX_STDERR_BYTES: return DrainResult(True, "capture_overflow") for index in range(n): out.append(UInt8(Int(buf[index]))) @@ -285,15 +338,25 @@ def file_sha256(path: String, mut guard: CleanupGuard) raises -> String: out = alt^ var text = String(out.stdout.strip()) var digest = String(text[byte=0:64]) + _require_hex(digest, 64, "sha256 for " + path) + return digest + + +def _require_hex(digest: String, expected_len: Int, label: String) raises: + if digest.byte_length() != expected_len: + raise Error( + "measurement: " + + label + + " is not a " + + String(expected_len) + + "-character digest" + ) for byte in digest.as_bytes(): var b = Int(byte) var is_digit = b >= 48 and b <= 57 var is_hex = (b >= 97 and b <= 102) or (b >= 65 and b <= 70) if not is_digit and not is_hex: - raise Error( - "measurement: sha256 output not hexadecimal for " + path - ) - return digest + raise Error("measurement: " + label + " is not hexadecimal") def host_platform(mut guard: CleanupGuard) raises -> String: @@ -324,89 +387,167 @@ def toolchain_version(mut guard: CleanupGuard) raises -> String: @fieldwise_init -struct MeasurementIdentity(Movable): - """Exact source/binary/toolchain/host identity and launch profile.""" +struct SourceTreeIdentity(Movable): + """Verified source/tree identity for the measured build inputs. + + ``revision``/``tree`` are the exact commit and tree, ``manifest_sha256`` is + a deterministic content digest of the tracked ``src`` build inputs and + ``dirty_status`` is the exact ``git status --porcelain`` output for the + measured paths (``""`` proves a clean tree). A mismatch between capture, + build and run is a bounded drift failure, never a silent success. + """ var source_root: String - var source_revision: String - var cwd: String - var binary_path: String - var binary_sha256: String - var pixi_toml_sha256: String - var pixi_lock_sha256: String - var toolchain_version: String - var host_platform: String - var argv_profile: String - var env_profile: String + var revision: String + var tree: String + var manifest_sha256: String + var dirty_status: String def describe(self) -> String: return ( - "source_revision=" - + self.source_revision - + " cwd=" - + self.cwd - + " binary_sha256=" - + self.binary_sha256 - + " pixi_toml_sha256=" - + self.pixi_toml_sha256 - + " pixi_lock_sha256=" - + self.pixi_lock_sha256 - + " toolchain=" - + self.toolchain_version - + " host=" - + self.host_platform - + " argv=" - + self.argv_profile - + " env=" - + self.env_profile + "revision=" + + self.revision + + " tree=" + + self.tree + + " manifest_sha256=" + + self.manifest_sha256 + + " tree_state=" + + (self.dirty_status if self.dirty_status == "" else "dirty") ) + def same_as(self, other: SourceTreeIdentity) -> Bool: + return ( + self.revision == other.revision + and self.tree == other.tree + and self.manifest_sha256 == other.manifest_sha256 + and self.dirty_status == other.dirty_status + ) -comptime MEASUREMENT_BUILD_DEADLINE_MS = 600000 +def git_capture( + source_root: String, var args: List[String], mut guard: CleanupGuard +) raises -> CommandOutput: + return run_capture( + "git", args^, MEASUREMENT_SAMPLE_DEADLINE_MS, guard, source_root + ) -def build_product_binary( - temp_dir: String, mut guard: CleanupGuard -) raises -> String: - """Build the existing product entry point once, outside measured frames. - The build is a separate bounded phase; the measurement itself never - compiles or starts a process per frame. - """ - var output = temp_dir + "/hyfd" +def source_revision( + source_root: String, mut guard: CleanupGuard +) raises -> String: + """Exact capsule source revision; a measurement must not guess it.""" var args = List[String]() - args.append("build") - args.append("-I") - args.append("src") - args.append("src/main.mojo") - args.append("-o") - args.append(output) - var out = run_capture("mojo", args^, MEASUREMENT_BUILD_DEADLINE_MS, guard) - if out.exit_code != 0: + args.append("rev-parse") + args.append("HEAD") + var out = git_capture(source_root, args^, guard) + var text = String(out.stdout.strip()) + if out.exit_code != 0 or text.byte_length() != 40: raise Error( - "measurement: product build failed (" + out.describe() + ")" + "measurement: source revision unavailable in " + source_root ) - return output^ + _require_hex(text, 40, "source revision") + return text^ -def source_revision( +def source_tree_id( source_root: String, mut guard: CleanupGuard ) raises -> String: - """Exact capsule source revision; a measurement must not guess it.""" + """Exact commit tree object id; distinguishes a content-identical tree.""" var args = List[String]() args.append("rev-parse") - args.append("HEAD") + args.append("HEAD^{tree}") + var out = git_capture(source_root, args^, guard) + var text = String(out.stdout.strip()) + if out.exit_code != 0 or text.byte_length() != 40: + raise Error("measurement: source tree id unavailable in " + source_root) + _require_hex(text, 40, "source tree id") + return text^ + + +def source_manifest_sha256( + source_root: String, mut guard: CleanupGuard +) raises -> String: + """Deterministic content digest of the tracked ``src`` build inputs. + + The digest covers the index entries (mode/blob/path) of every tracked file + under ``src`` — the exact inputs of ``mojo build -I src src/main.mojo``. + It never reads or hashes secrets or arbitrary workspace files. + """ + var args = List[String]() + args.append("-c") + args.append( + "git ls-files -s -- src | LC_ALL=C sort | shasum -a 256 | cut -d' ' -f1" + ) var out = run_capture( - "git", args^, MEASUREMENT_SAMPLE_DEADLINE_MS, guard, source_root + "sh", args^, MEASUREMENT_SAMPLE_DEADLINE_MS, guard, source_root ) var text = String(out.stdout.strip()) - if out.exit_code != 0 or text.byte_length() != 40: + if out.exit_code != 0: raise Error( - "measurement: source revision unavailable in " + source_root + "measurement: source content manifest unavailable in " + source_root ) + _require_hex(text, 64, "source content manifest") return text^ +def source_dirty_status( + source_root: String, mut guard: CleanupGuard +) raises -> String: + """Exact porcelain status for the measured build inputs and toolchain.""" + var args = List[String]() + args.append("status") + args.append("--porcelain") + args.append("--") + args.append("src") + args.append("pixi.toml") + args.append("pixi.lock") + var out = git_capture(source_root, args^, guard) + if out.exit_code != 0: + raise Error( + "measurement: source tree status unavailable in " + source_root + ) + return String(out.stdout.strip()) + + +def source_identity( + source_root: String, mut guard: CleanupGuard +) raises -> SourceTreeIdentity: + return SourceTreeIdentity( + source_root=source_root, + revision=source_revision(source_root, guard), + tree=source_tree_id(source_root, guard), + manifest_sha256=source_manifest_sha256(source_root, guard), + dirty_status=source_dirty_status(source_root, guard), + ) + + +def require_clean_source(identity: SourceTreeIdentity, phase: String) raises: + """Reject a dirty measured build input at capture, build or run time.""" + if identity.dirty_status != "": + raise Error( + "measurement: measured source tree is dirty at " + + phase + + " (" + + identity.dirty_status + + ")" + ) + + +def require_identity_drift_free( + before: SourceTreeIdentity, after: SourceTreeIdentity, phase: String +) raises: + if not before.same_as(after): + raise Error( + "measurement: source drifted during " + + phase + + " (before " + + before.describe() + + " after " + + after.describe() + + ")" + ) + + def working_directory(mut guard: CleanupGuard) raises -> String: var args = List[String]() var out = run_capture_simple("pwd", args^, guard) @@ -420,7 +561,8 @@ def verified_environment_profile() -> String: Recorded from the live process environment rather than asserted, so the reproduced profile is truthful: the measured child inherits exactly these - values through ``execvp``. + values through ``execvp``. The profile is deliberately bounded to the two + governed, secret-free HYF_PATHS variables. """ var profile = std.os.getenv("HYF_PATHS_PROFILE") var root = std.os.getenv("HYF_PATHS_REPO_LOCAL_ROOT") @@ -432,18 +574,76 @@ def verified_environment_profile() -> String: ) +@fieldwise_init +struct MeasurementIdentity(Movable): + """Exact source/binary/toolchain/host identity and launch profile.""" + + var source_root: String + var source_revision: String + var source_tree: String + var source_manifest_sha256: String + var source_tree_state: String + var binding: String + var cwd: String + var binary_path: String + var binary_sha256: String + var pixi_toml_sha256: String + var pixi_lock_sha256: String + var toolchain_version: String + var host_platform: String + var argv_profile: String + var env_profile: String + + def describe(self) -> String: + return ( + "binding=" + + self.binding + + " source_revision=" + + self.source_revision + + " source_tree=" + + self.source_tree + + " source_manifest_sha256=" + + self.source_manifest_sha256 + + " source_tree_state=" + + self.source_tree_state + + " cwd=" + + self.cwd + + " binary_sha256=" + + self.binary_sha256 + + " pixi_toml_sha256=" + + self.pixi_toml_sha256 + + " pixi_lock_sha256=" + + self.pixi_lock_sha256 + + " toolchain=" + + self.toolchain_version + + " host=" + + self.host_platform + + " argv=" + + self.argv_profile + + " env=" + + self.env_profile + ) + + def measurement_identity( source_root: String, binary_path: String, + binary_sha256: String, + binding: String, argv_profile: String, mut guard: CleanupGuard, ) raises -> MeasurementIdentity: + var observed = source_identity(source_root, guard) return MeasurementIdentity( source_root=source_root, - source_revision=source_revision(source_root, guard), + source_revision=observed.revision, + source_tree=observed.tree, + source_manifest_sha256=observed.manifest_sha256, + source_tree_state=("dirty" if observed.dirty_status != "" else "clean"), + binding=binding, cwd=working_directory(guard), binary_path=binary_path, - binary_sha256=file_sha256(binary_path, guard), + binary_sha256=binary_sha256, pixi_toml_sha256=file_sha256(source_root + "/pixi.toml", guard), pixi_lock_sha256=file_sha256(source_root + "/pixi.lock", guard), toolchain_version=toolchain_version(guard), @@ -453,106 +653,461 @@ def measurement_identity( ) +@fieldwise_init +struct BuiltProduct(Movable): + """A product binary bound to the clean source/tree it was built from. + + The build is a separate, explicitly bounded phase: the pre-build source + identity is captured, the binary is built, the binary digest is recorded + and the post-build source identity must be byte-for-byte the same. A + later run re-verifies the same identity and binary before spawning. + """ + + var binary_path: String + var binary_sha256: String + var source: SourceTreeIdentity + var build_ms: Int + + +def build_product_binary( + source_root: String, temp_dir: String, mut guard: CleanupGuard +) raises -> BuiltProduct: + """Build the existing product entry point once, outside measured frames. + + The measurement itself never compiles or starts a process per frame. The + build input tree must be clean and must not drift across the build. + """ + var before = source_identity(source_root, guard) + require_clean_source(before, "build capture") + var output = temp_dir + "/hyfd" + var args = List[String]() + args.append("build") + args.append("-I") + args.append("src") + args.append("src/main.mojo") + args.append("-o") + args.append(output) + var start = now_ms() + var out = run_capture( + "mojo", args^, MEASUREMENT_BUILD_DEADLINE_MS, guard, source_root + ) + var build_ms = now_ms() - start + if out.exit_code != 0: + raise Error( + "measurement: product build failed (" + out.describe() + ")" + ) + var binary_sha256 = file_sha256(output, guard) + var after = source_identity(source_root, guard) + require_identity_drift_free(before, after, "build") + return BuiltProduct( + binary_path=output, + binary_sha256=binary_sha256, + source=after^, + build_ms=build_ms, + ) + + # ── Persistent measured process ───────────────────────────────────────────── @fieldwise_init +struct MeasurementPoll(Movable): + """One bounded ``poll(2)`` result with an explicit EINTR classification.""" + + var count: Int + var r0: Int + var r1: Int + var r2: Int + var interrupted: Bool + + +def measurement_poll( + fd0: Int, + events0: Int, + fd1: Int, + events1: Int, + fd2: Int, + events2: Int, + timeout_ms: Int, + fault_eintr_count: Int = 0, +) -> MeasurementPoll: + """Poll three descriptors, classifying a real ``EINTR`` as retryable. + + ``fault_eintr_count`` is a bounded test-only seam: any nonzero value makes + every attempt report ``EINTR`` so the deadline-bounded retry branch is + deterministically executable without changing host signal state. + """ + if fault_eintr_count != 0: + return MeasurementPoll(-1, 0, 0, 0, True) + var cell = InlineArray[Int32, 12](fill=0) + cell[0] = Int32(fd0) + cell[1] = Int32(events0) + cell[2] = Int32(fd1) + cell[3] = Int32(events1) + cell[4] = Int32(fd2) + cell[5] = Int32(events2) + var n = Int( + external_call["poll", c_int]( + cell.unsafe_ptr(), c_uint(3), c_int(timeout_ms) + ) + ) + if n < 0: + if get_errno() == ErrNo.EINTR: + return MeasurementPoll(-1, 0, 0, 0, True) + return MeasurementPoll(-1, 0, 0, 0, False) + if n == 0: + return MeasurementPoll(0, 0, 0, 0, False) + return MeasurementPoll( + n, + (Int(cell[1]) >> 16) & 0xFFFF, + (Int(cell[3]) >> 16) & 0xFFFF, + (Int(cell[5]) >> 16) & 0xFFFF, + False, + ) + + +def measurement_poll_retry( + fd0: Int, + events0: Int, + fd1: Int, + events1: Int, + fd2: Int, + events2: Int, + max_wait_ms: Int, + fault_eintr_count: Int = 0, +) -> MeasurementPoll: + """Poll with a deadline-bounded retry of a real ``EINTR``. + + The retry can never outlive ``max_wait_ms``: once the budget is exhausted + the last interrupted result is returned to the deadline owner instead of + spinning. ``fault_eintr_count`` is the bounded test-only seam. + """ + var start = now_ms() + var slice = min(MEASUREMENT_POLL_SLICE_MS, max_wait_ms) + if slice < 1: + slice = 1 + while True: + var pr = measurement_poll( + fd0, events0, fd1, events1, fd2, events2, slice, fault_eintr_count + ) + if not pr.interrupted: + return pr^ + if now_ms() - start >= max_wait_ms: + return pr^ + sleep_ms(1) + + +@fieldwise_init struct MeasurementProcess(Movable): - """Exactly one persistent HYF stdio process owned by the measurement.""" + """Exactly one persistent HYF stdio process owned by the measurement. + + stdout/stderr are consumed concurrently with request writes and exit + observation and are drained to EOF. Every owned descriptor carries an + explicit close-once flag: EOF is a stream property and never implies that + the parent descriptor was closed, so no later cleanup can target a reused + descriptor number. + """ var pid: Int var stdin_fd: Int - var reader: PipedChildState + var stdout_fd: Int var stderr_fd: Int + var stdout_pending: List[UInt8] var stderr_bytes: List[UInt8] + var stdout_eof: Bool var stderr_eof: Bool - var deadline_ms: Int - var start_ms: Int var stdin_closed: Bool + var stdout_closed: Bool + var stderr_closed: Bool var reaped: Bool var status: ProcessStatus + var spawn_ms: Int + var deadline_at_ms: Int + var stdout_total_bytes: Int + var guard: UnsafePointer[CleanupGuard, MutAnyOrigin] + + def work_remaining_ms(self) -> Int: + """Remaining part of the one spawn-relative measurement work budget. + + A nonpositive result must never be turned into a new successful + interval; callers fail and use only the separate bounded cleanup + allowance. + """ + return self.deadline_at_ms - now_ms() + + def close_stdin(mut self): + if self.stdin_closed: + return + close_fd(self.stdin_fd) + self.stdin_fd = -1 + self.stdin_closed = True - def remaining_ms(mut self) -> Int: - return self.deadline_ms - (now_ms() - self.start_ms) + def close_stdout(mut self): + if self.stdout_closed: + return + # A descriptor retained by the guard for recovery is closed by the + # guard exactly once, so the handle can never close a reused number. + if not self.guard[].close_retained_fd(self.stdout_fd): + close_fd(self.stdout_fd) + self.stdout_closed = True + + def close_stderr(mut self): + if self.stderr_closed: + return + close_fd(self.stderr_fd) + self.stderr_closed = True - def send_frame(mut self, frame: String, deadline_ms: Int) raises: - """Write one newline-delimited request frame with a finite deadline.""" + def stdout_pending_len(self) -> Int: + return len(self.stdout_pending) + + def stderr_len(self) -> Int: + return len(self.stderr_bytes) + + def drain_stdout(mut self, revents: Int, budget_ms: Int) raises: + """Drain ready stdout bytes with a bounded size cap and cause.""" + if (revents & (POLLIN | POLLHUP | POLLERR | POLLNVAL)) == 0: + return + if budget_ms < 1: + raise Error( + "measurement: work deadline expired during stdout drain" + " (work_deadline_expired)" + ) + var buf = InlineArray[Byte, 4096](fill=0) + var n = read_fd(self.stdout_fd, buf.unsafe_ptr(), 4096, budget_ms) + if n == IO_DEADLINE_EXPIRED: + raise Error( + "measurement: stdout read deadline expired" + " (work_deadline_expired)" + ) + if n < 0: + raise Error("measurement: stdout read_error") + if n == 0: + self.stdout_eof = True + return + self.stdout_total_bytes += n + if len(self.stdout_pending) + n > MEASUREMENT_FRAME_BYTES: + raise Error("measurement: stdout_overflow") + for index in range(n): + self.stdout_pending.append(UInt8(Int(buf[index]))) + + def drain_stderr(mut self, revents: Int, budget_ms: Int) raises: + """Drain ready stderr bytes; errors and overflow fail explicitly.""" + if (revents & (POLLIN | POLLHUP | POLLERR | POLLNVAL)) == 0: + return + if budget_ms < 1: + raise Error( + "measurement: work deadline expired during stderr drain" + " (work_deadline_expired)" + ) + var buf = InlineArray[Byte, 4096](fill=0) + var n = read_fd(self.stderr_fd, buf.unsafe_ptr(), 4096, budget_ms) + if n == IO_DEADLINE_EXPIRED: + raise Error( + "measurement: stderr read deadline expired" + " (work_deadline_expired)" + ) + if n < 0: + raise Error("measurement: stderr read_error") + if n == 0: + self.stderr_eof = True + return + if len(self.stderr_bytes) + n > MEASUREMENT_MAX_STDERR_BYTES: + raise Error("measurement: stderr_overflow") + for index in range(n): + self.stderr_bytes.append(UInt8(Int(buf[index]))) + + def poll_and_drain( + mut self, + want_stdin_write: Bool, + max_wait_ms: Int, + fault_eintr_count: Int = 0, + ) raises -> Bool: + """Poll stdin/stdout/stderr and drain whatever is ready. + + Returns True when the stdin write end is writable. All waits are sliced + by ``MEASUREMENT_POLL_SLICE_MS`` inside the caller's remaining budget; + a real ``EINTR`` is retried until the budget expires. + """ + if max_wait_ms <= 0: + raise Error( + "measurement: work deadline expired (work_deadline_expired)" + ) + var slice = min(MEASUREMENT_POLL_SLICE_MS, max_wait_ms) + var fd0 = self.stdin_fd if ( + want_stdin_write and not self.stdin_closed + ) else -1 + var fd1 = self.stdout_fd if ( + not self.stdout_eof and not self.stdout_closed + ) else -1 + var fd2 = self.stderr_fd if ( + not self.stderr_eof and not self.stderr_closed + ) else -1 + if fd0 < 0 and fd1 < 0 and fd2 < 0: + sleep_ms(slice) + return False + var start = now_ms() + var pr = measurement_poll_retry( + fd0, + POLLOUT if fd0 >= 0 else 0, + fd1, + POLLIN if fd1 >= 0 else 0, + fd2, + POLLIN if fd2 >= 0 else 0, + max_wait_ms, + fault_eintr_count, + ) + if pr.count > 0: + var remaining = max_wait_ms - (now_ms() - start) + if fd0 >= 0 and (pr.r0 & (POLLERR | POLLHUP | POLLNVAL)) != 0: + raise Error("measurement: request write pipe closed") + if ( + fd1 >= 0 + and (pr.r1 & (POLLIN | POLLHUP | POLLERR | POLLNVAL)) != 0 + ): + self.drain_stdout(pr.r1, remaining) + if ( + fd2 >= 0 + and (pr.r2 & (POLLIN | POLLHUP | POLLERR | POLLNVAL)) != 0 + ): + self.drain_stderr(pr.r2, remaining) + return fd0 >= 0 and (pr.r0 & POLLOUT) != 0 + return False + + def send_frame(mut self, frame: String, index: Int) raises: + """Write one newline-delimited request frame within the work budget. + + stdout/stderr are drained concurrently while the request is written, so + a child that emits output during the write can never deadlock the + parent. There is no fixed idle wait in the request latency. + """ + if self.stdin_closed: + raise Error("measurement: request write after stdin closed") var payload = frame + "\n" var sent = 0 - var start = now_ms() - while sent < payload.byte_length(): - var remaining = deadline_ms - (now_ms() - start) + var total = payload.byte_length() + while sent < total: + var remaining = self.work_remaining_ms() if remaining <= 0: - raise Error("measurement: request write deadline expired") - var ev = poll_three(self.stdin_fd, POLLOUT, -1, 0, -1, 0, 25) - if ev.count < 0: - raise Error("measurement: request write poll failed") - if ev.count == 0: - continue - if (ev.r0 & (POLLERR | POLLHUP | POLLNVAL)) != 0: + raise Error( + "measurement: work deadline expired during request write" + " (work_deadline_expired)" + ) + var writable = self.poll_and_drain(True, remaining) + if self.stdin_closed: raise Error("measurement: request write pipe closed") - var cw = write_fd_chunk(self.stdin_fd, payload, sent, remaining) + if not writable: + continue + var remaining_write = self.work_remaining_ms() + if remaining_write <= 0: + raise Error( + "measurement: work deadline expired during request write" + " (work_deadline_expired)" + ) + var cw = write_fd_chunk( + self.stdin_fd, payload, sent, remaining_write + ) + if cw.reason == "write_deadline_expired": + raise Error("measurement: request write deadline expired") if cw.reason != "": raise Error( "measurement: request write failed (" + cw.reason + ")" ) + if cw.written <= 0: + continue sent += cw.written - self.drain_stderr(25) - - def read_response(mut self, deadline_ms: Int) raises -> String: - """Read one complete newline-terminated response with bounded I/O.""" - self.drain_stderr(1) - var line = self.reader.read_line(MEASUREMENT_FRAME_BYTES, deadline_ms) - if not self.reader.last_terminated: - if line == "": + + def read_response(mut self, index: Int) raises -> String: + """Read one complete newline-terminated response within the budget.""" + while True: + var found = _newline_index(self.stdout_pending) + if found >= 0: + if found > MEASUREMENT_FRAME_BYTES: + raise Error("measurement: response frame overflow") + var line_bytes = _list_prefix(self.stdout_pending, found) + self.stdout_pending = _list_suffix( + self.stdout_pending, found + 1 + ) + return _decode_frame(line_bytes^) + if self.stdout_eof: + if len(self.stdout_pending) == 0: + raise Error( + "measurement: response stream reached EOF with no" + " frame (early_eof)" + ) raise Error( - "measurement: response stream reached EOF with no frame" - " (early_eof)" + "measurement: response frame was not newline-terminated" ) + var remaining = self.work_remaining_ms() + if remaining <= 0: + raise Error( + "measurement: work deadline expired while awaiting" + " response (work_deadline_expired)" + ) + _ = self.poll_and_drain(False, remaining) + + def assert_no_surplus(mut self, index: Int) raises: + """Exactly one response per request: surplus bytes are a failure.""" + if len(self.stdout_pending) > 0: raise Error( - "measurement: response frame was not newline-terminated" + "measurement: unexpected trailing stdout (" + + String(len(self.stdout_pending)) + + " bytes) after response " + + String(index) ) - return line^ - def drain_stderr(mut self, slice_ms: Int): - if self.stderr_eof: - return - var pr = poll_three(-1, 0, -1, 0, self.stderr_fd, POLLIN, slice_ms) - if pr.count <= 0: - return - var buf = InlineArray[Byte, 4096](fill=0) - var n = read_fd(self.stderr_fd, buf.unsafe_ptr(), 4096, 0) - if n <= 0: - if n == 0: - self.stderr_eof = True - return - if len(self.stderr_bytes) + n <= MEASUREMENT_MAX_CAPTURE_BYTES: - for index in range(n): - self.stderr_bytes.append(UInt8(Int(buf[index]))) - - def finish(mut self) raises -> ProcessStatus: - """Close stdin, require a clean bounded child exit and prove cleanup.""" - if not self.stdin_closed: - close_fd(self.stdin_fd) - self.stdin_fd = -1 - self.stdin_closed = True - var remaining = self.remaining_ms() - if remaining < 1: - remaining = 1 - var st = wait_bounded(self.pid, remaining) + def finish_expected(mut self, total: Int) raises -> ProcessStatus: + """Close stdin, account the whole stream through EOF and reap. + + The remaining stdout stream must be exactly empty after the expected + responses, stderr must reach a clean EOF and the child must exit within + the remaining work budget. The terminal state is cached before any + further wait or signal. + """ + if len(self.stdout_pending) > 0: + raise Error( + "measurement: unexpected trailing stdout (" + + String(len(self.stdout_pending)) + + " bytes) after " + + String(total) + + " expected responses" + ) + self.close_stdin() + while not (self.stdout_eof and self.stderr_eof): + var remaining = self.work_remaining_ms() + if remaining <= 0: + raise Error( + "measurement: work deadline expired during output drain" + " (work_deadline_expired)" + ) + _ = self.poll_and_drain(False, remaining) + if len(self.stdout_pending) > 0: + raise Error( + "measurement: unexpected trailing stdout (" + + String(len(self.stdout_pending)) + + " bytes) after " + + String(total) + + " expected responses" + ) + var remaining_exit = self.work_remaining_ms() + if remaining_exit <= 0: + raise Error( + "measurement: work deadline expired before child exit" + " (work_deadline_expired)" + ) + var st = wait_bounded(self.pid, remaining_exit) if not st.cleanup_proved(): + self.status = st.copy() var term = terminate_owned(self.pid, TERMINATION_GRACE_MS) - self.reader.guard[].retain( - self.pid, - self.reader.report_fd, - "measurement process cleanup unproved " + term.describe(), - ) self.status = term.copy() if term.cleanup_proved(): - self.reader.guard[].resolve_pid(self.pid) self.reaped = True - self.reader.close_reader() - close_fd(self.stderr_fd) + self.guard[].resolve_pid(self.pid) + else: + self.guard[].retain( + self.pid, + self.stdout_fd, + "measurement process cleanup unproved " + term.describe(), + ) raise Error( "measurement: process did not exit within its budget (" + term.describe() @@ -560,36 +1115,61 @@ struct MeasurementProcess(Movable): ) self.status = st.copy() self.reaped = True - self.reader.guard[].resolve_pid(self.pid) - self.reader.close_reader() - close_fd(self.stderr_fd) + self.guard[].resolve_pid(self.pid) + self.close_stdout() + self.close_stderr() return st^ def cleanup(mut self): """Non-raising cleanup so a failing assertion still reaps the child.""" - if self.stdin_fd >= 0 and not self.stdin_closed: - close_fd(self.stdin_fd) - self.stdin_fd = -1 - self.stdin_closed = True + self.close_stdin() if not self.reaped: var st = terminate_owned(self.pid, TERMINATION_GRACE_MS) self.status = st.copy() if st.cleanup_proved(): self.reaped = True - self.reader.guard[].resolve_pid(self.pid) + self.guard[].resolve_pid(self.pid) else: - self.reader.guard[].retain( + self.guard[].retain( self.pid, - self.reader.report_fd, + self.stdout_fd, "measurement process cleanup unproved " + st.describe(), ) - self.reader.close_reader() - if not self.stderr_eof: - close_fd(self.stderr_fd) - self.stderr_eof = True + self.close_stdout() + self.close_stderr() + + def stderr_text(self) raises -> String: + return bytes_to_text( + _list_prefix(self.stderr_bytes, len(self.stderr_bytes)) + ) + + +def _newline_index(bytes: List[UInt8]) -> Int: + for index in range(len(bytes)): + if Int(bytes[index]) == 10: + return index + return -1 + + +def _list_prefix(bytes: List[UInt8], count: Int) -> List[UInt8]: + var out = List[UInt8]() + for index in range(count): + out.append(bytes[index]) + return out^ + + +def _list_suffix(bytes: List[UInt8], start: Int) -> List[UInt8]: + var out = List[UInt8]() + for index in range(start, len(bytes)): + out.append(bytes[index]) + return out^ - def stderr_text(mut self) raises -> String: - return bytes_to_text(self.stderr_bytes) + +def _decode_frame(var bytes: List[UInt8]) raises -> String: + try: + return bytes_to_text(bytes^) + except: + raise Error("measurement: response frame is not valid UTF-8") def spawn_measurement_process( @@ -622,7 +1202,7 @@ def spawn_measurement_process( var command_ptr = elements[0].as_c_string_slice().unsafe_ptr() var argv_ptr = argv_c.unsafe_ptr() - var start = now_ms() + var spawn_ms = now_ms() var pid = fork_owned_or_close3(pipes) if pid == 0: if dup2_fd(stdin_read_fd, 0) < 0: @@ -644,21 +1224,24 @@ def spawn_measurement_process( close_fd(stdin_read_fd) close_fd(stdout_write_fd) close_fd(stderr_write_fd) - var reader = piped_child_state( - pid, stdout_read_fd, deadline_ms, 0, UnsafePointer(to=guard) - ) return MeasurementProcess( pid=pid, stdin_fd=stdin_write_fd, - reader=reader^, + stdout_fd=stdout_read_fd, stderr_fd=stderr_read_fd, + stdout_pending=List[UInt8](), stderr_bytes=List[UInt8](), + stdout_eof=False, stderr_eof=False, - deadline_ms=deadline_ms, - start_ms=start, stdin_closed=False, + stdout_closed=False, + stderr_closed=False, reaped=False, status=ProcessStatus("pending", False, -1, 0, 0, ""), + spawn_ms=spawn_ms, + deadline_at_ms=spawn_ms + deadline_ms, + stdout_total_bytes=0, + guard=UnsafePointer(to=guard), ) @@ -768,15 +1351,27 @@ def build_status_frame(index: Int) -> List[String]: def sample_rss_kb( - pid: Int, command: String, mut guard: CleanupGuard + pid: Int, + command: String, + mut guard: CleanupGuard, + deadline_ms: Int = MEASUREMENT_SAMPLE_DEADLINE_MS, ) raises -> Int: - """Numeric resident-set size in kB, or an explicit sampling failure.""" + """Numeric resident-set size in kB, or an explicit sampling failure. + + The sampling subprocess consumes the caller's remaining work budget, so + instrumentation can never extend the measured window. + """ + if deadline_ms <= 0: + raise Error( + "measurement: work deadline expired before rss sampling" + " (work_deadline_expired)" + ) var args = List[String]() args.append("-o") args.append("rss=") args.append("-p") args.append(String(pid)) - var out = run_capture(command, args^, MEASUREMENT_SAMPLE_DEADLINE_MS, guard) + var out = run_capture(command, args^, deadline_ms, guard) if out.exit_code != 0: raise Error( "measurement: rss sampling unavailable (" @@ -796,15 +1391,23 @@ def sample_rss_kb( def sample_fd_count( - pid: Int, command: String, mut guard: CleanupGuard + pid: Int, + command: String, + mut guard: CleanupGuard, + deadline_ms: Int = MEASUREMENT_SAMPLE_DEADLINE_MS, ) raises -> Int: """Numeric descriptor count, excluding lsof cwd/txt/mem pseudo rows.""" + if deadline_ms <= 0: + raise Error( + "measurement: work deadline expired before descriptor sampling" + " (work_deadline_expired)" + ) var args = List[String]() args.append("-p") args.append(String(pid)) args.append("-F") args.append("f") - var out = run_capture(command, args^, MEASUREMENT_SAMPLE_DEADLINE_MS, guard) + var out = run_capture(command, args^, deadline_ms, guard) if out.exit_code != 0: raise Error( "measurement: descriptor sampling unavailable (" @@ -834,6 +1437,34 @@ def sample_fd_count( return count +def child_process_count( + pid: Int, + mut guard: CleanupGuard, + deadline_ms: Int = MEASUREMENT_SAMPLE_DEADLINE_MS, +) raises -> Int: + """Numeric count of live direct children of ``pid`` (0 when none).""" + if deadline_ms <= 0: + raise Error( + "measurement: work deadline expired before child census" + " (work_deadline_expired)" + ) + var args = List[String]() + args.append("-P") + args.append(String(pid)) + var out = run_capture("pgrep", args^, deadline_ms, guard) + if out.exit_code == 1: + return 0 + if out.exit_code != 0: + raise Error( + "measurement: child census unavailable (" + out.describe() + ")" + ) + var count = 0 + for line in out.stdout.split("\n"): + if String(line).strip().byte_length() > 0: + count += 1 + return count + + # ── Measurement session ───────────────────────────────────────────────────── @@ -848,8 +1479,15 @@ struct MeasurementSession(Movable): var failed_frames: Int var first_failure: String var startup_ms: Int + var startup_wall_ms: Int + var startup_sampling_ms: Int var warmup_ms: Int + var measured_wall_ms: Int + var measured_sampling_ms: Int var measured_ms: Int + var request_total_ms: Int + var request_min_ms: Int + var request_max_ms: Int var rss_kb_before_warmup: Int var rss_kb_after_warmup: Int var rss_kb_after_measured: Int @@ -861,6 +1499,7 @@ struct MeasurementSession(Movable): var child_exit: String var stderr_excerpt: String var sampling_method: String + var sampling_cadence: String def summary(self) -> String: return ( @@ -872,10 +1511,24 @@ struct MeasurementSession(Movable): + String(self.failed_frames) + " startup_ms=" + String(self.startup_ms) + + " startup_wall_ms=" + + String(self.startup_wall_ms) + + " startup_sampling_ms=" + + String(self.startup_sampling_ms) + " warmup_ms=" + String(self.warmup_ms) + " measured_ms=" + String(self.measured_ms) + + " measured_wall_ms=" + + String(self.measured_wall_ms) + + " measured_sampling_ms=" + + String(self.measured_sampling_ms) + + " request_total_ms=" + + String(self.request_total_ms) + + " request_min_ms=" + + String(self.request_min_ms) + + " request_max_ms=" + + String(self.request_max_ms) + " rss_kb=" + String(self.rss_kb_before_warmup) + "/" @@ -892,6 +1545,8 @@ struct MeasurementSession(Movable): + String(self.fd_peak) + " child=" + self.child_exit + + " stderr_bytes=" + + String(self.stderr_excerpt.byte_length()) + " first_failure=" + (self.first_failure if self.first_failure != "" else "none") ) @@ -915,31 +1570,80 @@ def measure_persistent_process( mut guard: CleanupGuard, rss_sampler: String = "ps", fd_sampler: String = "lsof", + expected_revision: String = "", + expected_manifest_sha256: String = "", + expected_binary_sha256: String = "", ) raises -> MeasurementSession: """Drive warmup then measured frames over one persistent owned process. - Every frame is written and read with a finite deadline and validated for - parsed envelope, correlation and outcome; every sample is a real numeric - value or an explicit error; the child must exit cleanly after EOF and its - cleanup is always proved through the shared guard. + Identity capture is a separate bounded phase before the spawn; the single + spawn-relative work deadline then covers sampling, request writes, response + reads, the EOF drain and the child exit. When ``expected_revision`` is + provided (product measurement) the measured source tree must be clean and + the observed revision/tree/manifest must match the build exactly; a + mismatch before or after the run is a bounded drift failure. Controlled + test children record the observed identity as ``controlled_child`` and do + not claim a build binding. """ if warmup_frames < 0 or measured_frames < 1: raise Error("measurement: invalid warmup/measured frame counts") + if deadline_ms <= 0: + raise Error("measurement: invalid work deadline") + var binding = ( + "clean_product_tree" if expected_revision != "" else "controlled_child" + ) + var observed = source_identity(source_root, guard) + if expected_revision != "": + require_clean_source(observed, "run capture") + if observed.revision != expected_revision: + raise Error( + "measurement: source revision drift before run (expected " + + expected_revision + + " observed " + + observed.revision + + ")" + ) + if observed.manifest_sha256 != expected_manifest_sha256: + raise Error( + "measurement: source manifest drift before run (expected " + + expected_manifest_sha256 + + " observed " + + observed.manifest_sha256 + + ")" + ) + var binary_sha256 = file_sha256(binary_path, guard) + if expected_binary_sha256 != "" and binary_sha256 != expected_binary_sha256: + raise Error( + "measurement: binary drift before run (expected " + + expected_binary_sha256 + + " observed " + + binary_sha256 + + ")" + ) var identity = measurement_identity( source_root, binary_path, + binary_sha256, + binding, argv_profile_of(binary_path, argv.copy()), guard, ) var process = spawn_measurement_process( - identity.binary_path, argv^, deadline_ms, guard + binary_path, argv^, deadline_ms, guard ) var ok_frames = 0 var failed_frames = 0 var first_failure = "" var startup_ms = -1 + var startup_wall_ms = -1 + var startup_sampling_ms = 0 var warmup_ms = 0 + var measured_wall_ms = 0 + var measured_sampling_ms = 0 var measured_ms = 0 + var request_total_ms = 0 + var request_min_ms = -1 + var request_max_ms = 0 var rss_before = 0 var rss_after_warmup = 0 var rss_after_measured = 0 @@ -952,11 +1656,22 @@ def measure_persistent_process( var stderr_excerpt = "" try: - rss_before = sample_rss_kb(process.pid, rss_sampler, guard) - fd_before = sample_fd_count(process.pid, fd_sampler, guard) + # Instrumentation before the first frame; timed as overhead so the + # recorded startup latency is the child's own work. + var pre_sample_start = now_ms() + var pre_remaining = process.work_remaining_ms() + rss_before = sample_rss_kb( + process.pid, rss_sampler, guard, pre_remaining + ) + var pre_remaining_fd = process.work_remaining_ms() + fd_before = sample_fd_count( + process.pid, fd_sampler, guard, pre_remaining_fd + ) + var pre_sample_ms = now_ms() - pre_sample_start rss_peak = rss_before fd_peak = fd_before var warmup_start = now_ms() + var measured_start = -1 var total = warmup_frames + measured_frames for index in range(total): var pair = build_status_frame(index) @@ -964,11 +1679,18 @@ def measure_persistent_process( var request_id = pair[1] var trace_id = pair[2] var frame_start = now_ms() - process.send_frame(frame, deadline_ms) - var response = process.read_response(deadline_ms) + process.send_frame(frame, index) + var response = process.read_response(index) + process.assert_no_surplus(index) var latency = now_ms() - frame_start if index == 0: - startup_ms = latency + # Startup is the true spawn-to-first-validated-response + # latency; any sampling that overlapped the child's own + # startup is recorded separately as instrumentation rather than + # subtracted, so the value can never understate real startup. + startup_wall_ms = now_ms() - process.spawn_ms + startup_sampling_ms = pre_sample_ms + startup_ms = startup_wall_ms var verdict = validate_status_frame(response, request_id, trace_id) if verdict.declared_max_requests >= 0: declared = verdict.declared_max_requests @@ -990,43 +1712,85 @@ def measure_persistent_process( raise Error("measurement: frame invalid " + first_failure) if index == warmup_frames - 1: warmup_ms = now_ms() - warmup_start + var s_start = now_ms() + var s_remaining = process.work_remaining_ms() rss_after_warmup = sample_rss_kb( - process.pid, rss_sampler, guard + process.pid, rss_sampler, guard, s_remaining ) + var s_remaining_fd = process.work_remaining_ms() fd_peak = max( - fd_peak, sample_fd_count(process.pid, fd_sampler, guard) + fd_peak, + sample_fd_count( + process.pid, fd_sampler, guard, s_remaining_fd + ), ) if rss_after_warmup > rss_peak: rss_peak = rss_after_warmup + measured_start = now_ms() + measured_sampling_ms -= now_ms() - s_start if index >= warmup_frames: var measured_index = index - warmup_frames + request_total_ms += latency + if request_min_ms < 0 or latency < request_min_ms: + request_min_ms = latency + if latency > request_max_ms: + request_max_ms = latency if measured_index % MEASUREMENT_SAMPLE_INTERVAL == 0: - var mid_rss = sample_rss_kb(process.pid, rss_sampler, guard) - var mid_fd = sample_fd_count(process.pid, fd_sampler, guard) + var s_start = now_ms() + var s_remaining = process.work_remaining_ms() + var mid_rss = sample_rss_kb( + process.pid, rss_sampler, guard, s_remaining + ) + var s_remaining_fd = process.work_remaining_ms() + var mid_fd = sample_fd_count( + process.pid, fd_sampler, guard, s_remaining_fd + ) + measured_sampling_ms += now_ms() - s_start if mid_rss > rss_peak: rss_peak = mid_rss if mid_fd > fd_peak: fd_peak = mid_fd - if measured_index == measured_frames - 1: - measured_ms = now_ms() - warmup_start - warmup_ms - rss_after_measured = sample_rss_kb(process.pid, rss_sampler, guard) - fd_after = sample_fd_count(process.pid, fd_sampler, guard) + # Post-measured instrumentation (still inside the measured window). + var post_start = now_ms() + var post_remaining = process.work_remaining_ms() + rss_after_measured = sample_rss_kb( + process.pid, rss_sampler, guard, post_remaining + ) + var post_remaining_fd = process.work_remaining_ms() + fd_after = sample_fd_count( + process.pid, fd_sampler, guard, post_remaining_fd + ) + measured_sampling_ms += now_ms() - post_start if rss_after_measured > rss_peak: rss_peak = rss_after_measured if fd_after > fd_peak: fd_peak = fd_after - var st = process.finish() + measured_wall_ms = now_ms() - measured_start + measured_ms = measured_wall_ms - measured_sampling_ms + var st = process.finish_expected(total) child_exit = st.describe() + " cleanup=proved" stderr_excerpt = process.stderr_text() if not st.exited or st.exit_code != 0: - # A failed child is a measurement failure, never a success with a - # missing or partial outcome. raise Error( "measurement: child exited nonzero (" + st.describe() + ") stderr=" + stderr_excerpt ) + # Post-run drift rejection: the measured binary and source must be the + # same ones captured before the run. + var after = source_identity(source_root, guard) + if expected_revision != "": + require_identity_drift_free(observed, after, "run") + var binary_after = file_sha256(binary_path, guard) + if binary_after != binary_sha256: + raise Error( + "measurement: binary drifted during run (before " + + binary_sha256 + + " after " + + binary_after + + ")" + ) except e: process.cleanup() raise Error("measurement: persistent measurement failed: " + String(e)) @@ -1039,8 +1803,15 @@ def measure_persistent_process( failed_frames=failed_frames, first_failure=first_failure, startup_ms=startup_ms, + startup_wall_ms=startup_wall_ms, + startup_sampling_ms=startup_sampling_ms, warmup_ms=warmup_ms, + measured_wall_ms=measured_wall_ms, + measured_sampling_ms=measured_sampling_ms, measured_ms=measured_ms, + request_total_ms=request_total_ms, + request_min_ms=request_min_ms, + request_max_ms=request_max_ms, rss_kb_before_warmup=rss_before, rss_kb_after_warmup=rss_after_warmup, rss_kb_after_measured=rss_after_measured, @@ -1056,4 +1827,9 @@ def measure_persistent_process( + fd_sampler + " -p <pid> -F f numeric f<digits> rows" ), + sampling_cadence=( + "before warmup, after warmup, every " + + String(MEASUREMENT_SAMPLE_INTERVAL) + + " measured frames, after measured phase" + ), ) diff --git a/tests/measurement_runner.mojo b/tests/measurement_runner.mojo @@ -4,6 +4,11 @@ Runs one persistent HYF stdio process through warmup and measured frames and prints the exact identity, per-phase timing, numeric RSS/FD samples, frame counts and child exit. Exits nonzero on any failed mandatory guarantee. +ADR-0019 D39 MR04 additionally records the numeric direct-local-provider +elapsed time, request/connection counts and the deterministic client/schema +construction characterization, so the pre-migration baseline is source-backed +rather than a test duration or an unprinted assertion. + Usage (governed lane): ``cargo extbuild run -- pixi run --frozen measure-h005a`` """ @@ -11,7 +16,7 @@ from std.collections import List from safe_tempdir import SafeTempDir -from parent_lifecycle import CleanupGuard +from parent_lifecycle import CleanupGuard, now_ms from stdio_process_helper import ( HYF_PATHS_PROFILE_ENV, HYF_PATHS_REPO_LOCAL_ROOT_ENV, @@ -21,35 +26,145 @@ from measurement_process_helper import ( build_product_binary, measure_persistent_process, ) +from max_local_process_helper import spawn_max_local_stub + +from json import Value, loads + +from hyf_core.request_context import default_request_context +from hyf_provider.client import post_max_local_chat_completion +from hyf_provider.config import MaxLocalProviderConfig +from hyf_provider.result import parse_query_analysis_from_chat_completion +from hyf_provider.schema import build_query_rewrite_request_body comptime MEASUREMENT_DEADLINE_MS = 180000 comptime WARMUP_FRAMES = 100 comptime MEASURED_FRAMES = 1000 +comptime ANALYSIS_JSON_TEXT = ( + '{"original_text":"eggs near me",' + '"normalized_text":"eggs near me",' + '"rewritten_text":"eggs",' + '"query_terms":["eggs"],' + '"normalization_signals":["local_intent_detected"],' + '"ranking_hints":["prefer_local_results"],' + '"extracted_filters":{' + '"local_intent":true,' + '"fulfillment":"unspecified",' + '"time_window":"unspecified"' + "}}" +) + + +def characterize_direct_provider(mut guard: CleanupGuard) raises: + """Numeric direct-local-provider and client/schema characterization. + + One direct provider request over one verified connection against the local + stub: the elapsed stub-startup, schema-construction and request times and + the request/connection counts are printed, and the deterministic request + body/response parsing is proved. The Morph daemon-assisted path is out of + scope here and remains the explicit H024 obligation. + """ + var stub_start = now_ms() + with spawn_max_local_stub(0, "count_requests", 1, guard) as stub: + var stub_startup_ms = now_ms() - stub_start + var config = MaxLocalProviderConfig( + base_url="http://127.0.0.1:" + String(stub.port) + "/v1/", + health_url="http://127.0.0.1:" + String(stub.port) + "/health", + model="max-local-query-rewrite", + request_timeout_ms=15000, + ) + var context = default_request_context() + context.return_provenance = True + var construct_start = now_ms() + var body = build_query_rewrite_request_body( + config, "eggs near me", context + ) + var construct_ms = now_ms() - construct_start + var response_format = body["response_format"] + print( + "h005a.client_schema", + "fields=" + String(body.object_count()), + "messages=" + String(body["messages"].array_count()), + "response_format_type=" + response_format["type"].string_value(), + "json_schema_name=" + + response_format["json_schema"]["name"].string_value(), + "model=" + body["model"].string_value(), + ) + var request_start = now_ms() + var outcome = post_max_local_chat_completion(config, body) + var request_ms = now_ms() - request_start + stub.wait() + var failure_text = "false" + if outcome.failure: + failure_text = "true" + print( + "h005a.direct_provider", + "stub_startup_ms=" + String(stub_startup_ms), + "construct_ms=" + String(construct_ms), + "request_ms=" + String(request_ms), + "requests=" + String(stub.request_count()), + "connections=" + String(stub.connection_count()), + "failure=" + failure_text, + ) + if outcome.failure: + raise Error("measurement: direct provider request failed") + if stub.request_count() != 1: + raise Error("measurement: direct provider request count mismatch") + if stub.connection_count() != 1: + raise Error( + "measurement: direct provider connection count mismatch" + ) + var response = loads("{}") + var choices = loads("[]") + var choice = loads("{}") + var message = loads("{}") + message.set("content", Value(ANALYSIS_JSON_TEXT)) + choice.set("message", message) + choices.append(choice) + response.set("choices", choices) + var analysis = parse_query_analysis_from_chat_completion(response) + if analysis.original_text != "eggs near me": + raise Error("measurement: provider response parsing mismatch") + if len(analysis.query_terms) != 1: + raise Error("measurement: provider response term count mismatch") + def main() raises: var guard = CleanupGuard() + var source_root = "." with SafeTempDir() as temp_dir: with ScopedEnvVar(HYF_PATHS_PROFILE_ENV, "repo_local"): with ScopedEnvVar(HYF_PATHS_REPO_LOCAL_ROOT_ENV, temp_dir): - var binary = build_product_binary(temp_dir, guard) + var built = build_product_binary(source_root, temp_dir, guard) + print( + "h005a.build", + "build_ms=" + String(built.build_ms), + built.source.describe(), + ) var argv = List[String]() var measured = measure_persistent_process( - ".", - binary, + source_root, + built.binary_path, argv^, WARMUP_FRAMES, MEASURED_FRAMES, MEASUREMENT_DEADLINE_MS, guard, + "ps", + "lsof", + built.source.revision, + built.source.manifest_sha256, + built.binary_sha256, ) print("h005a.identity", measured.identity.describe()) print("h005a.measurement", measured.summary()) print("h005a.sampling_method", measured.sampling_method) + print("h005a.sampling_cadence", measured.sampling_cadence) print( "h005a.stderr_bytes", measured.stderr_excerpt.byte_length(), ) + characterize_direct_provider(guard) guard.assert_clean() print("h005a_measurement: ok") diff --git a/tests/test_measurement_contract.mojo b/tests/test_measurement_contract.mojo @@ -1,11 +1,19 @@ -"""H005A measurement contract tests (ADR-0012 D29, ADR-0014 D34). +"""H005A measurement contract tests (ADR-0012 D29, ADR-0014 D34, ADR-0019 D39). Drives the governed persistent-process measurement tooling against a real build of the existing product entry point and against controlled child -processes, proving that the reproduced H005 defects (R56/R57) now fail the -measurement instead of reporting success: not-json output, wrong correlation, -unterminated output, early EOF, a failed child and unavailable ``ps``/``lsof`` -sampling. +processes, proving that the reproduced H005/R56/R57 defects and the period-10 +counterexamples now fail the measurement instead of reporting success: + +* MR01 — extra, coalesced, split, unterminated and malformed trailing stdout, + early EOF and a nonzero child exit; +* MR02 — a valid response or exit that arrives after the one work budget, a + sampling subprocess that would outlive it, stderr overflow and a + deadline-bounded EINTR retry; +* MR03 — repeated failing public measurement calls leave no descriptor or + child behind, with successful recovery afterward; +* MR04 — source/binary identity drift rejection, delayed-startup timing and + truthful per-request/instrumentation accounting. This module is test-only tooling. It changes no product policy, schema, dependency or lock. @@ -16,7 +24,12 @@ from std.testing import TestSuite, assert_equal, assert_true from safe_tempdir import SafeTempDir -from parent_lifecycle import CleanupGuard, owned_pid, now_ms +from parent_lifecycle import ( + CleanupGuard, + now_ms, + open_fd_count_checked, + owned_pid, +) from stdio_process_helper import ( HYF_PATHS_PROFILE_ENV, HYF_PATHS_REPO_LOCAL_ROOT_ENV, @@ -26,8 +39,11 @@ from measurement_process_helper import ( MeasurementSession, build_product_binary, build_status_frame, + child_process_count, file_sha256, measure_persistent_process, + measurement_poll, + measurement_poll_retry, run_capture, sample_fd_count, sample_rss_kb, @@ -40,16 +56,34 @@ from hyf_provider.client import post_max_local_chat_completion from hyf_provider.config import MaxLocalProviderConfig from hyf_provider.result import parse_query_analysis_from_chat_completion from hyf_provider.schema import build_query_rewrite_request_body -from max_local_process_helper import ( - reserve_loopback_port, - spawn_max_local_stub, -) +from max_local_process_helper import spawn_max_local_stub comptime MEASUREMENT_DEADLINE_MS = 120000 comptime WARMUP_FRAMES = 20 comptime MEASURED_FRAMES = 200 +# One valid sys.status response for request/trace index 0. +comptime STATUS0 = ( + '{"version":1,"request_id":"meas-status-0","trace_id":"meas-trace-0",' + '"ok":true,"output":{"daemon":"hyfd"}}' +) +comptime WRONG_REVISION = "0000000000000000000000000000000000000000" +comptime WRONG_DIGEST = ( + "0000000000000000000000000000000000000000000000000000000000000000" +) + + +def _one_response() -> String: + """A well-behaved responder: answer exactly one request, then consume the + rest of stdin until EOF so the child exits cleanly and closes stdout. + """ + return ( + "IFS= read -r line; printf '%s\\n' '" + + STATUS0 + + "'; while IFS= read -r line; do :; done" + ) + def _sh(args_text: String) -> List[String]: var args = List[String]() @@ -65,6 +99,7 @@ def _run_sh_measurement( mut guard: CleanupGuard, rss_sampler: String = "ps", fd_sampler: String = "lsof", + deadline_ms: Int = MEASUREMENT_DEADLINE_MS, ) raises -> MeasurementSession: var argv = _sh(args_text) return measure_persistent_process( @@ -73,44 +108,85 @@ def _run_sh_measurement( argv^, warmup, measured, - MEASUREMENT_DEADLINE_MS, + deadline_ms, guard, rss_sampler, fd_sampler, ) +def _run_sh_failure( + args_text: String, + warmup: Int, + measured: Int, + mut guard: CleanupGuard, + rss_sampler: String = "ps", + fd_sampler: String = "lsof", + deadline_ms: Int = MEASUREMENT_DEADLINE_MS, +) -> String: + try: + _ = _run_sh_measurement( + args_text, + warmup, + measured, + guard, + rss_sampler, + fd_sampler, + deadline_ms, + ) + except e: + return String(e) + return "" + + # ── Positive persistent measurement ───────────────────────────────────────── def test_persistent_measurement_validates_every_frame() raises: - # D29: one process serves warmup and measured frames; every frame has a - # parsed envelope, matching correlation and expected outcome, with numeric - # RSS/FD samples, a checked child exit and proved cleanup. The recorded - # environment profile is the verified live profile of the measured child. + # D29/MR01: one process serves warmup and measured frames; every frame has + # a parsed envelope, matching correlation and expected outcome, the whole + # stream is accounted through EOF, with numeric RSS/FD samples, a checked + # child exit and proved cleanup. The recorded environment profile is the + # verified live profile of the measured child. var guard = CleanupGuard() with SafeTempDir() as temp_dir: with ScopedEnvVar(HYF_PATHS_PROFILE_ENV, "repo_local"): with ScopedEnvVar(HYF_PATHS_REPO_LOCAL_ROOT_ENV, temp_dir): - var binary = build_product_binary(temp_dir, guard) - assert_true(file_sha256(binary, guard).byte_length() == 64) + var built = build_product_binary(".", temp_dir, guard) + assert_equal(built.binary_sha256.byte_length(), 64) + assert_equal(built.source.dirty_status, "") + assert_equal(built.source.manifest_sha256.byte_length(), 64) + assert_true(built.build_ms > 0) var argv = List[String]() var session = measure_persistent_process( ".", - binary, + built.binary_path, argv^, WARMUP_FRAMES, MEASURED_FRAMES, MEASUREMENT_DEADLINE_MS, guard, + "ps", + "lsof", + built.source.revision, + built.source.manifest_sha256, + built.binary_sha256, ) assert_equal(session.ok_frames, WARMUP_FRAMES + MEASURED_FRAMES) assert_equal(session.failed_frames, 0) assert_equal(session.first_failure, "") - # Startup, warmup and measured timing are recorded separately (ms). + # Startup is measured from spawn to the first validated + # response; instrumentation is recorded separately. assert_true(session.startup_ms >= 0) + assert_true(session.startup_wall_ms >= session.startup_ms) + assert_true(session.startup_sampling_ms >= 0) + # Measured-phase wall time excludes sampling instrumentation. assert_true(session.measured_ms >= 0) - # Numeric sampling with recorded units and method. + assert_true(session.measured_wall_ms >= session.measured_ms) + assert_true(session.measured_sampling_ms >= 0) + assert_true(session.request_total_ms > 0) + assert_true(session.request_max_ms >= session.request_min_ms) + # Numeric sampling with recorded units, method and cadence. assert_true(session.rss_kb_before_warmup > 0) assert_true(session.rss_kb_after_warmup > 0) assert_true(session.rss_kb_after_measured > 0) @@ -120,17 +196,35 @@ def test_persistent_measurement_validates_every_frame() raises: assert_true(session.fd_peak >= session.fd_after_measured) assert_true(session.sampling_method.find("kB") >= 0) assert_true(session.sampling_method.find("-F f") >= 0) + assert_true(session.sampling_cadence.find("every") >= 0) # The declared per-process request policy is characterized, not # assumed: the persistent loop does not enforce it. assert_equal(session.declared_max_requests_per_process, 1) assert_true(session.child_exit.find("exited=0") >= 0) assert_true(session.stderr_excerpt == "") - # Exact, truthful identity: source revision, cwd, binary, pixi - # files, toolchain, host and the verified environment profile. - assert_equal(len(session.identity.binary_sha256), 64) - assert_equal(len(session.identity.pixi_lock_sha256), 64) - assert_equal(len(session.identity.pixi_toml_sha256), 64) - assert_equal(len(session.identity.source_revision), 40) + # Exact, truthful identity: clean verified source/tree plus a + # deterministic content manifest, binary, pixi files, toolchain, + # host and the verified environment profile. + assert_equal(session.identity.binding, "clean_product_tree") + assert_equal( + session.identity.binary_sha256, built.binary_sha256 + ) + assert_equal( + session.identity.source_revision, built.source.revision + ) + assert_equal( + session.identity.source_manifest_sha256, + built.source.manifest_sha256, + ) + assert_equal(session.identity.source_tree_state, "clean") + assert_equal(session.identity.source_tree.byte_length(), 40) + assert_equal( + session.identity.pixi_lock_sha256.byte_length(), 64 + ) + assert_equal( + session.identity.pixi_toml_sha256.byte_length(), 64 + ) + assert_equal(session.identity.source_revision.byte_length(), 40) assert_true(session.identity.cwd.find("oss/hyf") >= 0) assert_true( session.identity.env_profile.find( @@ -153,6 +247,77 @@ def test_persistent_measurement_validates_every_frame() raises: guard.assert_clean() +def test_measurement_identity_drift_is_rejected() raises: + # MR04: a product measurement must reject a binary/source binding that does + # not match the observed clean source identity. + var guard = CleanupGuard() + var message = "" + with SafeTempDir() as temp_dir: + var built = build_product_binary(".", temp_dir, guard) + var wrong_revision = WRONG_REVISION + var argv = List[String]() + try: + _ = measure_persistent_process( + ".", + built.binary_path, + argv^, + 0, + 1, + MEASUREMENT_DEADLINE_MS, + guard, + "ps", + "lsof", + wrong_revision, + built.source.manifest_sha256, + built.binary_sha256, + ) + except e: + message = String(e) + assert_true(message.find("drift") >= 0) + guard.assert_clean() + + +def test_measurement_rejects_wrong_binary_digest() raises: + # MR04: the recorded binary digest is enforced, not merely recorded. + var guard = CleanupGuard() + var message = "" + with SafeTempDir() as temp_dir: + var built = build_product_binary(".", temp_dir, guard) + var argv = List[String]() + try: + _ = measure_persistent_process( + ".", + built.binary_path, + argv^, + 0, + 1, + MEASUREMENT_DEADLINE_MS, + guard, + "ps", + "lsof", + built.source.revision, + built.source.manifest_sha256, + WRONG_DIGEST, + ) + except e: + message = String(e) + assert_true(message.find("binary drift") >= 0) + guard.assert_clean() + + +def test_measurement_reports_delayed_startup() raises: + # MR04: startup timing is spawn-relative and truthful, so a deliberately + # delayed child is observed as such rather than as a near-zero value. + var guard = CleanupGuard() + var session = _run_sh_measurement( + "sleep 0.3; " + _one_response(), 0, 1, guard + ) + assert_equal(session.ok_frames, 1) + assert_true(session.startup_ms >= 200) + assert_true(session.startup_wall_ms >= session.startup_ms) + guard.assert_clean() + + def test_measurement_sampling_is_numeric_and_units_are_recorded() raises: # D29: the sampler must return real numeric values with explicit # unavailability, never a placeholder. @@ -181,15 +346,16 @@ comptime ANALYSIS_JSON_TEXT = ( def test_direct_provider_request_and_client_schema_characterization() raises: - # D29: distinguish startup, a deterministic daemon request, a direct local - # provider request and connection counts, and characterize client/schema - # construction with source evidence. The Morph daemon-assisted path is out - # of scope here and remains an explicit H024 obligation. + # D29/MR04: distinguish startup, a deterministic daemon request, a direct + # local provider request and connection counts, and characterize + # client/schema construction with numeric recorded values and source + # evidence. The Morph daemon-assisted path is out of scope here and remains + # an explicit H024 obligation. var guard = CleanupGuard() - var startup_start = now_ms() + var stub_start = now_ms() with spawn_max_local_stub(0, "count_requests", 1, guard) as stub: - var startup_ms = now_ms() - startup_start - assert_true(startup_ms >= 0) + var stub_startup_ms = now_ms() - stub_start + assert_true(stub_startup_ms >= 0) var config = MaxLocalProviderConfig( base_url="http://127.0.0.1:" + String(stub.port) + "/v1/", health_url="http://127.0.0.1:" + String(stub.port) + "/health", @@ -198,10 +364,17 @@ def test_direct_provider_request_and_client_schema_characterization() raises: ) var context = default_request_context() context.return_provenance = True + var construct_start = now_ms() var body = build_query_rewrite_request_body( config, "eggs near me", context ) - # Schema construction is deterministic and source-verifiable. + var construct_ms = now_ms() - construct_start + # Schema construction is deterministic and source-verifiable, with a + # numeric field/message count that is recorded rather than asserted as + # a bare non-negative duration. + assert_true(construct_ms >= 0) + assert_true(body.object_count() > 0) + assert_equal(body["messages"].array_count(), 2) assert_equal(body["model"].string_value(), "max-local-query-rewrite") assert_equal(body["messages"][0]["role"].string_value(), "system") assert_equal(body["messages"][1]["role"].string_value(), "user") @@ -212,9 +385,14 @@ def test_direct_provider_request_and_client_schema_characterization() raises: body["response_format"]["json_schema"]["name"].string_value(), "query_rewrite", ) - # One direct provider request over one verified connection. + # One direct provider request over one verified connection, with a + # numeric elapsed time that is asserted to be a plausible positive + # measurement. + var request_start = now_ms() var outcome = post_max_local_chat_completion(config, body) + var request_ms = now_ms() - request_start assert_true(not outcome.failure) + assert_true(request_ms >= 0) stub.wait() assert_equal(stub.request_count(), 1) assert_equal(stub.connection_count(), 1) @@ -234,83 +412,163 @@ def test_direct_provider_request_and_client_schema_characterization() raises: guard.assert_clean() -# ── R56/R57 counterexamples must fail the measurement ─────────────────────── +def test_measurement_poll_eintr_is_deadline_bounded() raises: + # MR02: a real poll EINTR is retried but can never outlive the budget. + var start = now_ms() + var pr = measurement_poll_retry(-1, 0, -1, 0, -1, 0, 60, 1) + var elapsed = now_ms() - start + assert_true(pr.interrupted) + assert_true(elapsed >= 50) + assert_true(elapsed < 2000) + # An ordinary no-readiness poll is not misclassified as interrupted. + var quiet = measurement_poll(-1, 0, -1, 0, -1, 0, 0) + assert_equal(quiet.count, 0) + assert_true(not quiet.interrupted) + + +# ── R56/R57/R69 counterexamples must fail the measurement ─────────────────── def test_measurement_rejects_not_json_response() raises: var guard = CleanupGuard() - var message = "" - try: - _ = _run_sh_measurement( - "while IFS= read -r line; do printf 'not-json\\n'; done", - 0, - 2, - guard, - ) - except e: - message = String(e) + var message = _run_sh_failure( + "while IFS= read -r line; do printf 'not-json\\n'; done", + 0, + 2, + guard, + ) assert_true(message.find("not_json") >= 0) guard.assert_clean() def test_measurement_rejects_wrong_correlation() raises: var guard = CleanupGuard() - var message = "" - try: - _ = _run_sh_measurement( - ( - "while IFS= read -r line; do printf '%s\\n' " - '\'{"version":1,"request_id":"wrong","trace_id":"wrong",' - '"ok":true,"output":{"daemon":"hyfd"}}\'; done' - ), - 0, - 2, - guard, - ) - except e: - message = String(e) + var message = _run_sh_failure( + ( + "while IFS= read -r line; do printf '%s\\n' " + '\'{"version":1,"request_id":"wrong","trace_id":"wrong",' + '"ok":true,"output":{"daemon":"hyfd"}}\'; done' + ), + 0, + 2, + guard, + ) assert_true(message.find("correlation_mismatch") >= 0) guard.assert_clean() def test_measurement_rejects_unterminated_response() raises: var guard = CleanupGuard() - var message = "" - try: - _ = _run_sh_measurement( - ( - "IFS= read -r line; printf '%s' " - '\'{"version":1,"request_id":"meas-status-0",' - '"trace_id":"meas-trace-0","ok":true,' - '"output":{"daemon":"hyfd"}}\'; exit 0' - ), - 0, - 1, - guard, - ) - except e: - message = String(e) + var message = _run_sh_failure( + ( + "IFS= read -r line; printf '%s' " + '\'{"version":1,"request_id":"meas-status-0",' + '"trace_id":"meas-trace-0","ok":true,' + '"output":{"daemon":"hyfd"}}\'; exit 0' + ), + 0, + 1, + guard, + ) assert_true(message.find("newline-terminated") >= 0) guard.assert_clean() +def test_measurement_rejects_extra_response_frame() raises: + # MR01/R69: an extra unvalidated frame after the expected response must + # fail, not be reported as success. + var guard = CleanupGuard() + var message = _run_sh_failure( + ( + "IFS= read -r line; printf '%s\\n%s\\n' '" + + STATUS0 + + "' 'EXTRA_UNVALIDATED_FRAME'; while IFS= read -r line; do :; done" + ), + 0, + 1, + guard, + ) + assert_true(message.find("unexpected trailing stdout") >= 0) + guard.assert_clean() + + +def test_measurement_rejects_coalesced_trailing_frame() raises: + # MR01: two frames arriving in the same read chunk are still two frames. + var guard = CleanupGuard() + var message = _run_sh_failure( + ( + "IFS= read -r line; printf '%s\\n%s\\n' '" + + STATUS0 + + "' '" + + STATUS0 + + "'; while IFS= read -r line; do :; done" + ), + 0, + 1, + guard, + ) + assert_true( + message.find("unexpected trailing stdout") >= 0 + or message.find("correlation_mismatch") >= 0 + or message.find("not_json") >= 0 + ) + guard.assert_clean() + + +def test_measurement_rejects_trailing_malformed_bytes() raises: + # MR01: trailing bytes without a complete frame are a bounded failure. + var guard = CleanupGuard() + var message = _run_sh_failure( + ( + "IFS= read -r line; printf '%s\\n' '" + + STATUS0 + + "'; printf 'garbage'; while IFS= read -r line; do :; done" + ), + 0, + 1, + guard, + ) + assert_true( + message.find("unexpected trailing stdout") >= 0 + or message.find("newline-terminated") >= 0 + ) + guard.assert_clean() + + +def test_measurement_rejects_unterminated_trailing_frame() raises: + # MR01: a second frame that never terminates is not silently ignored. + var guard = CleanupGuard() + var message = _run_sh_failure( + ( + "IFS= read -r line; printf '%s\\n' '" + + STATUS0 + + "'; printf '{\"partial\":true';" + " while IFS= read -r line; do :; done" + ), + 0, + 1, + guard, + ) + assert_true( + message.find("unexpected trailing stdout") >= 0 + or message.find("newline-terminated") >= 0 + ) + guard.assert_clean() + + def test_measurement_rejects_failed_child() raises: var guard = CleanupGuard() - var message = "" - try: - _ = _run_sh_measurement( - ( - "while IFS= read -r line; do printf '%s\\n' " - '\'{"version":1,"request_id":"meas-status-0",' - '"trace_id":"meas-trace-0","ok":true,' - '"output":{"daemon":"hyfd"}}\'; done; exit 17' - ), - 0, - 1, - guard, - ) - except e: - message = String(e) + var message = _run_sh_failure( + ( + "while IFS= read -r line; do printf '%s\\n' " + '\'{"version":1,"request_id":"meas-status-0",' + '"trace_id":"meas-trace-0","ok":true,' + '"output":{"daemon":"hyfd"}}\'; done; exit 17' + ), + 0, + 1, + guard, + ) assert_true(message.find("nonzero") >= 0) guard.assert_clean() @@ -319,66 +577,159 @@ def test_measurement_rejects_early_eof_child() raises: # A child that never answers and closes its stream must fail as an early # EOF, not be read as a successful empty response. var guard = CleanupGuard() - var message = "" - try: - _ = _run_sh_measurement("sleep 1; exit 0", 0, 1, guard) - except e: - message = String(e) + var message = _run_sh_failure("sleep 1; exit 0", 0, 1, guard) assert_true(message.find("measurement") >= 0) assert_true( message.find("early_eof") >= 0 or message.find("write") >= 0 - or message.find("descriptor sampling") >= 0 + or message.find("sampling") >= 0 ) guard.assert_clean() -def test_measurement_rejects_unavailable_rss_sampler() raises: +# ── MR02 work-budget and pressure controls ────────────────────────────────── + + +def test_measurement_rejects_late_response_after_budget() raises: + # MR02: a valid response that arrives after the one work budget is not + # accepted as a late success. var guard = CleanupGuard() - var message = "" - try: - _ = _run_sh_measurement( - "while IFS= read -r line; do printf '%s\\n' ok; done", - 0, - 1, - guard, - "hyf-no-such-rss-sampler", + var message = _run_sh_failure( + "sleep 2; " + _one_response(), 0, 1, guard, "ps", "lsof", 400 + ) + assert_true(message.find("work_deadline_expired") >= 0) + guard.assert_clean() + + +def test_measurement_rejects_late_exit_after_budget() raises: + # MR02: complete and valid output does not rescue a child that keeps the + # stream open past the work budget. + var guard = CleanupGuard() + var message = _run_sh_failure( + ( + "while IFS= read -r line; do printf '%s\\n' '" + + STATUS0 + + "'; done; sleep 2" + ), + 0, + 1, + guard, + "ps", + "lsof", + 400, + ) + assert_true(message.find("work_deadline_expired") >= 0) + guard.assert_clean() + + +def test_measurement_rejects_slow_sampling_past_budget() raises: + # MR02: a sampling subprocess consumes the remaining work budget and can + # never extend the measured window. + var guard = CleanupGuard() + with SafeTempDir() as temp_dir: + var script = temp_dir + "/slow_sampler.sh" + var mk = List[String]() + mk.append("-c") + mk.append( + "printf '#!/bin/sh\\nsleep 4\\necho 17\\n' > '" + + script + + "'; chmod +x '" + + script + + "'" ) - except e: - message = String(e) + var made = run_capture("sh", mk^, 10000, guard) + assert_equal(made.exit_code, 0) + var message = _run_sh_failure( + _one_response(), 0, 1, guard, script, "lsof", 1500 + ) + assert_true( + message.find("sample_deadline_expired") >= 0 + or message.find("work_deadline_expired") >= 0 + ) + guard.assert_clean() + + +def test_measurement_rejects_stderr_overflow() raises: + # MR02: stderr pressure past the cap fails explicitly instead of being + # silently dropped. + var guard = CleanupGuard() + var message = _run_sh_failure( + ( + "IFS= read -r line; printf '%s\\n' '" + + STATUS0 + + "'; head -c 200000 /dev/zero | tr '\\0' 'x' 1>&2;" + " while IFS= read -r line; do :; done" + ), + 0, + 1, + guard, + ) + assert_true(message.find("stderr_overflow") >= 0) + guard.assert_clean() + + +def test_measurement_rejects_unavailable_rss_sampler() raises: + var guard = CleanupGuard() + var message = _run_sh_failure( + "while IFS= read -r line; do printf '%s\\n' ok; done", + 0, + 1, + guard, + "hyf-no-such-rss-sampler", + ) assert_true(message.find("rss sampling unavailable") >= 0) guard.assert_clean() def test_measurement_rejects_unavailable_fd_sampler() raises: var guard = CleanupGuard() - var message = "" - try: - _ = _run_sh_measurement( - "while IFS= read -r line; do printf '%s\\n' ok; done", - 0, - 1, - guard, - "ps", - "hyf-no-such-fd-sampler", - ) - except e: - message = String(e) + var message = _run_sh_failure( + "while IFS= read -r line; do printf '%s\\n' ok; done", + 0, + 1, + guard, + "ps", + "hyf-no-such-fd-sampler", + ) assert_true(message.find("descriptor sampling unavailable") >= 0) guard.assert_clean() def test_measurement_rejects_invalid_frame_counts() raises: var guard = CleanupGuard() - var message = "" - try: - _ = _run_sh_measurement("exit 0", 0, 0, guard) - except e: - message = String(e) + var message = _run_sh_failure("exit 0", 0, 0, guard) assert_true(message.find("invalid warmup/measured") >= 0) guard.assert_clean() +# ── MR03 exact resource ownership ─────────────────────────────────────────── + + +def test_measurement_repeated_failures_leak_nothing() raises: + # MR03: repeated failing public measurement calls must leave no descriptor + # or child behind, and a subsequent supported call must still succeed. + var guard = CleanupGuard() + var self_pid = owned_pid() + var child_before = child_process_count(self_pid, guard) + var fd_before = open_fd_count_checked() + for index in range(4): + var message = _run_sh_failure( + "while IFS= read -r line; do printf 'not-json\\n'; done", + 0, + 1, + guard, + ) + assert_true(message.find("not_json") >= 0) + guard.assert_clean() + var fd_after = open_fd_count_checked() + assert_equal(fd_after - fd_before, 0) + assert_equal(child_process_count(self_pid, guard), child_before) + var recovered = _run_sh_measurement(_one_response(), 0, 1, guard) + assert_equal(recovered.ok_frames, 1) + assert_equal(recovered.failed_frames, 0) + guard.assert_clean() + + # ── Direct correlation validation unit controls ─────────────────────────────