diff --git a/docs/13_LIDAR_WORKER_PRODUCT_AND_ROADMAP.md b/docs/13_LIDAR_WORKER_PRODUCT_AND_ROADMAP.md index f8af7fa..10be97e 100644 --- a/docs/13_LIDAR_WORKER_PRODUCT_AND_ROADMAP.md +++ b/docs/13_LIDAR_WORKER_PRODUCT_AND_ROADMAP.md @@ -666,9 +666,25 @@ and exclusive PointSlab ownership is enforced. Conflict count remains `38`; the result is accepted only as the source-scoped diagnostic/shadow input for E33 and is not promoted as a detector-accuracy or staleness improvement. -Execution is strictly sequential through E33: E30 determines what E31 is -allowed to change; E31 determines the E32 profile; E32 determines the E33 -runtime input. E34 and E35 may proceed only after E33 closes exact accounting. +E33 accepted immutable worker result +`e33-worker-shadow-05cc0bb264410fd49536df90e94067ac39731aff0322a8873700d40008a8bb3a` +from package +`e33-worker-package-c8609151e71a35b3be3857aef168c0fe98380ff9055ce37560233bb332508823`. +The final E32 TrackGeometry stage delivered all `4,489` frames at recorded +`1.0×` pace on physical worker `DESKTOP-OPJ8J04`: effective rate was +`10.0062 FPS`, wall/source ratio `1.000142`, result-age p95 `2.668 ms` and +process RSS maximum `92.34 MiB`. Work and result queues remained bounded at +depth two with zero replacement, drop or deadline miss. The pinned container, +worker node, complete per-frame outcomes and resource samples are retained in +the immutable result. GPU board counters include co-tenant persistent +workloads and are evidence of node visibility, not E33 process consumption. +E33 qualifies downstream publication against the immutable E10→E32 chain; it +does not rerun or independently requalify upstream inference. + +Execution was strictly sequential through E33: E30 determined what E31 was +allowed to change; E31 determined the E32 profile; E32 determined the E33 +runtime input. Exact accounting is now closed, so E34 and E35 are the active +critical path. E36 is the first generalization gate. A separate product decision follows: either keep the result as operator/shadow evidence, or start L5 occupied-space integration. No LAB in this cycle can enable navigation, commands or safety diff --git a/docs/16_ARCHITECTURE_AUDIT_EXECUTION_ROADMAP.md b/docs/16_ARCHITECTURE_AUDIT_EXECUTION_ROADMAP.md index d888667..bdd9bb3 100644 --- a/docs/16_ARCHITECTURE_AUDIT_EXECUTION_ROADMAP.md +++ b/docs/16_ARCHITECTURE_AUDIT_EXECUTION_ROADMAP.md @@ -184,7 +184,15 @@ unqualified ranges and resolves 4,461 overlapping point claims. Of 709 arbitrated semantic observations, 563 retain `agree` and 146 conservatively become `unknown`. The 38 E29 conflicts remain 38; A6 is accepted as a diagnostic/shadow contract, not as detector-accuracy improvement. A7/E33 -recorded-source-paced worker execution is now the critical path. +is complete in immutable result +`e33-worker-shadow-05cc0bb264410fd49536df90e94067ac39731aff0322a8873700d40008a8bb3a`. +All 4,489 frames were delivered at `10.0062 FPS` over a `448.623 s` source +span with zero queue replacement, result drop or deadline miss. Result age +p95 was `2.668 ms`, process RSS peaked at `92.34 MiB`, both queues remained +bounded at depth two and the exact physical worker/container identity is +retained. This qualifies the final TrackGeometry publication stage against the +immutable upstream E10→E32 chain; it is not a new model-inference benchmark. +A8/E34+E35 is now the critical path. - [x] Reproduce all 4,489 immutable E29 frames with the exact frozen profile before applying E31/E30 changes. @@ -196,6 +204,18 @@ recorded-source-paced worker execution is now the critical path. cause while preserving source availability and unknown/free-space policy. - [x] Add a digest-bound compact binary PointSlab encoding and strict TrackGeometryFrame reconstruction/validation. +- [x] Package the exact E32 result, frozen runtime and profile as immutable + E33 worker input with complete artifact digests. +- [x] Replay all 4,489 frames at recorded `1.0×` pace through independent + bounded latest-wins work/result queues. +- [x] Close every frame as delivered, input-superseded or result-superseded + and publish per-frame health/timing plus resource telemetry. +- [x] Bind the accepted result to physical worker `DESKTOP-OPJ8J04` and the + pinned container image; independently verify all result artifacts. +- [ ] Build E34 as a separate short-TTL occupied/unknown temporal layer over + accepted E32/E33 evidence. +- [ ] Run E35 deterministic degradation/recovery variants without changing + the immutable source or accepted E32/E33 results. ### A3 residual and human-exception policy diff --git a/docs/adr/0026-e33-recorded-source-paced-track-geometry-shadow.md b/docs/adr/0026-e33-recorded-source-paced-track-geometry-shadow.md new file mode 100644 index 0000000..5ef1dcc --- /dev/null +++ b/docs/adr/0026-e33-recorded-source-paced-track-geometry-shadow.md @@ -0,0 +1,97 @@ +# ADR 0026 — E33 recorded-source-paced TrackGeometry worker shadow + +Date: 2026-07-27 + +Status: accepted for source-scoped diagnostic/shadow use + +## Context + +E32 closed source, observation and exclusive point-ownership accounting for all +4,489 immutable RAVNOVES00 frames. That result proved the data contract, but it +did not prove that the final TrackGeometry publication stage could sustain the +recorded source timeline on the designated worker with bounded channels and +explicit terminal outcomes. + +The upstream E10 result already contains the full source-paced detector and +semantic worker execution. E33 therefore qualifies the final E32 +TrackGeometry transport/publication stage. It does not rerun, replace or tune +the detector, segmentation, tracking, calibration or geometry algorithms. + +## Decision + +1. E33 receives one content-addressed package containing the complete immutable + E32 result, the frozen E33 profile and the minimum runtime required to + validate and replay it. +2. The package and result bind the exact E10, E29, E30, E31 and E32 identities. + Every package artifact and every published result artifact has a byte length + and SHA-256 digest. +3. The producer releases all frames at `1.0×` recorded host pace. Work and + result channels are independent latest-wins queues with capacity two. +4. Queue replacement is never silent. Every source frame has exactly one + terminal outcome: delivered, input-superseded or result-superseded. +5. Delivered frames include release lag, work wait, processing time, result + wait, result age, health and a digest over the referenced E32 point slab. + A late result becomes explicit health state. +6. The runtime samples process CPU and RSS plus GPU board visibility, + utilization and memory. GPU board values are node-level observations and + may include co-tenant workloads; they are not process-level attribution to + the E33 publication stage. +7. The immutable identity binds both the expected physical worker node and the + container hostname. A correct container on an unexpected node fails closed. +8. The Windows wrapper admits only a real `D:` package/output directory, + enforces a 300 GiB free-space floor, verifies the pinned container image, + disables networking, uses a read-only filesystem, drops capabilities and + publishes only one atomically completed result directory. +9. E33 is diagnostic/shadow only. Commands, navigation and safety acceptance + remain false. No persistent reconstruction or source artifact is modified. + +## Accepted result + +Package: +`e33-worker-package-c8609151e71a35b3be3857aef168c0fe98380ff9055ce37560233bb332508823`. + +Result: +`e33-worker-shadow-05cc0bb264410fd49536df90e94067ac39731aff0322a8873700d40008a8bb3a`. + +Runtime: + +- physical worker: `DESKTOP-OPJ8J04`; +- container hostname: `ebca59deeafd`; +- image: + `nvcr.io/nvidia/tritonserver:26.06-py3@sha256:58df7489c3f2276f9591d500a012dee03e23d35543ce3c390b4c001e6bf90794`; +- source span: `448.623 s`; +- replay wall time: `448.687 s`; +- wall/source ratio: `1.000142`; +- effective delivery rate: `10.0062 FPS`. + +Accounting and timing: + +- 4,489 / 4,489 frames delivered; +- zero input replacements, result replacements and deadline misses; +- work/result queue maximum depth `2/2`; +- release lag p50/p95/max: + `0.075 / 0.462 / 8.594 ms`; +- processing p50/p95/max: + `0.174 / 0.750 / 7.698 ms`; +- result age p50/p95/max: + `0.468 / 2.668 / 13.638 ms`; +- process RSS p50/p95/max: + `79.28 / 90.33 / 92.34 MiB`. + +All 13 predeclared acceptance requirements pass. The independently copied +result archive has SHA-256 +`9336f8df72ed43c2b46c68c544400db2fda52cf22f7ed8520d23fa0a9fc21081`. + +## Consequences + +- A7 is complete for the final TrackGeometry worker publication stage. +- E33 proves source-paced bounded execution and closed accounting for this + exact source, profile, package, worker and container identity. It does not + independently requalify the already accepted upstream model-inference + stages. +- An earlier accepted runtime result without a physical worker-node binding is + retained as historical evidence and is not the accepted A7 result. +- E34 may consume the accepted E32 evidence under the E33 timing/accounting + envelope to build a separate short-TTL occupied/unknown layer. +- E35 must still prove deterministic degradation and recovery. E33 does not + establish second-source transfer, free space, traversability or control. diff --git a/experiments/perception/LAB_E33_REPORT_2026-07-27.md b/experiments/perception/LAB_E33_REPORT_2026-07-27.md new file mode 100644 index 0000000..aa3aa88 --- /dev/null +++ b/experiments/perception/LAB_E33_REPORT_2026-07-27.md @@ -0,0 +1,161 @@ +# LAB E33 — recorded-source-paced TrackGeometry worker shadow + +Date: 2026-07-27 +Status: accepted for source-scoped diagnostic/shadow use; navigation, safety +and command authority are false +Immutable result: +`e33-worker-shadow-05cc0bb264410fd49536df90e94067ac39731aff0322a8873700d40008a8bb3a` + +## Objective and architecture stage + +E33 closes A7 of the architecture-audit roadmap. It asks whether the final E32 +TrackGeometry publication stage can replay every immutable RAVNOVES00 frame at +the original recorded host pace on the designated worker while keeping queues, +deadlines, health, resource use and every terminal frame outcome explicit. + +The upstream E10 result already ran the complete detector and semantic worker +pipeline over the same source at recorded pace. E33 does not rerun or tune +those models. It qualifies the downstream E32 TrackGeometry +validation/publication stage and preserves the complete upstream identity +chain. + +## Source evidence + +- E32 input: + `e32-track-geometry-a14ca0e7fb3850ca0dfa3c41634e1b490a2d58ab74d101afc6d6921fbdb0e6fd`. +- Source frames: `4,489`. +- Recorded timeline: `35.421857292–484.044857292 s`. +- Source span: `448.623 s`. +- Upstream camera result: + `e10-integrated-perception-459aac93918d8f6414b342986ccc6968fefcef6c1f3a78a5254df0b565255ad2`. +- E29 result: + `e29-camera-geometry-421a9d930638bef12cd5eb10979a477917fa4a389e655ed95f73ba4bd62e13dc`. +- E31 source profile: + `e31-source-qualification-b2460a5eb143688c7eea6821b2277e13aea79868abe81d83f7e78548c119159a`. + +The E33 package includes the complete E32 manifest, frame JSONL, PointSlab +arrays, change journal and report. The package builder verifies every input +digest before copying it and hashes every packaged runtime/source artifact. + +## Method, runtime and worker + +Package: +`e33-worker-package-c8609151e71a35b3be3857aef168c0fe98380ff9055ce37560233bb332508823`. + +Archive SHA-256: +`1495fd04a383f169098c249778bf9234cdf853e53e97d2bd4cdf500e9a921b2c`. + +Execution: + +- ID: `e33-full-1x-20260727-a14ca0e7-v4`; +- physical worker: `DESKTOP-OPJ8J04`; +- container hostname: `ebca59deeafd`; +- Python `3.12.3`, NumPy `1.26.4`, Linux `x86_64`; +- pinned image: + `nvcr.io/nvidia/tritonserver:26.06-py3@sha256:58df7489c3f2276f9591d500a012dee03e23d35543ce3c390b4c001e6bf90794`; +- `1.0×` recorded pacing; +- work queue capacity two; +- result queue capacity two; +- result deadline `100 ms`, stale boundary `150 ms`, unavailable boundary + `500 ms`; +- worker free-space floor `300 GiB`. + +The container runs without network access, with a read-only filesystem, +capabilities dropped, `no-new-privileges`, a bounded PID limit and only +read-only package plus writable result mounts. The result appears only through +an atomic final directory rename. + +## Implementation and validation + +The runtime: + +1. verifies package and E32 content identities; +2. binds the expected physical worker and pinned container image; +3. memory-maps the E32 PointSlab arrays; +4. releases each frame on its recorded timeline; +5. validates the TrackGeometry record and hashes its point-slab byte ranges; +6. transports work and results through independent bounded latest-wins queues; +7. records exactly one terminal outcome per source frame; +8. samples process and GPU-board telemetry once per second; +9. evaluates predeclared gates; +10. hashes and atomically publishes the immutable result. + +Automated coverage includes the complete positive lifecycle, deterministic +overload replacement, artifact tamper rejection, container identity rejection +and physical worker-node rejection. The shared queue test now also verifies +the exact evicted item so the runtime can close per-frame replacement +accounting. + +## Result + +Result: +`e33-worker-shadow-05cc0bb264410fd49536df90e94067ac39731aff0322a8873700d40008a8bb3a`. + +Result archive SHA-256: +`9336f8df72ed43c2b46c68c544400db2fda52cf22f7ed8520d23fa0a9fc21081`. + +| Measure | Result | Gate | +| --- | ---: | ---: | +| delivered frames | 4,489 / 4,489 | closed accounting | +| input superseded | 0 | ≤ 0.1% | +| result superseded | 0 | 0 | +| deadline miss fraction | 0 | ≤ 1% | +| effective delivery | 10.0062 FPS | ≥ 9.5 FPS | +| wall/source ratio | 1.000142 | recorded speed `1.0×` | +| release lag p95 | 0.462 ms | ≤ 25 ms | +| result age p95 | 2.668 ms | ≤ 100 ms | +| process RSS p95 / max | 90.33 / 92.34 MiB | max ≤ 1,024 MiB | +| work/result max queue depth | 2 / 2 | bounded by 2 | + +All 4,489 outcomes are `healthy`. Processing p50/p95/max is +`0.174 / 0.750 / 7.698 ms`; result age p50/p95/max is +`0.468 / 2.668 / 13.638 ms`. + +The GPU is visible. Across 436 node samples, board utilization is +`48%` p50, `52%` p95 and `55%` maximum; used board memory is approximately +`13,090 MiB`. These board counters include persistent co-tenant inference +workloads. They prove node visibility only and must not be described as E33 +process GPU consumption. The E33 final publication stage is predominantly CPU +validation and hashing. + +All 13 predeclared requirements pass. The result is +`accepted-recorded-source-paced-shadow`. + +## Regressions, failed attempts and limitations + +Two package attempts failed before result publication: + +- one because the container mount basename did not preserve the + content-addressed package ID; +- one because the package reader assumed a non-contractual `kind` field in the + E32 artifact descriptors. + +Both defects were fixed in new package identities; neither produced a result. +An intermediate accepted replay proved the timing path but recorded only the +container hostname. It remains immutable history but is not the accepted A7 +evidence. The final profile and runtime now require +`worker_node=DESKTOP-OPJ8J04` and were replayed in full. + +E33 does not: + +- rerun or independently time the upstream detector/semantic model stages; +- improve detector accuracy, calibration or E32 geometry quality; +- attribute shared GPU-board use to this process; +- prove a second source or changed mount; +- infer free space, traversability or safety; +- enable commands. + +## Decision and next stage + +A7 is complete. E33 proves that the final immutable TrackGeometry stage keeps +pace with this exact recorded source on the designated worker with no hidden +frame loss and bounded timing/resources. + +The next critical-path stage is A8: + +- E34 builds a separate short-TTL occupied/unknown temporal layer over the + accepted E32/E33 evidence without modifying persistent reconstruction; +- E35 injects deterministic source loss, staleness, delay, bounded drops and + timing offsets and requires explicit safe degradation/recovery. + +Planner, navigation, command and safety authority remain false. diff --git a/experiments/perception/prepare_e33_worker_package.py b/experiments/perception/prepare_e33_worker_package.py new file mode 100644 index 0000000..163df9b --- /dev/null +++ b/experiments/perception/prepare_e33_worker_package.py @@ -0,0 +1,276 @@ +#!/usr/bin/env python3 +"""Build a minimal immutable E33 worker package around one accepted E32 result.""" + +from __future__ import annotations + +import argparse +import hashlib +import json +import os +import shutil +import uuid +from datetime import UTC, datetime +from pathlib import Path +from typing import Any + +from k1link.compute.e32_track_geometry_replay import read_e32_track_geometry_replay +from k1link.compute.e33_worker_shadow import E33_PACKAGE_SCHEMA, E33_PROFILE_SCHEMA + +_RUNTIME_FILES = { + "runtime/k1link/__init__.py": "src/k1link/__init__.py", + "runtime/k1link/compute/__init__.py": None, + "runtime/k1link/compute/e33_worker_shadow.py": ( + "src/k1link/compute/e33_worker_shadow.py" + ), + "runtime/k1link/compute/live_perception.py": "src/k1link/compute/live_perception.py", + "runtime/k1link/data_plane/__init__.py": "src/k1link/data_plane/__init__.py", + "runtime/k1link/data_plane/views.py": "src/k1link/data_plane/views.py", + "runtime/run_e33_worker_shadow.py": ( + "experiments/perception/worker/run_e33_worker_shadow.py" + ), + "runtime/Invoke-E33WorkerShadow.ps1": ( + "experiments/perception/worker/Invoke-E33WorkerShadow.ps1" + ), +} +_GENERATED_COMPUTE_INIT = ( + '"""Minimal E33 worker projection; import e33_worker_shadow explicitly."""\n' +) + + +class E33WorkerPackageError(RuntimeError): + """The E33 package source or immutable package is invalid.""" + + +def build_e33_worker_package( + *, + repository_root: Path, + e32_result_root: Path, + profile_path: Path, + output_root: Path, +) -> Path: + """Build or verify one content-addressed E33 worker package.""" + + repository = repository_root.resolve(strict=True) + e32 = read_e32_track_geometry_replay(e32_result_root) + profile_source = profile_path.resolve(strict=True) + profile = _read_json(profile_source) + if ( + profile.get("schema_version") != E33_PROFILE_SCHEMA + or profile.get("expected_e32_result_id") != e32.result_id + ): + raise E33WorkerPackageError("E33 package profile does not bind the E32 result") + runtime_sources: dict[str, Path | None] = {} + for target, relative in _RUNTIME_FILES.items(): + source = None if relative is None else repository / relative + if source is not None and (not source.is_file() or source.is_symlink()): + raise E33WorkerPackageError(f"E33 runtime source is invalid: {relative}") + runtime_sources[target] = source + e32_files = sorted( + path + for path in e32.result_root.iterdir() + if path.is_file() and not path.is_symlink() + ) + if not e32_files or any(path.is_symlink() for path in e32.result_root.iterdir()): + raise E33WorkerPackageError("E32 result contains an invalid package member") + e32_prefix = f"input/e32/{e32.result_id}" + sources: dict[str, Path | None] = { + **runtime_sources, + "profile.json": profile_source, + **{f"{e32_prefix}/{path.name}": path for path in e32_files}, + } + source_descriptors = [] + for relative, source in sorted(sources.items()): + payload = ( + _GENERATED_COMPUTE_INIT.encode() + if source is None + else source.read_bytes() + ) + source_descriptors.append( + { + "path": relative, + "byte_length": len(payload), + "sha256": hashlib.sha256(payload).hexdigest(), + } + ) + identity = { + "schema_version": E33_PACKAGE_SCHEMA, + "classification": "immutable-recorded-shadow-worker-input", + "e32_result_id": e32.result_id, + "e32_identity_sha256": e32.manifest["identity_sha256"], + "profile_sha256": _sha256(profile_source), + "artifact_paths": [item["path"] for item in source_descriptors], + "source_artifacts": source_descriptors, + "authority": { + "commands_enabled": False, + "navigation_or_safety_accepted": False, + }, + } + identity_sha256 = hashlib.sha256(_canonical_json(identity)).hexdigest() + package_id = f"e33-worker-package-{identity_sha256}" + output = output_root.expanduser().absolute() + output.mkdir(mode=0o700, parents=True, exist_ok=True) + destination = output / package_id + if destination.exists(): + validate_e33_worker_package(destination) + return destination + staging = output / f".{package_id}.{uuid.uuid4().hex}.tmp" + staging.mkdir(mode=0o700, exist_ok=False) + try: + for relative, source in sources.items(): + target_path = staging / relative + target_path.parent.mkdir(mode=0o700, parents=True, exist_ok=True) + if source is None: + target_path.write_text(_GENERATED_COMPUTE_INIT, encoding="utf-8") + else: + shutil.copyfile(source, target_path) + artifacts = [ + { + "kind": relative, + "path": relative, + "byte_length": (staging / relative).stat().st_size, + "sha256": _sha256(staging / relative), + } + for relative in sorted(sources) + ] + manifest = { + "schema_version": E33_PACKAGE_SCHEMA, + "package_id": package_id, + "identity_sha256": identity_sha256, + "identity": identity, + "created_at_utc": datetime.now(UTC) + .isoformat(timespec="milliseconds") + .replace("+00:00", "Z"), + "artifacts": artifacts, + } + _write_json(staging / "manifest.json", manifest) + validate_e33_worker_package(staging, allow_staging=True) + os.replace(staging, destination) + except BaseException: + shutil.rmtree(staging, ignore_errors=True) + raise + validate_e33_worker_package(destination) + return destination + + +def validate_e33_worker_package( + root: Path, + *, + allow_staging: bool = False, +) -> dict[str, Any]: + """Validate package identity, exact file set, and every member digest.""" + + resolved = root.resolve(strict=True) + manifest = _read_json(resolved / "manifest.json") + identity = manifest.get("identity") + identity_sha256 = manifest.get("identity_sha256") + package_id = manifest.get("package_id") + artifacts = manifest.get("artifacts") + expected_name = ( + isinstance(package_id, str) + and ( + resolved.name == package_id + or ( + allow_staging + and resolved.name.startswith(f".{package_id}.") + and resolved.name.endswith(".tmp") + ) + ) + ) + if ( + manifest.get("schema_version") != E33_PACKAGE_SCHEMA + or not isinstance(identity, dict) + or not isinstance(identity_sha256, str) + or hashlib.sha256(_canonical_json(identity)).hexdigest() != identity_sha256 + or package_id != f"e33-worker-package-{identity_sha256}" + or not expected_name + or not isinstance(artifacts, list) + ): + raise E33WorkerPackageError("E33 worker package identity is invalid") + expected_paths = set(identity.get("artifact_paths", [])) + actual_paths = { + path.relative_to(resolved).as_posix() + for path in resolved.rglob("*") + if path.is_file() + } + if ( + not expected_paths + or actual_paths != expected_paths | {"manifest.json"} + or len(artifacts) != len(expected_paths) + ): + raise E33WorkerPackageError("E33 worker package file set changed") + observed: set[str] = set() + for item in artifacts: + if not isinstance(item, dict): + raise E33WorkerPackageError("E33 worker package artifact is invalid") + relative = item.get("path") + path = resolved / str(relative) + if ( + not isinstance(relative, str) + or relative not in expected_paths + or relative in observed + or Path(relative).is_absolute() + or ".." in Path(relative).parts + or not path.is_file() + or path.is_symlink() + or item.get("kind") != relative + or item.get("byte_length") != path.stat().st_size + or item.get("sha256") != _sha256(path) + ): + raise E33WorkerPackageError("E33 worker package artifact changed") + observed.add(relative) + if observed != expected_paths: + raise E33WorkerPackageError("E33 worker package artifact coverage changed") + return manifest + + +def _canonical_json(value: object) -> bytes: + return json.dumps( + value, + sort_keys=True, + separators=(",", ":"), + allow_nan=False, + ).encode() + + +def _sha256(path: Path) -> str: + digest = hashlib.sha256() + with path.open("rb") as stream: + while chunk := stream.read(1024 * 1024): + digest.update(chunk) + return digest.hexdigest() + + +def _read_json(path: Path) -> dict[str, Any]: + value = json.loads(path.read_text(encoding="utf-8-sig")) + if not isinstance(value, dict): + raise E33WorkerPackageError(f"JSON object expected: {path.name}") + return value + + +def _write_json(path: Path, value: object) -> None: + with path.open("x", encoding="utf-8", newline="\n") as stream: + json.dump(value, stream, indent=2, sort_keys=True) + stream.write("\n") + stream.flush() + os.fsync(stream.fileno()) + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("--repository-root", type=Path, required=True) + parser.add_argument("--e32-result", type=Path, required=True) + parser.add_argument("--profile", type=Path, required=True) + parser.add_argument("--output-root", type=Path, required=True) + args = parser.parse_args() + package = build_e33_worker_package( + repository_root=args.repository_root, + e32_result_root=args.e32_result, + profile_path=args.profile, + output_root=args.output_root, + ) + print(package) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/experiments/perception/worker/Invoke-E33WorkerShadow.ps1 b/experiments/perception/worker/Invoke-E33WorkerShadow.ps1 new file mode 100644 index 0000000..f7b94d5 --- /dev/null +++ b/experiments/perception/worker/Invoke-E33WorkerShadow.ps1 @@ -0,0 +1,141 @@ +[CmdletBinding()] +param( + [Parameter(Mandatory = $true)] + [string]$PackageRoot, + [Parameter(Mandatory = $true)] + [ValidatePattern("^[a-z0-9][a-z0-9._-]{2,95}$")] + [string]$ExecutionId, + [string]$OutputRoot = "D:\NDC_MISSIONCORE\runtime\derived\e33-worker-shadow", + [string]$ContainerImage = "nvcr.io/nvidia/tritonserver:26.06-py3@sha256:58df7489c3f2276f9591d500a012dee03e23d35543ce3c390b4c001e6bf90794", + [ValidateRange(1, 1000)] + [int]$FreeGiBFloor = 300 +) + +$ErrorActionPreference = "Stop" +$ProgressPreference = "SilentlyContinue" + +function Assert-LastExitCode([string]$Operation) { + if ($LASTEXITCODE -ne 0) { + throw "$Operation failed with exit code $LASTEXITCODE" + } +} + +function Resolve-DDirectory([string]$Path, [string]$Label) { + $item = Get-Item -LiteralPath (Resolve-Path -LiteralPath $Path).Path -Force + $root = [IO.Path]::GetPathRoot($item.FullName).TrimEnd("\") + if ( + -not $item.PSIsContainer -or + ($item.Attributes -band [IO.FileAttributes]::ReparsePoint) -or + $root -ine "D:" + ) { + throw "$Label must be a real D: directory" + } + return $item.FullName +} + +function Convert-ToDockerPath([string]$Path) { + return $Path.Replace("\", "/") +} + +function Assert-FreeSpace([string]$Phase) { + $free = [int64](Get-PSDrive -Name D).Free + $floor = [int64]$FreeGiBFloor * 1GB + Write-Host ( + "DISK_GUARD PHASE={0} DRIVE=D FREE_BYTES={1} FREE_GIB={2} FLOOR_GIB={3}" -f + $Phase, $free, [math]::Round($free / 1GB, 3), $FreeGiBFloor + ) + if ($free -lt ($floor + 1GB)) { + throw "D: lacks the guarded E33 reserve during $Phase" + } + return $free +} + +$package = Resolve-DDirectory $PackageRoot "E33 package" +$packageManifestPath = Join-Path $package "manifest.json" +if (-not (Test-Path -LiteralPath $packageManifestPath -PathType Leaf)) { + throw "E33 package manifest is missing" +} +$packageManifest = Get-Content -LiteralPath $packageManifestPath -Raw | ConvertFrom-Json +if ( + $packageManifest.schema_version -ne "missioncore.e33-worker-package/v1" -or + $packageManifest.package_id -ne (Split-Path $package -Leaf) -or + $packageManifest.package_id -notmatch "^e33-worker-package-[a-f0-9]{64}$" +) { + throw "E33 package manifest is incompatible" +} + +if (-not (Test-Path -LiteralPath $OutputRoot)) { + $null = New-Item -ItemType Directory -Path $OutputRoot +} +$output = Resolve-DDirectory $OutputRoot "E33 output root" +$freeBefore = Assert-FreeSpace "preflight" + +& docker image inspect $ContainerImage *> $null +Assert-LastExitCode "Pinned E33 container image inspection" + +$dockerPackage = Convert-ToDockerPath $package +$dockerOutput = Convert-ToDockerPath $output +$packageName = Split-Path $package -Leaf +$containerPackage = "/opt/e33-input/$packageName" +$command = @( + "run", "--rm", + "--network", "none", + "--read-only", + "--security-opt", "no-new-privileges:true", + "--cap-drop", "ALL", + "--pids-limit", "128", + "--gpus", "all", + "--tmpfs", "/tmp:rw,noexec,nosuid,size=256m", + "-e", "PYTHONDONTWRITEBYTECODE=1", + "-e", ("PYTHONPATH={0}/runtime" -f $containerPackage), + "-e", ("E33_CONTAINER_IMAGE={0}" -f $ContainerImage), + "-e", ("E33_WORKER_NODE={0}" -f $env:COMPUTERNAME), + "-v", ("{0}:{1}:ro" -f $dockerPackage, $containerPackage), + "-v", ("{0}:/output:rw" -f $dockerOutput), + "--entrypoint", "python3", + $ContainerImage, + ("{0}/runtime/run_e33_worker_shadow.py" -f $containerPackage), + "--package", $containerPackage, + "--output-root", "/output", + "--execution-id", $ExecutionId +) + +Write-Output ("PACKAGE_ID={0}" -f $packageManifest.package_id) +Write-Output ("PACKAGE_IDENTITY_SHA256={0}" -f $packageManifest.identity_sha256) +Write-Output ("CONTAINER_IMAGE={0}" -f $ContainerImage) +Write-Output ("EXECUTION_ID={0}" -f $ExecutionId) +& docker @command +Assert-LastExitCode "E33 worker shadow" + +$matches = @( + Get-ChildItem -LiteralPath $output -Directory -Filter "e33-worker-shadow-*" | + Where-Object { + $resultPath = Join-Path $_.FullName "result.json" + if (-not (Test-Path -LiteralPath $resultPath -PathType Leaf)) { + return $false + } + $result = Get-Content -LiteralPath $resultPath -Raw | ConvertFrom-Json + return ( + $result.identity.execution_id -eq $ExecutionId -and + $result.identity.package_id -eq $packageManifest.package_id + ) + } +) +if ($matches.Count -ne 1) { + throw "E33 immutable result could not be resolved uniquely" +} +$resultRoot = $matches[0].FullName +$resultManifest = Get-Content -LiteralPath (Join-Path $resultRoot "result.json") -Raw | + ConvertFrom-Json +if ( + $resultManifest.schema_version -ne "missioncore.e33-worker-shadow-result/v1" -or + $resultManifest.publication_scope -ne "worker-shadow-diagnostic-only" +) { + throw "E33 result manifest is incompatible" +} +$freeAfter = Assert-FreeSpace "completed" +Write-Output ("RESULT_ROOT={0}" -f $resultRoot) +Write-Output ("RESULT_ID={0}" -f $resultManifest.result_id) +Write-Output ("ACCEPTANCE_STATE={0}" -f $resultManifest.acceptance_state) +Write-Output ("DISK_FREE_BYTES_BEFORE={0}" -f $freeBefore) +Write-Output ("DISK_FREE_BYTES_AFTER={0}" -f $freeAfter) diff --git a/experiments/perception/worker/e33_worker_shadow_profile.json b/experiments/perception/worker/e33_worker_shadow_profile.json new file mode 100644 index 0000000..35bbc8a --- /dev/null +++ b/experiments/perception/worker/e33_worker_shadow_profile.json @@ -0,0 +1,37 @@ +{ + "schema_version": "missioncore.e33-worker-shadow-profile/v1", + "mode": "full-session-qualification", + "expected_e32_result_id": "e32-track-geometry-a14ca0e7fb3850ca0dfa3c41634e1b490a2d58ab74d101afc6d6921fbdb0e6fd", + "expected_worker_node": "DESKTOP-OPJ8J04", + "container_image": "nvcr.io/nvidia/tritonserver:26.06-py3@sha256:58df7489c3f2276f9591d500a012dee03e23d35543ce3c390b4c001e6bf90794", + "pacing": { + "speed": 1.0, + "start_delay_ms": 250.0 + }, + "queues": { + "work_capacity": 2, + "result_capacity": 2 + }, + "deadlines": { + "result_ms": 100.0, + "stale_ms": 150.0, + "unavailable_ms": 500.0 + }, + "resources": { + "sample_interval_seconds": 1.0 + }, + "acceptance": { + "minimum_effective_fps": 9.5, + "maximum_input_drop_fraction": 0.001, + "maximum_result_drop_fraction": 0.0, + "maximum_result_age_p95_ms": 100.0, + "maximum_deadline_miss_fraction": 0.01, + "maximum_release_lag_p95_ms": 25.0, + "maximum_process_rss_mib": 1024.0, + "require_gpu_visible": true + }, + "authority": { + "commands_enabled": false, + "navigation_or_safety_accepted": false + } +} diff --git a/experiments/perception/worker/run_e33_worker_shadow.py b/experiments/perception/worker/run_e33_worker_shadow.py new file mode 100644 index 0000000..46c7bdf --- /dev/null +++ b/experiments/perception/worker/run_e33_worker_shadow.py @@ -0,0 +1,40 @@ +#!/usr/bin/env python3 +"""Execute one immutable E33 package inside the bounded worker container.""" + +from __future__ import annotations + +import argparse +import json +from pathlib import Path + +from k1link.compute.e33_worker_shadow import run_e33_worker_shadow + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("--package", type=Path, required=True) + parser.add_argument("--output-root", type=Path, required=True) + parser.add_argument("--execution-id", required=True) + args = parser.parse_args() + result = run_e33_worker_shadow( + args.package, + args.output_root, + execution_id=args.execution_id, + ) + print( + json.dumps( + { + "result_id": result.result_id, + "result_root": str(result.result_root), + "accepted": result.accepted, + "metrics": result.report["metrics"], + }, + ensure_ascii=False, + sort_keys=True, + ) + ) + return 0 if result.accepted else 2 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/src/k1link/compute/__init__.py b/src/k1link/compute/__init__.py index 3c14017..c0645d2 100644 --- a/src/k1link/compute/__init__.py +++ b/src/k1link/compute/__init__.py @@ -27,6 +27,19 @@ from .e32_track_geometry_replay import ( e32_track_geometry_frame, read_e32_track_geometry_replay, ) +from .e33_worker_shadow import ( + E33_FRAME_SCHEMA, + E33_OUTCOME_SCHEMA, + E33_PACKAGE_SCHEMA, + E33_PROFILE_SCHEMA, + E33_REPORT_SCHEMA, + E33_RESOURCE_SCHEMA, + E33_RESULT_SCHEMA, + E33WorkerShadowError, + E33WorkerShadowResult, + read_e33_worker_shadow_result, + run_e33_worker_shadow, +) from .evaluation_pack import ( ANNOTATION_CONTRACT_SCHEMA, EVALUATION_PACK_SCHEMA, @@ -372,6 +385,17 @@ __all__ = [ "E32_TRACK_GEOMETRY_RECORD_SCHEMA", "E32TrackGeometryReplay", "E32TrackGeometryReplayError", + "E33_FRAME_SCHEMA", + "E33_OUTCOME_SCHEMA", + "E33_PACKAGE_SCHEMA", + "E33_PROFILE_SCHEMA", + "E33_REPORT_SCHEMA", + "E33_RESOURCE_SCHEMA", + "E33_RESULT_SCHEMA", + "E33WorkerShadowError", + "E33WorkerShadowResult", + "read_e33_worker_shadow_result", + "run_e33_worker_shadow", "build_lidar_ground_annotation_template", "build_lidar_ground_benchmark", "build_k1_local_surface", diff --git a/src/k1link/compute/e33_worker_shadow.py b/src/k1link/compute/e33_worker_shadow.py new file mode 100644 index 0000000..d4fdd9c --- /dev/null +++ b/src/k1link/compute/e33_worker_shadow.py @@ -0,0 +1,1290 @@ +"""Recorded-source-paced worker shadow for immutable E32 TrackGeometry. + +The E33 runtime is deliberately downstream of E32. It does not rerun or tune +the detector, semantics, calibration, or geometry policy. It verifies one +content-addressed E32 package, releases every TrackGeometry frame on the +recorded timeline, and exercises bounded work/result channels with complete +per-frame outcome and resource accounting. +""" + +from __future__ import annotations + +import hashlib +import json +import math +import os +import platform +import re +import resource +import shutil +import subprocess +import sys +import threading +import time +from collections import Counter +from collections.abc import Mapping, Sequence +from dataclasses import dataclass +from datetime import UTC, datetime +from pathlib import Path +from typing import Any, Final + +import numpy as np +import numpy.typing as npt + +from .live_perception import LatestWinsQueue + +E33_PACKAGE_SCHEMA: Final = "missioncore.e33-worker-package/v1" +E33_PROFILE_SCHEMA: Final = "missioncore.e33-worker-shadow-profile/v1" +E33_RESULT_SCHEMA: Final = "missioncore.e33-worker-shadow-result/v1" +E33_REPORT_SCHEMA: Final = "missioncore.e33-worker-shadow-report/v1" +E33_FRAME_SCHEMA: Final = "missioncore.e33-worker-shadow-frame/v1" +E33_OUTCOME_SCHEMA: Final = "missioncore.e33-worker-shadow-outcome/v1" +E33_RESOURCE_SCHEMA: Final = "missioncore.e33-worker-resource-sample/v1" + +_E32_RESULT_SCHEMA = "missioncore.e32-track-geometry-replay/v1" +_E32_RECORD_SCHEMA = "missioncore.e32-track-geometry-record/v1" +_E32_RESULT_ID = re.compile(r"^e32-track-geometry-[a-f0-9]{64}$") +_E33_PACKAGE_ID = re.compile(r"^e33-worker-package-[a-f0-9]{64}$") +_E33_RESULT_ID = re.compile(r"^e33-worker-shadow-[a-f0-9]{64}$") +_EXECUTION_ID = re.compile(r"^[a-z0-9][a-z0-9._-]{2,95}$") +_SHA256 = re.compile(r"^[a-f0-9]{64}$") + +_FRAMES_NAME = "track-geometry-frames.jsonl" +_OFFSETS_NAME = "frame-point-offsets.npy" +_SOURCE_INDICES_NAME = "point-source-indices.npy" +_POINTS_NAME = "point-coordinates-map-f32.npy" +_OWNERS_NAME = "point-owner-indices.npy" + + +class E33WorkerShadowError(RuntimeError): + """The E33 package, runtime, accounting, or immutable result is invalid.""" + + +@dataclass(frozen=True, slots=True) +class E33WorkerShadowResult: + """One validated E33 worker-shadow result.""" + + result_root: Path + result_id: str + result: dict[str, Any] + report: dict[str, Any] + + @property + def accepted(self) -> bool: + return bool(self.report["acceptance"]["accepted"]) + + +@dataclass(frozen=True, slots=True) +class _InputEnvelope: + frame_index: int + source_frame_index: int + session_seconds: float + source_available: bool + record: dict[str, Any] + row_start: int + row_end: int + scheduled_monotonic: float + enqueued_monotonic: float + release_lag_ms: float + + +@dataclass(frozen=True, slots=True) +class _OutputEnvelope: + frame_index: int + source_frame_index: int + session_seconds: float + source_available: bool + scheduled_monotonic: float + input_enqueued_monotonic: float + work_dequeued_monotonic: float + result_enqueued_monotonic: float + release_lag_ms: float + work_wait_ms: float + processing_ms: float + document: dict[str, Any] + + +def run_e33_worker_shadow( + package_root: Path, + output_root: Path, + *, + execution_id: str, +) -> E33WorkerShadowResult: + """Run or verify one immutable recorded-source-paced E33 qualification.""" + + if _EXECUTION_ID.fullmatch(execution_id) is None: + raise E33WorkerShadowError("E33 execution ID is invalid") + package, package_artifacts = _validate_package(package_root) + profile = _read_profile(package_artifacts["profile.json"]) + e32_root = _e32_root(package_root, package) + e32_manifest, e32_artifacts = _validate_e32_input(e32_root, profile) + frame_count = _positive_int( + _object(e32_manifest.get("identity"), "E32 identity").get("frame_count"), + "E32 frame count", + ) + worker = _worker_identity(profile) + identity = { + "schema_version": E33_RESULT_SCHEMA, + "execution_id": execution_id, + "pipeline": "e32-track-geometry-recorded-source-paced-worker-shadow/v1", + "package_id": package["package_id"], + "package_identity_sha256": package["identity_sha256"], + "e32_result_id": e32_manifest["result_id"], + "e32_identity_sha256": e32_manifest["identity_sha256"], + "frame_count": frame_count, + "timeline_start_seconds": _object( + e32_manifest.get("identity"), + "E32 identity", + )["timeline_start_seconds"], + "timeline_end_seconds": _object( + e32_manifest.get("identity"), + "E32 identity", + )["timeline_end_seconds"], + "profile": profile, + "worker": worker, + "upstream": _object(e32_manifest["identity"], "E32 identity")["source"], + "authority": _authority(), + } + identity_sha256 = hashlib.sha256(_canonical_json(identity)).hexdigest() + result_id = f"e33-worker-shadow-{identity_sha256}" + destination = output_root.expanduser().absolute() + destination.mkdir(mode=0o700, parents=True, exist_ok=True) + result_root = destination / result_id + if result_root.exists(): + return read_e33_worker_shadow_result(result_root) + staging = destination / f".{result_id}.{os.getpid()}.incomplete" + staging.mkdir(mode=0o700, exist_ok=False) + try: + report = _execute( + staging=staging, + result_id=result_id, + identity=identity, + profile=profile, + frame_count=frame_count, + e32_artifacts=e32_artifacts, + ) + _write_json(staging / "run-report.json", report) + artifacts = [ + _artifact(staging / "shadow-frames.jsonl", "shadow-frames", E33_FRAME_SCHEMA), + _artifact( + staging / "frame-outcomes.jsonl", + "frame-outcomes", + E33_OUTCOME_SCHEMA, + ), + _artifact( + staging / "resource-telemetry.jsonl", + "resource-telemetry", + E33_RESOURCE_SCHEMA, + ), + _artifact(staging / "run-report.json", "run-report", E33_REPORT_SCHEMA), + ] + result = { + "schema_version": E33_RESULT_SCHEMA, + "result_id": result_id, + "identity_sha256": identity_sha256, + "identity": identity, + "created_at_utc": report["created_at_utc"], + "acceptance_state": ( + "accepted-recorded-source-paced-shadow" + if report["acceptance"]["accepted"] + else "rejected-fail-closed" + ), + "publication_scope": "worker-shadow-diagnostic-only", + "ground_truth": False, + "artifacts": artifacts, + "authority": _authority(), + } + _write_json(staging / "result.json", result) + os.replace(staging, result_root) + except BaseException: + shutil.rmtree(staging, ignore_errors=True) + raise + return read_e33_worker_shadow_result(result_root) + + +def read_e33_worker_shadow_result(root: Path) -> E33WorkerShadowResult: + """Read and fully validate one immutable E33 result.""" + + resolved = root.expanduser().resolve(strict=True) + if not resolved.is_dir() or _E33_RESULT_ID.fullmatch(resolved.name) is None: + raise E33WorkerShadowError("E33 result root is invalid") + result = _read_json(resolved / "result.json") + identity = _object(result.get("identity"), "E33 result identity") + identity_sha256 = result.get("identity_sha256") + if ( + result.get("schema_version") != E33_RESULT_SCHEMA + or result.get("result_id") != resolved.name + or not isinstance(identity_sha256, str) + or _SHA256.fullmatch(identity_sha256) is None + or resolved.name != f"e33-worker-shadow-{identity_sha256}" + or hashlib.sha256(_canonical_json(identity)).hexdigest() != identity_sha256 + or result.get("publication_scope") != "worker-shadow-diagnostic-only" + or result.get("ground_truth") is not False + or result.get("authority") != _authority() + ): + raise E33WorkerShadowError("E33 result identity is inconsistent") + artifacts = _verified_result_artifacts(resolved, result.get("artifacts")) + report = _read_json(artifacts["run-report"]) + metrics = _object(report.get("metrics"), "E33 report metrics") + accounting = _object(metrics.get("accounting"), "E33 accounting") + frame_count = _positive_int(identity.get("frame_count"), "E33 frame count") + if ( + report.get("schema_version") != E33_REPORT_SCHEMA + or report.get("result_id") != resolved.name + or report.get("identity") != identity + or not isinstance(report.get("acceptance"), dict) + or result.get("acceptance_state") + != ( + "accepted-recorded-source-paced-shadow" + if report["acceptance"].get("accepted") + else "rejected-fail-closed" + ) + or accounting.get("source_frames") != frame_count + or accounting.get("source_frames") + != ( + _nonnegative_int(accounting.get("delivered"), "delivered frames") + + _nonnegative_int( + accounting.get("input_superseded"), + "input superseded frames", + ) + + _nonnegative_int( + accounting.get("result_superseded"), + "result superseded frames", + ) + ) + ): + raise E33WorkerShadowError("E33 report accounting is inconsistent") + _validate_outcomes( + artifacts["frame-outcomes"], + frame_count=frame_count, + accounting=accounting, + ) + _validate_jsonl_schema( + artifacts["shadow-frames"], + E33_FRAME_SCHEMA, + expected_count=_nonnegative_int(accounting.get("delivered"), "delivered"), + ) + resource_count = _validate_jsonl_schema( + artifacts["resource-telemetry"], + E33_RESOURCE_SCHEMA, + expected_count=None, + ) + if resource_count < 1: + raise E33WorkerShadowError("E33 resource telemetry is empty") + return E33WorkerShadowResult( + result_root=resolved, + result_id=resolved.name, + result=result, + report=report, + ) + + +def _execute( + *, + staging: Path, + result_id: str, + identity: dict[str, Any], + profile: dict[str, Any], + frame_count: int, + e32_artifacts: Mapping[str, Path], +) -> dict[str, Any]: + offsets, source_indices, points, owners = _load_point_arrays( + e32_artifacts, + frame_count=frame_count, + ) + work_queue = LatestWinsQueue[_InputEnvelope]( + _positive_int(profile["queues"]["work_capacity"], "work queue capacity") + ) + result_queue = LatestWinsQueue[_OutputEnvelope]( + _positive_int(profile["queues"]["result_capacity"], "result queue capacity") + ) + errors: list[BaseException] = [] + outcomes: dict[int, dict[str, Any]] = {} + outcome_lock = threading.Lock() + release_lag_ms: list[float] = [] + work_wait_ms: list[float] = [] + processing_ms: list[float] = [] + result_wait_ms: list[float] = [] + result_age_ms: list[float] = [] + resource_sampler = _ResourceSampler( + staging / "resource-telemetry.jsonl", + interval_seconds=float(profile["resources"]["sample_interval_seconds"]), + ) + frames_path = e32_artifacts["track-geometry-frames"] + speed = float(profile["pacing"]["speed"]) + start_delay = float(profile["pacing"]["start_delay_ms"]) / 1_000.0 + test_delay = float(profile.get("test_controls", {}).get("processing_delay_ms", 0.0)) + recorded_start = float(identity["timeline_start_seconds"]) + replay_started = time.perf_counter() + start_delay + + def record_dropped(envelope: _InputEnvelope | _OutputEnvelope, status: str) -> None: + with outcome_lock: + if envelope.frame_index in outcomes: + raise E33WorkerShadowError("E33 frame received multiple terminal outcomes") + outcomes[envelope.frame_index] = { + "schema_version": E33_OUTCOME_SCHEMA, + "frame_index": envelope.frame_index, + "source_frame_index": envelope.source_frame_index, + "session_seconds": envelope.session_seconds, + "source_available": envelope.source_available, + "status": status, + "health": "unavailable", + "reason": status, + "result_age_ms": None, + } + + def produce() -> None: + try: + previous_seconds: float | None = None + with frames_path.open("r", encoding="utf-8") as stream: + for frame_index, line in enumerate(stream): + if frame_index >= frame_count: + raise E33WorkerShadowError("E32 frame stream has extra rows") + record = _frame_record(line, expected_frame_index=frame_index) + session_seconds = _finite_number( + record.get("session_seconds"), + "E32 session seconds", + ) + if previous_seconds is not None and session_seconds <= previous_seconds: + raise E33WorkerShadowError("E32 session timeline is not increasing") + previous_seconds = session_seconds + scheduled = replay_started + (session_seconds - recorded_start) / speed + remaining = scheduled - time.perf_counter() + if remaining > 0: + time.sleep(remaining) + enqueued = time.perf_counter() + envelope = _InputEnvelope( + frame_index=frame_index, + source_frame_index=_nonnegative_int( + record.get("source_frame_index"), + "source frame index", + ), + session_seconds=session_seconds, + source_available=_boolean( + record.get("source_available"), + "source availability", + ), + record=record, + row_start=int(offsets[frame_index]), + row_end=int(offsets[frame_index + 1]), + scheduled_monotonic=scheduled, + enqueued_monotonic=enqueued, + release_lag_ms=max(0.0, (enqueued - scheduled) * 1_000.0), + ) + dropped = work_queue.publish(envelope) + if dropped is not None: + record_dropped(dropped, "input-superseded") + if frame_index + 1 != frame_count: + raise E33WorkerShadowError("E32 frame stream is incomplete") + except BaseException as exc: # noqa: BLE001 - relayed to owner thread + errors.append(exc) + finally: + work_queue.close() + + def process() -> None: + try: + while True: + envelope = work_queue.take_next(timeout=0.5) + if envelope is None: + if work_queue.snapshot().closed: + break + continue + dequeued = time.perf_counter() + if test_delay > 0: + time.sleep(test_delay / 1_000.0) + row = slice(envelope.row_start, envelope.row_end) + document = _shadow_frame( + envelope, + source_indices=source_indices[row], + points=points[row], + owners=owners[row], + ) + completed = time.perf_counter() + output = _OutputEnvelope( + frame_index=envelope.frame_index, + source_frame_index=envelope.source_frame_index, + session_seconds=envelope.session_seconds, + source_available=envelope.source_available, + scheduled_monotonic=envelope.scheduled_monotonic, + input_enqueued_monotonic=envelope.enqueued_monotonic, + work_dequeued_monotonic=dequeued, + result_enqueued_monotonic=completed, + release_lag_ms=envelope.release_lag_ms, + work_wait_ms=(dequeued - envelope.enqueued_monotonic) * 1_000.0, + processing_ms=(completed - dequeued) * 1_000.0, + document=document, + ) + dropped = result_queue.publish(output) + if dropped is not None: + record_dropped(dropped, "result-superseded") + except BaseException as exc: # noqa: BLE001 - relayed to owner thread + errors.append(exc) + finally: + result_queue.close() + + def publish() -> None: + try: + with (staging / "shadow-frames.jsonl").open( + "x", + encoding="utf-8", + newline="\n", + ) as stream: + while True: + envelope = result_queue.take_next(timeout=0.5) + if envelope is None: + if result_queue.snapshot().closed: + break + continue + published = time.perf_counter() + age_ms = (published - envelope.scheduled_monotonic) * 1_000.0 + health = _health(age_ms, profile["deadlines"]) + document = dict(envelope.document) + document["delivery"] = { + "health": health, + "release_lag_ms": envelope.release_lag_ms, + "work_wait_ms": envelope.work_wait_ms, + "processing_ms": envelope.processing_ms, + "result_wait_ms": ( + published - envelope.result_enqueued_monotonic + ) + * 1_000.0, + "result_age_ms": age_ms, + } + _write_jsonl(stream, document) + with outcome_lock: + if envelope.frame_index in outcomes: + raise E33WorkerShadowError( + "E33 frame received multiple terminal outcomes" + ) + outcomes[envelope.frame_index] = { + "schema_version": E33_OUTCOME_SCHEMA, + "frame_index": envelope.frame_index, + "source_frame_index": envelope.source_frame_index, + "session_seconds": envelope.session_seconds, + "source_available": envelope.source_available, + "status": "delivered", + "health": health, + "reason": None if health == "healthy" else "result-age", + "result_age_ms": age_ms, + } + release_lag_ms.append(envelope.release_lag_ms) + work_wait_ms.append(envelope.work_wait_ms) + processing_ms.append(envelope.processing_ms) + result_wait_ms.append( + (published - envelope.result_enqueued_monotonic) * 1_000.0 + ) + result_age_ms.append(age_ms) + except BaseException as exc: # noqa: BLE001 - relayed to owner thread + errors.append(exc) + + resource_sampler.start() + producer = threading.Thread(target=produce, name="e33-source-pacer", daemon=True) + processor = threading.Thread(target=process, name="e33-track-geometry", daemon=True) + publisher = threading.Thread(target=publish, name="e33-result-publisher", daemon=True) + producer.start() + processor.start() + publisher.start() + producer.join() + processor.join() + publisher.join() + resource_sampler.stop() + if errors: + raise E33WorkerShadowError("E33 worker thread failed") from errors[0] + if len(outcomes) != frame_count: + raise E33WorkerShadowError("E33 frame outcome accounting is incomplete") + + ordered_outcomes = [outcomes[index] for index in range(frame_count)] + with (staging / "frame-outcomes.jsonl").open( + "x", + encoding="utf-8", + newline="\n", + ) as stream: + for outcome in ordered_outcomes: + _write_jsonl(stream, outcome) + + finished = time.perf_counter() + source_span = float(identity["timeline_end_seconds"]) - recorded_start + replay_wall = max(0.0, finished - replay_started) + status_counts = Counter(str(item["status"]) for item in ordered_outcomes) + health_counts = Counter(str(item["health"]) for item in ordered_outcomes) + work = work_queue.snapshot() + results = result_queue.snapshot() + delivered = status_counts["delivered"] + metrics = { + "accounting": { + "source_frames": frame_count, + "delivered": delivered, + "input_superseded": status_counts["input-superseded"], + "result_superseded": status_counts["result-superseded"], + "closed": sum(status_counts.values()) == frame_count, + }, + "health_counts": dict(sorted(health_counts.items())), + "source_span_seconds": source_span, + "replay_wall_seconds": replay_wall, + "wall_to_ideal_ratio": replay_wall / (source_span / speed), + "effective_delivery_fps": delivered / source_span, + "deadline_miss_fraction": ( + sum( + count + for state, count in health_counts.items() + if state != "healthy" + ) + / frame_count + ), + "work_queue": _queue_metrics(work), + "result_queue": _queue_metrics(results), + "release_lag_ms": _distribution(release_lag_ms), + "work_wait_ms": _distribution(work_wait_ms), + "processing_ms": _distribution(processing_ms), + "result_wait_ms": _distribution(result_wait_ms), + "result_age_ms": _distribution(result_age_ms), + "resources": resource_sampler.summary(), + } + acceptance = _acceptance(profile, metrics) + return { + "schema_version": E33_REPORT_SCHEMA, + "result_id": result_id, + "created_at_utc": datetime.now(UTC) + .isoformat(timespec="milliseconds") + .replace("+00:00", "Z"), + "identity": identity, + "metrics": metrics, + "acceptance": acceptance, + "authority": _authority(), + } + + +class _ResourceSampler: + def __init__(self, path: Path, *, interval_seconds: float) -> None: + if not 0.1 <= interval_seconds <= 10.0: + raise E33WorkerShadowError("E33 resource sample interval is invalid") + self._path = path + self._interval = interval_seconds + self._stop = threading.Event() + self._thread = threading.Thread( + target=self._sample, + name="e33-resource-sampler", + daemon=True, + ) + self._samples: list[dict[str, Any]] = [] + self._error: BaseException | None = None + + def start(self) -> None: + self._thread.start() + + def stop(self) -> None: + self._stop.set() + self._thread.join(timeout=max(5.0, self._interval * 2)) + if self._thread.is_alive(): + raise E33WorkerShadowError("E33 resource sampler did not stop") + if self._error is not None: + raise E33WorkerShadowError("E33 resource sampler failed") from self._error + if not self._samples: + raise E33WorkerShadowError("E33 resource sampler emitted no samples") + + def summary(self) -> dict[str, Any]: + rss = [float(item["process_rss_mib"]) for item in self._samples] + cpu = [float(item["process_cpu_percent"]) for item in self._samples] + gpu = [ + float(item["gpu"]["utilization_percent"]) + for item in self._samples + if isinstance(item.get("gpu"), dict) + and item["gpu"].get("state") == "available" + ] + memory = [ + float(item["gpu"]["memory_used_mib"]) + for item in self._samples + if isinstance(item.get("gpu"), dict) + and item["gpu"].get("state") == "available" + ] + return { + "sample_count": len(self._samples), + "process_rss_mib": _distribution(rss), + "process_cpu_percent": _distribution(cpu), + "gpu_visible": bool(gpu), + "gpu_utilization_percent": _distribution(gpu), + "gpu_memory_used_mib": _distribution(memory), + } + + def _sample(self) -> None: + previous_wall = time.perf_counter() + previous_cpu = time.process_time() + try: + with self._path.open("x", encoding="utf-8", newline="\n") as stream: + while True: + now_wall = time.perf_counter() + now_cpu = time.process_time() + elapsed = max(1e-9, now_wall - previous_wall) + sample = { + "schema_version": E33_RESOURCE_SCHEMA, + "sample_index": len(self._samples), + "elapsed_seconds": now_wall, + "process_rss_mib": _process_rss_mib(), + "process_cpu_percent": ( + (now_cpu - previous_cpu) / elapsed * 100.0 + ), + "gpu": _gpu_sample(), + } + self._samples.append(sample) + _write_jsonl(stream, sample) + previous_wall = now_wall + previous_cpu = now_cpu + if self._stop.wait(self._interval): + break + except BaseException as exc: # noqa: BLE001 - relayed to owner thread + self._error = exc + + +def _shadow_frame( + envelope: _InputEnvelope, + *, + source_indices: npt.NDArray[np.int64], + points: npt.NDArray[np.float32], + owners: npt.NDArray[np.uint32], +) -> dict[str, Any]: + geometries = envelope.record.get("geometries") + slab = envelope.record.get("point_slab") + if not isinstance(geometries, list) or not isinstance(slab, dict): + raise E33WorkerShadowError("E32 frame payload is invalid") + if slab.get("row_count") != source_indices.size: + raise E33WorkerShadowError("E32 frame PointSlab row count changed") + evidence = Counter( + str(item.get("evidence_state")) + for item in geometries + if isinstance(item, dict) + ) + return { + "schema_version": E33_FRAME_SCHEMA, + "frame_index": envelope.frame_index, + "source_frame_index": envelope.source_frame_index, + "session_seconds": envelope.session_seconds, + "source_available": envelope.source_available, + "geometry_count": len(geometries), + "point_count": int(source_indices.size), + "evidence_state_counts": dict(sorted(evidence.items())), + "frame_record_sha256": hashlib.sha256( + _canonical_json(envelope.record) + ).hexdigest(), + "point_slab_digests": { + "source_indices_sha256": hashlib.sha256( + source_indices.tobytes(order="C") + ).hexdigest(), + "points_map_sha256": hashlib.sha256(points.tobytes(order="C")).hexdigest(), + "owner_indices_sha256": hashlib.sha256( + owners.tobytes(order="C") + ).hexdigest(), + }, + "delivery": None, + "authority": _authority(), + } + + +def _acceptance( + profile: Mapping[str, Any], + metrics: Mapping[str, Any], +) -> dict[str, Any]: + gate = _object(profile.get("acceptance"), "E33 acceptance profile") + accounting = _object(metrics.get("accounting"), "E33 accounting") + work = _object(metrics.get("work_queue"), "E33 work queue metrics") + results = _object(metrics.get("result_queue"), "E33 result queue metrics") + resources = _object(metrics.get("resources"), "E33 resource metrics") + source_frames = _positive_int(accounting.get("source_frames"), "source frames") + requirements = { + "complete_frame_accounting": accounting.get("closed") is True, + "work_queue_bounded": work.get("maximum_depth") + <= profile["queues"]["work_capacity"], + "result_queue_bounded": results.get("maximum_depth") + <= profile["queues"]["result_capacity"], + "input_drop_fraction_within_gate": ( + accounting.get("input_superseded", source_frames) / source_frames + <= float(gate["maximum_input_drop_fraction"]) + ), + "result_drop_fraction_within_gate": ( + accounting.get("result_superseded", source_frames) / source_frames + <= float(gate["maximum_result_drop_fraction"]) + ), + "delivery_rate_within_gate": float(metrics["effective_delivery_fps"]) + >= float(gate["minimum_effective_fps"]), + "result_age_p95_within_gate": _required_distribution_value( + metrics["result_age_ms"], + "p95", + ) + <= float(gate["maximum_result_age_p95_ms"]), + "deadline_miss_fraction_within_gate": float( + metrics["deadline_miss_fraction"] + ) + <= float(gate["maximum_deadline_miss_fraction"]), + "producer_release_p95_within_gate": _required_distribution_value( + metrics["release_lag_ms"], + "p95", + ) + <= float(gate["maximum_release_lag_p95_ms"]), + "process_rss_within_gate": _required_distribution_value( + resources["process_rss_mib"], + "maximum", + ) + <= float(gate["maximum_process_rss_mib"]), + "gpu_visible": ( + resources.get("gpu_visible") is True + if gate.get("require_gpu_visible") is True + else True + ), + "recorded_speed_is_one": float(profile["pacing"]["speed"]) == 1.0, + "navigation_or_safety_authority_false": True, + } + failed = [name for name, passed in requirements.items() if not passed] + return { + "accepted": not failed, + "requirements": requirements, + "rejection_reasons": failed, + "scope": "recorded-source-paced-track-geometry-worker-shadow", + "commands_enabled": False, + "navigation_or_safety_accepted": False, + } + + +def _validate_package(root: Path) -> tuple[dict[str, Any], dict[str, Path]]: + resolved = root.expanduser().resolve(strict=True) + manifest = _read_json(resolved / "manifest.json") + identity = _object(manifest.get("identity"), "E33 package identity") + identity_sha256 = manifest.get("identity_sha256") + package_id = manifest.get("package_id") + if ( + manifest.get("schema_version") != E33_PACKAGE_SCHEMA + or not isinstance(identity_sha256, str) + or _SHA256.fullmatch(identity_sha256) is None + or hashlib.sha256(_canonical_json(identity)).hexdigest() != identity_sha256 + or package_id != f"e33-worker-package-{identity_sha256}" + or not isinstance(package_id, str) + or _E33_PACKAGE_ID.fullmatch(package_id) is None + or resolved.name != package_id + ): + raise E33WorkerShadowError("E33 package identity is invalid") + artifacts = _verified_artifacts(resolved, manifest.get("artifacts")) + if set(artifacts) != set(identity.get("artifact_paths", [])): + raise E33WorkerShadowError("E33 package artifact set changed") + return manifest, artifacts + + +def _read_profile(path: Path) -> dict[str, Any]: + profile = _read_json(path) + queues = profile.get("queues") + pacing = profile.get("pacing") + deadlines = profile.get("deadlines") + resources = profile.get("resources") + acceptance = profile.get("acceptance") + authority = profile.get("authority") + if not all( + isinstance(item, dict) + for item in (queues, pacing, deadlines, resources, acceptance) + ): + raise E33WorkerShadowError("E33 profile sections are invalid") + queue_profile = _object(queues, "E33 queue profile") + pacing_profile = _object(pacing, "E33 pacing profile") + deadline_profile = _object(deadlines, "E33 deadline profile") + resource_profile = _object(resources, "E33 resource profile") + if ( + profile.get("schema_version") != E33_PROFILE_SCHEMA + or profile.get("mode") + not in {"full-session-qualification", "contract-test"} + or authority != _authority() + or not 1 + <= _positive_int(queue_profile["work_capacity"], "work capacity") + <= 8 + or not 1 + <= _positive_int(queue_profile["result_capacity"], "result capacity") + <= 8 + or not 0.1 <= float(pacing_profile["speed"]) <= 100.0 + or not 0 <= float(pacing_profile["start_delay_ms"]) <= 5_000 + or not 0 + < float(deadline_profile["result_ms"]) + < float(deadline_profile["stale_ms"]) + or not float(deadline_profile["stale_ms"]) + < float(deadline_profile["unavailable_ms"]) + or not 0.1 <= float(resource_profile["sample_interval_seconds"]) <= 10.0 + or not isinstance(profile.get("expected_e32_result_id"), str) + or _E32_RESULT_ID.fullmatch(profile["expected_e32_result_id"]) is None + or not isinstance(profile.get("expected_worker_node"), str) + or not profile["expected_worker_node"] + ): + raise E33WorkerShadowError("E33 profile contract is invalid") + test_controls = profile.get("test_controls", {}) + if ( + not isinstance(test_controls, dict) + or not 0 <= float(test_controls.get("processing_delay_ms", 0.0)) <= 5_000 + or ( + profile["mode"] == "full-session-qualification" + and ( + float(pacing_profile["speed"]) != 1.0 + or float(test_controls.get("processing_delay_ms", 0.0)) != 0.0 + ) + ) + ): + raise E33WorkerShadowError("E33 test controls are invalid") + return profile + + +def _validate_e32_input( + root: Path, + profile: Mapping[str, Any], +) -> tuple[dict[str, Any], dict[str, Path]]: + manifest = _read_json(root / "manifest.json") + result_id = manifest.get("result_id") + identity = _object(manifest.get("identity"), "E32 identity") + identity_sha256 = manifest.get("identity_sha256") + if ( + manifest.get("schema_version") != _E32_RESULT_SCHEMA + or result_id != root.name + or result_id != profile["expected_e32_result_id"] + or not isinstance(identity_sha256, str) + or _SHA256.fullmatch(identity_sha256) is None + or result_id != f"e32-track-geometry-{identity_sha256}" + or hashlib.sha256(_canonical_json(identity)).hexdigest() != identity_sha256 + or manifest.get("authority") != _authority() + ): + raise E33WorkerShadowError("E32 worker input identity is invalid") + raw_artifacts = _verified_artifacts(root, manifest.get("artifacts")) + storage_names = { + "track-geometry-frames": _FRAMES_NAME, + "frame-point-offsets": _OFFSETS_NAME, + "point-source-indices": _SOURCE_INDICES_NAME, + "point-coordinates-map-f32": _POINTS_NAME, + "point-owner-indices": _OWNERS_NAME, + } + artifacts: dict[str, Path] = {} + for kind, name in storage_names.items(): + matches = { + path + for key, path in raw_artifacts.items() + if key in {kind, name} or path.name == name + } + if len(matches) != 1: + raise E33WorkerShadowError( + f"E32 worker input artifact is not unique: {kind}" + ) + artifacts[kind] = matches.pop() + return manifest, artifacts + + +def _load_point_arrays( + artifacts: Mapping[str, Path], + *, + frame_count: int, +) -> tuple[ + npt.NDArray[np.int64], + npt.NDArray[np.int64], + npt.NDArray[np.float32], + npt.NDArray[np.uint32], +]: + try: + offsets = np.load( + artifacts["frame-point-offsets"], + allow_pickle=False, + mmap_mode="r", + ) + source_indices = np.load( + artifacts["point-source-indices"], + allow_pickle=False, + mmap_mode="r", + ) + points = np.load( + artifacts["point-coordinates-map-f32"], + allow_pickle=False, + mmap_mode="r", + ) + owners = np.load( + artifacts["point-owner-indices"], + allow_pickle=False, + mmap_mode="r", + ) + except (OSError, ValueError, KeyError) as exc: + raise E33WorkerShadowError("E32 worker point storage is unreadable") from exc + if ( + offsets.dtype != np.dtype(" dict[str, Any]: + try: + value = json.loads(line) + except json.JSONDecodeError as exc: + raise E33WorkerShadowError("E32 frame record is invalid JSON") from exc + if ( + not isinstance(value, dict) + or value.get("schema_version") != _E32_RECORD_SCHEMA + or value.get("frame_index") != expected_frame_index + or not isinstance(value.get("geometries"), list) + or not isinstance(value.get("point_slab"), dict) + or value.get("authority") != _authority() + ): + raise E33WorkerShadowError("E32 frame record contract changed") + return value + + +def _worker_identity(profile: Mapping[str, Any]) -> dict[str, Any]: + expected_image = profile.get("container_image") + actual_image = os.environ.get("E33_CONTAINER_IMAGE") + expected_worker_node = profile.get("expected_worker_node") + actual_worker_node = os.environ.get("E33_WORKER_NODE") + if not isinstance(expected_image, str) or not expected_image: + raise E33WorkerShadowError("E33 container image identity is missing") + if actual_image != expected_image: + raise E33WorkerShadowError("E33 container image identity changed") + if ( + not isinstance(expected_worker_node, str) + or not expected_worker_node + or actual_worker_node != expected_worker_node + ): + raise E33WorkerShadowError("E33 worker node identity changed") + return { + "worker_node": actual_worker_node, + "container_hostname": platform.node(), + "system": platform.system(), + "machine": platform.machine(), + "python": platform.python_version(), + "numpy": np.__version__, + "container_image": actual_image, + } + + +def _e32_root(package_root: Path, package: Mapping[str, Any]) -> Path: + identity = _object(package.get("identity"), "E33 package identity") + result_id = identity.get("e32_result_id") + if not isinstance(result_id, str) or _E32_RESULT_ID.fullmatch(result_id) is None: + raise E33WorkerShadowError("E33 package E32 result ID is invalid") + return package_root.resolve(strict=True) / "input" / "e32" / result_id + + +def _health(age_ms: float, deadlines: Mapping[str, Any]) -> str: + if age_ms <= float(deadlines["result_ms"]): + return "healthy" + if age_ms <= float(deadlines["stale_ms"]): + return "late" + if age_ms <= float(deadlines["unavailable_ms"]): + return "stale" + return "unavailable" + + +def _gpu_sample() -> dict[str, Any]: + try: + completed = subprocess.run( + [ + "nvidia-smi", + "--query-gpu=name,utilization.gpu,memory.used,power.draw,temperature.gpu", + "--format=csv,noheader,nounits", + ], + check=False, + capture_output=True, + text=True, + timeout=3, + ) + except (OSError, subprocess.SubprocessError): + return {"state": "unavailable"} + lines = [line.strip() for line in completed.stdout.splitlines() if line.strip()] + if completed.returncode != 0 or len(lines) != 1: + return {"state": "unavailable"} + parts = [item.strip() for item in lines[0].split(",")] + if len(parts) != 5: + return {"state": "unavailable"} + try: + return { + "state": "available", + "name": parts[0], + "utilization_percent": float(parts[1]), + "memory_used_mib": float(parts[2]), + "power_w": float(parts[3]), + "temperature_c": float(parts[4]), + } + except ValueError: + return {"state": "unavailable"} + + +def _process_rss_mib() -> float: + proc_statm = Path("/proc/self/statm") + if proc_statm.is_file(): + fields = proc_statm.read_text(encoding="ascii").split() + if len(fields) >= 2: + return float(int(fields[1]) * os.sysconf("SC_PAGE_SIZE") / (1024 * 1024)) + value = resource.getrusage(resource.RUSAGE_SELF).ru_maxrss + return float(value / (1024 * 1024) if sys.platform == "darwin" else value / 1024) + + +def _queue_metrics(snapshot: Any) -> dict[str, Any]: + return { + "capacity": snapshot.capacity, + "depth": snapshot.depth, + "maximum_depth": snapshot.maximum_depth, + "published": snapshot.published, + "consumed": snapshot.consumed, + "dropped_overflow": snapshot.dropped_overflow, + "dropped_superseded": snapshot.dropped_superseded, + "closed": snapshot.closed, + } + + +def _distribution(values: Sequence[float]) -> dict[str, float | int | None]: + if not values: + return { + "sample_count": 0, + "minimum": None, + "mean": None, + "p50": None, + "p95": None, + "maximum": None, + } + array = np.asarray(values, dtype=np.float64) + if not np.isfinite(array).all(): + raise E33WorkerShadowError("E33 metric distribution is non-finite") + return { + "sample_count": int(array.size), + "minimum": float(np.min(array)), + "mean": float(np.mean(array)), + "p50": float(np.percentile(array, 50)), + "p95": float(np.percentile(array, 95)), + "maximum": float(np.max(array)), + } + + +def _required_distribution_value(value: object, key: str) -> float: + distribution = _object(value, "E33 metric distribution") + candidate = distribution.get(key) + if ( + not isinstance(candidate, (int, float)) + or isinstance(candidate, bool) + or not math.isfinite(float(candidate)) + ): + raise E33WorkerShadowError("E33 metric distribution is incomplete") + return float(candidate) + + +def _validate_outcomes( + path: Path, + *, + frame_count: int, + accounting: Mapping[str, Any], +) -> None: + counts: Counter[str] = Counter() + with path.open("r", encoding="utf-8") as stream: + for expected_index, line in enumerate(stream): + value = _json_line(line, "E33 outcome") + if ( + expected_index >= frame_count + or value.get("schema_version") != E33_OUTCOME_SCHEMA + or value.get("frame_index") != expected_index + or value.get("status") + not in {"delivered", "input-superseded", "result-superseded"} + or value.get("health") + not in {"healthy", "late", "stale", "unavailable"} + ): + raise E33WorkerShadowError("E33 outcome stream changed") + counts[str(value["status"])] += 1 + if ( + sum(counts.values()) != frame_count + or counts["delivered"] != accounting.get("delivered") + or counts["input-superseded"] != accounting.get("input_superseded") + or counts["result-superseded"] != accounting.get("result_superseded") + ): + raise E33WorkerShadowError("E33 outcome accounting changed") + + +def _validate_jsonl_schema( + path: Path, + schema: str, + *, + expected_count: int | None, +) -> int: + count = 0 + with path.open("r", encoding="utf-8") as stream: + for line in stream: + value = _json_line(line, schema) + if value.get("schema_version") != schema: + raise E33WorkerShadowError("E33 artifact row schema changed") + count += 1 + if expected_count is not None and count != expected_count: + raise E33WorkerShadowError("E33 artifact row count changed") + return count + + +def _verified_result_artifacts( + root: Path, + raw: object, +) -> dict[str, Path]: + expected = { + "shadow-frames": ("shadow-frames.jsonl", E33_FRAME_SCHEMA), + "frame-outcomes": ("frame-outcomes.jsonl", E33_OUTCOME_SCHEMA), + "resource-telemetry": ("resource-telemetry.jsonl", E33_RESOURCE_SCHEMA), + "run-report": ("run-report.json", E33_REPORT_SCHEMA), + } + artifacts = _verified_artifacts(root, raw) + if set(artifacts) != set(expected): + raise E33WorkerShadowError("E33 result artifact set changed") + if not isinstance(raw, list): + raise E33WorkerShadowError("E33 result artifact descriptors are invalid") + descriptors = { + str(item["kind"]): item + for item in raw + if isinstance(item, dict) and isinstance(item.get("kind"), str) + } + for kind, (name, schema) in expected.items(): + descriptor = descriptors[kind] + if descriptor.get("path") != name or descriptor.get("schema_version") != schema: + raise E33WorkerShadowError("E33 result artifact descriptor changed") + return artifacts + + +def _verified_artifacts(root: Path, raw: object) -> dict[str, Path]: + if not isinstance(raw, list) or not raw: + raise E33WorkerShadowError("artifact list is invalid") + artifacts: dict[str, Path] = {} + for item in raw: + if not isinstance(item, dict): + raise E33WorkerShadowError("artifact descriptor is invalid") + kind = item.get("kind", item.get("path")) + relative = item.get("path") + if ( + not isinstance(kind, str) + or not isinstance(relative, str) + or kind in artifacts + or Path(relative).is_absolute() + or ".." in Path(relative).parts + ): + raise E33WorkerShadowError("artifact path is invalid") + path = (root / relative).resolve(strict=True) + if ( + not path.is_relative_to(root) + or not path.is_file() + or path.is_symlink() + or item.get("byte_length") != path.stat().st_size + or not isinstance(item.get("sha256"), str) + or _SHA256.fullmatch(item["sha256"]) is None + or _sha256(path) != item["sha256"] + ): + raise E33WorkerShadowError("artifact identity changed") + artifacts[kind] = path + return artifacts + + +def _artifact(path: Path, kind: str, schema: str) -> dict[str, Any]: + return { + "kind": kind, + "path": path.name, + "media_type": ( + "application/json" if path.suffix == ".json" else "application/x-ndjson" + ), + "schema_version": schema, + "byte_length": path.stat().st_size, + "sha256": _sha256(path), + } + + +def _authority() -> dict[str, bool]: + return { + "commands_enabled": False, + "navigation_or_safety_accepted": False, + } + + +def _canonical_json(value: object) -> bytes: + return json.dumps( + value, + ensure_ascii=False, + sort_keys=True, + separators=(",", ":"), + allow_nan=False, + ).encode() + + +def _sha256(path: Path) -> str: + digest = hashlib.sha256() + with path.open("rb") as stream: + while chunk := stream.read(1024 * 1024): + digest.update(chunk) + return digest.hexdigest() + + +def _read_json(path: Path) -> dict[str, Any]: + try: + value = json.loads(path.read_text(encoding="utf-8-sig")) + except (OSError, UnicodeError, json.JSONDecodeError) as exc: + raise E33WorkerShadowError(f"invalid JSON artifact: {path.name}") from exc + if not isinstance(value, dict): + raise E33WorkerShadowError(f"JSON artifact is not an object: {path.name}") + return value + + +def _json_line(line: str, label: str) -> dict[str, Any]: + try: + value = json.loads(line) + except json.JSONDecodeError as exc: + raise E33WorkerShadowError(f"{label} row is invalid JSON") from exc + if not isinstance(value, dict): + raise E33WorkerShadowError(f"{label} row is not an object") + return value + + +def _write_json(path: Path, value: object) -> None: + with path.open("x", encoding="utf-8", newline="\n") as stream: + json.dump(value, stream, ensure_ascii=False, indent=2, sort_keys=True, allow_nan=False) + stream.write("\n") + stream.flush() + os.fsync(stream.fileno()) + + +def _write_jsonl(stream: Any, value: object) -> None: + stream.write( + json.dumps( + value, + ensure_ascii=False, + sort_keys=True, + separators=(",", ":"), + allow_nan=False, + ) + + "\n" + ) + + +def _object(value: object, label: str) -> dict[str, Any]: + if not isinstance(value, dict): + raise E33WorkerShadowError(f"{label} is invalid") + return value + + +def _positive_int(value: object, label: str) -> int: + if not isinstance(value, int) or isinstance(value, bool) or value < 1: + raise E33WorkerShadowError(f"{label} is invalid") + return value + + +def _nonnegative_int(value: object, label: str) -> int: + if not isinstance(value, int) or isinstance(value, bool) or value < 0: + raise E33WorkerShadowError(f"{label} is invalid") + return value + + +def _finite_number(value: object, label: str) -> float: + if ( + not isinstance(value, (int, float)) + or isinstance(value, bool) + or not math.isfinite(float(value)) + ): + raise E33WorkerShadowError(f"{label} is invalid") + return float(value) + + +def _boolean(value: object, label: str) -> bool: + if not isinstance(value, bool): + raise E33WorkerShadowError(f"{label} is invalid") + return value diff --git a/src/k1link/compute/live_perception.py b/src/k1link/compute/live_perception.py index 033b34f..825327d 100644 --- a/src/k1link/compute/live_perception.py +++ b/src/k1link/compute/live_perception.py @@ -592,17 +592,21 @@ class LatestWinsQueue[T]: self._dropped_superseded = 0 self._closed = False - def publish(self, item: T) -> None: + def publish(self, item: T) -> T | None: + """Publish fresh work and return the evicted oldest item, if any.""" + with self._condition: if self._closed: raise RuntimeError("cannot publish to a closed latest-wins queue") + dropped: T | None = None if len(self._items) == self._capacity: - self._items.popleft() + dropped = self._items.popleft() self._dropped_overflow += 1 self._items.append(item) self._published += 1 self._maximum_depth = max(self._maximum_depth, len(self._items)) self._condition.notify() + return dropped def take_next(self, timeout: float | None = None) -> T | None: with self._condition: diff --git a/tests/test_e33_worker_shadow.py b/tests/test_e33_worker_shadow.py new file mode 100644 index 0000000..cfcc55f --- /dev/null +++ b/tests/test_e33_worker_shadow.py @@ -0,0 +1,333 @@ +from __future__ import annotations + +import hashlib +import json +from pathlib import Path +from typing import Any + +import numpy as np +import pytest + +from k1link.compute.e33_worker_shadow import ( + E33_PACKAGE_SCHEMA, + E33_PROFILE_SCHEMA, + E33WorkerShadowError, + read_e33_worker_shadow_result, + run_e33_worker_shadow, +) + +_IMAGE = "test/e33-worker@sha256:" + "a" * 64 +_AUTHORITY = { + "commands_enabled": False, + "navigation_or_safety_accepted": False, +} + + +def test_e33_worker_shadow_closes_every_frame_and_reuses_existing_result( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + package = _package(tmp_path / "package-source", frame_count=8) + monkeypatch.setenv("E33_CONTAINER_IMAGE", _IMAGE) + monkeypatch.setenv("E33_WORKER_NODE", "test-worker") + + result = run_e33_worker_shadow( + package, + tmp_path / "results", + execution_id="contract-pass", + ) + reused = run_e33_worker_shadow( + package, + tmp_path / "results", + execution_id="contract-pass", + ) + + assert result.accepted is True + assert reused.result_id == result.result_id + accounting = result.report["metrics"]["accounting"] + assert accounting == { + "source_frames": 8, + "delivered": 8, + "input_superseded": 0, + "result_superseded": 0, + "closed": True, + } + assert result.report["acceptance"]["requirements"][ + "navigation_or_safety_authority_false" + ] + assert result.result["authority"] == _AUTHORITY + + +def test_e33_worker_shadow_records_overload_instead_of_hiding_loss( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + package = _package( + tmp_path / "package-source", + frame_count=24, + interval_seconds=0.005, + processing_delay_ms=30.0, + work_capacity=1, + ) + monkeypatch.setenv("E33_CONTAINER_IMAGE", _IMAGE) + monkeypatch.setenv("E33_WORKER_NODE", "test-worker") + + result = run_e33_worker_shadow( + package, + tmp_path / "results", + execution_id="contract-overload", + ) + + accounting = result.report["metrics"]["accounting"] + assert result.accepted is False + assert accounting["source_frames"] == 24 + assert accounting["delivered"] + accounting["input_superseded"] == 24 + assert accounting["input_superseded"] > 0 + assert "input_drop_fraction_within_gate" in result.report["acceptance"][ + "rejection_reasons" + ] + outcomes = [ + json.loads(line) + for line in (result.result_root / "frame-outcomes.jsonl").read_text().splitlines() + ] + assert [row["frame_index"] for row in outcomes] == list(range(24)) + assert any(row["status"] == "input-superseded" for row in outcomes) + + +def test_e33_reader_rejects_tampered_artifact( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + package = _package(tmp_path / "package-source", frame_count=4) + monkeypatch.setenv("E33_CONTAINER_IMAGE", _IMAGE) + monkeypatch.setenv("E33_WORKER_NODE", "test-worker") + result = run_e33_worker_shadow( + package, + tmp_path / "results", + execution_id="tamper-check", + ) + with (result.result_root / "frame-outcomes.jsonl").open("a", encoding="utf-8") as stream: + stream.write("{}\n") + + with pytest.raises(E33WorkerShadowError, match="artifact identity changed"): + read_e33_worker_shadow_result(result.result_root) + + +def test_e33_rejects_wrong_container_identity( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + package = _package(tmp_path / "package-source", frame_count=4) + monkeypatch.setenv("E33_CONTAINER_IMAGE", "different") + monkeypatch.setenv("E33_WORKER_NODE", "test-worker") + + with pytest.raises(E33WorkerShadowError, match="container image identity changed"): + run_e33_worker_shadow( + package, + tmp_path / "results", + execution_id="wrong-container", + ) + + +def test_e33_rejects_wrong_worker_node_identity( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + package = _package(tmp_path / "package-source", frame_count=4) + monkeypatch.setenv("E33_CONTAINER_IMAGE", _IMAGE) + monkeypatch.setenv("E33_WORKER_NODE", "unexpected-worker") + + with pytest.raises(E33WorkerShadowError, match="worker node identity changed"): + run_e33_worker_shadow( + package, + tmp_path / "results", + execution_id="wrong-worker-node", + ) + + +def _package( + root: Path, + *, + frame_count: int, + interval_seconds: float = 0.02, + processing_delay_ms: float = 0.0, + work_capacity: int = 2, +) -> Path: + root.mkdir(parents=True) + e32_identity = { + "frame_count": frame_count, + "timeline_start_seconds": 0.0, + "timeline_end_seconds": (frame_count - 1) * interval_seconds, + "source": { + "camera_result_id": "e10-integrated-perception-" + "1" * 64, + "source_pack_id": "e10-lidar-pack-" + "2" * 64, + }, + "authority": _AUTHORITY, + } + e32_identity_sha = hashlib.sha256(_canonical(e32_identity)).hexdigest() + e32_id = f"e32-track-geometry-{e32_identity_sha}" + e32_root = root / "input" / "e32" / e32_id + e32_root.mkdir(parents=True) + offsets = np.arange(frame_count + 1, dtype=" dict[str, Any]: + return { + "path": path.name, + "byte_length": path.stat().st_size, + "sha256": _sha256(path), + } + + +def _write_json(path: Path, value: object) -> None: + path.write_text(json.dumps(value, sort_keys=True), encoding="utf-8") + + +def _canonical(value: object) -> bytes: + return json.dumps( + value, + sort_keys=True, + separators=(",", ":"), + allow_nan=False, + ).encode() + + +def _sha256(path: Path) -> str: + return hashlib.sha256(path.read_bytes()).hexdigest() diff --git a/tests/test_live_perception.py b/tests/test_live_perception.py index 7af23cb..10c90e0 100644 --- a/tests/test_live_perception.py +++ b/tests/test_live_perception.py @@ -210,9 +210,9 @@ def test_live_result_rejects_tampering_and_partial_cuboid() -> None: def test_latest_wins_queue_never_exceeds_capacity() -> None: queue = LatestWinsQueue[int](capacity=2) - queue.publish(1) - queue.publish(2) - queue.publish(3) + assert queue.publish(1) is None + assert queue.publish(2) is None + assert queue.publish(3) == 1 assert queue.take_next(timeout=0) == 2 assert queue.take_next(timeout=0) == 3