measurement_process_helper.mojo (76811B)
1 """Governed test-only persistent-process measurement tooling (ADR-0012 D29). 2 3 H005A requires repo-owned, standalone measurement tooling rather than a 4 docs-resident script. This helper owns exactly one persistent HYF stdio process 5 and drives warmup plus measured frames over that single process, validating 6 every frame's parsed envelope, correlation, outcome and count with finite 7 bounded I/O, a checked child exit and exception-safe cleanup. 8 9 It exists only for tests. It changes no HYF product policy, no schema and no 10 dependency: the process is launched from a build of the existing product entry 11 point and the samples are read back with the bounded lifecycle primitives. 12 13 ADR-0019 D39 MR01-MR05 requalification: 14 15 * MR01 — the *whole* stdout stream is accounted for through EOF. Exactly the 16 expected warmup/measured frames may succeed; a retained surplus byte, a 17 trailing frame in the same or a later chunk, an unterminated tail, early EOF 18 and a nonzero exit all fail with a bounded cause. stdout/stderr are drained 19 concurrently with request writes and exit observation. 20 * MR02 — one spawn-relative measurement work deadline covers sampling, request 21 writes, response reads, the EOF drain and the child exit. An expired budget 22 is never clamped into a fresh successful interval and a late valid response 23 is never accepted. Build and pre-spawn identity capture are separate, 24 explicitly bounded phases. Byte/output caps distinguish timeout, read error, 25 EOF and overflow; stderr read errors and overflow fail explicitly. 26 * MR03 — every owned descriptor (stdin, stdout, stderr) is closed exactly once 27 on success, invalid output, sampling failure, early EOF, nonzero exit and 28 timeout. EOF is not descriptor closure. The terminal child state is cached 29 before any further wait/signal, and an unproved cleanup retains the exact 30 child/descriptor ownership in the caller-held guard. 31 * MR04 — the product binary is bound to a clean, verified source/tree identity 32 plus a deterministic content manifest of the tracked build inputs; capture, 33 build and run drift is rejected. Startup is measured from spawn to the first 34 validated response, sampling is recorded separately as instrumentation 35 overhead, and the secret-free HYF_PATHS profile is read back from the live 36 environment. 37 38 Explicit failure policy (R56/R57, R69/R70/R71): 39 40 * a frame that is not JSON, has the wrong correlation or outcome, or a child 41 that exits nonzero, fails the measurement instead of reporting success; 42 * unavailable ``ps``/``lsof`` sampling raises a sampling error rather than 43 recording a placeholder value; 44 * every read/write/wait is parent-bounded by the one finite work deadline and 45 the owned child is always terminated and reaped through the shared ownership 46 guard. 47 """ 48 49 import std.os 50 from std.collections import List 51 from std.ffi import ( 52 ErrNo, 53 CStringSlice, 54 c_int, 55 c_ssize_t, 56 c_size_t, 57 c_uint, 58 external_call, 59 get_errno, 60 ) 61 62 from json import Value, loads 63 from std.pathlib import Path 64 65 from safe_tempdir import SafeTempDir 66 67 from parent_lifecycle import ( 68 IO_DEADLINE_EXPIRED, 69 POLLERR, 70 POLLHUP, 71 POLLIN, 72 POLLNVAL, 73 POLLOUT, 74 TERMINATION_GRACE_MS, 75 CleanupGuard, 76 ProcessStatus, 77 child_exit, 78 close_fd, 79 dup2_fd, 80 fork_owned_or_close3, 81 make_three_pipes, 82 now_ms, 83 poll_three, 84 read_fd, 85 set_alarm, 86 sleep_ms, 87 terminate_owned, 88 wait_bounded, 89 write_fd_chunk, 90 ) 91 92 93 comptime MEASUREMENT_FRAME_BYTES = 1048576 94 comptime MEASUREMENT_CHILD_ALARM_SECONDS = 900 95 comptime MEASUREMENT_MAX_STDERR_BYTES = 65536 96 comptime MEASUREMENT_SAMPLE_DEADLINE_MS = 10000 97 comptime MEASUREMENT_BUILD_DEADLINE_MS = 600000 98 comptime MEASUREMENT_SAMPLE_INTERVAL = 50 99 # Bounded polling granularity. Every wait is sliced by this value so an 100 # idle-but-live child is observed incrementally; the recorded request latency 101 # therefore carries at most one polling slice of tolerance rather than a fixed 102 # idle delay. Instrumentation (sampling subprocesses) is timed separately. 103 comptime MEASUREMENT_POLL_SLICE_MS = 10 104 105 106 struct MeasurementFaults(ImplicitlyCopyable, Movable): 107 """Bounded test-only fault seam for the measurement lifecycle. 108 109 Each counter is consumed once. ``stderr_read_errors`` makes the next ready 110 stderr drain report a real read-error cause without touching host state, so 111 the explicit stderr-error branch is deterministically executable. 112 ``cleanup_failures`` makes the next cleanup attempt report an unproved result 113 while the real forked child keeps running, so the retained-ownership and 114 recovery path is exercised against a real exact-owned child. 115 """ 116 117 var stderr_read_errors: Int 118 var cleanup_failures: Int 119 120 def __init__( 121 out self, stderr_read_errors: Int = 0, cleanup_failures: Int = 0 122 ): 123 self.stderr_read_errors = stderr_read_errors 124 self.cleanup_failures = cleanup_failures 125 126 def __copyinit__(out self, existing: Self): 127 self.stderr_read_errors = existing.stderr_read_errors 128 self.cleanup_failures = existing.cleanup_failures 129 130 131 # ── Bounded captured commands (identity and sampling) ─────────────────────── 132 133 134 @fieldwise_init 135 struct CommandOutput(Movable): 136 var exit_code: Int 137 var signal: Int 138 var stdout: String 139 var stderr: String 140 141 def describe(self) -> String: 142 return ( 143 "exited=" 144 + String(self.exit_code) 145 + " signal=" 146 + String(self.signal) 147 ) 148 149 150 def run_capture( 151 command: String, 152 var args: List[String], 153 deadline_ms: Int, 154 mut guard: CleanupGuard, 155 cwd: String = "", 156 ) raises -> CommandOutput: 157 """Run one bounded child command and capture its stdout/stderr. 158 159 The child is owned by this test: it is forked once, its stdio is closed or 160 captured, every drain is bounded by ``deadline_ms`` and a timeout terminates 161 and reaps the exact owned pid. A non-``execvp`` failure is reported as an 162 exit rather than silently succeeding. 163 """ 164 var parts = List[String]() 165 parts.append(command) 166 for index in range(len(args)): 167 parts.append(args[index]) 168 var argv = List[Optional[CStringSlice[ImmutAnyOrigin]]]( 169 length=len(parts) + 1, fill={} 170 ) 171 var elements = parts.unsafe_ptr() 172 for index in range(len(parts)): 173 argv[index] = rebind[CStringSlice[ImmutAnyOrigin]]( 174 elements[index].as_c_string_slice() 175 ) 176 177 var pipes = make_three_pipes() 178 var stdin_read_fd = pipes.stdin_pipe.read_fd 179 var stdin_write_fd = pipes.stdin_pipe.write_fd 180 var stdout_read_fd = pipes.stdout_pipe.read_fd 181 var stdout_write_fd = pipes.stdout_pipe.write_fd 182 var stderr_read_fd = pipes.stderr_pipe.read_fd 183 var stderr_write_fd = pipes.stderr_pipe.write_fd 184 var command_ptr = elements[0].as_c_string_slice().unsafe_ptr() 185 var argv_ptr = argv.unsafe_ptr() 186 var cwd_local = String(cwd) 187 var cwd_ptr = cwd_local.as_c_string_slice().unsafe_ptr() 188 189 var pid = fork_owned_or_close3(pipes) 190 if pid == 0: 191 if dup2_fd(stdin_read_fd, 0) < 0: 192 child_exit(126) 193 if dup2_fd(stdout_write_fd, 1) < 0: 194 child_exit(126) 195 if dup2_fd(stderr_write_fd, 2) < 0: 196 child_exit(126) 197 close_fd(stdin_read_fd) 198 close_fd(stdin_write_fd) 199 close_fd(stdout_read_fd) 200 close_fd(stdout_write_fd) 201 close_fd(stderr_read_fd) 202 close_fd(stderr_write_fd) 203 if cwd != "": 204 if Int(external_call["chdir", c_int](cwd_ptr)) != 0: 205 child_exit(126) 206 _ = set_alarm(MEASUREMENT_CHILD_ALARM_SECONDS) 207 _ = external_call["execvp", c_int](command_ptr, argv_ptr) 208 child_exit(127) 209 210 close_fd(stdin_read_fd) 211 close_fd(stdin_write_fd) 212 close_fd(stdout_write_fd) 213 close_fd(stderr_write_fd) 214 215 var stdout = List[UInt8]() 216 var stderr_bytes = List[UInt8]() 217 var stdout_eof = False 218 var stderr_eof = False 219 var read_reason = "" 220 var start = now_ms() 221 while not (stdout_eof and stderr_eof): 222 var elapsed = now_ms() - start 223 if elapsed >= deadline_ms: 224 read_reason = "sample_deadline_expired" 225 break 226 var slice_ms = min(MEASUREMENT_POLL_SLICE_MS * 2, deadline_ms - elapsed) 227 if slice_ms < 1: 228 slice_ms = 1 229 var pr = poll_three( 230 -1, 231 0, 232 stdout_read_fd, 233 0 if stdout_eof else POLLIN, 234 stderr_read_fd, 235 0 if stderr_eof else POLLIN, 236 slice_ms, 237 ) 238 if pr.count < 0: 239 read_reason = "sample_poll_failed" 240 break 241 if pr.count == 0: 242 continue 243 var remaining = deadline_ms - (now_ms() - start) 244 if not stdout_eof: 245 var d = drain_capture(stdout_read_fd, stdout, pr.r1, remaining) 246 stdout_eof = d.eof 247 if d.reason != "": 248 read_reason = "stdout_" + d.reason 249 break 250 if not stderr_eof: 251 var d = drain_capture( 252 stderr_read_fd, stderr_bytes, pr.r2, remaining 253 ) 254 stderr_eof = d.eof 255 if d.reason != "": 256 read_reason = "stderr_" + d.reason 257 break 258 259 close_fd(stdout_read_fd) 260 close_fd(stderr_read_fd) 261 if read_reason != "": 262 var term = terminate_owned(pid, TERMINATION_GRACE_MS) 263 guard.retain( 264 pid, -1, "measurement sample command " + command + " " + read_reason 265 ) 266 if term.cleanup_proved(): 267 guard.resolve_pid(pid) 268 raise Error( 269 "measurement: sample command " + command + " " + read_reason 270 ) 271 272 var remaining = deadline_ms - (now_ms() - start) 273 if remaining <= 0: 274 # A bounded sampling phase never converts an expired budget into a new 275 # successful interval: the exact owned child is reaped or retained. 276 var expired = terminate_owned(pid, TERMINATION_GRACE_MS) 277 guard.retain( 278 pid, -1, "measurement sample command " + command + " deadline" 279 ) 280 if expired.cleanup_proved(): 281 guard.resolve_pid(pid) 282 raise Error( 283 "measurement: sample command " 284 + command 285 + " sample_deadline_expired" 286 ) 287 var st = wait_bounded(pid, remaining) 288 if not st.cleanup_proved(): 289 var term = terminate_owned(pid, TERMINATION_GRACE_MS) 290 guard.retain( 291 pid, -1, "measurement sample command " + command + " timeout" 292 ) 293 if term.cleanup_proved(): 294 guard.resolve_pid(pid) 295 raise Error("measurement: sample command " + command + " timeout") 296 guard.resolve_pid(pid) 297 298 return CommandOutput( 299 exit_code=st.exit_code if st.exited else -1, 300 signal=st.signal, 301 stdout=bytes_to_text(stdout), 302 stderr=bytes_to_text(stderr_bytes), 303 ) 304 305 306 @fieldwise_init 307 struct DrainResult(Movable): 308 var eof: Bool 309 var reason: String 310 311 312 def drain_capture( 313 fd: Int, mut out: List[UInt8], revents: Int, deadline_ms: Int 314 ) -> DrainResult: 315 """Drain one captured descriptor with a bounded cause.""" 316 if (revents & (POLLIN | POLLHUP | POLLERR | POLLNVAL)) == 0: 317 return DrainResult(False, "") 318 var buf = InlineArray[Byte, 4096](fill=0) 319 var budget = deadline_ms 320 if budget < 1: 321 budget = 1 322 var n = read_fd(fd, buf.unsafe_ptr(), 4096, budget) 323 if n == IO_DEADLINE_EXPIRED: 324 return DrainResult(True, "read_deadline_expired") 325 if n < 0: 326 return DrainResult(True, "read_error") 327 if n == 0: 328 return DrainResult(True, "") 329 if len(out) + n > MEASUREMENT_MAX_STDERR_BYTES: 330 return DrainResult(True, "capture_overflow") 331 for index in range(n): 332 out.append(UInt8(Int(buf[index]))) 333 return DrainResult(False, "") 334 335 336 def bytes_to_text(bytes: List[UInt8]) raises -> String: 337 if len(bytes) == 0: 338 return "" 339 return String(from_utf8=Span(ptr=bytes.unsafe_ptr(), length=len(bytes))) 340 341 342 # ── Identity ──────────────────────────────────────────────────────────────── 343 344 345 def run_capture_simple( 346 command: String, var args: List[String], mut guard: CleanupGuard 347 ) raises -> CommandOutput: 348 return run_capture(command, args^, MEASUREMENT_SAMPLE_DEADLINE_MS, guard) 349 350 351 def file_sha256(path: String, mut guard: CleanupGuard) raises -> String: 352 """Exact sha256 of ``path`` via a bounded ``shasum``/``sha256sum`` call.""" 353 var shasum_args = List[String]() 354 shasum_args.append("-a") 355 shasum_args.append("256") 356 shasum_args.append(path) 357 var out = run_capture_simple("shasum", shasum_args^, guard) 358 if out.exit_code != 0 or out.stdout.strip().byte_length() < 64: 359 var sum_args = List[String]() 360 sum_args.append(path) 361 var alt = run_capture( 362 "sha256sum", sum_args^, MEASUREMENT_SAMPLE_DEADLINE_MS, guard 363 ) 364 if alt.exit_code != 0 or alt.stdout.strip().byte_length() < 64: 365 raise Error("measurement: sha256 unavailable for " + path) 366 out = alt^ 367 var text = String(out.stdout.strip()) 368 var digest = String(text[byte=0:64]) 369 _require_hex(digest, 64, "sha256 for " + path) 370 return digest 371 372 373 def _require_hex(digest: String, expected_len: Int, label: String) raises: 374 if digest.byte_length() != expected_len: 375 raise Error( 376 "measurement: " 377 + label 378 + " is not a " 379 + String(expected_len) 380 + "-character digest" 381 ) 382 for byte in digest.as_bytes(): 383 var b = Int(byte) 384 var is_digit = b >= 48 and b <= 57 385 var is_hex = (b >= 97 and b <= 102) or (b >= 65 and b <= 70) 386 if not is_digit and not is_hex: 387 raise Error("measurement: " + label + " is not hexadecimal") 388 389 390 def host_platform(mut guard: CleanupGuard) raises -> String: 391 var args = List[String]() 392 args.append("-s") 393 args.append("-m") 394 var out = run_capture_simple("uname", args^, guard) 395 if out.exit_code != 0: 396 raise Error("measurement: host platform unavailable") 397 return String(out.stdout.strip()) 398 399 400 def toolchain_version(mut guard: CleanupGuard) raises -> String: 401 var args = List[String]() 402 args.append("--version") 403 var out = run_capture_simple("mojo", args^, guard) 404 if out.exit_code != 0: 405 raise Error("measurement: toolchain version unavailable") 406 var text = String(out.stdout.strip()) 407 if text == "": 408 text = String(out.stderr.strip()) 409 if text.byte_length() == 0: 410 raise Error("measurement: toolchain version output empty") 411 var newline = text.find("\n") 412 if newline >= 0: 413 return String(text[byte=0:newline]) 414 return text^ 415 416 417 @fieldwise_init 418 struct SourceTreeIdentity(Movable): 419 """Verified source/tree identity for the measured build inputs. 420 421 ``revision``/``tree`` are the exact commit and tree, ``manifest_sha256`` is 422 a deterministic content digest of the tracked ``src`` build inputs and 423 ``dirty_status`` is the exact ``git status --porcelain`` output for the 424 measured paths (``""`` proves a clean tree). A mismatch between capture, 425 build and run is a bounded drift failure, never a silent success. 426 """ 427 428 var source_root: String 429 var revision: String 430 var tree: String 431 var manifest_sha256: String 432 var dirty_status: String 433 434 def describe(self) -> String: 435 return ( 436 "revision=" 437 + self.revision 438 + " tree=" 439 + self.tree 440 + " manifest_sha256=" 441 + self.manifest_sha256 442 + " tree_state=" 443 + (self.dirty_status if self.dirty_status == "" else "dirty") 444 ) 445 446 def same_as(self, other: SourceTreeIdentity) -> Bool: 447 return ( 448 self.revision == other.revision 449 and self.tree == other.tree 450 and self.manifest_sha256 == other.manifest_sha256 451 and self.dirty_status == other.dirty_status 452 ) 453 454 455 def git_capture( 456 source_root: String, var args: List[String], mut guard: CleanupGuard 457 ) raises -> CommandOutput: 458 return run_capture( 459 "git", args^, MEASUREMENT_SAMPLE_DEADLINE_MS, guard, source_root 460 ) 461 462 463 def source_revision( 464 source_root: String, mut guard: CleanupGuard 465 ) raises -> String: 466 """Exact capsule source revision; a measurement must not guess it.""" 467 var args = List[String]() 468 args.append("rev-parse") 469 args.append("HEAD") 470 var out = git_capture(source_root, args^, guard) 471 var text = String(out.stdout.strip()) 472 if out.exit_code != 0 or text.byte_length() != 40: 473 raise Error( 474 "measurement: source revision unavailable in " + source_root 475 ) 476 _require_hex(text, 40, "source revision") 477 return text^ 478 479 480 def source_tree_id( 481 source_root: String, mut guard: CleanupGuard 482 ) raises -> String: 483 """Exact commit tree object id; distinguishes a content-identical tree.""" 484 var args = List[String]() 485 args.append("rev-parse") 486 args.append("HEAD^{tree}") 487 var out = git_capture(source_root, args^, guard) 488 var text = String(out.stdout.strip()) 489 if out.exit_code != 0 or text.byte_length() != 40: 490 raise Error("measurement: source tree id unavailable in " + source_root) 491 _require_hex(text, 40, "source tree id") 492 return text^ 493 494 495 # ── Fail-closed checked digest pipelines (ADR-0020 MC02) ──────────────────── 496 497 498 def _path_basename(path: String) -> String: 499 """Last path component, used to keep manifest text path-independent.""" 500 var slash = -1 501 var index = 0 502 for byte in path.as_bytes(): 503 if Int(byte) == 47: 504 slash = index 505 index += 1 506 if slash < 0: 507 return String(path) 508 return String(path[byte = slash + 1 :]) 509 510 511 def _checked_file_digests( 512 var files: List[String], 513 mut guard: CleanupGuard, 514 primary_hasher: String = "shasum", 515 fallback_hasher: String = "sha256sum", 516 ) raises -> List[String]: 517 """Per-file sha256 of explicit argv paths, fail-closed. 518 519 Every path is passed as argv data, so apostrophes, spaces and other shell 520 characters are handled literally and no interpolation is performed. A 521 missing or unreadable input, an empty input list, a failed hasher stage or 522 an unavailable hasher raises instead of yielding a valid empty digest. This 523 replaces the former status-masking shell pipeline whose final stage could 524 succeed on empty input and report the empty-input digest as success. 525 526 ``primary_hasher``/``fallback_hasher`` name the two checked stages. The 527 defaults are the supported host hashers; the parameters exist only as a 528 bounded test seam so the failed-hasher-stage and fallback controls can run 529 against present, valid regular-file inputs without altering live tools or 530 host settings (ADR-0021 MP01). 531 """ 532 if len(files) == 0: 533 raise Error("measurement: refusing to digest an empty input list") 534 var shasum_args = List[String]() 535 shasum_args.append("-a") 536 shasum_args.append("256") 537 for index in range(len(files)): 538 shasum_args.append(files[index]) 539 var out = run_capture( 540 primary_hasher, shasum_args^, MEASUREMENT_SAMPLE_DEADLINE_MS, guard 541 ) 542 if out.exit_code != 0: 543 var sum_args = List[String]() 544 for index in range(len(files)): 545 sum_args.append(files[index]) 546 out = run_capture( 547 fallback_hasher, sum_args^, MEASUREMENT_SAMPLE_DEADLINE_MS, guard 548 ) 549 if out.exit_code != 0: 550 raise Error( 551 "measurement: sha256 unavailable for " 552 + String(len(files)) 553 + " input(s) (" 554 + out.describe() 555 + ")" 556 ) 557 var digests = List[String]() 558 for line in out.stdout.split("\n"): 559 var entry = String(line).strip() 560 if entry.byte_length() == 0: 561 continue 562 if entry.byte_length() < 64: 563 raise Error("measurement: truncated sha256 output line") 564 var digest = String(entry[byte=0:64]) 565 _require_hex(digest, 64, "sha256 input") 566 digests.append(digest) 567 if len(digests) != len(files): 568 raise Error( 569 "measurement: sha256 digest count mismatch (" 570 + String(len(digests)) 571 + " of " 572 + String(len(files)) 573 + ")" 574 ) 575 return digests^ 576 577 578 def sha256_text( 579 text: String, label: String, mut guard: CleanupGuard 580 ) raises -> String: 581 """sha256 of exact text via a private temp file and a checked argv call.""" 582 with SafeTempDir() as temp_dir: 583 var path = temp_dir + "/measurement-manifest.txt" 584 try: 585 Path(path).write_text(text) 586 except: 587 raise Error("measurement: " + label + " staging failed") 588 var files = List[String]() 589 files.append(path) 590 var digests = _checked_file_digests(files^, guard) 591 _require_hex(digests[0], 64, label) 592 return digests[0] 593 594 595 def sha256_file_set( 596 label: String, 597 var files: List[String], 598 mut guard: CleanupGuard, 599 primary_hasher: String = "shasum", 600 fallback_hasher: String = "sha256sum", 601 ) raises -> String: 602 """Deterministic, path-independent digest of an explicit file set. 603 604 Each declared file's exact content digest is checked first, then the 605 manifest text ``<basename> <digest>`` (in declared order) is hashed, so the 606 result depends only on the declared files' content and names, never on the 607 checkout location. ``primary_hasher``/``fallback_hasher`` are the bounded 608 test seam documented on ``_checked_file_digests``. 609 """ 610 var names = List[String]() 611 for index in range(len(files)): 612 names.append(_path_basename(files[index])) 613 var digests = _checked_file_digests( 614 files^, guard, primary_hasher, fallback_hasher 615 ) 616 var canonical = "" 617 for index in range(len(names)): 618 canonical += names[index] + " " + digests[index] + "\n" 619 return sha256_text(canonical, label, guard) 620 621 622 def measurement_tooling_files(source_root: String) -> List[String]: 623 """Bounded closure of the test-only tooling that produced the evidence. 624 625 The set is enumerated explicitly — never a workspace scan — and covers the 626 measurement helper/runner, the contract test and the shared helpers those 627 import (lifecycle, stdio, temp, fixtures). A change to an imported helper 628 therefore changes the recorded tooling identity (ADR-0020 MC02). Product 629 ``src`` inputs are bound separately by the clean source content manifest. 630 """ 631 var files = List[String]() 632 files.append(source_root + "/tests/measurement_process_helper.mojo") 633 files.append(source_root + "/tests/measurement_runner.mojo") 634 files.append(source_root + "/tests/test_measurement_contract.mojo") 635 files.append(source_root + "/tests/parent_lifecycle.mojo") 636 files.append(source_root + "/tests/safe_tempdir.mojo") 637 files.append(source_root + "/tests/stdio_process_helper.mojo") 638 files.append(source_root + "/tests/max_local_process_helper.mojo") 639 files.append(source_root + "/tests/strict_fixture.mojo") 640 return files^ 641 642 643 def source_manifest_sha256( 644 source_root: String, mut guard: CleanupGuard 645 ) raises -> String: 646 """Deterministic content digest of the tracked ``src`` build inputs. 647 648 The digest covers the exact ``git ls-files -s -- src`` index listing 649 (mode/blob/path) — the inputs of ``mojo build -I src src/main.mojo``. The 650 git stage and the digest stage are both checked, so a missing repository, a 651 failed git command or an empty listing is an error, never a valid empty 652 digest. It never reads or hashes secrets or arbitrary workspace files. 653 """ 654 var args = List[String]() 655 args.append("ls-files") 656 args.append("-s") 657 args.append("--") 658 args.append("src") 659 var out = git_capture(source_root, args^, guard) 660 if out.exit_code != 0: 661 raise Error( 662 "measurement: source content manifest unavailable in " + source_root 663 ) 664 var listing = String(out.stdout) 665 if listing.strip().byte_length() == 0: 666 raise Error( 667 "measurement: source content manifest is empty in " + source_root 668 ) 669 return sha256_text(listing, "source content manifest", guard) 670 671 672 def source_dirty_status( 673 source_root: String, mut guard: CleanupGuard 674 ) raises -> String: 675 """Exact porcelain status for the measured build inputs and toolchain.""" 676 var args = List[String]() 677 args.append("status") 678 args.append("--porcelain") 679 args.append("--") 680 args.append("src") 681 args.append("pixi.toml") 682 args.append("pixi.lock") 683 var out = git_capture(source_root, args^, guard) 684 if out.exit_code != 0: 685 raise Error( 686 "measurement: source tree status unavailable in " + source_root 687 ) 688 return String(out.stdout.strip()) 689 690 691 def tooling_manifest_sha256( 692 source_root: String, 693 mut guard: CleanupGuard, 694 primary_hasher: String = "shasum", 695 fallback_hasher: String = "sha256sum", 696 ) raises -> String: 697 """Content digest of the measurement tooling that produced the evidence. 698 699 The tooling is test-only source outside the product build inputs, so it is 700 not covered by the clean-tree build binding. Recording its exact content 701 digest ties the emitted evidence to the reviewed tooling revision and its 702 imported helper closure without requiring the working tree to be committed 703 at capture time. A missing input or failed hasher stage is a bounded error, 704 never a valid empty digest. The optional hasher parameters are the bounded 705 test seam described on ``_checked_file_digests``. 706 """ 707 var files = measurement_tooling_files(source_root) 708 return sha256_file_set( 709 "tooling manifest", files^, guard, primary_hasher, fallback_hasher 710 ) 711 712 713 def source_identity( 714 source_root: String, mut guard: CleanupGuard 715 ) raises -> SourceTreeIdentity: 716 return SourceTreeIdentity( 717 source_root=source_root, 718 revision=source_revision(source_root, guard), 719 tree=source_tree_id(source_root, guard), 720 manifest_sha256=source_manifest_sha256(source_root, guard), 721 dirty_status=source_dirty_status(source_root, guard), 722 ) 723 724 725 def require_clean_source(identity: SourceTreeIdentity, phase: String) raises: 726 """Reject a dirty measured build input at capture, build or run time.""" 727 if identity.dirty_status != "": 728 raise Error( 729 "measurement: measured source tree is dirty at " 730 + phase 731 + " (" 732 + identity.dirty_status 733 + ")" 734 ) 735 736 737 def require_identity_drift_free( 738 before: SourceTreeIdentity, after: SourceTreeIdentity, phase: String 739 ) raises: 740 if not before.same_as(after): 741 raise Error( 742 "measurement: source drifted during " 743 + phase 744 + " (before " 745 + before.describe() 746 + " after " 747 + after.describe() 748 + ")" 749 ) 750 751 752 def working_directory(mut guard: CleanupGuard) raises -> String: 753 var args = List[String]() 754 var out = run_capture_simple("pwd", args^, guard) 755 if out.exit_code != 0: 756 raise Error("measurement: working directory unavailable") 757 return String(out.stdout.strip()) 758 759 760 def verified_environment_profile() -> String: 761 """The environment actually inherited by the measured child, read back. 762 763 Recorded from the live process environment rather than asserted, so the 764 reproduced profile is truthful: the measured child inherits exactly these 765 values through ``execvp``. The profile is deliberately bounded to the two 766 governed, secret-free HYF_PATHS variables. 767 """ 768 var profile = std.os.getenv("HYF_PATHS_PROFILE") 769 var root = std.os.getenv("HYF_PATHS_REPO_LOCAL_ROOT") 770 return ( 771 "HYF_PATHS_PROFILE=" 772 + (profile if profile != "" else "<unset>") 773 + " HYF_PATHS_REPO_LOCAL_ROOT=" 774 + (root if root != "" else "<unset>") 775 ) 776 777 778 @fieldwise_init 779 struct MeasurementIdentity(Movable): 780 """Exact source/binary/toolchain/host identity and launch profile.""" 781 782 var source_root: String 783 var source_revision: String 784 var source_tree: String 785 var source_manifest_sha256: String 786 var source_tree_state: String 787 var tooling_manifest_sha256: String 788 var binding: String 789 var cwd: String 790 var binary_path: String 791 var binary_sha256: String 792 var pixi_toml_sha256: String 793 var pixi_lock_sha256: String 794 var toolchain_version: String 795 var host_platform: String 796 var argv_profile: String 797 var env_profile: String 798 799 def describe(self) -> String: 800 return ( 801 "binding=" 802 + self.binding 803 + " source_revision=" 804 + self.source_revision 805 + " source_tree=" 806 + self.source_tree 807 + " source_manifest_sha256=" 808 + self.source_manifest_sha256 809 + " source_tree_state=" 810 + self.source_tree_state 811 + " tooling_manifest_sha256=" 812 + self.tooling_manifest_sha256 813 + " cwd=" 814 + self.cwd 815 + " binary_sha256=" 816 + self.binary_sha256 817 + " pixi_toml_sha256=" 818 + self.pixi_toml_sha256 819 + " pixi_lock_sha256=" 820 + self.pixi_lock_sha256 821 + " toolchain=" 822 + self.toolchain_version 823 + " host=" 824 + self.host_platform 825 + " argv=" 826 + self.argv_profile 827 + " env=" 828 + self.env_profile 829 ) 830 831 832 def measurement_identity( 833 source_root: String, 834 binary_path: String, 835 binary_sha256: String, 836 binding: String, 837 argv_profile: String, 838 mut guard: CleanupGuard, 839 ) raises -> MeasurementIdentity: 840 var observed = source_identity(source_root, guard) 841 return MeasurementIdentity( 842 source_root=source_root, 843 source_revision=observed.revision, 844 source_tree=observed.tree, 845 source_manifest_sha256=observed.manifest_sha256, 846 source_tree_state=("dirty" if observed.dirty_status != "" else "clean"), 847 tooling_manifest_sha256=tooling_manifest_sha256(source_root, guard), 848 binding=binding, 849 cwd=working_directory(guard), 850 binary_path=binary_path, 851 binary_sha256=binary_sha256, 852 pixi_toml_sha256=file_sha256(source_root + "/pixi.toml", guard), 853 pixi_lock_sha256=file_sha256(source_root + "/pixi.lock", guard), 854 toolchain_version=toolchain_version(guard), 855 host_platform=host_platform(guard), 856 argv_profile=argv_profile, 857 env_profile=verified_environment_profile(), 858 ) 859 860 861 @fieldwise_init 862 struct BuiltProduct(Movable): 863 """A product binary bound to the clean source/tree it was built from. 864 865 The build is a separate, explicitly bounded phase: the pre-build source 866 identity is captured, the binary is built, the binary digest is recorded 867 and the post-build source identity must be byte-for-byte the same. A 868 later run re-verifies the same identity and binary before spawning. 869 """ 870 871 var binary_path: String 872 var binary_sha256: String 873 var source: SourceTreeIdentity 874 var build_ms: Int 875 876 877 def build_product_binary( 878 source_root: String, temp_dir: String, mut guard: CleanupGuard 879 ) raises -> BuiltProduct: 880 """Build the existing product entry point once, outside measured frames. 881 882 The measurement itself never compiles or starts a process per frame. The 883 build input tree must be clean and must not drift across the build. 884 """ 885 var before = source_identity(source_root, guard) 886 require_clean_source(before, "build capture") 887 var output = temp_dir + "/hyfd" 888 var args = List[String]() 889 args.append("build") 890 args.append("-I") 891 args.append("src") 892 args.append("src/main.mojo") 893 args.append("-o") 894 args.append(output) 895 var start = now_ms() 896 var out = run_capture( 897 "mojo", args^, MEASUREMENT_BUILD_DEADLINE_MS, guard, source_root 898 ) 899 var build_ms = now_ms() - start 900 if out.exit_code != 0: 901 raise Error( 902 "measurement: product build failed (" + out.describe() + ")" 903 ) 904 var binary_sha256 = file_sha256(output, guard) 905 var after = source_identity(source_root, guard) 906 require_identity_drift_free(before, after, "build") 907 return BuiltProduct( 908 binary_path=output, 909 binary_sha256=binary_sha256, 910 source=after^, 911 build_ms=build_ms, 912 ) 913 914 915 # ── Persistent measured process ───────────────────────────────────────────── 916 917 918 @fieldwise_init 919 struct MeasurementPoll(Movable): 920 """One bounded ``poll(2)`` result with an explicit EINTR classification.""" 921 922 var count: Int 923 var r0: Int 924 var r1: Int 925 var r2: Int 926 var interrupted: Bool 927 928 929 def measurement_poll( 930 fd0: Int, 931 events0: Int, 932 fd1: Int, 933 events1: Int, 934 fd2: Int, 935 events2: Int, 936 timeout_ms: Int, 937 fault_eintr_count: Int = 0, 938 ) -> MeasurementPoll: 939 """Poll three descriptors, classifying a real ``EINTR`` as retryable. 940 941 ``fault_eintr_count`` is a bounded test-only seam: any nonzero value makes 942 every attempt report ``EINTR`` so the deadline-bounded retry branch is 943 deterministically executable without changing host signal state. 944 """ 945 if fault_eintr_count != 0: 946 return MeasurementPoll(-1, 0, 0, 0, True) 947 var cell = InlineArray[Int32, 12](fill=0) 948 cell[0] = Int32(fd0) 949 cell[1] = Int32(events0) 950 cell[2] = Int32(fd1) 951 cell[3] = Int32(events1) 952 cell[4] = Int32(fd2) 953 cell[5] = Int32(events2) 954 var n = Int( 955 external_call["poll", c_int]( 956 cell.unsafe_ptr(), c_uint(3), c_int(timeout_ms) 957 ) 958 ) 959 if n < 0: 960 if get_errno() == ErrNo.EINTR: 961 return MeasurementPoll(-1, 0, 0, 0, True) 962 return MeasurementPoll(-1, 0, 0, 0, False) 963 if n == 0: 964 return MeasurementPoll(0, 0, 0, 0, False) 965 return MeasurementPoll( 966 n, 967 (Int(cell[1]) >> 16) & 0xFFFF, 968 (Int(cell[3]) >> 16) & 0xFFFF, 969 (Int(cell[5]) >> 16) & 0xFFFF, 970 False, 971 ) 972 973 974 def measurement_poll_retry( 975 fd0: Int, 976 events0: Int, 977 fd1: Int, 978 events1: Int, 979 fd2: Int, 980 events2: Int, 981 max_wait_ms: Int, 982 fault_eintr_count: Int = 0, 983 ) -> MeasurementPoll: 984 """Poll with a deadline-bounded retry of a real ``EINTR``. 985 986 The retry can never outlive ``max_wait_ms``: once the budget is exhausted 987 the last interrupted result is returned to the deadline owner instead of 988 spinning. ``fault_eintr_count`` is the bounded test-only seam. 989 """ 990 var start = now_ms() 991 var slice = min(MEASUREMENT_POLL_SLICE_MS, max_wait_ms) 992 if slice < 1: 993 slice = 1 994 while True: 995 var pr = measurement_poll( 996 fd0, events0, fd1, events1, fd2, events2, slice, fault_eintr_count 997 ) 998 if not pr.interrupted: 999 return pr^ 1000 if now_ms() - start >= max_wait_ms: 1001 return pr^ 1002 sleep_ms(1) 1003 1004 1005 @fieldwise_init 1006 struct MeasurementProcess(Movable): 1007 """Exactly one persistent HYF stdio process owned by the measurement. 1008 1009 stdout/stderr are consumed concurrently with request writes and exit 1010 observation and are drained to EOF. Every owned descriptor carries an 1011 explicit close-once flag: EOF is a stream property and never implies that 1012 the parent descriptor was closed, so no later cleanup can target a reused 1013 descriptor number. 1014 """ 1015 1016 var pid: Int 1017 var stdin_fd: Int 1018 var stdout_fd: Int 1019 var stderr_fd: Int 1020 var stdout_pending: List[UInt8] 1021 var stderr_bytes: List[UInt8] 1022 var stdout_eof: Bool 1023 var stderr_eof: Bool 1024 var stdin_closed: Bool 1025 var stdout_closed: Bool 1026 var stderr_closed: Bool 1027 var reaped: Bool 1028 var status: ProcessStatus 1029 var spawn_ms: Int 1030 var deadline_at_ms: Int 1031 var stdout_total_bytes: Int 1032 var faults: MeasurementFaults 1033 var guard: UnsafePointer[CleanupGuard, MutAnyOrigin] 1034 1035 def work_remaining_ms(self) -> Int: 1036 """Remaining part of the one spawn-relative measurement work budget. 1037 1038 A nonpositive result must never be turned into a new successful 1039 interval; callers fail and use only the separate bounded cleanup 1040 allowance. 1041 """ 1042 return self.deadline_at_ms - now_ms() 1043 1044 def close_stdin(mut self): 1045 if self.stdin_closed: 1046 return 1047 close_fd(self.stdin_fd) 1048 self.stdin_fd = -1 1049 self.stdin_closed = True 1050 1051 def close_stdout(mut self): 1052 if self.stdout_closed: 1053 return 1054 # A descriptor retained by the guard for recovery is closed by the 1055 # guard exactly once, so the handle can never close a reused number. 1056 if not self.guard[].close_retained_fd(self.stdout_fd): 1057 close_fd(self.stdout_fd) 1058 self.stdout_closed = True 1059 1060 def close_stderr(mut self): 1061 if self.stderr_closed: 1062 return 1063 close_fd(self.stderr_fd) 1064 self.stderr_closed = True 1065 1066 def stdout_pending_len(self) -> Int: 1067 return len(self.stdout_pending) 1068 1069 def stderr_len(self) -> Int: 1070 return len(self.stderr_bytes) 1071 1072 def drain_stdout(mut self, revents: Int, budget_ms: Int) raises: 1073 """Drain ready stdout bytes with a bounded size cap and cause.""" 1074 if (revents & (POLLIN | POLLHUP | POLLERR | POLLNVAL)) == 0: 1075 return 1076 if budget_ms < 1: 1077 raise Error( 1078 "measurement: work deadline expired during stdout drain" 1079 " (work_deadline_expired)" 1080 ) 1081 var buf = InlineArray[Byte, 4096](fill=0) 1082 var n = read_fd(self.stdout_fd, buf.unsafe_ptr(), 4096, budget_ms) 1083 if n == IO_DEADLINE_EXPIRED: 1084 raise Error( 1085 "measurement: stdout read deadline expired" 1086 " (work_deadline_expired)" 1087 ) 1088 if n < 0: 1089 raise Error("measurement: stdout read_error") 1090 if n == 0: 1091 self.stdout_eof = True 1092 return 1093 self.stdout_total_bytes += n 1094 if len(self.stdout_pending) + n > MEASUREMENT_FRAME_BYTES: 1095 raise Error("measurement: stdout_overflow") 1096 for index in range(n): 1097 self.stdout_pending.append(UInt8(Int(buf[index]))) 1098 1099 def drain_stderr(mut self, revents: Int, budget_ms: Int) raises: 1100 """Drain ready stderr bytes; errors and overflow fail explicitly.""" 1101 if (revents & (POLLIN | POLLHUP | POLLERR | POLLNVAL)) == 0: 1102 return 1103 if self.faults.stderr_read_errors > 0: 1104 self.faults.stderr_read_errors -= 1 1105 raise Error("measurement: stderr read_error") 1106 if budget_ms < 1: 1107 raise Error( 1108 "measurement: work deadline expired during stderr drain" 1109 " (work_deadline_expired)" 1110 ) 1111 var buf = InlineArray[Byte, 4096](fill=0) 1112 var n = read_fd(self.stderr_fd, buf.unsafe_ptr(), 4096, budget_ms) 1113 if n == IO_DEADLINE_EXPIRED: 1114 raise Error( 1115 "measurement: stderr read deadline expired" 1116 " (work_deadline_expired)" 1117 ) 1118 if n < 0: 1119 raise Error("measurement: stderr read_error") 1120 if n == 0: 1121 self.stderr_eof = True 1122 return 1123 if len(self.stderr_bytes) + n > MEASUREMENT_MAX_STDERR_BYTES: 1124 raise Error("measurement: stderr_overflow") 1125 for index in range(n): 1126 self.stderr_bytes.append(UInt8(Int(buf[index]))) 1127 1128 def poll_and_drain( 1129 mut self, 1130 want_stdin_write: Bool, 1131 max_wait_ms: Int, 1132 fault_eintr_count: Int = 0, 1133 ) raises -> Bool: 1134 """Poll stdin/stdout/stderr and drain whatever is ready. 1135 1136 Returns True when the stdin write end is writable. All waits are sliced 1137 by ``MEASUREMENT_POLL_SLICE_MS`` inside the caller's remaining budget; 1138 a real ``EINTR`` is retried until the budget expires. 1139 """ 1140 if max_wait_ms <= 0: 1141 raise Error( 1142 "measurement: work deadline expired (work_deadline_expired)" 1143 ) 1144 var slice = min(MEASUREMENT_POLL_SLICE_MS, max_wait_ms) 1145 var fd0 = self.stdin_fd if ( 1146 want_stdin_write and not self.stdin_closed 1147 ) else -1 1148 var fd1 = self.stdout_fd if ( 1149 not self.stdout_eof and not self.stdout_closed 1150 ) else -1 1151 var fd2 = self.stderr_fd if ( 1152 not self.stderr_eof and not self.stderr_closed 1153 ) else -1 1154 if fd0 < 0 and fd1 < 0 and fd2 < 0: 1155 sleep_ms(slice) 1156 return False 1157 var start = now_ms() 1158 var pr = measurement_poll_retry( 1159 fd0, 1160 POLLOUT if fd0 >= 0 else 0, 1161 fd1, 1162 POLLIN if fd1 >= 0 else 0, 1163 fd2, 1164 POLLIN if fd2 >= 0 else 0, 1165 max_wait_ms, 1166 fault_eintr_count, 1167 ) 1168 if pr.count > 0: 1169 var remaining = max_wait_ms - (now_ms() - start) 1170 if fd0 >= 0 and (pr.r0 & (POLLERR | POLLHUP | POLLNVAL)) != 0: 1171 raise Error("measurement: request write pipe closed") 1172 if ( 1173 fd1 >= 0 1174 and (pr.r1 & (POLLIN | POLLHUP | POLLERR | POLLNVAL)) != 0 1175 ): 1176 self.drain_stdout(pr.r1, remaining) 1177 if ( 1178 fd2 >= 0 1179 and (pr.r2 & (POLLIN | POLLHUP | POLLERR | POLLNVAL)) != 0 1180 ): 1181 self.drain_stderr(pr.r2, remaining) 1182 return fd0 >= 0 and (pr.r0 & POLLOUT) != 0 1183 return False 1184 1185 def send_frame(mut self, frame: String, index: Int) raises: 1186 """Write one newline-delimited request frame within the work budget. 1187 1188 stdout/stderr are drained concurrently while the request is written, so 1189 a child that emits output during the write can never deadlock the 1190 parent. There is no fixed idle wait in the request latency. 1191 """ 1192 if self.stdin_closed: 1193 raise Error("measurement: request write after stdin closed") 1194 var payload = frame + "\n" 1195 var sent = 0 1196 var total = payload.byte_length() 1197 while sent < total: 1198 var remaining = self.work_remaining_ms() 1199 if remaining <= 0: 1200 raise Error( 1201 "measurement: work deadline expired during request write" 1202 " (work_deadline_expired)" 1203 ) 1204 var writable = self.poll_and_drain(True, remaining) 1205 if self.stdin_closed: 1206 raise Error("measurement: request write pipe closed") 1207 if not writable: 1208 continue 1209 var remaining_write = self.work_remaining_ms() 1210 if remaining_write <= 0: 1211 raise Error( 1212 "measurement: work deadline expired during request write" 1213 " (work_deadline_expired)" 1214 ) 1215 var cw = write_fd_chunk( 1216 self.stdin_fd, payload, sent, remaining_write 1217 ) 1218 if cw.reason == "write_deadline_expired": 1219 raise Error("measurement: request write deadline expired") 1220 if cw.reason != "": 1221 raise Error( 1222 "measurement: request write failed (" + cw.reason + ")" 1223 ) 1224 if cw.written <= 0: 1225 continue 1226 sent += cw.written 1227 1228 def read_response(mut self, index: Int) raises -> String: 1229 """Read one complete newline-terminated response within the budget.""" 1230 while True: 1231 var found = _newline_index(self.stdout_pending) 1232 if found >= 0: 1233 if found > MEASUREMENT_FRAME_BYTES: 1234 raise Error("measurement: response frame overflow") 1235 var line_bytes = _list_prefix(self.stdout_pending, found) 1236 self.stdout_pending = _list_suffix( 1237 self.stdout_pending, found + 1 1238 ) 1239 return _decode_frame(line_bytes^) 1240 if self.stdout_eof: 1241 if len(self.stdout_pending) == 0: 1242 raise Error( 1243 "measurement: response stream reached EOF with no" 1244 " frame (early_eof)" 1245 ) 1246 raise Error( 1247 "measurement: response frame was not newline-terminated" 1248 ) 1249 var remaining = self.work_remaining_ms() 1250 if remaining <= 0: 1251 raise Error( 1252 "measurement: work deadline expired while awaiting" 1253 " response (work_deadline_expired)" 1254 ) 1255 _ = self.poll_and_drain(False, remaining) 1256 1257 def assert_no_surplus(mut self, index: Int) raises: 1258 """Exactly one response per request: surplus bytes are a failure.""" 1259 if len(self.stdout_pending) > 0: 1260 raise Error( 1261 "measurement: unexpected trailing stdout (" 1262 + String(len(self.stdout_pending)) 1263 + " bytes) after response " 1264 + String(index) 1265 ) 1266 1267 def finish_expected(mut self, total: Int) raises -> ProcessStatus: 1268 """Close stdin, account the whole stream through EOF and reap. 1269 1270 The remaining stdout stream must be exactly empty after the expected 1271 responses, stderr must reach a clean EOF and the child must exit within 1272 the remaining work budget. The terminal state is cached before any 1273 further wait or signal. 1274 """ 1275 if len(self.stdout_pending) > 0: 1276 raise Error( 1277 "measurement: unexpected trailing stdout (" 1278 + String(len(self.stdout_pending)) 1279 + " bytes) after " 1280 + String(total) 1281 + " expected responses" 1282 ) 1283 self.close_stdin() 1284 while not (self.stdout_eof and self.stderr_eof): 1285 var remaining = self.work_remaining_ms() 1286 if remaining <= 0: 1287 raise Error( 1288 "measurement: work deadline expired during output drain" 1289 " (work_deadline_expired)" 1290 ) 1291 _ = self.poll_and_drain(False, remaining) 1292 if len(self.stdout_pending) > 0: 1293 raise Error( 1294 "measurement: unexpected trailing stdout (" 1295 + String(len(self.stdout_pending)) 1296 + " bytes) after " 1297 + String(total) 1298 + " expected responses" 1299 ) 1300 var remaining_exit = self.work_remaining_ms() 1301 if remaining_exit <= 0: 1302 raise Error( 1303 "measurement: work deadline expired before child exit" 1304 " (work_deadline_expired)" 1305 ) 1306 var st = wait_bounded(self.pid, remaining_exit) 1307 if not st.cleanup_proved(): 1308 self.status = st.copy() 1309 var term = terminate_owned(self.pid, TERMINATION_GRACE_MS) 1310 self.status = term.copy() 1311 if term.cleanup_proved(): 1312 self.reaped = True 1313 self.guard[].resolve_pid(self.pid) 1314 else: 1315 self.guard[].retain( 1316 self.pid, 1317 self.stdout_fd, 1318 "measurement process cleanup unproved " + term.describe(), 1319 ) 1320 raise Error( 1321 "measurement: process did not exit within its budget (" 1322 + term.describe() 1323 + ")" 1324 ) 1325 self.status = st.copy() 1326 self.reaped = True 1327 self.guard[].resolve_pid(self.pid) 1328 self.close_stdout() 1329 self.close_stderr() 1330 return st^ 1331 1332 def cleanup(mut self): 1333 """Non-raising cleanup so a failing assertion still reaps the child.""" 1334 self.close_stdin() 1335 if not self.reaped: 1336 var st: ProcessStatus 1337 if self.faults.cleanup_failures > 0: 1338 self.faults.cleanup_failures -= 1 1339 st = ProcessStatus( 1340 "wait_error", False, -1, 0, -1, "injected_cleanup_failure" 1341 ) 1342 else: 1343 st = terminate_owned(self.pid, TERMINATION_GRACE_MS) 1344 self.status = st.copy() 1345 if st.cleanup_proved(): 1346 self.reaped = True 1347 self.guard[].resolve_pid(self.pid) 1348 else: 1349 self.guard[].retain( 1350 self.pid, 1351 self.stdout_fd, 1352 "measurement process cleanup unproved " + st.describe(), 1353 ) 1354 self.close_stdout() 1355 self.close_stderr() 1356 1357 def stderr_text(self) raises -> String: 1358 return bytes_to_text( 1359 _list_prefix(self.stderr_bytes, len(self.stderr_bytes)) 1360 ) 1361 1362 1363 def _newline_index(bytes: List[UInt8]) -> Int: 1364 for index in range(len(bytes)): 1365 if Int(bytes[index]) == 10: 1366 return index 1367 return -1 1368 1369 1370 def _list_prefix(bytes: List[UInt8], count: Int) -> List[UInt8]: 1371 var out = List[UInt8]() 1372 for index in range(count): 1373 out.append(bytes[index]) 1374 return out^ 1375 1376 1377 def _list_suffix(bytes: List[UInt8], start: Int) -> List[UInt8]: 1378 var out = List[UInt8]() 1379 for index in range(start, len(bytes)): 1380 out.append(bytes[index]) 1381 return out^ 1382 1383 1384 def _decode_frame(var bytes: List[UInt8]) raises -> String: 1385 try: 1386 return bytes_to_text(bytes^) 1387 except: 1388 raise Error("measurement: response frame is not valid UTF-8") 1389 1390 1391 def spawn_measurement_process( 1392 binary_path: String, 1393 var argv: List[String], 1394 deadline_ms: Int, 1395 mut guard: CleanupGuard, 1396 faults: MeasurementFaults = MeasurementFaults(), 1397 ) raises -> MeasurementProcess: 1398 """Fork exactly one persistent process for warmup and measured frames.""" 1399 var parts = List[String]() 1400 parts.append(binary_path) 1401 for index in range(len(argv)): 1402 parts.append(argv[index]) 1403 var argv_c = List[Optional[CStringSlice[ImmutAnyOrigin]]]( 1404 length=len(parts) + 1, fill={} 1405 ) 1406 var elements = parts.unsafe_ptr() 1407 for index in range(len(parts)): 1408 argv_c[index] = rebind[CStringSlice[ImmutAnyOrigin]]( 1409 elements[index].as_c_string_slice() 1410 ) 1411 1412 var pipes = make_three_pipes() 1413 var stdin_read_fd = pipes.stdin_pipe.read_fd 1414 var stdin_write_fd = pipes.stdin_pipe.write_fd 1415 var stdout_read_fd = pipes.stdout_pipe.read_fd 1416 var stdout_write_fd = pipes.stdout_pipe.write_fd 1417 var stderr_read_fd = pipes.stderr_pipe.read_fd 1418 var stderr_write_fd = pipes.stderr_pipe.write_fd 1419 var command_ptr = elements[0].as_c_string_slice().unsafe_ptr() 1420 var argv_ptr = argv_c.unsafe_ptr() 1421 1422 var spawn_ms = now_ms() 1423 var pid = fork_owned_or_close3(pipes) 1424 if pid == 0: 1425 if dup2_fd(stdin_read_fd, 0) < 0: 1426 child_exit(126) 1427 if dup2_fd(stdout_write_fd, 1) < 0: 1428 child_exit(126) 1429 if dup2_fd(stderr_write_fd, 2) < 0: 1430 child_exit(126) 1431 close_fd(stdin_read_fd) 1432 close_fd(stdin_write_fd) 1433 close_fd(stdout_read_fd) 1434 close_fd(stdout_write_fd) 1435 close_fd(stderr_read_fd) 1436 close_fd(stderr_write_fd) 1437 _ = set_alarm(MEASUREMENT_CHILD_ALARM_SECONDS) 1438 _ = external_call["execvp", c_int](command_ptr, argv_ptr) 1439 child_exit(127) 1440 1441 close_fd(stdin_read_fd) 1442 close_fd(stdout_write_fd) 1443 close_fd(stderr_write_fd) 1444 return MeasurementProcess( 1445 pid=pid, 1446 stdin_fd=stdin_write_fd, 1447 stdout_fd=stdout_read_fd, 1448 stderr_fd=stderr_read_fd, 1449 stdout_pending=List[UInt8](), 1450 stderr_bytes=List[UInt8](), 1451 stdout_eof=False, 1452 stderr_eof=False, 1453 stdin_closed=False, 1454 stdout_closed=False, 1455 stderr_closed=False, 1456 reaped=False, 1457 status=ProcessStatus("pending", False, -1, 0, 0, ""), 1458 spawn_ms=spawn_ms, 1459 deadline_at_ms=spawn_ms + deadline_ms, 1460 stdout_total_bytes=0, 1461 faults=faults, 1462 guard=UnsafePointer(to=guard), 1463 ) 1464 1465 1466 # ── Frame validation ──────────────────────────────────────────────────────── 1467 1468 1469 @fieldwise_init 1470 struct FrameVerdict(Movable): 1471 """Bounded per-frame validation result.""" 1472 1473 var ok: Bool 1474 var reason: String 1475 var outcome: String 1476 var declared_max_requests: Int 1477 1478 def describe(self) -> String: 1479 return ( 1480 "ok=" 1481 + String(self.ok) 1482 + " outcome=" 1483 + self.outcome 1484 + " reason=" 1485 + self.reason 1486 ) 1487 1488 1489 def _has_key(value: Value, key: String) -> Bool: 1490 for candidate in value.object_keys(): 1491 if String(candidate) == key: 1492 return True 1493 return False 1494 1495 1496 def _field_string(value: Value, key: String) raises -> String: 1497 if not value.is_object(): 1498 return "" 1499 if not _has_key(value, key): 1500 return "" 1501 var field = value[key] 1502 if not field.is_string(): 1503 return "" 1504 return String(field.string_value()) 1505 1506 1507 def validate_status_frame( 1508 response_text: String, 1509 expected_request_id: String, 1510 expected_trace_id: String, 1511 ) raises -> FrameVerdict: 1512 """Validate one status frame's parsed envelope, correlation and outcome.""" 1513 var parsed = Value(None) 1514 try: 1515 parsed = loads(response_text) 1516 except: 1517 return FrameVerdict(False, "not_json", "", -1) 1518 if not parsed.is_object(): 1519 return FrameVerdict(False, "not_object", "", -1) 1520 if not _has_key(parsed, "version"): 1521 return FrameVerdict(False, "missing_version", "", -1) 1522 if Int(parsed["version"].int_value()) != 1: 1523 return FrameVerdict(False, "version_mismatch", "", -1) 1524 var request_id = _field_string(parsed, "request_id") 1525 if request_id != expected_request_id: 1526 return FrameVerdict(False, "correlation_mismatch", "", -1) 1527 var trace_id = _field_string(parsed, "trace_id") 1528 if trace_id != expected_trace_id: 1529 return FrameVerdict(False, "trace_mismatch", "", -1) 1530 if not _has_key(parsed, "ok"): 1531 return FrameVerdict(False, "missing_outcome", "", -1) 1532 if not parsed["ok"].bool_value(): 1533 return FrameVerdict(False, "not_ok", "", -1) 1534 if _has_key(parsed, "error"): 1535 return FrameVerdict(False, "unexpected_error", "", -1) 1536 if not _has_key(parsed, "output") or not parsed["output"].is_object(): 1537 return FrameVerdict(False, "missing_output", "", -1) 1538 var daemon = _field_string(parsed["output"], "daemon") 1539 if daemon != "hyfd": 1540 return FrameVerdict(False, "outcome_mismatch", "", -1) 1541 var declared = -1 1542 var output = parsed["output"] 1543 if _has_key(output, "limits") and output["limits"].is_object(): 1544 var limits = output["limits"] 1545 if _has_key(limits, "max_requests_per_process"): 1546 declared = Int(limits["max_requests_per_process"].int_value()) 1547 return FrameVerdict(True, "", "sys.status_ok", declared) 1548 1549 1550 def build_status_frame(index: Int) -> List[String]: 1551 """One status request frame plus its expected request/trace correlation.""" 1552 var request_id = "meas-status-" + String(index) 1553 var trace_id = "meas-trace-" + String(index) 1554 var frame = ( 1555 '{"version":1,"request_id":"' 1556 + request_id 1557 + '","trace_id":"' 1558 + trace_id 1559 + '","capability":"sys.status","input":{}}' 1560 ) 1561 var pair = List[String]() 1562 pair.append(frame) 1563 pair.append(request_id) 1564 pair.append(trace_id) 1565 return pair^ 1566 1567 1568 # ── Sampling ──────────────────────────────────────────────────────────────── 1569 1570 1571 def sample_rss_kb( 1572 pid: Int, 1573 command: String, 1574 mut guard: CleanupGuard, 1575 deadline_ms: Int = MEASUREMENT_SAMPLE_DEADLINE_MS, 1576 ) raises -> Int: 1577 """Numeric resident-set size in kB, or an explicit sampling failure. 1578 1579 The sampling subprocess consumes the caller's remaining work budget, so 1580 instrumentation can never extend the measured window. 1581 """ 1582 if deadline_ms <= 0: 1583 raise Error( 1584 "measurement: work deadline expired before rss sampling" 1585 " (work_deadline_expired)" 1586 ) 1587 var args = List[String]() 1588 args.append("-o") 1589 args.append("rss=") 1590 args.append("-p") 1591 args.append(String(pid)) 1592 var out = run_capture(command, args^, deadline_ms, guard) 1593 if out.exit_code != 0: 1594 raise Error( 1595 "measurement: rss sampling unavailable (" 1596 + command 1597 + " " 1598 + out.describe() 1599 + ")" 1600 ) 1601 var text = String(out.stdout.strip()) 1602 if text.byte_length() == 0: 1603 raise Error("measurement: rss sampling returned no value") 1604 for byte in text.as_bytes(): 1605 var b = Int(byte) 1606 if b < 48 or b > 57: 1607 raise Error("measurement: rss sampling value is not numeric") 1608 return Int(text) 1609 1610 1611 def sample_fd_count( 1612 pid: Int, 1613 command: String, 1614 mut guard: CleanupGuard, 1615 deadline_ms: Int = MEASUREMENT_SAMPLE_DEADLINE_MS, 1616 ) raises -> Int: 1617 """Numeric descriptor count, excluding lsof cwd/txt/mem pseudo rows.""" 1618 if deadline_ms <= 0: 1619 raise Error( 1620 "measurement: work deadline expired before descriptor sampling" 1621 " (work_deadline_expired)" 1622 ) 1623 var args = List[String]() 1624 args.append("-p") 1625 args.append(String(pid)) 1626 args.append("-F") 1627 args.append("f") 1628 var out = run_capture(command, args^, deadline_ms, guard) 1629 if out.exit_code != 0: 1630 raise Error( 1631 "measurement: descriptor sampling unavailable (" 1632 + command 1633 + " " 1634 + out.describe() 1635 + ")" 1636 ) 1637 var count = 0 1638 for line in out.stdout.split("\n"): 1639 var entry = String(line).strip() 1640 if entry.byte_length() < 2: 1641 continue 1642 if not entry.startswith("f"): 1643 continue 1644 var digits = String(entry[byte=1:]) 1645 var numeric = True 1646 for byte in digits.as_bytes(): 1647 var b = Int(byte) 1648 if b < 48 or b > 57: 1649 numeric = False 1650 break 1651 if numeric: 1652 count += 1 1653 if count == 0: 1654 raise Error("measurement: descriptor sampling counted no descriptors") 1655 return count 1656 1657 1658 def child_process_count( 1659 pid: Int, 1660 mut guard: CleanupGuard, 1661 deadline_ms: Int = MEASUREMENT_SAMPLE_DEADLINE_MS, 1662 ) raises -> Int: 1663 """Numeric count of live direct children of ``pid`` (0 when none).""" 1664 if deadline_ms <= 0: 1665 raise Error( 1666 "measurement: work deadline expired before child census" 1667 " (work_deadline_expired)" 1668 ) 1669 var args = List[String]() 1670 args.append("-P") 1671 args.append(String(pid)) 1672 var out = run_capture("pgrep", args^, deadline_ms, guard) 1673 if out.exit_code == 1: 1674 return 0 1675 if out.exit_code != 0: 1676 raise Error( 1677 "measurement: child census unavailable (" + out.describe() + ")" 1678 ) 1679 var count = 0 1680 for line in out.stdout.split("\n"): 1681 if String(line).strip().byte_length() > 0: 1682 count += 1 1683 return count 1684 1685 1686 # ── Measurement session ───────────────────────────────────────────────────── 1687 1688 1689 @fieldwise_init 1690 struct MeasurementSession(Movable): 1691 """Validated result of one persistent warmup/measured process.""" 1692 1693 var identity: MeasurementIdentity 1694 var warmup_frames: Int 1695 var measured_frames: Int 1696 var ok_frames: Int 1697 var failed_frames: Int 1698 var first_failure: String 1699 var startup_ms: Int 1700 var startup_wall_ms: Int 1701 var startup_sampling_ms: Int 1702 var warmup_ms: Int 1703 var measured_wall_ms: Int 1704 var measured_sampling_ms: Int 1705 var measured_ms: Int 1706 var run_wall_ms: Int 1707 var request_total_ms: Int 1708 var request_min_ms: Int 1709 var request_max_ms: Int 1710 var rss_kb_before_warmup: Int 1711 var rss_kb_after_warmup: Int 1712 var rss_kb_after_measured: Int 1713 var rss_kb_peak: Int 1714 var fd_before_warmup: Int 1715 var fd_after_measured: Int 1716 var fd_peak: Int 1717 var declared_max_requests_per_process: Int 1718 var child_exit: String 1719 var stderr_excerpt: String 1720 var sampling_method: String 1721 var sampling_cadence: String 1722 1723 def accounting_error_ms(self) -> Int: 1724 """Coherent accounting identity: wall = measured + instrumentation (ms). 1725 1726 ``measured_ms`` and ``measured_sampling_ms`` are both milliseconds 1727 inside the same measured wall interval, so a truthful session returns 1728 exactly zero here (ADR-0020 MC01). 1729 """ 1730 return ( 1731 self.measured_wall_ms - self.measured_ms - self.measured_sampling_ms 1732 ) 1733 1734 def summary(self) -> String: 1735 return ( 1736 "frames=" 1737 + String(self.warmup_frames + self.measured_frames) 1738 + " ok=" 1739 + String(self.ok_frames) 1740 + " failed=" 1741 + String(self.failed_frames) 1742 + " startup_ms=" 1743 + String(self.startup_ms) 1744 + " startup_wall_ms=" 1745 + String(self.startup_wall_ms) 1746 + " startup_sampling_ms=" 1747 + String(self.startup_sampling_ms) 1748 + " warmup_ms=" 1749 + String(self.warmup_ms) 1750 + " measured_ms=" 1751 + String(self.measured_ms) 1752 + " measured_wall_ms=" 1753 + String(self.measured_wall_ms) 1754 + " measured_sampling_ms=" 1755 + String(self.measured_sampling_ms) 1756 + " run_wall_ms=" 1757 + String(self.run_wall_ms) 1758 + " request_total_ms=" 1759 + String(self.request_total_ms) 1760 + " request_min_ms=" 1761 + String(self.request_min_ms) 1762 + " request_max_ms=" 1763 + String(self.request_max_ms) 1764 + " rss_kb=" 1765 + String(self.rss_kb_before_warmup) 1766 + "/" 1767 + String(self.rss_kb_after_warmup) 1768 + "/" 1769 + String(self.rss_kb_after_measured) 1770 + " peak=" 1771 + String(self.rss_kb_peak) 1772 + " fd=" 1773 + String(self.fd_before_warmup) 1774 + "/" 1775 + String(self.fd_after_measured) 1776 + " peak=" 1777 + String(self.fd_peak) 1778 + " child=" 1779 + self.child_exit 1780 + " stderr_bytes=" 1781 + String(self.stderr_excerpt.byte_length()) 1782 + " first_failure=" 1783 + (self.first_failure if self.first_failure != "" else "none") 1784 ) 1785 1786 1787 def argv_profile_of(binary_path: String, var argv: List[String]) -> String: 1788 """Machine-derived argv profile of the measured child (not asserted).""" 1789 var profile = String("<" + binary_path + ">") 1790 for index in range(len(argv)): 1791 profile += " " + argv[index] 1792 return profile^ 1793 1794 1795 def sample_warmup_boundary( 1796 mut process: MeasurementProcess, 1797 rss_sampler: String, 1798 fd_sampler: String, 1799 mut guard: CleanupGuard, 1800 mut rss_peak: Int, 1801 mut fd_peak: Int, 1802 ) raises -> Int: 1803 """Instrumentation between warmup and the measured wall interval. 1804 1805 Both samples are taken *before* the measured interval opens, so they are 1806 recorded as after-warmup/peak observations and are neither subtracted from 1807 nor added to the measured instrumentation total. Subtracting this interval 1808 afterwards (the period-11 defect) biased the measured time upward; ADR-0020 1809 MC01 requires it to be excluded exactly once by staying outside the 1810 interval. Returns the recorded after-warmup RSS in kB. 1811 """ 1812 var remaining = process.work_remaining_ms() 1813 var rss = sample_rss_kb(process.pid, rss_sampler, guard, remaining) 1814 var remaining_fd = process.work_remaining_ms() 1815 var fds = sample_fd_count(process.pid, fd_sampler, guard, remaining_fd) 1816 if rss > rss_peak: 1817 rss_peak = rss 1818 if fds > fd_peak: 1819 fd_peak = fds 1820 return rss 1821 1822 1823 def measure_persistent_process( 1824 source_root: String, 1825 binary_path: String, 1826 var argv: List[String], 1827 warmup_frames: Int, 1828 measured_frames: Int, 1829 deadline_ms: Int, 1830 mut guard: CleanupGuard, 1831 rss_sampler: String = "ps", 1832 fd_sampler: String = "lsof", 1833 expected_revision: String = "", 1834 expected_manifest_sha256: String = "", 1835 expected_binary_sha256: String = "", 1836 faults: MeasurementFaults = MeasurementFaults(), 1837 ) raises -> MeasurementSession: 1838 """Drive warmup then measured frames over one persistent owned process. 1839 1840 Identity capture is a separate bounded phase before the spawn; the single 1841 spawn-relative work deadline then covers sampling, request writes, response 1842 reads, the EOF drain and the child exit. When ``expected_revision`` is 1843 provided (product measurement) the measured source tree must be clean and 1844 the observed revision/tree/manifest must match the build exactly; a 1845 mismatch before or after the run is a bounded drift failure. Controlled 1846 test children record the observed identity as ``controlled_child`` and do 1847 not claim a build binding. 1848 """ 1849 if warmup_frames < 0 or measured_frames < 1: 1850 raise Error("measurement: invalid warmup/measured frame counts") 1851 if deadline_ms <= 0: 1852 raise Error("measurement: invalid work deadline") 1853 var binding = ( 1854 "clean_product_tree" if expected_revision != "" else "controlled_child" 1855 ) 1856 var observed = source_identity(source_root, guard) 1857 if expected_revision != "": 1858 require_clean_source(observed, "run capture") 1859 if observed.revision != expected_revision: 1860 raise Error( 1861 "measurement: source revision drift before run (expected " 1862 + expected_revision 1863 + " observed " 1864 + observed.revision 1865 + ")" 1866 ) 1867 if observed.manifest_sha256 != expected_manifest_sha256: 1868 raise Error( 1869 "measurement: source manifest drift before run (expected " 1870 + expected_manifest_sha256 1871 + " observed " 1872 + observed.manifest_sha256 1873 + ")" 1874 ) 1875 var binary_sha256 = file_sha256(binary_path, guard) 1876 if expected_binary_sha256 != "" and binary_sha256 != expected_binary_sha256: 1877 raise Error( 1878 "measurement: binary drift before run (expected " 1879 + expected_binary_sha256 1880 + " observed " 1881 + binary_sha256 1882 + ")" 1883 ) 1884 var identity = measurement_identity( 1885 source_root, 1886 binary_path, 1887 binary_sha256, 1888 binding, 1889 argv_profile_of(binary_path, argv.copy()), 1890 guard, 1891 ) 1892 var process = spawn_measurement_process( 1893 binary_path, argv^, deadline_ms, guard, faults 1894 ) 1895 var ok_frames = 0 1896 var failed_frames = 0 1897 var first_failure = "" 1898 var startup_ms = -1 1899 var startup_wall_ms = -1 1900 var startup_sampling_ms = 0 1901 var warmup_ms = 0 1902 var measured_wall_ms = 0 1903 var measured_sampling_ms = 0 1904 var measured_ms = 0 1905 var run_wall_ms = 0 1906 var request_total_ms = 0 1907 var request_min_ms = -1 1908 var request_max_ms = 0 1909 var rss_before = 0 1910 var rss_after_warmup = 0 1911 var rss_after_measured = 0 1912 var rss_peak = 0 1913 var fd_before = 0 1914 var fd_after = 0 1915 var fd_peak = 0 1916 var declared = -1 1917 var child_exit = "" 1918 var stderr_excerpt = "" 1919 1920 try: 1921 # Instrumentation before the first frame; timed as overhead so the 1922 # recorded startup latency is the child's own work. 1923 var pre_sample_start = now_ms() 1924 var pre_remaining = process.work_remaining_ms() 1925 rss_before = sample_rss_kb( 1926 process.pid, rss_sampler, guard, pre_remaining 1927 ) 1928 var pre_remaining_fd = process.work_remaining_ms() 1929 fd_before = sample_fd_count( 1930 process.pid, fd_sampler, guard, pre_remaining_fd 1931 ) 1932 var pre_sample_ms = now_ms() - pre_sample_start 1933 rss_peak = rss_before 1934 fd_peak = fd_before 1935 var warmup_start = now_ms() 1936 var measured_start = -1 1937 var total = warmup_frames + measured_frames 1938 if warmup_frames == 0: 1939 # ADR-0020 MC01: a zero-warmup measurement must open the measured 1940 # wall interval before its first measured request. Leaving it unset 1941 # reported host uptime as a duration (the period-11 defect: an 1942 # 804 ms run reported 347472213 ms). No warmup ran, so the 1943 # after-warmup checkpoint is the pre-warmup observation and the 1944 # measured instrumentation total starts empty. 1945 warmup_ms = 0 1946 rss_after_warmup = rss_before 1947 measured_start = now_ms() 1948 for index in range(total): 1949 var pair = build_status_frame(index) 1950 var frame = pair[0] 1951 var request_id = pair[1] 1952 var trace_id = pair[2] 1953 var frame_start = now_ms() 1954 process.send_frame(frame, index) 1955 var response = process.read_response(index) 1956 process.assert_no_surplus(index) 1957 var latency = now_ms() - frame_start 1958 var verdict = validate_status_frame(response, request_id, trace_id) 1959 if verdict.declared_max_requests >= 0: 1960 declared = verdict.declared_max_requests 1961 if verdict.ok: 1962 ok_frames += 1 1963 if index == 0: 1964 # Startup is the true spawn-to-first-validated-response 1965 # latency, captured only after the first response has been 1966 # validated. Any sampling that overlapped the child's own 1967 # startup is recorded separately as instrumentation rather 1968 # than subtracted, so the value cannot understate startup. 1969 startup_wall_ms = now_ms() - process.spawn_ms 1970 startup_sampling_ms = pre_sample_ms 1971 startup_ms = startup_wall_ms 1972 else: 1973 failed_frames += 1 1974 if first_failure == "": 1975 first_failure = ( 1976 "frame=" 1977 + String(index) 1978 + " reason=" 1979 + verdict.reason 1980 + " response=" 1981 + response 1982 ) 1983 # A mandatory per-frame guarantee failed: fail loudly instead 1984 # of reporting a partially successful measurement. 1985 raise Error("measurement: frame invalid " + first_failure) 1986 if index == warmup_frames - 1: 1987 warmup_ms = now_ms() - warmup_start 1988 # Boundary instrumentation is taken before the measured wall 1989 # interval opens, so it is excluded exactly once by never 1990 # entering that interval; it is never subtracted afterwards. 1991 rss_after_warmup = sample_warmup_boundary( 1992 process, rss_sampler, fd_sampler, guard, rss_peak, fd_peak 1993 ) 1994 measured_start = now_ms() 1995 if index >= warmup_frames: 1996 var measured_index = index - warmup_frames 1997 request_total_ms += latency 1998 if request_min_ms < 0 or latency < request_min_ms: 1999 request_min_ms = latency 2000 if latency > request_max_ms: 2001 request_max_ms = latency 2002 if measured_index % MEASUREMENT_SAMPLE_INTERVAL == 0: 2003 var s_start = now_ms() 2004 var s_remaining = process.work_remaining_ms() 2005 var mid_rss = sample_rss_kb( 2006 process.pid, rss_sampler, guard, s_remaining 2007 ) 2008 var s_remaining_fd = process.work_remaining_ms() 2009 var mid_fd = sample_fd_count( 2010 process.pid, fd_sampler, guard, s_remaining_fd 2011 ) 2012 measured_sampling_ms += now_ms() - s_start 2013 if mid_rss > rss_peak: 2014 rss_peak = mid_rss 2015 if mid_fd > fd_peak: 2016 fd_peak = mid_fd 2017 # Post-measured instrumentation (still inside the measured window). 2018 var post_start = now_ms() 2019 var post_remaining = process.work_remaining_ms() 2020 rss_after_measured = sample_rss_kb( 2021 process.pid, rss_sampler, guard, post_remaining 2022 ) 2023 var post_remaining_fd = process.work_remaining_ms() 2024 fd_after = sample_fd_count( 2025 process.pid, fd_sampler, guard, post_remaining_fd 2026 ) 2027 measured_sampling_ms += now_ms() - post_start 2028 if rss_after_measured > rss_peak: 2029 rss_peak = rss_after_measured 2030 if fd_after > fd_peak: 2031 fd_peak = fd_after 2032 if measured_start < 0: 2033 raise Error("measurement: measured phase was never initialized") 2034 measured_wall_ms = now_ms() - measured_start 2035 if measured_wall_ms < 0 or measured_sampling_ms < 0: 2036 raise Error("measurement: inconsistent measured timing window") 2037 measured_ms = measured_wall_ms - measured_sampling_ms 2038 if measured_ms < 0: 2039 raise Error( 2040 "measurement: instrumentation exceeds the measured wall" 2041 " interval" 2042 ) 2043 # The measured interval is bounded by the observed run, not by host 2044 # uptime: an uninitialized start would report the latter (ADR-0020 MC01). 2045 run_wall_ms = now_ms() - process.spawn_ms 2046 if measured_wall_ms > run_wall_ms: 2047 raise Error( 2048 "measurement: measured interval exceeds the observed run" 2049 ) 2050 var st = process.finish_expected(total) 2051 child_exit = st.describe() + " cleanup=proved" 2052 stderr_excerpt = process.stderr_text() 2053 if not st.exited or st.exit_code != 0: 2054 raise Error( 2055 "measurement: child exited nonzero (" 2056 + st.describe() 2057 + ") stderr=" 2058 + stderr_excerpt 2059 ) 2060 # Post-run drift rejection: the measured binary and source must be the 2061 # same ones captured before the run. 2062 var after = source_identity(source_root, guard) 2063 if expected_revision != "": 2064 require_identity_drift_free(observed, after, "run") 2065 var binary_after = file_sha256(binary_path, guard) 2066 if binary_after != binary_sha256: 2067 raise Error( 2068 "measurement: binary drifted during run (before " 2069 + binary_sha256 2070 + " after " 2071 + binary_after 2072 + ")" 2073 ) 2074 except e: 2075 process.cleanup() 2076 raise Error("measurement: persistent measurement failed: " + String(e)) 2077 2078 return MeasurementSession( 2079 identity=identity^, 2080 warmup_frames=warmup_frames, 2081 measured_frames=measured_frames, 2082 ok_frames=ok_frames, 2083 failed_frames=failed_frames, 2084 first_failure=first_failure, 2085 startup_ms=startup_ms, 2086 startup_wall_ms=startup_wall_ms, 2087 startup_sampling_ms=startup_sampling_ms, 2088 warmup_ms=warmup_ms, 2089 measured_wall_ms=measured_wall_ms, 2090 measured_sampling_ms=measured_sampling_ms, 2091 measured_ms=measured_ms, 2092 run_wall_ms=run_wall_ms, 2093 request_total_ms=request_total_ms, 2094 request_min_ms=request_min_ms, 2095 request_max_ms=request_max_ms, 2096 rss_kb_before_warmup=rss_before, 2097 rss_kb_after_warmup=rss_after_warmup, 2098 rss_kb_after_measured=rss_after_measured, 2099 rss_kb_peak=rss_peak, 2100 fd_before_warmup=fd_before, 2101 fd_after_measured=fd_after, 2102 fd_peak=fd_peak, 2103 declared_max_requests_per_process=declared, 2104 child_exit=child_exit, 2105 stderr_excerpt=stderr_excerpt, 2106 sampling_method=( 2107 "ps -o rss= -p <pid> (kB) / " 2108 + fd_sampler 2109 + " -p <pid> -F f numeric f<digits> rows" 2110 ), 2111 sampling_cadence=( 2112 "before warmup, after warmup, every " 2113 + String(MEASUREMENT_SAMPLE_INTERVAL) 2114 + " measured frames, after measured phase" 2115 ), 2116 )