From 7720594790d45367b730d6164e191f9fe86742b7 Mon Sep 17 00:00:00 2001 From: DCCONSTRUCTIONS Date: Thu, 27 Aug 2026 11:27:28 +0300 Subject: [PATCH] feat(perception): add integrated TGS graph shadow gate --- .../m49-tgs-integrated-graph-shadow-v1.json | 73 +++ .../run_m48s_reference_graph_shadow_worker.py | 41 ++ .../Invoke-M49TgsIntegratedGraphShadow.ps1 | 450 ++++++++++++++++++ .../build_tgs_full_shadow_binary.sh | 15 + .../build_tgs_integrated_graph_evidence.py | 339 +++++++++++++ .../m49_t3_travel/run_tgs_full_shadow.cpp | 52 +- .../run_tgs_integrated_shadow.sh | 26 + ...49_tgs_integrated_graph_worker_artifact.py | 257 ++++++++++ tests/test_m48s_reference_graph_timing.py | 31 ++ tests/test_m49_tgs_integrated_graph_shadow.py | 208 ++++++++ 10 files changed, 1489 insertions(+), 3 deletions(-) create mode 100644 config/perception/m49-tgs-integrated-graph-shadow-v1.json create mode 100644 experiments/perception/worker/Invoke-M49TgsIntegratedGraphShadow.ps1 create mode 100644 experiments/perception/worker/m49_t3_travel/build_tgs_full_shadow_binary.sh create mode 100644 experiments/perception/worker/m49_t3_travel/build_tgs_integrated_graph_evidence.py create mode 100644 experiments/perception/worker/m49_t3_travel/run_tgs_integrated_shadow.sh create mode 100644 scripts/build_m49_tgs_integrated_graph_worker_artifact.py create mode 100644 tests/test_m49_tgs_integrated_graph_shadow.py diff --git a/config/perception/m49-tgs-integrated-graph-shadow-v1.json b/config/perception/m49-tgs-integrated-graph-shadow-v1.json new file mode 100644 index 0000000..038e3f9 --- /dev/null +++ b/config/perception/m49-tgs-integrated-graph-shadow-v1.json @@ -0,0 +1,73 @@ +{ + "schema_version": "missioncore.m49-tgs-integrated-graph-shadow-profile/v1", + "profile_id": "m49-ravnoves00-tgs-native-risk-integrated-shadow/v1", + "source": { + "source_id": "RAVNOVES00", + "expected_timeline_frames": 4489, + "expected_available_lidar_frames": 3928, + "source_pack_sha256": "0685d24219d8236caf8b7f1685e93f6d6b59e7fd015a768d88a92bbe8b154944", + "requested_source_rate_hz": 12.0, + "shared_start_barrier": true + }, + "stages": { + "reference_graph": { + "graph_config": "m48n-rf-detr-native-reference-graph-shadow-v0.json", + "graph_config_sha256": "1db6f6fa0561d819505a5e4f62fe256dab6fe4d9da6b4f1a9a5073bfcf772c90", + "detector_profile": "rf-detr-large-native-kb4-risk-shadow-v0.json", + "detector_profile_sha256": "dbf4da5dbad6c3c22b1280b46ffcad81719bd183c81c263a4859847d829019b6", + "detector_provider_id": "triton-rf-detr-large-coco-native-kb4-risk-fp16-shadow/v0", + "single_inference_pass": true + }, + "tgs": { + "profile": "m49-tgs-full-shadow-v1.json", + "profile_sha256": "c2e07010aaee78259d36c057962d6bfb885349251ff7356d867e5813e632881c", + "candidate_id": "travel-tgs-only-gravity-aligned", + "linked_accepted_result_id": "m49-tgs-full-shadow-ef98de7db7596d48e8c8c0549ce68e6704ee03e87c8c4bcf1e3e748b7ccb032e", + "aos_allowed": false, + "gpu_allowed": false, + "states_remain_separate": true + } + }, + "acceptance": { + "minimum_delivery_ratio": 1.0, + "reference_world_state_fps": 11.79902, + "maximum_world_state_fps_regression_fraction": 0.05, + "minimum_effective_world_state_fps": 11.209069, + "maximum_world_state_completion_p95_ms": 125.0, + "candidate_stage_p95_ms_max": 25.0, + "candidate_stage_p99_ms_max": 50.0, + "combined_output_age_p99_ms_max": 125.0, + "capacity_drop_count_max": 0, + "unaccounted_frame_count_max": 0 + }, + "telemetry": { + "sample_interval_seconds": 1.0, + "required_products": [ + "graph_pipeline_timing", + "graph_queue_high_watermarks", + "tgs_stage_timing", + "host_container_cpu_memory", + "gpu_utilization_memory_power_temperature" + ] + }, + "invariants": { + "native_fisheye_raster_unchanged": true, + "camera_rectification_allowed": false, + "camera_resize_allowed": false, + "tgs_candidate_parameters_unchanged": true, + "tgs_modifies_reference_graph_state": false, + "missing_lidar_means_unobserved": true, + "future_frames_used": false, + "low_step_used": false, + "gauss_or_playcanvas_in_scope": false + }, + "authority": { + "visual_quality_accepted": false, + "traversability_accepted": false, + "physical_free_space_accepted": false, + "commands_enabled": false, + "actuation_allowed": false, + "navigation_or_safety_accepted": false, + "production_accepted": false + } +} diff --git a/experiments/perception/run_m48s_reference_graph_shadow_worker.py b/experiments/perception/run_m48s_reference_graph_shadow_worker.py index e33d052..4c17e46 100644 --- a/experiments/perception/run_m48s_reference_graph_shadow_worker.py +++ b/experiments/perception/run_m48s_reference_graph_shadow_worker.py @@ -67,6 +67,27 @@ AUTHORITY: Final = { LOAD_PURPOSES: Final = ("production-rate", "reserve-gate", "limit-discovery") +def wait_for_shared_start( + *, + ready_file: Path, + start_file: Path, + timeout_seconds: float, +) -> None: + """Join an external source-admission barrier after local warmup is complete.""" + + if timeout_seconds <= 0: + raise RuntimeError("shared-start timeout must be positive") + if ready_file.exists(): + raise RuntimeError("shared-start ready file already exists") + ready_file.parent.mkdir(mode=0o700, parents=True, exist_ok=True) + ready_file.write_text("ready\n", encoding="utf-8") + deadline = time.monotonic() + timeout_seconds + while not start_file.is_file(): + if time.monotonic() >= deadline: + raise RuntimeError("shared-start barrier timed out") + time.sleep(0.01) + + class GpuTelemetry: def __init__(self, interval_seconds: float) -> None: self.interval_seconds = interval_seconds @@ -394,6 +415,9 @@ def main() -> int: parser.add_argument("--runtime-artifact-sha256", required=True) parser.add_argument("--runner-sha256", required=True) parser.add_argument("--telemetry-interval-seconds", type=float, default=1.0) + parser.add_argument("--shared-start-ready-file", type=Path) + parser.add_argument("--shared-start-file", type=Path) + parser.add_argument("--shared-start-timeout-seconds", type=float, default=600.0) parser.add_argument("--output", type=Path, required=True) parser.add_argument("--progress", type=Path, required=True) parser.add_argument("--frame-ledger", type=Path, required=True) @@ -404,6 +428,17 @@ def main() -> int: raise RuntimeError("maximum frame count must be positive") if arguments.telemetry_interval_seconds <= 0: raise RuntimeError("telemetry interval must be positive") + shared_start_requested = ( + arguments.shared_start_ready_file is not None or arguments.shared_start_file is not None + ) + if shared_start_requested and ( + arguments.shared_start_ready_file is None or arguments.shared_start_file is None + ): + raise RuntimeError("shared-start ready and start files must be configured together") + if shared_start_requested and arguments.loops != 1: + raise RuntimeError("shared-start barrier requires exactly one loop") + if arguments.shared_start_timeout_seconds <= 0: + raise RuntimeError("shared-start timeout must be positive") if arguments.source_rate_hz is not None and ( not np.isfinite(arguments.source_rate_hz) or arguments.source_rate_hz <= 0 ): @@ -506,6 +541,12 @@ def main() -> int: ) detector_warmup = runtime.warm_up_detector() source_prefetch = runtime.prepare_source() + if shared_start_requested: + wait_for_shared_start( + ready_file=arguments.shared_start_ready_file, + start_file=arguments.shared_start_file, + timeout_seconds=arguments.shared_start_timeout_seconds, + ) gc_policy = CyclicGcHotLoopPolicy() with gc_policy: runtime.mark_source_admission_started() diff --git a/experiments/perception/worker/Invoke-M49TgsIntegratedGraphShadow.ps1 b/experiments/perception/worker/Invoke-M49TgsIntegratedGraphShadow.ps1 new file mode 100644 index 0000000..5688d3a --- /dev/null +++ b/experiments/perception/worker/Invoke-M49TgsIntegratedGraphShadow.ps1 @@ -0,0 +1,450 @@ +[CmdletBinding()] +param( + [Parameter(Mandatory = $true)] + [string]$ReleaseRoot, + [Parameter(Mandatory = $true)] + [string]$CandidateRoot, + [Parameter(Mandatory = $true)] + [ValidatePattern("^[a-f0-9]{64}$")] + [string]$ExpectedArtifactSha256, + [Parameter(Mandatory = $true)] + [ValidatePattern("^[A-Za-z0-9._-]{1,96}$")] + [string]$RunId, + [ValidateRange(1.0, 120.0)] + [double]$SourceRateHz = 12.0, + [string]$OutputRoot = ( + "D:\NDC_MISSIONCORE\runtime\results\m49-tgs-integrated-graph-shadow" + ) +) + +$ErrorActionPreference = "Stop" +$ProgressPreference = "SilentlyContinue" +$TravelImageTag = "ndc/mission-core-m49-t3-travel:20260826" +$TravelImageId = "sha256:7b412020f4d8392d1d1ed1b33beadc44140f0ea8f781e62dd69796042334300f" +$ParityImageTag = "ndc-mission-core-m48t-upstream-parity:1.9.4-cu130" +$ParityImageId = "sha256:ceb13548617e4bd3f619766bfdff00af3fa5160946b367828da6d2233dcdcba0" +$RuntimeImage = ( + "nvcr.io/nvidia/tritonserver:26.06-py3@" + + "sha256:58df7489c3f2276f9591d500a012dee03e23d35543ce3c390b4c001e6bf90794" +) + +function Assert-LastExitCode([string]$Operation) { + if ($LASTEXITCODE -ne 0) { throw "$Operation failed with exit code $LASTEXITCODE" } +} + +function Get-Sha256([string]$Path) { + return (Get-FileHash -LiteralPath $Path -Algorithm SHA256).Hash.ToLowerInvariant() +} + +function Resolve-DDirectory([string]$Path, [string]$Label, [bool]$Create) { + if ($Create -and -not (Test-Path -LiteralPath $Path)) { + $null = New-Item -ItemType Directory -Path $Path + } + $item = Get-Item -LiteralPath (Resolve-Path -LiteralPath $Path).Path -Force + if ( + -not $item.PSIsContainer -or + ($item.Attributes -band [IO.FileAttributes]::ReparsePoint) -or + [IO.Path]::GetPathRoot($item.FullName).TrimEnd("\") -ine "D:" + ) { throw "$Label must be a real D: directory" } + return $item.FullName +} + +function Resolve-DFile([string]$Path, [string]$Label) { + $item = Get-Item -LiteralPath (Resolve-Path -LiteralPath $Path).Path -Force + if ( + $item.PSIsContainer -or + ($item.Attributes -band [IO.FileAttributes]::ReparsePoint) -or + [IO.Path]::GetPathRoot($item.FullName).TrimEnd("\") -ine "D:" + ) { throw "$Label must be a real D: file" } + return $item.FullName +} + +function Convert-ToDockerPath([string]$Path) { return ($Path -replace "\\", "/") } + +function Get-Container([string]$Name) { + $rows = @(((& docker inspect $Name) | ConvertFrom-Json)) + Assert-LastExitCode "Docker inspection for $Name" + if ($rows.Count -ne 1) { throw "Container identity for $Name is not unique" } + return $rows[0] +} + +function Assert-Image([string]$Tag, [string]$ExpectedId) { + $rows = @(((& docker image inspect $Tag) | ConvertFrom-Json)) + Assert-LastExitCode "Docker image inspection for $Tag" + if ($rows.Count -ne 1 -or [string]$rows[0].Id -cne $ExpectedId) { + throw "Pinned image identity changed for $Tag" + } +} + +function Remove-ExactContainer([string]$Name) { + if (& docker ps -a --format "{{.Names}}" --filter "name=^/$Name$") { + & docker rm --force $Name *> $null + } +} + +function Wait-Healthy([string]$Name) { + foreach ($attempt in 1..60) { + Start-Sleep -Seconds 2 + $container = Get-Container $Name + if (-not $container.State.Running) { + & docker logs $Name + throw "$Name stopped during startup" + } + if ($container.State.Health.Status -ceq "healthy") { return } + } + throw "$Name did not become healthy" +} + +function Wait-SharedReady( + [string]$GraphReady, + [string]$TgsReady, + [string]$GraphName, + [string]$TgsName +) { + $deadline = [DateTimeOffset]::UtcNow.AddMinutes(10) + while (-not ((Test-Path -LiteralPath $GraphReady) -and (Test-Path -LiteralPath $TgsReady))) { + foreach ($name in @($GraphName, $TgsName)) { + $container = Get-Container $name + if (-not $container.State.Running) { + & docker logs $name + throw "$name stopped before shared-start readiness" + } + } + if ([DateTimeOffset]::UtcNow -ge $deadline) { + throw "M49 integrated shared-start readiness timed out" + } + Start-Sleep -Milliseconds 250 + } +} + +if ($env:COMPUTERNAME -cne "DESKTOP-OPJ8J04") { + throw "M49 integrated TGS graph shadow is pinned to Worker 006" +} +$release = Resolve-DDirectory $ReleaseRoot "M49 integrated release root" $false +$payload = Resolve-DDirectory (Join-Path $release "payload") "M49 integrated payload" $false +$candidate = Resolve-DDirectory $CandidateRoot "M49 native candidate root" $false +$output = Resolve-DDirectory $OutputRoot "M49 integrated output root" $true +$runCandidate = Join-Path $output $RunId +if (Test-Path -LiteralPath $runCandidate) { throw "M49 integrated output already exists" } +$null = New-Item -ItemType Directory -Path $runCandidate +$runOutput = Resolve-DDirectory $runCandidate "M49 integrated run output" $false +foreach ($directory in @("bin", "control", "graph", "tgs")) { + $null = New-Item -ItemType Directory -Path (Join-Path $runOutput $directory) +} + +$releaseDocument = Get-Content -LiteralPath (Join-Path $payload "release.json") -Raw | ConvertFrom-Json +if ( + $releaseDocument.schema_version -cne "missioncore.m49-tgs-integrated-graph-worker-release/v1" -or + $releaseDocument.worker_id -cne "worker-006" -or + $releaseDocument.transition -cne "m49-tgs-native-risk-integrated-shadow/v1" +) { throw "M49 integrated release contract changed" } +foreach ($property in $releaseDocument.files.PSObject.Properties) { + $path = Join-Path $payload $property.Name + if ((Get-Sha256 $path) -cne [string]$property.Value.sha256) { + throw "M49 integrated payload digest changed: $($property.Name)" + } +} +$wheelSha256 = [string]$releaseDocument.files."nodedc_mission_core-0.1.0-py3-none-any.whl".sha256 +$runnerSha256 = [string]$releaseDocument.files."run_m48s_reference_graph_shadow_worker.py".sha256 + +$source = [ordered]@{ + CameraIndex = ( + "D:\NDC_MISSIONCORE\runtime\jobs\recorded-camera-602ac89026ed12978619801d" + + "\input\camera\sensor.camera.right\epoch-1\index.jsonl" + ) + SourcePack = ( + "D:\NDC_MISSIONCORE\runtime\derived" + + "\e10-lidar-pack-576c994a6c814e2592dd6240ace3902a5db94843312c759a73ba0c9166157d2b" + + "\lidar-pack.npz" + ) + LocalSurface = ( + "D:\NDC_MISSIONCORE\runtime\derived" + + "\k1-local-surface-23762244c8bdb97de26fb721ac957d7a00bc9a63571ac4cfa4be19c4effc7d55" + + "\local-surface.npz" + ) + Video = ( + "D:\NDC_MISSIONCORE\runtime\experiments\e46e\inputs" + + "\right-cadd1696ff000904eb78633a0a8418104b8024f178b91f3421789021ccb160e8.mp4" + ) + Mask = ( + "D:\NDC_MISSIONCORE\runtime\inputs\e2" + + "\valid-fov-mask-b4dd8ddf2b87c1d520ee8a0868c4fea062d7c14d1bae73ccabd3abe1f3acbac2" + + "\mask.png" + ) +} +foreach ($entry in $source.GetEnumerator()) { + $null = Resolve-DFile $entry.Value "M49 source $($entry.Key)" +} +if ((Get-Sha256 $source.SourcePack) -cne [string]$releaseDocument.source_pack_sha256) { + throw "RAVNOVES00 source pack digest changed" +} + +$nativeConfig = Resolve-DFile ( + (Join-Path $payload "rf_detr_large_native_kb4_config.pbtxt") +) "native RF-DETR config" +$nativeEngine = Resolve-DFile ( + (Join-Path $candidate "rf-detr-native-uint8.plan") +) "native RF-DETR engine" +if ((Get-Sha256 $nativeEngine) -cne "b8a40b3580edff001ec9680de68707242294ff590ab296000fae371f1083f695") { + throw "native RF-DETR engine SHA-256 changed" +} +$modelRoot = Join-Path $runOutput "triton-models" +$modelDirectory = Join-Path $modelRoot "rf_detr_large_native_kb4" +$modelVersionDirectory = Join-Path $modelDirectory "1" +$null = New-Item -ItemType Directory -Path $modelVersionDirectory +Copy-Item -LiteralPath $nativeConfig -Destination (Join-Path $modelDirectory "config.pbtxt") +Copy-Item -LiteralPath $nativeEngine -Destination (Join-Path $modelVersionDirectory "model.plan") + +$media = Resolve-DDirectory ( + "D:\NDC_MISSIONCORE\runtime\derived\perception-e15-media-pyav180-lz445-v1" +) "PyAV dependency" $false +$opencv = Resolve-DDirectory ( + "D:\NDC_MISSIONCORE\runtime\derived\perception-e3-opencv413092-v1\packages" +) "OpenCV dependency" $false +$pillow = Resolve-DDirectory ( + "D:\NDC_MISSIONCORE\runtime\derived\perception-p0-env-v1" +) "Pillow dependency" $false + +Assert-Image $TravelImageTag $TravelImageId +Assert-Image $ParityImageTag $ParityImageId +& docker image inspect $RuntimeImage *> $null +Assert-LastExitCode "pinned runtime image inspection" +$os = Get-CimInstance Win32_OperatingSystem +$freeMemoryGiB = [double]$os.FreePhysicalMemory / 1MB +if ($freeMemoryGiB -lt 24.0) { + throw ("M49 integrated shadow requires 24 GiB free memory; observed {0:N2} GiB" -f $freeMemoryGiB) +} +$canonicalBefore = Get-Container "ndc-mission-core-triton" +if (-not $canonicalBefore.State.Running -or $canonicalBefore.State.Health.Status -cne "healthy") { + throw "Canonical Mission Core Triton must remain healthy" +} +$canonicalId = [string]$canonicalBefore.Id + +$prepareName = "ndc-mission-core-m49-integrated-prepare-$RunId" +$compileName = "ndc-mission-core-m49-integrated-compile-$RunId" +$tritonName = "ndc-mission-core-m49-integrated-triton-$RunId" +$graphName = "ndc-mission-core-m49-integrated-graph-$RunId" +$tgsName = "ndc-mission-core-m49-integrated-tgs-$RunId" +$analyzeName = "ndc-mission-core-m49-integrated-analyze-$RunId" +$evidenceName = "ndc-mission-core-m49-integrated-evidence-$RunId" +$containers = @($prepareName, $compileName, $tritonName, $graphName, $tgsName, $analyzeName, $evidenceName) +foreach ($name in $containers) { + if (& docker ps -a --format "{{.Names}}" --filter "name=^/$name$") { + throw "M49 integrated container name already exists: $name" + } +} + +$started = [DateTimeOffset]::UtcNow +try { + & docker run --rm --name $prepareName --network none --cpus 8 --memory 16g ` + --entrypoint python3 ` + --volume ((Convert-ToDockerPath $source.SourcePack) + ":/source/lidar-pack.npz:ro") ` + --volume ((Convert-ToDockerPath $payload) + ":/release:ro") ` + --volume ((Convert-ToDockerPath (Join-Path $runOutput "tgs")) + ":/tgs") ` + $ParityImageTag /release/prepare_tgs_full_shadow_inputs.py ` + --source-pack /source/lidar-pack.npz ` + --config /release/m49-tgs-full-shadow-v1.json ` + --output-root /tgs/inputs + Assert-LastExitCode "M49 integrated TGS input preparation" + + & docker run --rm --name $compileName --network none --cpus 8 --memory 8g ` + --entrypoint /bin/bash ` + --volume ((Convert-ToDockerPath $payload) + ":/release:ro") ` + --volume ((Convert-ToDockerPath (Join-Path $runOutput "bin")) + ":/out") ` + $TravelImageTag /release/build_tgs_full_shadow_binary.sh /out/run_tgs_full_shadow + Assert-LastExitCode "M49 integrated TGS binary build" + + & docker create --name $tritonName ` + --read-only --security-opt "no-new-privileges:true" --cap-drop ALL ` + --pids-limit 512 --shm-size 1g --gpus all ` + --tmpfs "/tmp:rw,noexec,nosuid,size=2g" ` + --health-cmd "curl --fail --silent http://127.0.0.1:8000/v2/health/ready" ` + --health-interval 5s --health-timeout 3s --health-start-period 20s --health-retries 24 ` + -v ((Convert-ToDockerPath $modelRoot) + ":/models:ro") ` + $RuntimeImage tritonserver --model-repository=/models ` + --model-control-mode=explicit --load-model=rf_detr_large_native_kb4 ` + --disable-auto-complete-config --strict-readiness=true --exit-on-error=true ` + --allow-http=true --allow-grpc=false --allow-metrics=false *> $null + Assert-LastExitCode "M49 integrated Triton creation" + & docker start $tritonName *> $null + Assert-LastExitCode "M49 integrated Triton start" + Wait-Healthy $tritonName + + $dockerRelease = Convert-ToDockerPath $payload + $dockerRun = Convert-ToDockerPath $runOutput + $rate = [string]::Format([Globalization.CultureInfo]::InvariantCulture, "{0:R}", $SourceRateHz) + $graphArguments = @( + "create", "--name", $graphName, + "--network", ("container:{0}" -f $tritonName), + "--read-only", "--security-opt", "no-new-privileges:true", "--cap-drop", "ALL", + "--pids-limit", "256", "--gpus", "all", + "--tmpfs", "/tmp:rw,noexec,nosuid,size=2g", + "-e", "PYTHONDONTWRITEBYTECODE=1", + "-e", "PYTHONPATH=/release/nodedc_mission_core-0.1.0-py3-none-any.whl:/opt/media:/opt/opencv:/opt/pillow", + "-v", ("{0}:/release:ro" -f $dockerRelease), + "-v", ("{0}:/shared:rw" -f $dockerRun), + "-v", ((Convert-ToDockerPath $media) + ":/opt/media:ro"), + "-v", ((Convert-ToDockerPath $opencv) + ":/opt/opencv:ro"), + "-v", ((Convert-ToDockerPath $pillow) + ":/opt/pillow:ro"), + "-v", ((Convert-ToDockerPath $source.CameraIndex) + ":/source/camera-index.jsonl:ro"), + "-v", ((Convert-ToDockerPath $source.SourcePack) + ":/source/source-pack.npz:ro"), + "-v", ((Convert-ToDockerPath $source.LocalSurface) + ":/source/local-surface.npz:ro"), + "-v", ((Convert-ToDockerPath $source.Video) + ":/source/right.mp4:ro"), + "-v", ((Convert-ToDockerPath $source.Mask) + ":/source/mask.png:ro"), + "--entrypoint", "python3", $RuntimeImage, + "/release/run_m48s_reference_graph_shadow_worker.py", + "--graph-config", "/release/m48n-rf-detr-native-reference-graph-shadow-v0.json", + "--baseline-profile", "/release/m4-recorded-realtime-baseline-v1.json", + "--detector-profile", "/release/rf-detr-large-native-kb4-risk-shadow-v0.json", + "--geometry-profile", "/release/m4-geometry-association-v1.json", + "--temporal-motion-profile", "/release/m4-temporal-motion-v1.json", + "--rolling-map-profile", "/release/m4-rolling-local-map-v1.json", + "--threat-profile", "/release/m4-replay-threat-v3.json", + "--camera-index", "/source/camera-index.jsonl", + "--source-pack", "/source/source-pack.npz", + "--local-surface", "/source/local-surface.npz", + "--video", "/source/right.mp4", + "--valid-fov-mask", "/source/mask.png", + "--triton-origin", "http://127.0.0.1:8000", + "--loops", "1", "--maximum-frames", "4489", "--source-rate-hz", $rate, + "--minimum-delivery-ratio", "1.0", + "--minimum-effective-world-state-fps", "11.209069", + "--maximum-world-state-completion-p95-ms", "125.0", + "--load-purpose", "reserve-gate", + "--runtime-artifact-sha256", $wheelSha256, + "--runner-sha256", $runnerSha256, + "--shared-start-ready-file", "/shared/control/graph.ready", + "--shared-start-file", "/shared/control/start.signal", + "--output", "/shared/graph/result.json", + "--progress", "/shared/graph/progress.jsonl", + "--frame-ledger", "/shared/graph/frames.jsonl" + ) + & docker @graphArguments *> $null + Assert-LastExitCode "M49 integrated graph creation" + + & docker create --name $tgsName --network none --cpus 16 --memory 24g ` + --read-only --security-opt "no-new-privileges:true" --cap-drop ALL ` + --pids-limit 256 --tmpfs "/tmp:rw,noexec,nosuid,size=1g" ` + -e ("M49_SOURCE_RATE_HZ={0}" -f $rate) ` + --entrypoint /bin/bash ` + --volume ($dockerRelease + ":/release:ro") ` + --volume ($dockerRun + ":/shared:rw") ` + $TravelImageTag /release/run_tgs_integrated_shadow.sh *> $null + Assert-LastExitCode "M49 integrated TGS creation" + + & docker start $graphName *> $null + Assert-LastExitCode "M49 integrated graph start" + & docker start $tgsName *> $null + Assert-LastExitCode "M49 integrated TGS start" + $graphReady = Join-Path $runOutput "control\graph.ready" + $tgsReady = Join-Path $runOutput "control\tgs.ready" + Wait-SharedReady $graphReady $tgsReady $graphName $tgsName + [DateTimeOffset]::UtcNow.ToString("o") | Set-Content -LiteralPath ( + Join-Path $runOutput "control\start.signal" + ) -Encoding utf8 + + $telemetryPath = Join-Path $runOutput "container-telemetry.jsonl" + while ($true) { + $graphState = Get-Container $graphName + $tgsState = Get-Container $tgsName + $running = @() + if ($graphState.State.Running) { $running += $graphName } + if ($tgsState.State.Running) { $running += $tgsName } + if ((Get-Container $tritonName).State.Running) { $running += $tritonName } + if ($running.Count -gt 0) { + $stats = @((& docker stats --no-stream --format "{{json .}}" @running)) + Assert-LastExitCode "M49 integrated container telemetry" + foreach ($line in $stats) { + $value = $line | ConvertFrom-Json + $role = if ($value.Name -ceq $graphName) { + "graph" + } elseif ($value.Name -ceq $tgsName) { + "tgs" + } elseif ($value.Name -ceq $tritonName) { + "triton" + } else { + throw "Unknown M49 telemetry container" + } + [ordered]@{ + observed_utc = [DateTimeOffset]::UtcNow.ToString("o") + role = $role + name = [string]$value.Name + cpu_percent = [string]$value.CPUPerc + memory_usage = [string]$value.MemUsage + memory_percent = [string]$value.MemPerc + pids = [string]$value.PIDs + } | ConvertTo-Json -Compress | Out-File -LiteralPath $telemetryPath -Encoding utf8 -Append + } + } + if (-not $graphState.State.Running -and -not $tgsState.State.Running) { break } + Start-Sleep -Seconds 1 + } + $graphExit = [int](Get-Container $graphName).State.ExitCode + $tgsExit = [int](Get-Container $tgsName).State.ExitCode + (& docker logs $graphName 2>&1) | Set-Content -LiteralPath (Join-Path $runOutput "graph.log") -Encoding utf8 + (& docker logs $tgsName 2>&1) | Set-Content -LiteralPath (Join-Path $runOutput "tgs.log") -Encoding utf8 + if ($graphExit -ne 0) { throw "M49 integrated graph failed with exit code $graphExit" } + if ($tgsExit -ne 0) { throw "M49 integrated TGS failed with exit code $tgsExit" } + + & docker run --rm --name $analyzeName --network none --cpus 8 --memory 16g ` + --entrypoint python3 ` + --volume ($dockerRelease + ":/release:ro") ` + --volume ($dockerRun + ":/shared:rw") ` + $ParityImageTag /release/build_tgs_full_shadow_evidence.py ` + --run-root /shared/tgs ` + --config /release/m49-tgs-full-shadow-v1.json ` + --output-root /shared/tgs/evidence + Assert-LastExitCode "M49 integrated TGS evidence analysis" + + & docker run --rm --name $evidenceName --network none --cpus 4 --memory 8g ` + --entrypoint python3 ` + --volume ($dockerRelease + ":/release:ro") ` + --volume ($dockerRun + ":/shared:rw") ` + $ParityImageTag /release/build_tgs_integrated_graph_evidence.py ` + --profile /release/m49-tgs-integrated-graph-shadow-v1.json ` + --graph-result /shared/graph/result.json ` + --graph-frames /shared/graph/frames.jsonl ` + --tgs-result /shared/tgs/evidence/result.json ` + --tgs-timing /shared/tgs/tgs-full-timing.tsv ` + --telemetry /shared/container-telemetry.jsonl ` + --output /shared/result.json ` + --release-sha256 $ExpectedArtifactSha256 + Assert-LastExitCode "M49 integrated evidence gate" +} finally { + foreach ($name in $containers) { Remove-ExactContainer $name } + $canonicalAfter = Get-Container "ndc-mission-core-triton" + if ( + [string]$canonicalAfter.Id -cne $canonicalId -or + -not $canonicalAfter.State.Running -or + $canonicalAfter.State.Health.Status -cne "healthy" + ) { throw "Canonical Mission Core Triton changed during M49 integrated shadow" } +} + +$completed = [DateTimeOffset]::UtcNow +$resultPath = Join-Path $runOutput "result.json" +if (-not (Test-Path -LiteralPath $resultPath -PathType Leaf)) { + throw "M49 integrated result is missing" +} +$result = Get-Content -LiteralPath $resultPath -Raw | ConvertFrom-Json +$summary = [ordered]@{ + schema_version = "missioncore.m49-tgs-integrated-graph-worker-summary/v1" + worker_id = "worker-006" + run_id = $RunId + code_revision = [string]$releaseDocument.code_revision + source_rate_hz = $SourceRateHz + started_utc = $started.ToString("o") + completed_utc = $completed.ToString("o") + wall_seconds = [math]::Round(($completed - $started).TotalSeconds, 6) + free_memory_gib_before = [math]::Round($freeMemoryGiB, 6) + result_id = [string]$result.result_id + result_status = [string]$result.status + canonical_triton_id = $canonicalId + canonical_triton_health = "healthy" + gauss_or_playcanvas_action = "none" + durable_worker_action = "none" + navigation_or_actuation_allowed = $false +} +$summary | ConvertTo-Json -Depth 3 | Set-Content -LiteralPath ( + Join-Path $runOutput "worker-summary.json" +) -Encoding utf8 +$summary | ConvertTo-Json -Depth 3 diff --git a/experiments/perception/worker/m49_t3_travel/build_tgs_full_shadow_binary.sh b/experiments/perception/worker/m49_t3_travel/build_tgs_full_shadow_binary.sh new file mode 100644 index 0000000..7876c16 --- /dev/null +++ b/experiments/perception/worker/m49_t3_travel/build_tgs_full_shadow_binary.sh @@ -0,0 +1,15 @@ +#!/usr/bin/env bash +set -euo pipefail + +readonly TARGET=${1:?target binary path is required} + +test -f /release/run_tgs_full_shadow.cpp +test ! -e "${TARGET}" +mkdir -p "$(dirname "${TARGET}")" +g++ -std=c++17 -O3 -DNDEBUG -pthread \ + -I/opt/travel/src/TRAVEL/cpp/travel/core \ + -I/usr/include/eigen3 \ + /release/run_tgs_full_shadow.cpp \ + -o "${TARGET}" +chmod 0755 "${TARGET}" +sha256sum "${TARGET}" diff --git a/experiments/perception/worker/m49_t3_travel/build_tgs_integrated_graph_evidence.py b/experiments/perception/worker/m49_t3_travel/build_tgs_integrated_graph_evidence.py new file mode 100644 index 0000000..f44c2e8 --- /dev/null +++ b/experiments/perception/worker/m49_t3_travel/build_tgs_integrated_graph_evidence.py @@ -0,0 +1,339 @@ +#!/usr/bin/env python3 +"""Seal the synchronized TGS plus native RF-DETR reference-graph shadow.""" + +from __future__ import annotations + +import argparse +import csv +import hashlib +import json +import math +import re +from collections import defaultdict +from pathlib import Path +from typing import Any + +import numpy as np + +PROFILE_SCHEMA = "missioncore.m49-tgs-integrated-graph-shadow-profile/v1" +GRAPH_SCHEMA = "missioncore.m48s-reference-graph-shadow-load/v5" +TGS_SCHEMA = "missioncore.m49-tgs-full-shadow-result/v1" +RESULT_SCHEMA = "missioncore.m49-tgs-integrated-graph-shadow-result/v1" +FRAME_COUNT = 4_489 + + +class IntegratedShadowError(RuntimeError): + """The integrated shadow evidence is incomplete or incompatible.""" + + +def sha256_file(path: Path) -> str: + digest = hashlib.sha256() + with path.open("rb") as stream: + for chunk in iter(lambda: stream.read(1024 * 1024), b""): + digest.update(chunk) + return digest.hexdigest() + + +def canonical_json(value: object) -> bytes: + return json.dumps(value, sort_keys=True, separators=(",", ":")).encode("utf-8") + + +def load_json(path: Path, label: str) -> dict[str, Any]: + try: + value = json.loads(path.read_text(encoding="utf-8-sig")) + except (OSError, json.JSONDecodeError) as exc: + raise IntegratedShadowError(f"{label} is unreadable") from exc + if not isinstance(value, dict): + raise IntegratedShadowError(f"{label} is not an object") + return value + + +def distribution(values: list[float]) -> dict[str, float]: + if not values: + return {"mean": 0.0, "p50": 0.0, "p95": 0.0, "p99": 0.0, "maximum": 0.0} + array = np.asarray(values, dtype=np.float64) + return { + "mean": round(float(array.mean()), 6), + "p50": round(float(np.percentile(array, 50)), 6), + "p95": round(float(np.percentile(array, 95)), 6), + "p99": round(float(np.percentile(array, 99)), 6), + "maximum": round(float(array.max()), 6), + } + + +def graph_completion_ages(path: Path) -> list[float]: + ages: list[float] = [] + with path.open("r", encoding="utf-8") as stream: + for expected, line in enumerate(stream): + row = json.loads(line) + sequence = row.get("source_envelope", {}).get("sequence") + if sequence != expected: + raise IntegratedShadowError("graph frame ledger sequence changed") + age_ns = row.get("completion_age_ns") + if not isinstance(age_ns, int) or age_ns < 0: + raise IntegratedShadowError("graph completion age is invalid") + ages.append(age_ns / 1_000_000.0) + if len(ages) != FRAME_COUNT: + raise IntegratedShadowError("graph frame ledger is incomplete") + return ages + + +def tgs_completion_ages(path: Path) -> list[float]: + ages: list[float] = [] + with path.open("r", encoding="utf-8", newline="") as stream: + for expected, row in enumerate(csv.DictReader(stream, delimiter="\t")): + if int(row["timeline_frame_index"]) != expected: + raise IntegratedShadowError("TGS timing sequence changed") + value = float(row["completion_age_ms"]) + if not math.isfinite(value) or value < 0: + raise IntegratedShadowError("TGS completion age is invalid") + ages.append(value) + if len(ages) != FRAME_COUNT: + raise IntegratedShadowError("TGS timing ledger is incomplete") + return ages + + +_SIZE = re.compile(r"^\s*([0-9.]+)\s*([kmgt]?i?b)\s*$", re.IGNORECASE) + + +def size_mib(value: str) -> float: + match = _SIZE.fullmatch(value) + if match is None: + raise IntegratedShadowError("container memory telemetry is invalid") + number = float(match.group(1)) + unit = match.group(2).lower() + scale = { + "b": 1.0 / (1024.0 * 1024.0), + "kb": 1.0 / 1024.0, + "kib": 1.0 / 1024.0, + "mb": 1.0, + "mib": 1.0, + "gb": 1024.0, + "gib": 1024.0, + "tb": 1024.0 * 1024.0, + "tib": 1024.0 * 1024.0, + }[unit] + return number * scale + + +def host_telemetry(path: Path) -> dict[str, object]: + samples: dict[str, list[dict[str, float]]] = defaultdict(list) + with path.open("r", encoding="utf-8-sig") as stream: + for line in stream: + row = json.loads(line) + role = row.get("role") + if role not in {"graph", "tgs", "triton"}: + raise IntegratedShadowError("container telemetry role changed") + cpu_text = row.get("cpu_percent") + memory_text = row.get("memory_usage") + memory_percent_text = row.get("memory_percent") + if not all( + isinstance(value, str) for value in (cpu_text, memory_text, memory_percent_text) + ): + raise IntegratedShadowError("container telemetry row is incomplete") + used_text = memory_text.split("/", 1)[0].strip() + samples[role].append( + { + "cpu_percent": float(cpu_text.rstrip("%")), + "memory_used_mib": size_mib(used_text), + "memory_percent": float(memory_percent_text.rstrip("%")), + } + ) + if any(not samples[role] for role in ("graph", "tgs", "triton")): + raise IntegratedShadowError("container telemetry does not cover every runtime role") + return { + role: { + "sample_count": len(rows), + "cpu_percent": distribution([row["cpu_percent"] for row in rows]), + "memory_used_mib": distribution([row["memory_used_mib"] for row in rows]), + "memory_percent": distribution([row["memory_percent"] for row in rows]), + } + for role, rows in sorted(samples.items()) + } + + +def build( + *, + profile_path: Path, + graph_result_path: Path, + graph_frames_path: Path, + tgs_result_path: Path, + tgs_timing_path: Path, + telemetry_path: Path, + output_path: Path, + release_sha256: str, +) -> dict[str, object]: + if output_path.exists(): + raise IntegratedShadowError("integrated result already exists") + profile = load_json(profile_path, "integrated profile") + graph = load_json(graph_result_path, "reference graph result") + tgs = load_json(tgs_result_path, "TGS result") + if profile.get("schema_version") != PROFILE_SCHEMA: + raise IntegratedShadowError("integrated profile schema changed") + if graph.get("schema_version") != GRAPH_SCHEMA: + raise IntegratedShadowError("reference graph result schema changed") + if tgs.get("schema_version") != TGS_SCHEMA: + raise IntegratedShadowError("TGS result schema changed") + if len(release_sha256) != 64 or any( + value not in "0123456789abcdef" for value in release_sha256 + ): + raise IntegratedShadowError("release SHA-256 is invalid") + + graph_ages = graph_completion_ages(graph_frames_path) + tgs_ages = tgs_completion_ages(tgs_timing_path) + combined_ages = [ + max(graph_age, tgs_age) for graph_age, tgs_age in zip(graph_ages, tgs_ages, strict=True) + ] + combined = distribution(combined_ages) + telemetry = host_telemetry(telemetry_path) + acceptance_profile = profile["acceptance"] + execution = graph.get("execution", {}) + graph_metrics = graph.get("metrics", {}) + tgs_performance = tgs.get("performance", {}) + effective_fps = float(execution.get("effective_world_state_fps", 0.0)) + reference_fps = float(acceptance_profile["reference_world_state_fps"]) + fps_regression = max(0.0, (reference_fps - effective_fps) / reference_fps) + graph_p95 = float(graph_metrics.get("world_state_completion_age_ms", {}).get("p95", math.inf)) + tgs_p95 = float(tgs_performance.get("candidate_tgs_ms", {}).get("p95", math.inf)) + tgs_p99 = float(tgs_performance.get("candidate_tgs_ms", {}).get("p99", math.inf)) + tgs_drops = int(tgs_performance.get("capacity_drop_count", -1)) + terminal = execution.get("terminal_outcomes", {}) + superseded = int(terminal.get("superseded", 0)) if isinstance(terminal, dict) else -1 + gpu_samples = int(graph_metrics.get("gpu", {}).get("sample_count", 0)) + graph_inputs = graph.get("identity", {}).get("inputs", {}) + checks = { + "frozen_graph_identity": ( + graph_inputs.get("graph_config") + == profile["stages"]["reference_graph"]["graph_config_sha256"] + and graph_inputs.get("detector_profile") + == profile["stages"]["reference_graph"]["detector_profile_sha256"] + ), + "frozen_tgs_identity": ( + tgs.get("config_sha256") == profile["stages"]["tgs"]["profile_sha256"] + ), + "requested_source_rate_preserved": ( + execution.get("requested_source_rate_hz") + == profile["source"]["requested_source_rate_hz"] + ), + "all_graph_frames_delivered": ( + execution.get("admitted_frames") == FRAME_COUNT + and execution.get("delivered_world_states") == FRAME_COUNT + ), + "all_tgs_frames_accounted": tgs.get("timeline", {}).get("frame_count") == FRAME_COUNT, + "exact_sequence_join": len(combined_ages) == FRAME_COUNT, + "minimum_delivery_ratio": float(execution.get("delivery_ratio", 0.0)) + >= float(acceptance_profile["minimum_delivery_ratio"]), + "maximum_world_state_fps_regression": fps_regression + <= float(acceptance_profile["maximum_world_state_fps_regression_fraction"]), + "minimum_effective_world_state_fps": effective_fps + >= float(acceptance_profile["minimum_effective_world_state_fps"]), + "maximum_world_state_completion_p95_ms": graph_p95 + <= float(acceptance_profile["maximum_world_state_completion_p95_ms"]), + "candidate_stage_p95_ms": tgs_p95 + <= float(acceptance_profile["candidate_stage_p95_ms_max"]), + "candidate_stage_p99_ms": tgs_p99 + <= float(acceptance_profile["candidate_stage_p99_ms_max"]), + "combined_output_age_p99_ms": combined["p99"] + <= float(acceptance_profile["combined_output_age_p99_ms_max"]), + "zero_capacity_drops": tgs_drops <= int(acceptance_profile["capacity_drop_count_max"]) + and superseded <= int(acceptance_profile["capacity_drop_count_max"]), + "reference_graph_integrity": graph.get("evidence_integrity_gate_passed") is True, + "tgs_integrity": tgs.get("status") == "passed", + "host_resource_telemetry_complete": all( + telemetry[role]["sample_count"] > 0 for role in ("graph", "tgs", "triton") + ), + "gpu_telemetry_complete": gpu_samples > 0, + "authority_remains_false": all(value is False for value in profile["authority"].values()), + } + files = { + label: {"bytes": path.stat().st_size, "sha256": sha256_file(path)} + for label, path in ( + ("graph-result.json", graph_result_path), + ("graph-frames.jsonl", graph_frames_path), + ("tgs-result.json", tgs_result_path), + ("tgs-timing.tsv", tgs_timing_path), + ("container-telemetry.jsonl", telemetry_path), + ) + } + document: dict[str, object] = { + "schema_version": RESULT_SCHEMA, + "profile_id": profile["profile_id"], + "status": "passed" if all(checks.values()) else "failed", + "source": { + "source_id": profile["source"]["source_id"], + "source_pack_sha256": profile["source"]["source_pack_sha256"], + "requested_source_rate_hz": profile["source"]["requested_source_rate_hz"], + "joined_frame_count": len(combined_ages), + }, + "identity": { + "release_sha256": release_sha256, + "profile_sha256": sha256_file(profile_path), + "graph_config_sha256": profile["stages"]["reference_graph"]["graph_config_sha256"], + "tgs_profile_sha256": profile["stages"]["tgs"]["profile_sha256"], + "linked_accepted_tgs_result_id": profile["stages"]["tgs"]["linked_accepted_result_id"], + }, + "performance": { + "effective_world_state_fps": effective_fps, + "reference_world_state_fps": reference_fps, + "world_state_fps_regression_fraction": round(fps_regression, 9), + "world_state_completion_age_ms": graph_metrics.get("world_state_completion_age_ms"), + "tgs_candidate_stage_ms": tgs_performance.get("candidate_tgs_ms"), + "tgs_completion_age_ms": tgs_performance.get("completion_age_ms"), + "combined_output_age_ms": combined, + "gpu": graph_metrics.get("gpu"), + "host_containers": telemetry, + }, + "accounting": { + "graph_admitted": execution.get("admitted_frames"), + "graph_delivered": execution.get("delivered_world_states"), + "graph_terminal_outcomes": terminal, + "tgs_timeline_frames": tgs.get("timeline", {}).get("frame_count"), + "tgs_available_lidar_frames": tgs.get("timeline", {}).get( + "available_lidar_frame_count" + ), + "tgs_capacity_drops": tgs_drops, + }, + "checks": checks, + "integrated_runtime_gate_passed": all(checks.values()), + "visual_quality_accepted": False, + "traversability_accepted": False, + "production_accepted": False, + "authority": profile["authority"], + "files": files, + } + identity = hashlib.sha256(canonical_json(document)).hexdigest() + document["result_id"] = f"m49-tgs-integrated-graph-shadow-{identity}" + output_path.parent.mkdir(parents=True, exist_ok=True) + output_path.write_text(json.dumps(document, indent=2, sort_keys=True) + "\n", encoding="utf-8") + return document + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("--profile", type=Path, required=True) + parser.add_argument("--graph-result", type=Path, required=True) + parser.add_argument("--graph-frames", type=Path, required=True) + parser.add_argument("--tgs-result", type=Path, required=True) + parser.add_argument("--tgs-timing", type=Path, required=True) + parser.add_argument("--telemetry", type=Path, required=True) + parser.add_argument("--output", type=Path, required=True) + parser.add_argument("--release-sha256", required=True) + arguments = parser.parse_args() + result = build( + profile_path=arguments.profile, + graph_result_path=arguments.graph_result, + graph_frames_path=arguments.graph_frames, + tgs_result_path=arguments.tgs_result, + tgs_timing_path=arguments.tgs_timing, + telemetry_path=arguments.telemetry, + output_path=arguments.output, + release_sha256=arguments.release_sha256, + ) + print( + json.dumps({"result_id": result["result_id"], "status": result["status"]}, sort_keys=True) + ) + return 0 if result["status"] == "passed" else 2 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/experiments/perception/worker/m49_t3_travel/run_tgs_full_shadow.cpp b/experiments/perception/worker/m49_t3_travel/run_tgs_full_shadow.cpp index 8c5802c..c22e641 100644 --- a/experiments/perception/worker/m49_t3_travel/run_tgs_full_shadow.cpp +++ b/experiments/perception/worker/m49_t3_travel/run_tgs_full_shadow.cpp @@ -1,4 +1,5 @@ #include +#include #include #include #include @@ -76,11 +77,38 @@ double milliseconds(Clock::duration duration) { return std::chrono::duration(duration).count(); } +void waitForSharedStart(const std::string& ready_path, const std::string& start_path) { + if (ready_path.empty() != start_path.empty()) { + throw std::runtime_error("shared-start paths must be configured together"); + } + if (ready_path.empty()) { + return; + } + if (std::filesystem::exists(ready_path)) { + throw std::runtime_error("shared-start ready file already exists"); + } + { + std::ofstream ready(ready_path); + ready << "ready\n"; + if (!ready) { + throw std::runtime_error("cannot publish TGS shared-start readiness"); + } + } + const auto deadline = Clock::now() + std::chrono::minutes(10); + while (!std::filesystem::is_regular_file(start_path)) { + if (Clock::now() >= deadline) { + throw std::runtime_error("TGS shared-start barrier timed out"); + } + std::this_thread::sleep_for(std::chrono::milliseconds(10)); + } +} + } // namespace int main(int argc, char** argv) { - if (argc != 5) { - std::cerr << "Usage: run_tgs_full_shadow \n"; + if (argc != 5 && argc != 8) { + std::cerr << "Usage: run_tgs_full_shadow " + " [target_rate_hz ready_file start_file]\n"; return 1; } try { @@ -89,6 +117,19 @@ int main(int argc, char** argv) { const std::string output_dir = argv[3]; const std::string timing_path = argv[4]; const auto schedule = readSchedule(schedule_path); + const double target_rate_hz = argc == 8 ? std::stod(argv[5]) : 0.0; + if (target_rate_hz < 0.0 || !std::isfinite(target_rate_hz)) { + throw std::runtime_error("invalid TGS target rate"); + } + const double source_duration_seconds = + schedule.back().session_seconds - schedule.front().session_seconds; + if (!(source_duration_seconds > 0.0)) { + throw std::runtime_error("invalid TGS source duration"); + } + const double recorded_rate_hz = + static_cast(schedule.size() - 1) / source_duration_seconds; + const double pacing_scale = + target_rate_hz > 0.0 ? recorded_rate_hz / target_rate_hz : 1.0; KittiLoader loader(sequence_dir); if (loader.size() != 3928) { throw std::runtime_error("full-shadow available LiDAR frame count changed"); @@ -103,12 +144,14 @@ int main(int argc, char** argv) { << "\ttgs_ms\tstage_wall_ms\tqueue_delay_ms\tcompletion_age_ms\tcapacity_drop\n"; timing << std::fixed << std::setprecision(6); + waitForSharedStart(argc == 8 ? argv[6] : "", argc == 8 ? argv[7] : ""); const double first_source_seconds = schedule.front().session_seconds; const auto run_started = Clock::now(); std::size_t expected_slot = 0; for (const auto& row : schedule) { const auto target = run_started + std::chrono::duration_cast( - std::chrono::duration(row.session_seconds - first_source_seconds)); + std::chrono::duration( + (row.session_seconds - first_source_seconds) * pacing_scale)); const auto before_wait = Clock::now(); if (before_wait < target) { std::this_thread::sleep_until(target); @@ -173,6 +216,9 @@ int main(int argc, char** argv) { throw std::runtime_error("full-shadow available frame accounting changed"); } std::cout << "[TGS-FULL] complete timeline=4489 available=3928\n"; + std::cout << "[TGS-FULL] recorded_rate_hz=" << recorded_rate_hz + << " target_rate_hz=" << (target_rate_hz > 0.0 ? target_rate_hz : recorded_rate_hz) + << "\n"; return 0; } catch (const std::exception& error) { std::cerr << "[TGS-FULL] " << error.what() << '\n'; diff --git a/experiments/perception/worker/m49_t3_travel/run_tgs_integrated_shadow.sh b/experiments/perception/worker/m49_t3_travel/run_tgs_integrated_shadow.sh new file mode 100644 index 0000000..4554bd8 --- /dev/null +++ b/experiments/perception/worker/m49_t3_travel/run_tgs_integrated_shadow.sh @@ -0,0 +1,26 @@ +#!/usr/bin/env bash +set -euo pipefail + +readonly BINARY=/shared/bin/run_tgs_full_shadow +readonly INPUT_ROOT=/shared/tgs/inputs +readonly OUTPUT_ROOT=/shared/tgs/outputs/causal_rolling_1s +readonly TIMING_PATH=/shared/tgs/tgs-full-timing.tsv +readonly READY_FILE=/shared/control/tgs.ready +readonly START_FILE=/shared/control/start.signal +readonly SOURCE_RATE_HZ=${M49_SOURCE_RATE_HZ:-12.0} + +test -x "${BINARY}" +test -f "${INPUT_ROOT}/input-manifest.json" +test -f "${INPUT_ROOT}/schedule.tsv" +test ! -e /shared/tgs/outputs +test ! -e "${TIMING_PATH}" +test ! -e "${READY_FILE}" +mkdir -p "${OUTPUT_ROOT}" +exec /usr/bin/time -v "${BINARY}" \ + "${INPUT_ROOT}/profiles/causal_rolling_1s" \ + "${INPUT_ROOT}/schedule.tsv" \ + "${OUTPUT_ROOT}" \ + "${TIMING_PATH}" \ + "${SOURCE_RATE_HZ}" \ + "${READY_FILE}" \ + "${START_FILE}" diff --git a/scripts/build_m49_tgs_integrated_graph_worker_artifact.py b/scripts/build_m49_tgs_integrated_graph_worker_artifact.py new file mode 100644 index 0000000..91211ce --- /dev/null +++ b/scripts/build_m49_tgs_integrated_graph_worker_artifact.py @@ -0,0 +1,257 @@ +#!/usr/bin/env python3 +"""Build a clean-revision Worker 006 release for the integrated M4.9 shadow.""" + +from __future__ import annotations + +import argparse +import gzip +import hashlib +import io +import json +import os +import re +import subprocess +import tarfile +import tempfile +from pathlib import Path + +REPOSITORY_ROOT = Path(__file__).resolve().parents[1] +WHEEL_NAME = "nodedc_mission_core-0.1.0-py3-none-any.whl" +PATCH_ID = re.compile(r"^[A-Za-z0-9._-]{1,96}$") +SOURCES = ( + Path("experiments/perception/worker/Invoke-M49TgsIntegratedGraphShadow.ps1"), + Path("experiments/perception/run_m48s_reference_graph_shadow_worker.py"), + Path("experiments/perception/worker/rf_detr_large_native_kb4_config.pbtxt"), + Path("experiments/perception/worker/m49_t3_travel/prepare_tgs_fail_closed_inputs.py"), + Path("experiments/perception/worker/m49_t3_travel/prepare_tgs_full_shadow_inputs.py"), + Path("experiments/perception/worker/m49_t3_travel/run_tgs_full_shadow.cpp"), + Path("experiments/perception/worker/m49_t3_travel/build_tgs_full_shadow_binary.sh"), + Path("experiments/perception/worker/m49_t3_travel/run_tgs_integrated_shadow.sh"), + Path("experiments/perception/worker/m49_t3_travel/build_tgs_full_shadow_evidence.py"), + Path("experiments/perception/worker/m49_t3_travel/build_tgs_integrated_graph_evidence.py"), + Path("config/perception/m49-tgs-integrated-graph-shadow-v1.json"), + Path("config/perception/m49-tgs-full-shadow-v1.json"), + Path("config/perception/m48n-rf-detr-native-reference-graph-shadow-v0.json"), + Path("config/perception/m4-recorded-realtime-baseline-v1.json"), + Path("config/perception/rf-detr-large-native-kb4-risk-shadow-v0.json"), + Path("config/perception/m4-geometry-association-v1.json"), + Path("config/perception/m4-temporal-motion-v1.json"), + Path("config/perception/m4-rolling-local-map-v1.json"), + Path("config/perception/m4-replay-threat-v3.json"), +) + + +class ArtifactBuildError(RuntimeError): + """The integrated Worker release cannot be built from its declared revision.""" + + +def sha256_file(path: Path) -> str: + digest = hashlib.sha256() + with path.open("rb") as stream: + for chunk in iter(lambda: stream.read(1024 * 1024), b""): + digest.update(chunk) + return digest.hexdigest() + + +def git_revision() -> str: + result = subprocess.run( + ["git", "rev-parse", "HEAD"], + cwd=REPOSITORY_ROOT, + check=True, + capture_output=True, + text=True, + ) + revision = result.stdout.strip() + if re.fullmatch(r"[a-f0-9]{40}", revision) is None: + raise ArtifactBuildError("Git revision is not a full SHA-1") + return revision + + +def materialize_revision(revision: str, destination: Path) -> None: + archive_path = destination.parent / "source.tar" + subprocess.run( + ["git", "archive", "--format=tar", "--output", str(archive_path), revision], + cwd=REPOSITORY_ROOT, + check=True, + capture_output=True, + ) + destination.mkdir() + root = destination.resolve() + with tarfile.open(archive_path, "r:") as archive: + for member in archive.getmembers(): + target = (destination / member.name).resolve() + if target != root and root not in target.parents: + raise ArtifactBuildError("Git archive contains an unsafe path") + archive.extractall(destination) + + +def build_wheel(source_root: Path, output: Path) -> Path: + environment = os.environ.copy() + environment["SOURCE_DATE_EPOCH"] = "0" + result = subprocess.run( + ["uv", "build", "--wheel", "--out-dir", str(output)], + cwd=source_root, + env=environment, + check=False, + capture_output=True, + text=True, + ) + if result.returncode != 0: + detail = (result.stderr or result.stdout).strip() + raise ArtifactBuildError(f"wheel build failed: {detail}") + wheel = output / WHEEL_NAME + if not wheel.is_file() or wheel.is_symlink(): + raise ArtifactBuildError("expected Worker wheel was not built") + return wheel + + +def tar_info(path: Path, arcname: str) -> tarfile.TarInfo: + info = tarfile.TarInfo(arcname) + info.uid = info.gid = 0 + info.uname = info.gname = "root" + info.mtime = 0 + if path.is_dir(): + info.type = tarfile.DIRTYPE + info.mode = 0o755 + else: + info.type = tarfile.REGTYPE + info.mode = 0o755 if path.suffix in {".sh", ".ps1", ".py"} else 0o644 + info.size = path.stat().st_size + return info + + +def write_archive(stage: Path, target: Path) -> None: + members = [stage / "manifest.env", stage / "files.txt", stage / "payload"] + members.extend(sorted((stage / "payload").rglob("*"))) + target.parent.mkdir(parents=True, exist_ok=True) + with ( + target.open("wb") as raw, + gzip.GzipFile(filename="", mode="wb", fileobj=raw, compresslevel=9, mtime=0) as compressed, + tarfile.open(fileobj=compressed, mode="w", format=tarfile.PAX_FORMAT) as archive, + ): + for path in members: + info = tar_info(path, path.relative_to(stage).as_posix()) + if path.is_file(): + with path.open("rb") as stream: + archive.addfile(info, stream) + else: + archive.addfile(info, io.BytesIO()) + + +def build_artifact( + patch_id: str, + output_directory: Path, + *, + revision: str | None = None, + source_root: Path | None = None, +) -> dict[str, object]: + if PATCH_ID.fullmatch(patch_id) is None: + raise ArtifactBuildError("patch id is invalid") + selected_revision = revision or git_revision() + if re.fullmatch(r"[a-f0-9]{40}", selected_revision) is None: + raise ArtifactBuildError("artifact revision is invalid") + with tempfile.TemporaryDirectory(prefix="mission-core-m49-integrated-") as directory: + stage = Path(directory) + snapshot = source_root + if snapshot is None: + snapshot = stage / "source" + materialize_revision(selected_revision, snapshot) + sources = tuple(snapshot / relative for relative in SOURCES) + if any(path.is_symlink() or not path.is_file() for path in sources): + raise ArtifactBuildError("release input is not a regular file") + payload = stage / "payload" + payload.mkdir() + wheel = build_wheel(snapshot, stage / "wheel") + copied: list[Path] = [] + for source in sources: + destination = payload / source.name + if destination.exists(): + raise ArtifactBuildError("release payload file names are not unique") + destination.write_bytes(source.read_bytes()) + copied.append(destination) + wheel_destination = payload / WHEEL_NAME + wheel_destination.write_bytes(wheel.read_bytes()) + copied.append(wheel_destination) + release = { + "schema_version": "missioncore.m49-tgs-integrated-graph-worker-release/v1", + "patch_id": patch_id, + "transition": "m49-tgs-native-risk-integrated-shadow/v1", + "code_revision": selected_revision, + "worker_id": "worker-006", + "source_pack_sha256": ( + "0685d24219d8236caf8b7f1685e93f6d6b59e7fd015a768d88a92bbe8b154944" + ), + "expected_frames": 4489, + "requested_source_rate_hz": 12.0, + "native_engine_sha256": ( + "b8a40b3580edff001ec9680de68707242294ff590ab296000fae371f1083f695" + ), + "images": { + "travel": "sha256:7b412020f4d8392d1d1ed1b33beadc44140f0ea8f781e62dd69796042334300f", + "parity": "sha256:ceb13548617e4bd3f619766bfdff00af3fa5160946b367828da6d2233dcdcba0", + "runtime": ( + "sha256:58df7489c3f2276f9591d500a012dee03e23d35543ce3c390b4c001e6bf90794" + ), + }, + "authority": { + "visual_quality_accepted": False, + "traversability_accepted": False, + "physical_free_space_accepted": False, + "commands_enabled": False, + "actuation_allowed": False, + "navigation_or_safety_accepted": False, + "production_accepted": False, + }, + "scope": { + "gauss_or_playcanvas_action": "none", + "durable_worker_action": "none", + "canonical_triton_action": "none", + }, + "files": { + path.name: {"sha256": sha256_file(path), "bytes": path.stat().st_size} + for path in sorted(copied) + }, + } + release_path = payload / "release.json" + release_path.write_text( + json.dumps(release, indent=2, sort_keys=True) + "\n", encoding="utf-8" + ) + payload_files = sorted((*release["files"], release_path.name)) + (stage / "manifest.env").write_text( + f"id={patch_id}\ncomponent=mission-core-worker\ntype=shadow-release\n", + encoding="utf-8", + ) + (stage / "files.txt").write_text("\n".join(payload_files) + "\n", encoding="utf-8") + target = output_directory.resolve() / f"nodedc-{patch_id}.tgz" + write_archive(stage, target) + return { + "ok": True, + "patch_id": patch_id, + "artifact": str(target), + "sha256": sha256_file(target), + "code_revision": selected_revision, + "wheel_sha256": release["files"][WHEEL_NAME]["sha256"], + "payload_files": payload_files, + "transition": release["transition"], + } + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("patch_id") + parser.add_argument( + "--output-directory", + type=Path, + default=REPOSITORY_ROOT / ".runtime/worker-artifacts", + ) + arguments = parser.parse_args() + try: + result = build_artifact(arguments.patch_id, arguments.output_directory) + except (ArtifactBuildError, OSError, subprocess.SubprocessError) as exc: + parser.error(str(exc)) + print(json.dumps(result, indent=2, sort_keys=True)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/test_m48s_reference_graph_timing.py b/tests/test_m48s_reference_graph_timing.py index 64dea9a..75c651d 100644 --- a/tests/test_m48s_reference_graph_timing.py +++ b/tests/test_m48s_reference_graph_timing.py @@ -3,6 +3,7 @@ from __future__ import annotations import importlib.util import time from pathlib import Path +from threading import Thread from unittest.mock import patch from k1link.perception.detector import DetectorFrameTiming @@ -157,3 +158,33 @@ def test_cyclic_gc_policy_collects_outside_hot_loop_and_restores_state() -> None "pre_collected": 3, "post_collected": 0, } + + +def test_shared_start_barrier_publishes_readiness_and_waits_for_release( + tmp_path: Path, +) -> None: + ready = tmp_path / "graph.ready" + start = tmp_path / "start.signal" + completed: list[bool] = [] + + thread = Thread( + target=lambda: ( + RUNNER.wait_for_shared_start( + ready_file=ready, + start_file=start, + timeout_seconds=1.0, + ), + completed.append(True), + ) + ) + thread.start() + deadline = time.monotonic() + 1.0 + while not ready.exists() and time.monotonic() < deadline: + time.sleep(0.005) + assert ready.read_text(encoding="utf-8") == "ready\n" + assert completed == [] + + start.write_text("start\n", encoding="utf-8") + thread.join(timeout=1.0) + + assert completed == [True] diff --git a/tests/test_m49_tgs_integrated_graph_shadow.py b/tests/test_m49_tgs_integrated_graph_shadow.py new file mode 100644 index 0000000..d950c72 --- /dev/null +++ b/tests/test_m49_tgs_integrated_graph_shadow.py @@ -0,0 +1,208 @@ +from __future__ import annotations + +import importlib.util +import json +import tarfile +from pathlib import Path + +REPOSITORY_ROOT = Path(__file__).resolve().parents[1] +EVIDENCE_PATH = ( + REPOSITORY_ROOT + / "experiments/perception/worker/m49_t3_travel/build_tgs_integrated_graph_evidence.py" +) +ARTIFACT_PATH = REPOSITORY_ROOT / "scripts/build_m49_tgs_integrated_graph_worker_artifact.py" + + +def load_module(name: str, path: Path): + spec = importlib.util.spec_from_file_location(name, path) + assert spec is not None and spec.loader is not None + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + return module + + +EVIDENCE = load_module("m49_tgs_integrated_evidence", EVIDENCE_PATH) +ARTIFACT = load_module("m49_tgs_integrated_artifact", ARTIFACT_PATH) + + +def test_integrated_gate_joins_all_frames_and_preserves_false_authority(tmp_path: Path) -> None: + profile_path = tmp_path / "profile.json" + profile_path.write_text( + json.dumps( + { + "schema_version": EVIDENCE.PROFILE_SCHEMA, + "profile_id": "test", + "source": { + "source_id": "RAVNOVES00", + "source_pack_sha256": "a" * 64, + "requested_source_rate_hz": 12.0, + }, + "stages": { + "reference_graph": { + "graph_config_sha256": "b" * 64, + "detector_profile_sha256": "z" * 64, + }, + "tgs": { + "profile_sha256": "c" * 64, + "linked_accepted_result_id": "m49-tgs-full-shadow-" + "d" * 64, + }, + }, + "acceptance": { + "minimum_delivery_ratio": 1.0, + "reference_world_state_fps": 11.79902, + "maximum_world_state_fps_regression_fraction": 0.05, + "minimum_effective_world_state_fps": 11.209069, + "maximum_world_state_completion_p95_ms": 125.0, + "candidate_stage_p95_ms_max": 25.0, + "candidate_stage_p99_ms_max": 50.0, + "combined_output_age_p99_ms_max": 125.0, + "capacity_drop_count_max": 0, + }, + "authority": { + "visual_quality_accepted": False, + "traversability_accepted": False, + "physical_free_space_accepted": False, + "commands_enabled": False, + "actuation_allowed": False, + "navigation_or_safety_accepted": False, + "production_accepted": False, + }, + } + ), + encoding="utf-8", + ) + graph_result = tmp_path / "graph-result.json" + graph_result.write_text( + json.dumps( + { + "schema_version": EVIDENCE.GRAPH_SCHEMA, + "execution": { + "admitted_frames": EVIDENCE.FRAME_COUNT, + "delivered_world_states": EVIDENCE.FRAME_COUNT, + "effective_world_state_fps": 11.75, + "delivery_ratio": 1.0, + "terminal_outcomes": {"delivered": EVIDENCE.FRAME_COUNT}, + "requested_source_rate_hz": 12.0, + }, + "identity": { + "inputs": { + "graph_config": "b" * 64, + "detector_profile": "z" * 64, + } + }, + "metrics": { + "world_state_completion_age_ms": {"p95": 40.0}, + "gpu": {"sample_count": 2}, + }, + "evidence_integrity_gate_passed": True, + } + ), + encoding="utf-8", + ) + graph_frames = tmp_path / "graph-frames.jsonl" + graph_frames.write_text( + "".join( + json.dumps( + { + "source_envelope": {"sequence": index}, + "completion_age_ns": 40_000_000, + } + ) + + "\n" + for index in range(EVIDENCE.FRAME_COUNT) + ), + encoding="utf-8", + ) + tgs_result = tmp_path / "tgs-result.json" + tgs_result.write_text( + json.dumps( + { + "schema_version": EVIDENCE.TGS_SCHEMA, + "status": "passed", + "config_sha256": "c" * 64, + "timeline": { + "frame_count": EVIDENCE.FRAME_COUNT, + "available_lidar_frame_count": 3928, + }, + "performance": { + "candidate_tgs_ms": {"p95": 2.0, "p99": 3.0}, + "completion_age_ms": {"p99": 5.0}, + "capacity_drop_count": 0, + }, + } + ), + encoding="utf-8", + ) + tgs_timing = tmp_path / "tgs-timing.tsv" + tgs_timing.write_text( + "timeline_frame_index\tcompletion_age_ms\n" + + "".join(f"{index}\t5.0\n" for index in range(EVIDENCE.FRAME_COUNT)), + encoding="utf-8", + ) + telemetry = tmp_path / "telemetry.jsonl" + telemetry.write_text( + "".join( + json.dumps( + { + "role": role, + "cpu_percent": "10.0%", + "memory_usage": "1GiB / 64GiB", + "memory_percent": "1.56%", + } + ) + + "\n" + for role in ("graph", "tgs", "triton") + ), + encoding="utf-8", + ) + output = tmp_path / "result.json" + + result = EVIDENCE.build( + profile_path=profile_path, + graph_result_path=graph_result, + graph_frames_path=graph_frames, + tgs_result_path=tgs_result, + tgs_timing_path=tgs_timing, + telemetry_path=telemetry, + output_path=output, + release_sha256="e" * 64, + ) + + assert result["status"] == "passed" + assert result["source"]["joined_frame_count"] == EVIDENCE.FRAME_COUNT + assert result["performance"]["combined_output_age_ms"]["p99"] == 40.0 + assert result["checks"]["authority_remains_false"] is True + assert result["production_accepted"] is False + + +def test_worker_artifact_is_deterministic_and_excludes_gauss(monkeypatch, tmp_path: Path) -> None: + def fake_wheel(_source_root: Path, output: Path) -> Path: + output.mkdir(parents=True, exist_ok=True) + wheel = output / ARTIFACT.WHEEL_NAME + wheel.write_bytes(b"clean committed wheel\n") + return wheel + + monkeypatch.setattr(ARTIFACT, "build_wheel", fake_wheel) + revision = "f" * 40 + first = ARTIFACT.build_artifact( + "mission-core-m49-integrated-unit-001", + tmp_path / "first", + revision=revision, + source_root=REPOSITORY_ROOT, + ) + second = ARTIFACT.build_artifact( + "mission-core-m49-integrated-unit-001", + tmp_path / "second", + revision=revision, + source_root=REPOSITORY_ROOT, + ) + + assert Path(first["artifact"]).read_bytes() == Path(second["artifact"]).read_bytes() + with tarfile.open(first["artifact"], "r:gz") as archive: + names = set(archive.getnames()) + release_stream = archive.extractfile("payload/release.json") + assert release_stream is not None + release = json.loads(release_stream.read()) + assert not any("gauss" in name.lower() or "playcanvas" in name.lower() for name in names) + assert release["scope"]["gauss_or_playcanvas_action"] == "none" + assert all(value is False for value in release["authority"].values())