diff --git a/config/observatory-portable-run-definitions.json b/config/observatory-portable-run-definitions.json index 88e128b..ccf916d 100644 --- a/config/observatory-portable-run-definitions.json +++ b/config/observatory-portable-run-definitions.json @@ -158,8 +158,8 @@ { "setup_id": "m49-tgs-portable-v2", "definition_id": "m49-tgs-portable", - "version": 2, - "definition_sha256": "73611f24d70319ea1edca428726d6538a3cbad012a415cc0c1a7ecb7d9b4d910", + "version": 3, + "definition_sha256": "f56d6321bd794ccdfb7d2e3b05d044b11f616ffb81ee29517386cc253046d4eb", "source_requirements": { "plugin_id": "nodedc.device.xgrids-lixelkity-k1", "archive_id": "xgrids-k1.viewer-live.evidence", @@ -194,6 +194,11 @@ "component_id": "m49-tgs-portable-profile-v2", "kind": "configuration", "sha256": "6128d6af7e6137f9a9473db045e3b155e2105319159f17c32f344b4aedf823a9" + }, + { + "component_id": "m49-tgs-portable-runner-v1", + "kind": "runner", + "sha256": "e3bb2e91c70712eff74e8e69718075a407e042fb494616d65984500da42b21a9" } ], "models": [], @@ -218,12 +223,12 @@ }, "executor": { "contour_id": "worker-006", - "state": "not-installed", - "release_id": null, - "release_sha256": null, - "image_sha256": null, - "reason_code": "m49-portable-executor-not-installed", - "reason": "A source-independent M4.9 TGS executor release and image are not sealed or installed on Worker 006." + "state": "ready", + "release_id": "m49-tgs-portable-executor-v1", + "release_sha256": "c5b0670d943fe0452ef4bbfbc144ab2439a1a674f9ef164798ad9f8b1ecc29fa", + "image_sha256": "f9278ab21aa65045be993dd19bffc25f49955e19598893ac78cc4761ca63ecf3", + "reason_code": null, + "reason": null }, "authority": { "commands_enabled": false, diff --git a/config/observatory-worker-runtime-candidates.json b/config/observatory-worker-runtime-candidates.json index 9734d0e..16c5b71 100644 --- a/config/observatory-worker-runtime-candidates.json +++ b/config/observatory-worker-runtime-candidates.json @@ -181,15 +181,46 @@ "adapter_id": "m49-tgs-worker006-portable-v2", "setup_id": "m49-tgs-portable-v2", "definition_id": "m49-tgs-portable", - "definition_version": 2, - "definition_sha256": "73611f24d70319ea1edca428726d6538a3cbad012a415cc0c1a7ecb7d9b4d910", + "definition_version": 3, + "definition_sha256": "f56d6321bd794ccdfb7d2e3b05d044b11f616ffb81ee29517386cc253046d4eb", "source_adapter_sha256": "4e12be6d2503e2e237eddb290b28d6a7d16cf983855b6d8e6b4e3b65d2feb0de", "model_manifest_sha256": "489a43448f720a9b5c7993dc8279d167b77191a586f0d87b6d38b81cf728e2f1", "resource_profile_sha256": "49e373f1cbf314e7fab2d176db295f3d3d9ddbbcbd724bbde988ef60ad767dee", "result_contract_sha256": "9dd80c8e2504559d2156fca933de6eb27901e35305e6853aeb84707e1cb13892", - "state": "blocked", - "executor": null, + "state": "ready", + "executor": { + "release_id": "m49-tgs-portable-executor-v1", + "release_sha256": "c5b0670d943fe0452ef4bbfbc144ab2439a1a674f9ef164798ad9f8b1ecc29fa", + "image_sha256": "f9278ab21aa65045be993dd19bffc25f49955e19598893ac78cc4761ca63ecf3" + }, "reusable_assets": [ + { + "asset_id": "m49-portable-compiled-runner", + "kind": "local-file", + "sha256": "7be449392ef161fb8713b4c984705d2373bc3cd05645332d92b88a8bff0c7db3", + "byte_length": 274168, + "component_id": null, + "model_release_id": null, + "model_artifact_role": null + }, + { + "asset_id": "m49-portable-compiled-runner-build-seal", + "kind": "local-file", + "sha256": "e3bb2e91c70712eff74e8e69718075a407e042fb494616d65984500da42b21a9", + "byte_length": 1014, + "component_id": null, + "model_release_id": null, + "model_artifact_role": null + }, + { + "asset_id": "m49-portable-executor-release", + "kind": "local-file", + "sha256": "c5b0670d943fe0452ef4bbfbc144ab2439a1a674f9ef164798ad9f8b1ecc29fa", + "byte_length": 1693, + "component_id": null, + "model_release_id": null, + "model_artifact_role": null + }, { "asset_id": "m49-portable-profile", "kind": "definition-component", @@ -198,39 +229,12 @@ "component_id": "m49-tgs-portable-profile-v2", "model_release_id": null, "model_artifact_role": null - }, - { - "asset_id": "m49-portable-runner-manifest", - "kind": "local-file", - "sha256": "940aeb7bbcf1b28c15482e2d8bedd47ce4d3667137a5cd4702bbc7b754475baa", - "byte_length": 1825, - "component_id": null, - "model_release_id": null, - "model_artifact_role": null - }, - { - "asset_id": "m49-portable-runner-source", - "kind": "local-file", - "sha256": "52813392aabd02efc5c2b8f7c22ed88e3ef4cc8ad3aafeba2792efe503e29fe9", - "byte_length": 11052, - "component_id": null, - "model_release_id": null, - "model_artifact_role": null }, { - "asset_id": "m49-portable-runner-wrapper", + "asset_id": "m49-portable-worker-installation-receipt", "kind": "local-file", - "sha256": "2d6c32560682647f868e4ce4c2605749f17c60482a609f8c03ff951411f48ffb", - "byte_length": 858, - "component_id": null, - "model_release_id": null, - "model_artifact_role": null - }, - { - "asset_id": "m49-portable-smoke", - "kind": "local-file", - "sha256": "b8bc1cf9f21f69904bb64c6384db6fe301faf10f8041cc405a77b011602eaf04", - "byte_length": 1330, + "sha256": "b560ff9e02746cbb760502f2a3b4b7bd564f95ab52d1149d189927a326b1645a", + "byte_length": 2672, "component_id": null, "model_release_id": null, "model_artifact_role": null @@ -248,15 +252,15 @@ "phases": [ { "phase_id": "source-delivery", - "state": "missing" + "state": "implemented" }, { "phase_id": "camera-lidar-timeline-materializer", - "state": "missing" + "state": "implemented" }, { "phase_id": "portable-tgs-input-materializer", - "state": "missing" + "state": "implemented" }, { "phase_id": "portable-tgs-runner", @@ -264,29 +268,21 @@ }, { "phase_id": "result-v2-assembler", - "state": "missing" + "state": "implemented" }, { "phase_id": "observatory-result-publisher", - "state": "missing" + "state": "implemented" } ], - "blockers": [ - "executor-release-unsealed", - "portable-camera-lidar-timeline-unimplemented", - "portable-result-assembler-unimplemented", - "portable-source-delivery-unimplemented", - "portable-tgs-input-unimplemented", - "portable-tgs-runner-unsealed", - "result-publisher-integration-unaccepted" - ], + "blockers": [], "authority": { "commands_enabled": false, "actuation_allowed": false, "navigation_or_safety_accepted": false, "production_accepted": false }, - "candidate_sha256": "6b4625f2fdb4703efed852bca100f04398284bd6a894b65e4bfd51a6532430bf" + "candidate_sha256": "65cd2063146a1dd320e30d5f4e21e4bf0aab1ff683e846926cbbfbe25a9f8a5e" } ] } diff --git a/experiments/perception/worker/observatory_portable/Invoke-M49PortableExecutorPromotion.ps1 b/experiments/perception/worker/observatory_portable/Invoke-M49PortableExecutorPromotion.ps1 new file mode 100644 index 0000000..0876beb --- /dev/null +++ b/experiments/perception/worker/observatory_portable/Invoke-M49PortableExecutorPromotion.ps1 @@ -0,0 +1,401 @@ +[CmdletBinding()] +param() + +$ErrorActionPreference = "Stop" +$ProgressPreference = "SilentlyContinue" + +# This is one exact, additive Worker 006 transition. It never recompiles the +# runner, retags an image, starts a container, or changes the candidate files. +$WorkerId = "worker-006" +$ExpectedComputer = "DESKTOP-OPJ8J04" +$CandidateId = "beef090b79ab0e019cff976033d0446c45d4746ee921e02e8076a5a67f21817e" +$CandidateArchiveSha256 = "70cd95e611381e0f1d362f2e201b7615bb9530d69e26aafbf949b19aa523c630" +$CandidateReleaseSha256 = "a4fa00bf02c44e584667d4d48c6ec6a55832c704e7b47786dea638b59945c5db" +$CandidateBuildSealSha256 = "4b29155896db9f48461ea4af6b4c02a2c2b0224dd65748e827e5ca1125b8648d" +$CandidateReceiptSha256 = "a6b32fde8c09faa5419fb9c75a8207d4016a1d044ede3d6d6148da00b3092b8c" +$SourceRevision = "2f6e45bc96f2c9f65b8830dc536784e9aa85fd87" +$ReleaseId = "m49-tgs-portable-executor-v1" +$ReleaseParent = "D:\NDC_MISSIONCORE\runtime\releases\observatory-portable" +$CandidateRootName = "m49-tgs-portable-candidate-$CandidateId" +$RuntimeReleaseRoot = "/release" +$BaseImageSha256 = "7b412020f4d8392d1d1ed1b33beadc44140f0ea8f781e62dd69796042334300f" +$ExecutorImageSha256 = "f9278ab21aa65045be993dd19bffc25f49955e19598893ac78cc4761ca63ecf3" +$ExecutorImageSizeBytes = 1030212146 +$ExecutorImageTag = "ndc/mission-core-m49-tgs-portable-executor:beef090b79ab0e01-candidate" +$RunnerSha256 = "7be449392ef161fb8713b4c984705d2373bc3cd05645332d92b88a8bff0c7db3" +$RunnerBytes = 274168 +$ProfileSha256 = "6128d6af7e6137f9a9473db045e3b155e2105319159f17c32f344b4aedf823a9" +$ProfileBytes = 1683 +$RuntimeBuildSealSha256 = "e3bb2e91c70712eff74e8e69718075a407e042fb494616d65984500da42b21a9" +$RuntimeBuildSealBytes = 1014 +$InstalledReleaseSha256 = "c5b0670d943fe0452ef4bbfbc144ab2439a1a674f9ef164798ad9f8b1ecc29fa" +$InstalledReleaseBytes = 1693 +$ReadyReceiptSha256 = "b560ff9e02746cbb760502f2a3b4b7bd564f95ab52d1149d189927a326b1645a" +$ReadyReceiptBytes = 2672 +$ResultContractSha256 = "9dd80c8e2504559d2156fca933de6eb27901e35305e6853aeb84707e1cb13892" +$RunnerSourceSha256 = "52813392aabd02efc5c2b8f7c22ed88e3ef4cc8ad3aafeba2792efe503e29fe9" +$RunnerWrapperSha256 = "2d6c32560682647f868e4ce4c2605749f17c60482a609f8c03ff951411f48ffb" +$FixtureSmokeSha256 = "b8bc1cf9f21f69904bb64c6384db6fe301faf10f8041cc405a77b011602eaf04" +$ProtectedRuntime = @( + [ordered]@{ name = "ndc-mission-core-triton"; container_id = "4232fb040062a8384809e73612baa92b343ddcb60f00c23e7091eb4909959223" }, + [ordered]@{ name = "ndc-mission-core-perception-worker"; container_id = "db2024d05098a6beb6b73bbf43f02c88ede586a4664b2c91182f436a42b3e3ff" }, + [ordered]@{ name = "ndc-gaussian-pipeline-gaussian-gateway-1"; container_id = "67c6b7bdf37a3745a67126e5acc5929230b4f51e9e353dac4828df90c5f9df9f" }, + [ordered]@{ name = "ndc-gaussian-pipeline-gaussian-pipeline-1"; container_id = "ae855b3b3414ae94dcd9bdc2c8223fe278109484ce828cd82fe5ef0adecf17c9" }, + [ordered]@{ name = "ndc-gaussian-pipeline-gaussian-terrain-executor-1"; container_id = "4624a9fbcad8cda9d0ff0bc0a5b9bb09eb2340ec93f694033c5b3d2cd1806265" } +) + +$RuntimeBuildSealPayload = @' +{"authority":{"actuation_allowed":false,"commands_enabled":false,"navigation_or_safety_accepted":false,"production_accepted":false},"binary":{"byte_length":274168,"file_name":"run_m49_tgs_portable","format":"elf","sha256":"7be449392ef161fb8713b4c984705d2373bc3cd05645332d92b88a8bff0c7db3"},"build_image_sha256":"7b412020f4d8392d1d1ed1b33beadc44140f0ea8f781e62dd69796042334300f","compiler_contract":{"compiler":"g++","eigen_include":"/usr/include/eigen3","flags":["-O3","-DNDEBUG","-pthread"],"language_standard":"c++17","travel_include":"/opt/travel/src/TRAVEL/cpp/travel/core"},"profile_sha256":"6128d6af7e6137f9a9473db045e3b155e2105319159f17c32f344b4aedf823a9","runner_source_sha256":"52813392aabd02efc5c2b8f7c22ed88e3ef4cc8ad3aafeba2792efe503e29fe9","runner_wrapper_sha256":"2d6c32560682647f868e4ce4c2605749f17c60482a609f8c03ff951411f48ffb","schema_version":"missioncore.m49-tgs-portable-compiled-runner-build/v1","source_revision":"2f6e45bc96f2c9f65b8830dc536784e9aa85fd87","source_state":"committed-snapshot"} +'@.Trim() +$InstalledReleasePayload = @' +{"authority":{"actuation_allowed":false,"commands_enabled":false,"navigation_or_safety_accepted":false,"production_accepted":false},"base_image_sha256":"7b412020f4d8392d1d1ed1b33beadc44140f0ea8f781e62dd69796042334300f","blockers":[],"candidate":{"candidate_sha256":"beef090b79ab0e019cff976033d0446c45d4746ee921e02e8076a5a67f21817e","release_sha256":"a4fa00bf02c44e584667d4d48c6ec6a55832c704e7b47786dea638b59945c5db","worker_installation_receipt_sha256":"a6b32fde8c09faa5419fb9c75a8207d4016a1d044ede3d6d6148da00b3092b8c"},"compiled_runner":{"byte_length":274168,"relative_path":"run_m49_tgs_portable","sha256":"7be449392ef161fb8713b4c984705d2373bc3cd05645332d92b88a8bff0c7db3"},"compiled_runner_build_seal":{"byte_length":1014,"relative_path":"compiled-runner-build.json","sha256":"e3bb2e91c70712eff74e8e69718075a407e042fb494616d65984500da42b21a9"},"executor_image":{"sha256":"f9278ab21aa65045be993dd19bffc25f49955e19598893ac78cc4761ca63ecf3","size_bytes":1030212146},"fixture_smoke":{"classified_points_per_frame":124637,"frame_count":2,"input_points_per_frame":124668,"network":"none","result":"passed","source":"pinned TRAVEL KITTI fixture"},"profile":{"byte_length":1683,"relative_path":"m49-tgs-portable-v2.json","sha256":"6128d6af7e6137f9a9473db045e3b155e2105319159f17c32f344b4aedf823a9"},"release_id":"m49-tgs-portable-executor-v1","result_contract_sha256":"9dd80c8e2504559d2156fca933de6eb27901e35305e6853aeb84707e1cb13892","schema_version":"missioncore.m49-tgs-portable-executor-installed/v1","source_archive_sha256":"70cd95e611381e0f1d362f2e201b7615bb9530d69e26aafbf949b19aa523c630","source_revision":"2f6e45bc96f2c9f65b8830dc536784e9aa85fd87","state":"ready","worker_id":"worker-006"} +'@.Trim() +$ReadyReceiptPayload = @' +{"authority":{"actuation_allowed":false,"commands_enabled":false,"navigation_or_safety_accepted":false,"production_accepted":false},"base_image_sha256":"7b412020f4d8392d1d1ed1b33beadc44140f0ea8f781e62dd69796042334300f","blockers":[],"candidate_release_sha256":"a4fa00bf02c44e584667d4d48c6ec6a55832c704e7b47786dea638b59945c5db","candidate_sha256":"beef090b79ab0e019cff976033d0446c45d4746ee921e02e8076a5a67f21817e","candidate_worker_installation_receipt_sha256":"a6b32fde8c09faa5419fb9c75a8207d4016a1d044ede3d6d6148da00b3092b8c","computer_name":"DESKTOP-OPJ8J04","executor_image_sha256":"f9278ab21aa65045be993dd19bffc25f49955e19598893ac78cc4761ca63ecf3","files":[{"asset_id":"m49-portable-compiled-runner","byte_length":274168,"relative_path":"run_m49_tgs_portable","sha256":"7be449392ef161fb8713b4c984705d2373bc3cd05645332d92b88a8bff0c7db3"},{"asset_id":"m49-portable-compiled-runner-build-seal","byte_length":1014,"relative_path":"compiled-runner-build.json","sha256":"e3bb2e91c70712eff74e8e69718075a407e042fb494616d65984500da42b21a9"},{"asset_id":"m49-portable-executor-release","byte_length":1693,"relative_path":"executor-release.json","sha256":"c5b0670d943fe0452ef4bbfbc144ab2439a1a674f9ef164798ad9f8b1ecc29fa"},{"asset_id":"m49-portable-profile","byte_length":1683,"relative_path":"m49-tgs-portable-v2.json","sha256":"6128d6af7e6137f9a9473db045e3b155e2105319159f17c32f344b4aedf823a9"}],"fixture_smoke":"passed","legacy_m49_task_state":"Ready","protected_runtime":[{"container_id":"4232fb040062a8384809e73612baa92b343ddcb60f00c23e7091eb4909959223","name":"ndc-mission-core-triton"},{"container_id":"db2024d05098a6beb6b73bbf43f02c88ede586a4664b2c91182f436a42b3e3ff","name":"ndc-mission-core-perception-worker"},{"container_id":"67c6b7bdf37a3745a67126e5acc5929230b4f51e9e353dac4828df90c5f9df9f","name":"ndc-gaussian-pipeline-gaussian-gateway-1"},{"container_id":"ae855b3b3414ae94dcd9bdc2c8223fe278109484ce828cd82fe5ef0adecf17c9","name":"ndc-gaussian-pipeline-gaussian-pipeline-1"},{"container_id":"4624a9fbcad8cda9d0ff0bc0a5b9bb09eb2340ec93f694033c5b3d2cd1806265","name":"ndc-gaussian-pipeline-gaussian-terrain-executor-1"}],"receipt_state":"installed-ready","release_id":"m49-tgs-portable-executor-v1","release_root":"D:\\NDC_MISSIONCORE\\runtime\\releases\\observatory-portable\\m49-tgs-portable-candidate-beef090b79ab0e019cff976033d0446c45d4746ee921e02e8076a5a67f21817e\\ready","release_sha256":"c5b0670d943fe0452ef4bbfbc144ab2439a1a674f9ef164798ad9f8b1ecc29fa","runtime_release_root":"/release","schema_version":"missioncore.m49-tgs-portable-worker-installation-ready-receipt/v1","source_revision":"2f6e45bc96f2c9f65b8830dc536784e9aa85fd87","worker_id":"worker-006"} +'@.Trim() + +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 + 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 Get-Sha256([string]$Path) { + return (Get-FileHash -Algorithm SHA256 -LiteralPath $Path).Hash.ToLowerInvariant() +} + +function Get-PayloadSha256([string]$Value) { + $encoding = New-Object System.Text.UTF8Encoding($false) + $bytes = $encoding.GetBytes($Value) + $algorithm = [Security.Cryptography.SHA256]::Create() + try { + return ([BitConverter]::ToString($algorithm.ComputeHash($bytes))).Replace("-", "").ToLowerInvariant() + } finally { + $algorithm.Dispose() + } +} + +function Get-PayloadBytes([string]$Value) { + $encoding = New-Object System.Text.UTF8Encoding($false) + return $encoding.GetByteCount($Value) +} + +function Write-Utf8NoBom([string]$Path, [string]$Value) { + $encoding = New-Object System.Text.UTF8Encoding($false) + [IO.File]::WriteAllText($Path, $Value, $encoding) +} + +function Assert-ExactFile( + [string]$Path, + [string]$Sha256, + [long]$ByteLength, + [string]$Label +) { + $resolved = Resolve-DFile $Path $Label + $item = Get-Item -LiteralPath $resolved -Force + if ($item.Length -ne $ByteLength -or (Get-Sha256 $resolved) -cne $Sha256) { + throw "$Label identity changed" + } + return $resolved +} + +function Assert-ExactDigest([string]$Path, [string]$Sha256, [string]$Label) { + $resolved = Resolve-DFile $Path $Label + if ((Get-Sha256 $resolved) -cne $Sha256) { + throw "$Label identity changed" + } + return $resolved +} + +function Read-ExactJson( + [string]$Path, + [string]$Sha256, + [string]$Label +) { + $resolved = Resolve-DFile $Path $Label + if ((Get-Sha256 $resolved) -cne $Sha256) { + throw "$Label digest changed" + } + try { + return Get-Content -LiteralPath $resolved -Raw | ConvertFrom-Json + } catch { + throw "$Label is not valid JSON" + } +} + +function Get-ProtectedRuntime() { + $names = @($ProtectedRuntime | ForEach-Object { [string]$_.name }) + $rows = @(((& docker inspect $names) | ConvertFrom-Json)) + Assert-LastExitCode "protected runtime inspection" + if ($rows.Count -ne $ProtectedRuntime.Count) { + throw "protected runtime inventory is incomplete" + } + return $rows +} + +function Assert-ProtectedRuntime([object[]]$Rows) { + foreach ($expected in $ProtectedRuntime) { + $dockerName = "/$([string]$expected.name)" + $matched = @($Rows | Where-Object { [string]$_.Name -ceq $dockerName }) + if ( + $matched.Count -ne 1 -or + [string]$matched[0].Id -cne [string]$expected.container_id -or + -not $matched[0].State.Running + ) { + throw "protected runtime identity changed: $dockerName" + } + if ( + [string]$expected.name -cne "ndc-mission-core-perception-worker" -and + [string]$matched[0].State.Health.Status -cne "healthy" + ) { + throw "protected runtime health changed: $dockerName" + } + } +} + +function Assert-LegacyTaskReady() { + $task = Get-ScheduledTask -TaskName "MissionCore-M49TgsFullShadow" -ErrorAction SilentlyContinue + if ($null -eq $task -or [string]$task.State -cne "Ready") { + throw "legacy M49 scheduled task is not Ready" + } +} + +function Assert-ReadyRoot([string]$Path) { + $ready = Resolve-DDirectory $Path "portable M4.9 ready release" + $expectedNames = @( + "compiled-runner-build.json", + "executor-release.json", + "m49-tgs-portable-v2.json", + "run_m49_tgs_portable", + "worker-installation-receipt.json" + ) + $children = @(Get-ChildItem -LiteralPath $ready -Force) + $names = @($children | ForEach-Object { $_.Name } | Sort-Object) + if ( + $children.Count -ne $expectedNames.Count -or + (Compare-Object -ReferenceObject $expectedNames -DifferenceObject $names) + ) { + throw "portable M4.9 ready release contains unexpected files" + } + $null = Assert-ExactFile (Join-Path $ready "run_m49_tgs_portable") $RunnerSha256 $RunnerBytes "portable M4.9 ready runner" + $null = Assert-ExactFile (Join-Path $ready "compiled-runner-build.json") $RuntimeBuildSealSha256 $RuntimeBuildSealBytes "portable M4.9 runtime build seal" + $null = Assert-ExactFile (Join-Path $ready "m49-tgs-portable-v2.json") $ProfileSha256 $ProfileBytes "portable M4.9 ready profile" + $null = Assert-ExactFile (Join-Path $ready "executor-release.json") $InstalledReleaseSha256 $InstalledReleaseBytes "portable M4.9 installed release" + $null = Assert-ExactFile (Join-Path $ready "worker-installation-receipt.json") $ReadyReceiptSha256 $ReadyReceiptBytes "portable M4.9 ready receipt" + return $ready +} + +if ($env:COMPUTERNAME -cne $ExpectedComputer) { + throw "portable M4.9 promotion is pinned to Worker 006" +} +if ( + (Get-PayloadSha256 $RuntimeBuildSealPayload) -cne $RuntimeBuildSealSha256 -or + (Get-PayloadBytes $RuntimeBuildSealPayload) -ne $RuntimeBuildSealBytes -or + (Get-PayloadSha256 $InstalledReleasePayload) -cne $InstalledReleaseSha256 -or + (Get-PayloadBytes $InstalledReleasePayload) -ne $InstalledReleaseBytes -or + (Get-PayloadSha256 $ReadyReceiptPayload) -cne $ReadyReceiptSha256 -or + (Get-PayloadBytes $ReadyReceiptPayload) -ne $ReadyReceiptBytes +) { + throw "portable M4.9 promotion payload identity changed" +} + +$releaseParent = Resolve-DDirectory $ReleaseParent "portable Observatory release parent" +$candidateRoot = Resolve-DDirectory (Join-Path $releaseParent $CandidateRootName) "portable M4.9 candidate root" +if ((Split-Path -Parent $candidateRoot) -ine $releaseParent) { + throw "portable M4.9 candidate escaped its release parent" +} +$null = Assert-ExactDigest (Join-Path $candidateRoot "candidate.tgz") $CandidateArchiveSha256 "portable M4.9 candidate archive" +$installedCandidate = Resolve-DDirectory (Join-Path $candidateRoot "installed") "portable M4.9 installed candidate" +$candidateRelease = Read-ExactJson (Join-Path $installedCandidate "executor-release-candidate.json") $CandidateReleaseSha256 "portable M4.9 candidate release" +$candidateBuildSeal = Read-ExactJson (Join-Path $installedCandidate "compiled-runner-build-candidate.json") $CandidateBuildSealSha256 "portable M4.9 candidate build seal" +$candidateReceipt = Read-ExactJson (Join-Path $installedCandidate "worker-installation-receipt.json") $CandidateReceiptSha256 "portable M4.9 candidate receipt" +$candidateRunner = Assert-ExactFile (Join-Path $installedCandidate "run_m49_tgs_portable") $RunnerSha256 $RunnerBytes "portable M4.9 candidate runner" +$candidateProfile = Assert-ExactFile (Join-Path $installedCandidate "m49-tgs-portable-v2.json") $ProfileSha256 $ProfileBytes "portable M4.9 candidate profile" + +if ( + [string]$candidateRelease.schema_version -cne "missioncore.m49-tgs-portable-executor-installed-candidate/v1" -or + [string]$candidateRelease.state -cne "installed-candidate-blocked" -or + [string]$candidateRelease.worker_id -cne $WorkerId -or + [string]$candidateRelease.candidate_sha256 -cne $CandidateId -or + [string]$candidateRelease.source_revision -cne $SourceRevision -or + [string]$candidateRelease.source_state -cne "committed-snapshot" -or + [string]$candidateRelease.source_archive_sha256 -cne $CandidateArchiveSha256 -or + [string]$candidateRelease.base_image_sha256 -cne $BaseImageSha256 -or + [string]$candidateRelease.executor_image.sha256 -cne $ExecutorImageSha256 -or + [long]$candidateRelease.executor_image.size_bytes -ne $ExecutorImageSizeBytes -or + [string]$candidateRelease.compiled_runner.sha256 -cne $RunnerSha256 -or + [long]$candidateRelease.compiled_runner.byte_length -ne $RunnerBytes -or + [string]$candidateRelease.compiled_runner_build_seal.sha256 -cne $CandidateBuildSealSha256 -or + [string]$candidateRelease.profile.sha256 -cne $ProfileSha256 -or + [string]$candidateRelease.result_contract_sha256 -cne $ResultContractSha256 -or + [string]$candidateRelease.fixture_smoke.result -cne "passed" -or + @($candidateRelease.blockers).Count -ne 1 -or + [string]$candidateRelease.blockers[0] -cne "runtime-registry-promotion-pending" +) { + throw "portable M4.9 candidate release is not promotable" +} +if ( + [string]$candidateBuildSeal.schema_version -cne "missioncore.m49-tgs-portable-compiled-runner-build-candidate/v1" -or + [string]$candidateBuildSeal.source_revision -cne $SourceRevision -or + [string]$candidateBuildSeal.source_state -cne "committed-snapshot" -or + [string]$candidateBuildSeal.candidate_sha256 -cne $CandidateId -or + [string]$candidateBuildSeal.build_image_sha256 -cne $BaseImageSha256 -or + [string]$candidateBuildSeal.profile_sha256 -cne $ProfileSha256 -or + [string]$candidateBuildSeal.runner_source_sha256 -cne $RunnerSourceSha256 -or + [string]$candidateBuildSeal.runner_wrapper_sha256 -cne $RunnerWrapperSha256 -or + [string]$candidateBuildSeal.binary.sha256 -cne $RunnerSha256 -or + [long]$candidateBuildSeal.binary.byte_length -ne $RunnerBytes +) { + throw "portable M4.9 candidate build seal is not promotable" +} +if ( + [string]$candidateReceipt.schema_version -cne "missioncore.m49-tgs-portable-worker-installation-receipt/v1" -or + [string]$candidateReceipt.receipt_state -cne "installed-candidate-blocked" -or + [string]$candidateReceipt.worker_id -cne $WorkerId -or + [string]$candidateReceipt.computer_name -cne $ExpectedComputer -or + [string]$candidateReceipt.release_id -cne $ReleaseId -or + [string]$candidateReceipt.release_sha256 -cne $CandidateReleaseSha256 -or + [string]$candidateReceipt.executor_image_sha256 -cne $ExecutorImageSha256 -or + [string]$candidateReceipt.compiled_runner_sha256 -cne $RunnerSha256 -or + [string]$candidateReceipt.compiled_runner_build_seal_sha256 -cne $CandidateBuildSealSha256 -or + [string]$candidateReceipt.fixture_smoke -cne "passed" -or + @($candidateReceipt.blockers).Count -ne 1 -or + [string]$candidateReceipt.blockers[0] -cne "runtime-registry-promotion-pending" +) { + throw "portable M4.9 candidate receipt is not promotable" +} +if (@($candidateReceipt.protected_runtime).Count -ne $ProtectedRuntime.Count) { + throw "portable M4.9 candidate receipt protected runtime is incomplete" +} +foreach ($expected in $ProtectedRuntime) { + $matched = @( + $candidateReceipt.protected_runtime | + Where-Object { + [string]$_.name -ceq [string]$expected.name -and + [string]$_.container_id -ceq [string]$expected.container_id + } + ) + if ($matched.Count -ne 1) { + throw "portable M4.9 candidate receipt protected runtime changed" + } +} + +$baseImage = @(((& docker image inspect "sha256:$BaseImageSha256") | ConvertFrom-Json)) +Assert-LastExitCode "portable M4.9 base image inspection" +if ($baseImage.Count -ne 1 -or [string]$baseImage[0].Id -cne "sha256:$BaseImageSha256") { + throw "portable M4.9 base image identity changed" +} +$executorImage = @(((& docker image inspect $ExecutorImageTag) | ConvertFrom-Json)) +Assert-LastExitCode "portable M4.9 executor image inspection" +if ( + $executorImage.Count -ne 1 -or + [string]$executorImage[0].Id -cne "sha256:$ExecutorImageSha256" -or + [long]$executorImage[0].Size -ne $ExecutorImageSizeBytes -or + [string]$executorImage[0].Config.Labels.'com.nodedc.authority' -cne "observation-only" -or + [string]$executorImage[0].Config.Labels.'com.nodedc.product' -cne "mission-core" -or + [string]$executorImage[0].Config.Labels.'com.nodedc.stack' -cne "ndc-mission-core-observatory" -or + [string]$executorImage[0].Config.Labels.'com.nodedc.role' -cne "m49-tgs-portable-executor" -or + [string]$executorImage[0].Config.Labels.'com.nodedc.managed-by' -cne "mission-core-worker-release" -or + [string]$executorImage[0].Config.Labels.'com.nodedc.worker-contour' -cne $WorkerId -or + [string]$executorImage[0].Config.Labels.'com.nodedc.base-image.sha256' -cne $BaseImageSha256 -or + [string]$executorImage[0].Config.Labels.'com.nodedc.runner-source.sha256' -cne $RunnerSourceSha256 -or + [string]$executorImage[0].Config.Labels.'com.nodedc.fixture-smoke.sha256' -cne $FixtureSmokeSha256 -or + [string]$executorImage[0].Config.Labels.'com.nodedc.profile.sha256' -cne $ProfileSha256 +) { + throw "portable M4.9 executor image identity changed" +} + +$protectedBefore = @(Get-ProtectedRuntime) +Assert-ProtectedRuntime $protectedBefore +Assert-LegacyTaskReady + +$readyRoot = Join-Path $candidateRoot "ready" +if (Test-Path -LiteralPath $readyRoot) { + $readyRoot = Assert-ReadyRoot $readyRoot + [pscustomobject]@{ + ok = $true + already_promoted = $true + worker_id = $WorkerId + release_id = $ReleaseId + release_sha256 = $InstalledReleaseSha256 + executor_image_sha256 = $ExecutorImageSha256 + compiled_runner_sha256 = $RunnerSha256 + compiled_runner_build_seal_sha256 = $RuntimeBuildSealSha256 + worker_installation_receipt_sha256 = $ReadyReceiptSha256 + release_root = $readyRoot + runtime_release_root = $RuntimeReleaseRoot + } | ConvertTo-Json -Compress + exit 0 +} + +$imageUsers = @(& docker ps -aq --filter "ancestor=sha256:$ExecutorImageSha256") +Assert-LastExitCode "portable M4.9 executor image use inspection" +if ($imageUsers.Count -ne 0) { + throw "portable M4.9 executor image is already used by a container" +} + +$stagingRoot = Join-Path $candidateRoot ".ready-stage" +if (Test-Path -LiteralPath $stagingRoot) { + throw "portable M4.9 promotion staging requires reconciliation" +} +$null = New-Item -ItemType Directory -Path $stagingRoot +try { + $stagingRoot = Resolve-DDirectory $stagingRoot "portable M4.9 promotion staging" + Copy-Item -LiteralPath $candidateRunner -Destination (Join-Path $stagingRoot "run_m49_tgs_portable") + Copy-Item -LiteralPath $candidateProfile -Destination (Join-Path $stagingRoot "m49-tgs-portable-v2.json") + Write-Utf8NoBom (Join-Path $stagingRoot "compiled-runner-build.json") $RuntimeBuildSealPayload + Write-Utf8NoBom (Join-Path $stagingRoot "executor-release.json") $InstalledReleasePayload + Write-Utf8NoBom (Join-Path $stagingRoot "worker-installation-receipt.json") $ReadyReceiptPayload + $null = Assert-ReadyRoot $stagingRoot + + $protectedAfter = @(Get-ProtectedRuntime) + Assert-ProtectedRuntime $protectedAfter + Assert-LegacyTaskReady + if (Test-Path -LiteralPath $readyRoot) { + throw "portable M4.9 ready release collided during promotion" + } + Move-Item -LiteralPath $stagingRoot -Destination $readyRoot + $readyRoot = Assert-ReadyRoot $readyRoot +} finally { + if (Test-Path -LiteralPath $stagingRoot) { + Remove-Item -LiteralPath $stagingRoot -Recurse -Force + } +} + +[pscustomobject]@{ + ok = $true + already_promoted = $false + worker_id = $WorkerId + release_id = $ReleaseId + release_sha256 = $InstalledReleaseSha256 + executor_image_sha256 = $ExecutorImageSha256 + compiled_runner_sha256 = $RunnerSha256 + compiled_runner_build_seal_sha256 = $RuntimeBuildSealSha256 + worker_installation_receipt_sha256 = $ReadyReceiptSha256 + release_root = $readyRoot + runtime_release_root = $RuntimeReleaseRoot +} | ConvertTo-Json -Compress diff --git a/scripts/manage_mission_core_launch_agent.py b/scripts/manage_mission_core_launch_agent.py index cb47345..b9b1f0b 100644 --- a/scripts/manage_mission_core_launch_agent.py +++ b/scripts/manage_mission_core_launch_agent.py @@ -44,6 +44,14 @@ def main() -> int: ) parser.add_argument("--expected-current-sha256") parser.add_argument("--expected-desired-sha256") + parser.add_argument( + "--enable-local-observatory-worker", + action="store_true", + help=( + "enable the explicit local-only Observatory Worker transport and " + "bind its artifact roots below the preserved Mission Core data directory" + ), + ) arguments = parser.parse_args() if arguments.action == "status": print(json.dumps(_status(), sort_keys=True)) @@ -52,6 +60,7 @@ def main() -> int: repository_root=arguments.repository_root, agent_path=arguments.agent_path, expected_current_repository_root=arguments.expected_current_repository_root, + enable_local_observatory_worker=arguments.enable_local_observatory_worker, ) if arguments.action == "plan": print(json.dumps(plan.to_dict(), indent=2, sort_keys=True)) diff --git a/src/k1link/local_service_launchd.py b/src/k1link/local_service_launchd.py index f55e95a..82c44bf 100644 --- a/src/k1link/local_service_launchd.py +++ b/src/k1link/local_service_launchd.py @@ -12,6 +12,16 @@ from typing import Final, cast MISSION_CORE_LAUNCH_AGENT_LABEL: Final = "com.nodedc.mission-core.local" MISSION_CORE_LAUNCH_AGENT_SCHEMA: Final = "missioncore.local-launch-agent-plan/v1" +OBSERVATORY_LOCAL_WORKER_ENABLED_ENV: Final = ( + "MISSIONCORE_OBSERVATORY_WORKER_LOCAL_ENABLED" +) +OBSERVATORY_SOURCE_CAS_ROOT_ENV: Final = ( + "MISSIONCORE_OBSERVATORY_WORKER_SOURCE_CAS_ROOT" +) +OBSERVATORY_RESULT_STAGING_ROOT_ENV: Final = ( + "MISSIONCORE_OBSERVATORY_WORKER_RESULT_STAGING_ROOT" +) +ARTIFACT_STORE_ROOT_ENV: Final = "MISSIONCORE_ARTIFACT_STORE_ROOT" class MissionCoreLaunchAgentError(RuntimeError): @@ -28,6 +38,7 @@ class MissionCoreLaunchAgentPlan: preserved_data_directory: Path | None current_program_arguments: tuple[str, ...] desired_program_arguments: tuple[str, ...] + local_observatory_worker_enabled: bool desired_payload: bytes def to_dict(self) -> dict[str, object]: @@ -61,6 +72,9 @@ class MissionCoreLaunchAgentPlan: "bounded_launchd_exit_timeout_seconds": 20, "keep_alive": True, "process_group_owned": True, + "local_observatory_worker_enabled": ( + self.local_observatory_worker_enabled + ), }, } @@ -70,6 +84,7 @@ def plan_mission_core_launch_agent( repository_root: Path, agent_path: Path, expected_current_repository_root: Path | None = None, + enable_local_observatory_worker: bool = False, ) -> MissionCoreLaunchAgentPlan: repository = repository_root.expanduser().resolve(strict=True) expected_current_repository = ( @@ -131,6 +146,28 @@ def plan_mission_core_launch_agent( desired_environment["MISSIONCORE_SERVICE_WATCHDOG"] = "1" if preserved_data_directory is not None: desired_environment["MISSIONCORE_DATA_DIR"] = str(preserved_data_directory) + if enable_local_observatory_worker: + data_directory = _local_observatory_data_directory( + repository=repository, + environment=desired_environment, + ) + artifact_store = _private_local_worker_directory( + data_directory / "observatory-artifact-store", + "local Observatory artifact store", + ) + source_cas = _private_local_worker_directory( + data_directory / "observatory-worker-source-cas", + "local Observatory source CAS", + ) + result_staging = _private_local_worker_directory( + data_directory / "observatory-worker-result-staging", + "local Observatory result staging", + ) + desired_environment["MISSIONCORE_DATA_DIR"] = str(data_directory) + desired_environment[ARTIFACT_STORE_ROOT_ENV] = str(artifact_store) + desired_environment[OBSERVATORY_SOURCE_CAS_ROOT_ENV] = str(source_cas) + desired_environment[OBSERVATORY_RESULT_STAGING_ROOT_ENV] = str(result_staging) + desired_environment[OBSERVATORY_LOCAL_WORKER_ENABLED_ENV] = "1" log_path = repository / ".runtime/mission-core/k1link-serve-launchd.log" desired_program_arguments = ( str(uv_entrypoint), @@ -163,10 +200,44 @@ def plan_mission_core_launch_agent( preserved_data_directory=preserved_data_directory, current_program_arguments=current_arguments, desired_program_arguments=desired_program_arguments, + local_observatory_worker_enabled=enable_local_observatory_worker, desired_payload=desired_payload, ) +def _local_observatory_data_directory( + *, + repository: Path, + environment: dict[str, str], +) -> Path: + configured = environment.get("MISSIONCORE_DATA_DIR", "").strip() + candidate = Path(configured) if configured else repository / ".runtime" / "mission-core" + return _private_local_worker_directory( + candidate, + "Mission Core data directory", + ) + + +def _private_local_worker_directory(path: Path, label: str) -> Path: + if not path.is_absolute(): + raise MissionCoreLaunchAgentError(f"{label} is not an absolute directory") + try: + metadata = path.lstat() + resolved = path.resolve(strict=True) + except OSError as exc: + raise MissionCoreLaunchAgentError(f"{label} is unavailable") from exc + if ( + resolved != path + or not stat.S_ISDIR(metadata.st_mode) + or stat.S_IMODE(metadata.st_mode) != 0o700 + or metadata.st_uid != os.getuid() + ): + raise MissionCoreLaunchAgentError( + f"{label} is not a private canonical directory" + ) + return resolved + + def _program_arguments(document: dict[str, object]) -> tuple[str, ...]: value = document.get("ProgramArguments") if not isinstance(value, list) or not value or any(not isinstance(item, str) for item in value): diff --git a/src/k1link/observatory/m49_worker_main.py b/src/k1link/observatory/m49_worker_main.py new file mode 100644 index 0000000..3a0f66e --- /dev/null +++ b/src/k1link/observatory/m49_worker_main.py @@ -0,0 +1,6 @@ +"""Runnable fixed entrypoint for the portable M4.9 Observatory Worker.""" + +from k1link.observatory.m49_worker_service import main + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/src/k1link/observatory/m49_worker_service.py b/src/k1link/observatory/m49_worker_service.py new file mode 100644 index 0000000..cd01cb5 --- /dev/null +++ b/src/k1link/observatory/m49_worker_service.py @@ -0,0 +1,716 @@ +"""Fixed POSIX service composition for the sealed portable M4.9 executor. + +The entrypoint reads only immutable registries, one installation receipt and +path-only Worker service settings. Jobs cannot select commands, providers, +modules, images or filesystem locations. Every runtime asset is resolved from +the fixed release layout and re-verified against the ready runtime candidate +before the first queue claim. +""" + +from __future__ import annotations + +import argparse +import hashlib +import json +import os +import re +import signal +import stat +from collections.abc import Iterator, Mapping, Sequence +from contextlib import contextmanager +from dataclasses import dataclass +from datetime import UTC, datetime +from pathlib import Path, PurePosixPath +from threading import Event +from types import FrameType +from typing import Final, cast + +import httpx + +from k1link.observatory.m49_portable_executor import ( + M49_PORTABLE_COMPILED_RUNNER_ASSET_ID, + M49_PORTABLE_COMPILED_RUNNER_BUILD_SEAL_ASSET_ID, + M49_PORTABLE_PROFILE_ASSET_ID, + M49_PORTABLE_TRAVEL_IMAGE_ASSET_ID, + M49PortableRunnerInstallation, + compose_m49_portable_executor_adapter, +) +from k1link.observatory.portable_result_contract import OBSERVATION_ONLY_AUTHORITY +from k1link.observatory.portable_run_definitions import ( + PortableRunDefinition, + PortableRunDefinitionRegistry, +) +from k1link.observatory.portable_worker_runtime import ( + PortableWorkerLocalAssetBinding, + PortableWorkerRuntimeCandidate, + PortableWorkerRuntimeRegistry, + inspect_runtime_candidate, +) +from k1link.observatory.worker_agent import ( + WORKER_006_CONTOUR_ID, + ObservatoryWorkerAgent, + ObservatoryWorkerExecutorRegistration, + ObservatoryWorkerExecutorRegistry, +) +from k1link.observatory.worker_http_transport import ObservatoryWorkerHttpGateway +from k1link.observatory.worker_service import ( + InstalledObservatoryWorkerService, + ObservatoryWorkerServiceConfiguration, + load_observatory_worker_bearer_token, + require_ready_executor_coverage, +) + +M49_WORKER_DEFINITIONS_FILE_ENV: Final = "MISSIONCORE_OBSERVATORY_WORKER_DEFINITIONS_FILE" +M49_WORKER_RUNTIME_REGISTRY_FILE_ENV: Final = "MISSIONCORE_OBSERVATORY_WORKER_RUNTIME_REGISTRY_FILE" +M49_WORKER_INSTALLATION_RECEIPT_FILE_ENV: Final = ( + "MISSIONCORE_OBSERVATORY_M49_INSTALLATION_RECEIPT_FILE" +) + +M49_WORKER_INSTALLATION_RECEIPT_SCHEMA: Final = ( + "missioncore.m49-tgs-portable-worker-installation-ready-receipt/v1" +) +M49_WORKER_SETUP_ID: Final = "m49-tgs-portable-v2" +M49_WORKER_ADAPTER_ID: Final = "m49-tgs-worker006-portable-v2" +M49_WORKER_RELEASE_ID: Final = "m49-tgs-portable-executor-v1" + +M49_EXECUTOR_RELEASE_ASSET_ID: Final = "m49-portable-executor-release" +M49_WORKER_INSTALLATION_RECEIPT_ASSET_ID: Final = "m49-portable-worker-installation-receipt" + +_MAX_RECEIPT_BYTES: Final = 256 * 1024 +_SHA256: Final = re.compile(r"^[a-f0-9]{64}$") +_SOURCE_REVISION: Final = re.compile(r"^[a-f0-9]{40}$") +_WORKER_COMPUTER_NAME: Final = "DESKTOP-OPJ8J04" +_FIXED_RUNTIME_RELEASE_ROOT: Final = Path("/release") +_PROTECTED_RUNTIME_NAMES: Final = ( + "ndc-mission-core-triton", + "ndc-mission-core-perception-worker", + "ndc-gaussian-pipeline-gaussian-gateway-1", + "ndc-gaussian-pipeline-gaussian-pipeline-1", + "ndc-gaussian-pipeline-gaussian-terrain-executor-1", +) +_AUDIT_RELEASE_ROOT: Final = re.compile( + r"^D:\\NDC_MISSIONCORE\\runtime\\releases\\observatory-portable\\" + r"m49-tgs-portable-candidate-([a-f0-9]{64})\\ready$" +) + +_FIXED_FILE_ASSET_PATHS: Final = { + M49_PORTABLE_COMPILED_RUNNER_ASSET_ID: PurePosixPath("run_m49_tgs_portable"), + M49_PORTABLE_COMPILED_RUNNER_BUILD_SEAL_ASSET_ID: PurePosixPath("compiled-runner-build.json"), + M49_PORTABLE_PROFILE_ASSET_ID: PurePosixPath("m49-tgs-portable-v2.json"), + M49_EXECUTOR_RELEASE_ASSET_ID: PurePosixPath("executor-release.json"), +} +_REQUIRED_RUNTIME_ASSET_IDS: Final = frozenset( + { + *_FIXED_FILE_ASSET_PATHS, + M49_WORKER_INSTALLATION_RECEIPT_ASSET_ID, + M49_PORTABLE_TRAVEL_IMAGE_ASSET_ID, + } +) + + +class M49WorkerCompositionError(RuntimeError): + """The fixed M4.9 Worker installation cannot be composed safely.""" + + +@dataclass(frozen=True, slots=True) +class M49WorkerEntrypointConfiguration: + """Path-only service inputs selected before the Worker process starts.""" + + worker: ObservatoryWorkerServiceConfiguration + definitions_file: Path + runtime_registry_file: Path + installation_receipt_file: Path + + def __post_init__(self) -> None: + for path, label in ( + (self.definitions_file, "portable RunDefinition registry"), + (self.runtime_registry_file, "portable runtime registry"), + (self.installation_receipt_file, "M4.9 installation receipt"), + ): + _absolute_path(path, label) + + @classmethod + def from_environment( + cls, + environment: Mapping[str, str] | None = None, + ) -> M49WorkerEntrypointConfiguration: + values = os.environ if environment is None else environment + return cls( + worker=ObservatoryWorkerServiceConfiguration.from_environment(values), + definitions_file=_required_environment_path( + values, + M49_WORKER_DEFINITIONS_FILE_ENV, + ), + runtime_registry_file=_required_environment_path( + values, + M49_WORKER_RUNTIME_REGISTRY_FILE_ENV, + ), + installation_receipt_file=_required_environment_path( + values, + M49_WORKER_INSTALLATION_RECEIPT_FILE_ENV, + ), + ) + + +@dataclass(frozen=True, slots=True) +class M49WorkerReceiptFile: + asset_id: str + relative_path: PurePosixPath + byte_length: int + sha256: str + + def __post_init__(self) -> None: + expected = _FIXED_FILE_ASSET_PATHS.get(self.asset_id) + if expected is None or self.relative_path != expected: + raise M49WorkerCompositionError( + "M4.9 installation receipt contains an unknown release file" + ) + if isinstance(self.byte_length, bool) or not 1 <= self.byte_length <= 2**40: + raise M49WorkerCompositionError("M4.9 installation receipt file length is invalid") + _digest(self.sha256, "M4.9 installation receipt file SHA-256") + + +@dataclass(frozen=True, slots=True) +class M49WorkerInstallationReceipt: + """Strict installed-ready receipt; paths remain confined to one release root.""" + + path: Path + release_root: Path + source_revision: str + candidate_sha256: str + candidate_release_sha256: str + candidate_worker_installation_receipt_sha256: str + release_id: str + release_sha256: str + base_image_sha256: str + executor_image_sha256: str + files: tuple[M49WorkerReceiptFile, ...] + + def __post_init__(self) -> None: + if self.release_id != M49_WORKER_RELEASE_ID: + raise M49WorkerCompositionError("M4.9 installation receipt release changed") + if _SOURCE_REVISION.fullmatch(self.source_revision) is None: + raise M49WorkerCompositionError("M4.9 installation receipt source revision is invalid") + for value, label in ( + (self.candidate_sha256, "M4.9 source candidate SHA-256"), + (self.candidate_release_sha256, "M4.9 candidate release SHA-256"), + ( + self.candidate_worker_installation_receipt_sha256, + "M4.9 candidate installation receipt SHA-256", + ), + (self.release_sha256, "M4.9 executor release SHA-256"), + (self.base_image_sha256, "M4.9 base image SHA-256"), + (self.executor_image_sha256, "M4.9 executor image SHA-256"), + ): + _digest(value, label) + asset_ids = tuple(row.asset_id for row in self.files) + if asset_ids != tuple(sorted(_FIXED_FILE_ASSET_PATHS)): + raise M49WorkerCompositionError( + "M4.9 installation receipt file inventory is incomplete" + ) + release_file = next( + row for row in self.files if row.asset_id == M49_EXECUTOR_RELEASE_ASSET_ID + ) + if release_file.sha256 != self.release_sha256: + raise M49WorkerCompositionError( + "M4.9 executor release digest differs from its installation receipt" + ) + if self.path.parent != self.release_root: + raise M49WorkerCompositionError( + "M4.9 installation receipt is outside its fixed release root" + ) + + def asset_bindings( + self, + candidate: PortableWorkerRuntimeCandidate, + ) -> dict[str, PortableWorkerLocalAssetBinding]: + """Resolve every candidate asset or reject the unknown inventory.""" + + candidate_asset_ids = {asset.asset_id for asset in candidate.reusable_assets} + if candidate_asset_ids != _REQUIRED_RUNTIME_ASSET_IDS: + raise M49WorkerCompositionError( + "ready M4.9 runtime asset inventory differs from the fixed release" + ) + rows = {row.asset_id: row for row in self.files} + bindings: dict[str, PortableWorkerLocalAssetBinding] = {} + for requirement in candidate.reusable_assets: + if requirement.asset_id == M49_PORTABLE_TRAVEL_IMAGE_ASSET_ID: + bindings[requirement.asset_id] = PortableWorkerLocalAssetBinding( + asset_id=requirement.asset_id, + image_sha256=self.base_image_sha256, + ) + continue + if requirement.asset_id == M49_WORKER_INSTALLATION_RECEIPT_ASSET_ID: + bindings[requirement.asset_id] = PortableWorkerLocalAssetBinding( + asset_id=requirement.asset_id, + file_path=self.path, + ) + continue + row = rows.get(requirement.asset_id) + if row is None: + raise M49WorkerCompositionError( + "ready M4.9 runtime requires an unmapped local asset" + ) + path = _confined_release_file(self.release_root, row.relative_path) + metadata = path.stat() + if metadata.st_size != row.byte_length or _sha256_file(path) != row.sha256: + raise M49WorkerCompositionError( + "M4.9 release file differs from its installation receipt" + ) + bindings[requirement.asset_id] = PortableWorkerLocalAssetBinding( + asset_id=requirement.asset_id, + file_path=path, + ) + if set(bindings) != {asset.asset_id for asset in candidate.reusable_assets}: + raise M49WorkerCompositionError( + "M4.9 runtime asset bindings do not cover the exact candidate" + ) + return bindings + + +def load_m49_worker_installation_receipt( + path: Path, +) -> M49WorkerInstallationReceipt: + """Load one exact regular installed-ready receipt without following links.""" + + receipt_path = _regular_file(path, "M4.9 installation receipt") + try: + payload = receipt_path.read_bytes() + document: object = json.loads(payload.decode("utf-8")) + except (OSError, UnicodeDecodeError, json.JSONDecodeError) as exc: + raise M49WorkerCompositionError("M4.9 installation receipt is unreadable") from exc + if not 0 < len(payload) <= _MAX_RECEIPT_BYTES: + raise M49WorkerCompositionError("M4.9 installation receipt size is invalid") + row = _object(document, "M4.9 installation receipt") + _exact_keys( + row, + { + "schema_version", + "receipt_state", + "worker_id", + "computer_name", + "source_revision", + "candidate_sha256", + "candidate_release_sha256", + "candidate_worker_installation_receipt_sha256", + "release_id", + "release_sha256", + "release_root", + "runtime_release_root", + "base_image_sha256", + "executor_image_sha256", + "files", + "fixture_smoke", + "protected_runtime", + "legacy_m49_task_state", + "blockers", + "authority", + }, + "M4.9 installation receipt", + ) + blockers = _array(row["blockers"], "M4.9 installation blockers") + if ( + row["schema_version"] != M49_WORKER_INSTALLATION_RECEIPT_SCHEMA + or row["receipt_state"] != "installed-ready" + or row["worker_id"] != WORKER_006_CONTOUR_ID + or row["computer_name"] != _WORKER_COMPUTER_NAME + or row["fixture_smoke"] != "passed" + or blockers + or row["authority"] != OBSERVATION_ONLY_AUTHORITY + ): + raise M49WorkerCompositionError( + "M4.9 installation receipt is not an accepted ready installation" + ) + candidate_sha256 = _string( + row["candidate_sha256"], + "M4.9 source candidate SHA-256", + ) + audit_release_root = _string(row["release_root"], "M4.9 audit release root") + audit_match = _AUDIT_RELEASE_ROOT.fullmatch(audit_release_root) + if audit_match is None or audit_match.group(1) != candidate_sha256: + raise M49WorkerCompositionError("M4.9 audit release root identity changed") + release_root = _real_directory( + Path(_string(row["runtime_release_root"], "M4.9 runtime release root")), + "M4.9 runtime release root", + ) + if release_root != _FIXED_RUNTIME_RELEASE_ROOT: + raise M49WorkerCompositionError("M4.9 runtime release root is not fixed") + file_rows = tuple( + sorted( + (_receipt_file(value) for value in _array(row["files"], "M4.9 release files")), + key=lambda value: value.asset_id, + ) + ) + receipt = M49WorkerInstallationReceipt( + path=receipt_path, + release_root=release_root, + source_revision=_string(row["source_revision"], "M4.9 source revision"), + candidate_sha256=candidate_sha256, + candidate_release_sha256=_string( + row["candidate_release_sha256"], + "M4.9 candidate release SHA-256", + ), + candidate_worker_installation_receipt_sha256=_string( + row["candidate_worker_installation_receipt_sha256"], + "M4.9 candidate installation receipt SHA-256", + ), + release_id=_string(row["release_id"], "M4.9 release id"), + release_sha256=_string(row["release_sha256"], "M4.9 release SHA-256"), + base_image_sha256=_string( + row["base_image_sha256"], + "M4.9 base image SHA-256", + ), + executor_image_sha256=_string( + row["executor_image_sha256"], + "M4.9 executor image SHA-256", + ), + files=file_rows, + ) + protected = _array(row["protected_runtime"], "M4.9 protected runtime") + protected_names: list[str] = [] + for value in protected: + protected_row = _object(value, "M4.9 protected runtime row") + _exact_keys( + protected_row, + {"name", "container_id"}, + "M4.9 protected runtime row", + ) + protected_names.append( + _nonempty_string( + protected_row["name"], + "M4.9 protected runtime name", + ) + ) + container_id = _nonempty_string( + protected_row["container_id"], + "M4.9 protected runtime container id", + ) + _digest(container_id, "M4.9 protected runtime container id") + if tuple(protected_names) != _PROTECTED_RUNTIME_NAMES: + raise M49WorkerCompositionError("M4.9 protected runtime inventory changed") + legacy_state = row["legacy_m49_task_state"] + if legacy_state is not None: + _nonempty_string(legacy_state, "M4.9 legacy task state") + return receipt + + +def compose_installed_m49_worker_service( + configuration: M49WorkerEntrypointConfiguration, + *, + http_transport: httpx.BaseTransport | None = None, +) -> InstalledObservatoryWorkerService: + """Compose one fixed M4.9 executor around one shared HTTP gateway.""" + + if os.name != "posix": + raise M49WorkerCompositionError("portable M4.9 Worker requires a POSIX runtime") + definitions = PortableRunDefinitionRegistry.from_file(configuration.definitions_file) + definition = definitions.resolve_setup(M49_WORKER_SETUP_ID) + runtime_registry = PortableWorkerRuntimeRegistry.from_file( + configuration.runtime_registry_file, + definitions=definitions, + ) + candidate = runtime_registry.resolve( + definition.setup_id, + definition.definition_sha256, + ) + _verify_ready_identity(definitions, definition, candidate) + receipt = load_m49_worker_installation_receipt(configuration.installation_receipt_file) + _verify_receipt_identity(receipt, definition, candidate) + bindings = receipt.asset_bindings(candidate) + admission = inspect_runtime_candidate(candidate, bindings) + if not admission.ready: + raise M49WorkerCompositionError("portable M4.9 local asset admission is not ready") + profile = _bound_file(bindings, M49_PORTABLE_PROFILE_ASSET_ID) + runner = _bound_file(bindings, M49_PORTABLE_COMPILED_RUNNER_ASSET_ID) + build_seal = _bound_file( + bindings, + M49_PORTABLE_COMPILED_RUNNER_BUILD_SEAL_ASSET_ID, + ) + installation = M49PortableRunnerInstallation( + profile_path=profile, + runner_binary_path=runner, + runner_build_seal_path=build_seal, + runner_build_seal_sha256=_asset_sha256( + candidate, + M49_PORTABLE_COMPILED_RUNNER_BUILD_SEAL_ASSET_ID, + ), + output_parent=configuration.worker.work_root / "m49-portable" / "runner-output", + ) + token = load_observatory_worker_bearer_token(configuration.worker.bearer_token_file) + gateway: ObservatoryWorkerHttpGateway | None = None + try: + gateway = ObservatoryWorkerHttpGateway( + base_url=configuration.worker.base_url, + bearer_token=token, + work_root=configuration.worker.work_root, + transport=http_transport, + ) + adapter = compose_m49_portable_executor_adapter( + candidate=candidate, + definition=definition, + admission=admission, + source_transport=gateway, + result_transport=gateway, + installation=installation, + source_output_parent=( + configuration.worker.work_root / "m49-portable" / "source-output" + ), + created_at_utc=_utc_now, + ) + executors = ObservatoryWorkerExecutorRegistry( + ( + ObservatoryWorkerExecutorRegistration( + identity=candidate.executor_identity(), + adapter=adapter, + ), + ) + ) + require_ready_executor_coverage( + definitions=definitions, + executors=executors, + ) + return InstalledObservatoryWorkerService( + configuration=configuration.worker, + gateway=gateway, + agent=ObservatoryWorkerAgent(transport=gateway, executors=executors), + ) + except Exception: + if gateway is not None: + gateway.close() + raise + finally: + del token + + +def run_installed_m49_worker( + service: InstalledObservatoryWorkerService, + *, + stop: Event, + once: bool = False, +) -> None: + """Run one claim or the bounded-failure polling loop until POSIX shutdown.""" + + if once: + try: + service.agent.run_once() + finally: + service.close() + return + service.run(stop=stop) + + +def main(arguments: Sequence[str] | None = None) -> int: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument( + "--once", + action="store_true", + help="Run at most one claim cycle and exit.", + ) + options = parser.parse_args(arguments) + configuration = M49WorkerEntrypointConfiguration.from_environment() + service = compose_installed_m49_worker_service(configuration) + stop = Event() + with _posix_shutdown_signals(stop): + run_installed_m49_worker(service, stop=stop, once=cast(bool, options.once)) + return 0 + + +def _verify_ready_identity( + definitions: PortableRunDefinitionRegistry, + definition: PortableRunDefinition, + candidate: PortableWorkerRuntimeCandidate, +) -> None: + ready = definitions.ready_recorded_definitions() + if ( + len(ready) != 1 + or ready[0].setup_id != M49_WORKER_SETUP_ID + or definition.setup_id != M49_WORKER_SETUP_ID + or definition.executor.contour_id != WORKER_006_CONTOUR_ID + or candidate.adapter_id != M49_WORKER_ADAPTER_ID + or not definition.executor.ready + or not candidate.ready + ): + raise M49WorkerCompositionError( + "fixed M4.9 Worker requires exactly one ready M4.9 definition" + ) + + +def _verify_receipt_identity( + receipt: M49WorkerInstallationReceipt, + definition: PortableRunDefinition, + candidate: PortableWorkerRuntimeCandidate, +) -> None: + executor = candidate.executor + if ( + executor is None + or definition.executor.release_id != receipt.release_id + or definition.executor.release_sha256 != receipt.release_sha256 + or definition.executor.image_sha256 != receipt.executor_image_sha256 + or executor.release_id != receipt.release_id + or executor.release_sha256 != receipt.release_sha256 + or executor.image_sha256 != receipt.executor_image_sha256 + ): + raise M49WorkerCompositionError( + "M4.9 installation receipt differs from the sealed executor" + ) + + +def _receipt_file(value: object) -> M49WorkerReceiptFile: + row = _object(value, "M4.9 release file") + _exact_keys( + row, + {"asset_id", "relative_path", "byte_length", "sha256"}, + "M4.9 release file", + ) + relative = _safe_relative_path(_string(row["relative_path"], "M4.9 release relative path")) + byte_length = row["byte_length"] + if not isinstance(byte_length, int): + raise M49WorkerCompositionError("M4.9 release file length is invalid") + return M49WorkerReceiptFile( + asset_id=_string(row["asset_id"], "M4.9 release asset id"), + relative_path=relative, + byte_length=byte_length, + sha256=_string(row["sha256"], "M4.9 release file SHA-256"), + ) + + +def _bound_file( + bindings: Mapping[str, PortableWorkerLocalAssetBinding], + asset_id: str, +) -> Path: + binding = bindings.get(asset_id) + if binding is None or binding.file_path is None: + raise M49WorkerCompositionError("M4.9 fixed file asset is unavailable") + return binding.file_path + + +def _asset_sha256(candidate: PortableWorkerRuntimeCandidate, asset_id: str) -> str: + for asset in candidate.reusable_assets: + if asset.asset_id == asset_id: + return asset.sha256 + raise M49WorkerCompositionError("M4.9 fixed asset is absent from the runtime") + + +def _confined_release_file(root: Path, relative: PurePosixPath) -> Path: + candidate = _regular_file(root.joinpath(*relative.parts), "M4.9 release file") + if not candidate.is_relative_to(root): + raise M49WorkerCompositionError("M4.9 release file escapes its release root") + return candidate + + +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 _regular_file(path: Path, label: str) -> Path: + candidate = path.expanduser().absolute() + try: + metadata = candidate.lstat() + resolved = candidate.resolve(strict=True) + except OSError as exc: + raise M49WorkerCompositionError(f"{label} is unavailable") from exc + if ( + stat.S_ISLNK(metadata.st_mode) + or not stat.S_ISREG(metadata.st_mode) + or not os.path.samefile(candidate, resolved) + ): + raise M49WorkerCompositionError(f"{label} is unsafe") + return resolved + + +def _real_directory(path: Path, label: str) -> Path: + candidate = path.expanduser().absolute() + try: + metadata = candidate.lstat() + resolved = candidate.resolve(strict=True) + except OSError as exc: + raise M49WorkerCompositionError(f"{label} is unavailable") from exc + if ( + stat.S_ISLNK(metadata.st_mode) + or not stat.S_ISDIR(metadata.st_mode) + or not os.path.samefile(candidate, resolved) + ): + raise M49WorkerCompositionError(f"{label} is unsafe") + return resolved + + +def _safe_relative_path(value: str) -> PurePosixPath: + path = PurePosixPath(value) + if path.is_absolute() or not path.parts or any(part in {"", ".", ".."} for part in path.parts): + raise M49WorkerCompositionError("M4.9 release relative path is unsafe") + return path + + +def _absolute_path(path: Path, label: str) -> None: + if not path.is_absolute() or str(path) != str(path).strip(): + raise M49WorkerCompositionError(f"{label} must be an absolute path") + + +def _required_environment_path(values: Mapping[str, str], name: str) -> Path: + value = values.get(name, "") + if not value or value != value.strip(): + raise M49WorkerCompositionError(f"{name} is required") + path = Path(value) + _absolute_path(path, name) + return path + + +def _digest(value: str, label: str) -> None: + if _SHA256.fullmatch(value) is None: + raise M49WorkerCompositionError(f"{label} is invalid") + + +def _object(value: object, label: str) -> dict[str, object]: + if not isinstance(value, dict) or any(not isinstance(key, str) for key in value): + raise M49WorkerCompositionError(f"{label} must be an object") + return cast(dict[str, object], value) + + +def _array(value: object, label: str) -> list[object]: + if not isinstance(value, list): + raise M49WorkerCompositionError(f"{label} must be an array") + return value + + +def _exact_keys(row: Mapping[str, object], expected: set[str], label: str) -> None: + if set(row) != expected: + raise M49WorkerCompositionError(f"{label} fields are invalid") + + +def _string(value: object, label: str) -> str: + if not isinstance(value, str): + raise M49WorkerCompositionError(f"{label} must be a string") + return value + + +def _nonempty_string(value: object, label: str) -> str: + result = _string(value, label) + if not result or result != result.strip() or len(result) > 512: + raise M49WorkerCompositionError(f"{label} is invalid") + return result + + +def _utc_now() -> str: + return datetime.now(UTC).isoformat(timespec="milliseconds").replace("+00:00", "Z") + + +@contextmanager +def _posix_shutdown_signals(stop: Event) -> Iterator[None]: + def request_stop(_signal: int, _frame: FrameType | None) -> None: + stop.set() + + previous_int = signal.signal(signal.SIGINT, request_stop) + previous_term = signal.signal(signal.SIGTERM, request_stop) + try: + yield + finally: + signal.signal(signal.SIGINT, previous_int) + signal.signal(signal.SIGTERM, previous_term) + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/src/k1link/observatory/portable_worker_integration.py b/src/k1link/observatory/portable_worker_integration.py index f8c4638..7cf26a4 100644 --- a/src/k1link/observatory/portable_worker_integration.py +++ b/src/k1link/observatory/portable_worker_integration.py @@ -63,6 +63,9 @@ OBSERVATORY_WORKER_SOURCE_CAS_ROOT_ENV: Final = "MISSIONCORE_OBSERVATORY_WORKER_ OBSERVATORY_WORKER_RESULT_STAGING_ROOT_ENV: Final = ( "MISSIONCORE_OBSERVATORY_WORKER_RESULT_STAGING_ROOT" ) +OBSERVATORY_WORKER_LOCAL_ENABLED_ENV: Final = ( + "MISSIONCORE_OBSERVATORY_WORKER_LOCAL_ENABLED" +) class PortableWorkerIntegrationError(RuntimeError): @@ -197,7 +200,7 @@ _VALIDATOR_SPECS: Final = ( _ValidatorSpec( setup_id=PORTABLE_M49_SETUP_ID, definition_id="m49-tgs-portable", - definition_version=2, + definition_version=3, contract_id="m49-tgs-portable-review-v2", contract_version=2, result_schema=M49_PORTABLE_RESULT_SCHEMA, @@ -231,6 +234,27 @@ def portable_result_validator_registry( return PortableResultContractValidatorRegistry(tuple(registrations)) +def observatory_worker_local_enabled( + environment: Mapping[str, str] | None = None, +) -> bool: + """Return the explicit local-only Worker API gate. + + Absence is the fail-closed default. The single accepted enabled value is + deliberately exact so misspelled or whitespace-padded service settings + cannot expose the pull API. + """ + + values = os.environ if environment is None else environment + raw = values.get(OBSERVATORY_WORKER_LOCAL_ENABLED_ENV) + if raw is None or raw == "": + return False + if raw != "1": + raise PortableWorkerIntegrationError( + f"{OBSERVATORY_WORKER_LOCAL_ENABLED_ENV} must be exactly 1 when enabled" + ) + return True + + def build_portable_observatory_worker_integration( *, queue: ObservatoryRecordedJobQueue, diff --git a/src/k1link/web/app.py b/src/k1link/web/app.py index f3a17ba..3da4ec8 100644 --- a/src/k1link/web/app.py +++ b/src/k1link/web/app.py @@ -68,10 +68,12 @@ from k1link.observatory.portable_setup_projection import ( portable_calculation_profile_registry, ) from k1link.observatory.portable_worker_integration import ( + OBSERVATORY_WORKER_LOCAL_ENABLED_ENV, PortableObservatoryWorkerIntegration, PortableWorkerIntegrationError, PortableWorkerStorageRoots, build_portable_observatory_worker_integration, + observatory_worker_local_enabled, portable_result_validator_registry, ) from k1link.observatory.recorded_jobs import ( @@ -385,9 +387,15 @@ except (M49QueueBindingError, ObservatoryRecordedQueueError, OSError, ValueError OBSERVATORY_RECORDED_JOB_QUEUE = None OBSERVATORY_RECORDED_JOB_QUEUE_ERROR = str(exc) OBSERVATORY_WORKER_TOKEN_PATH = session_store.data_dir / "worker-auth" / "observatory-worker.token" -OBSERVATORY_WORKER_CLAIM_LEASE_READY = False -OBSERVATORY_WORKER_VERIFIED_RESULT_PUBLISHER_READY = False OBSERVATORY_WORKER_PRODUCTION_API_ENABLED = False +OBSERVATORY_WORKER_LOCAL_ENABLED: bool +OBSERVATORY_WORKER_LOCAL_ENABLED_ERROR: str | None +try: + OBSERVATORY_WORKER_LOCAL_ENABLED = observatory_worker_local_enabled() + OBSERVATORY_WORKER_LOCAL_ENABLED_ERROR = None +except PortableWorkerIntegrationError as exc: + OBSERVATORY_WORKER_LOCAL_ENABLED = False + OBSERVATORY_WORKER_LOCAL_ENABLED_ERROR = str(exc) OBSERVATORY_WORKER_AUTHENTICATION: ObservatoryWorkerAuthentication | None OBSERVATORY_WORKER_AUTHENTICATION_ERROR: str | None OBSERVATORY_WORKER_API_ERROR: str | None @@ -470,32 +478,47 @@ except (PortableWorkerIntegrationError, OSError, ValueError) as exc: # Failure remains isolated from K1, Simulation and legacy LAB. OBSERVATORY_PORTABLE_WORKER_INTEGRATION = None OBSERVATORY_PORTABLE_WORKER_INTEGRATION_ERROR = str(exc) -OBSERVATORY_WORKER_API_ERROR = ( - "Worker pull API is hard-disabled pending sealed installed executors, " - "a configured Worker credential, explicit integration acceptance and the " - "production gate" - + ( - "" - if OBSERVATORY_WORKER_AUTHENTICATION_ERROR is None - else ( - "; authentication unavailable: " - f"{OBSERVATORY_WORKER_AUTHENTICATION_ERROR}" - ) - ) - + ( - "" - if OBSERVATORY_PORTABLE_WORKER_INTEGRATION_ERROR is None - else f"; integration unavailable: {OBSERVATORY_PORTABLE_WORKER_INTEGRATION_ERROR}" - ) +OBSERVATORY_WORKER_API_GATE_ENABLED = OBSERVATORY_WORKER_LOCAL_ENABLED +OBSERVATORY_WORKER_CLAIM_LEASE_READY = ( + OBSERVATORY_WORKER_API_GATE_ENABLED + and OBSERVATORY_RECORDED_JOB_QUEUE is not None +) +OBSERVATORY_WORKER_VERIFIED_RESULT_PUBLISHER_READY = ( + OBSERVATORY_WORKER_API_GATE_ENABLED + and OBSERVATORY_PORTABLE_WORKER_INTEGRATION is not None ) OBSERVATORY_WORKER_DISPATCH_READY = ( - OBSERVATORY_WORKER_PRODUCTION_API_ENABLED - and OBSERVATORY_WORKER_CLAIM_LEASE_READY + OBSERVATORY_WORKER_CLAIM_LEASE_READY and OBSERVATORY_WORKER_VERIFIED_RESULT_PUBLISHER_READY and OBSERVATORY_RECORDED_JOB_QUEUE is not None and OBSERVATORY_WORKER_AUTHENTICATION is not None and OBSERVATORY_PORTABLE_WORKER_INTEGRATION is not None ) +if OBSERVATORY_WORKER_DISPATCH_READY: + OBSERVATORY_WORKER_API_ERROR = None +else: + worker_api_errors: list[str] = [] + if not OBSERVATORY_WORKER_API_GATE_ENABLED: + worker_api_errors.append( + OBSERVATORY_WORKER_LOCAL_ENABLED_ERROR + or ( + "local-only Worker API gate is disabled; set " + f"{OBSERVATORY_WORKER_LOCAL_ENABLED_ENV}=1 to enable it" + ) + ) + if OBSERVATORY_WORKER_AUTHENTICATION_ERROR is not None: + worker_api_errors.append( + "authentication unavailable: " + f"{OBSERVATORY_WORKER_AUTHENTICATION_ERROR}" + ) + if OBSERVATORY_PORTABLE_WORKER_INTEGRATION_ERROR is not None: + worker_api_errors.append( + "integration unavailable: " + f"{OBSERVATORY_PORTABLE_WORKER_INTEGRATION_ERROR}" + ) + OBSERVATORY_WORKER_API_ERROR = "Worker pull API is disabled; " + "; ".join( + worker_api_errors + ) OBSERVATORY_PORTABLE_BINDING_SERVICE: PortableRecordedQueueBindingService | None OBSERVATORY_PORTABLE_SETUP_PROJECTOR: PortableSetupProjector | None OBSERVATORY_PORTABLE_SETUP_PROJECTOR_ERROR: str | None diff --git a/tests/test_local_service_launchd.py b/tests/test_local_service_launchd.py index aa0ac84..fb01fe0 100644 --- a/tests/test_local_service_launchd.py +++ b/tests/test_local_service_launchd.py @@ -65,6 +65,72 @@ def test_launch_agent_plan_disables_sync_and_enables_watchdog(tmp_path: Path) -> assert "MISSIONCORE_DATA_DIR" not in desired["EnvironmentVariables"] assert plan.to_dict()["changes"]["repository_migration"] is False assert plan.to_dict()["changes"]["data_directory_preserved"] is False + assert plan.to_dict()["changes"]["local_observatory_worker_enabled"] is False + + +def test_launch_agent_plan_explicitly_enables_local_observatory_worker( + tmp_path: Path, +) -> None: + repository = tmp_path / "repo" + repository.mkdir() + data_directory = repository / ".runtime" / "mission-core" + data_directory.mkdir(parents=True, mode=0o700) + for name in ( + "observatory-artifact-store", + "observatory-worker-source-cas", + "observatory-worker-result-staging", + ): + (data_directory / name).mkdir(mode=0o700) + uv_entrypoint = tmp_path / "uv" + uv_entrypoint.write_text("#!/bin/sh\n") + uv_entrypoint.chmod(0o700) + agent = tmp_path / "agent.plist" + _write_agent(path=agent, repository=repository, uv_entrypoint=uv_entrypoint) + + plan = plan_mission_core_launch_agent( + repository_root=repository, + agent_path=agent, + enable_local_observatory_worker=True, + ) + desired = plistlib.loads(plan.desired_payload) + environment = desired["EnvironmentVariables"] + + assert environment["MISSIONCORE_DATA_DIR"] == str(data_directory) + assert environment["MISSIONCORE_OBSERVATORY_WORKER_LOCAL_ENABLED"] == "1" + assert environment["MISSIONCORE_ARTIFACT_STORE_ROOT"] == str( + data_directory / "observatory-artifact-store" + ) + assert environment["MISSIONCORE_OBSERVATORY_WORKER_SOURCE_CAS_ROOT"] == str( + data_directory / "observatory-worker-source-cas" + ) + assert environment["MISSIONCORE_OBSERVATORY_WORKER_RESULT_STAGING_ROOT"] == str( + data_directory / "observatory-worker-result-staging" + ) + assert plan.local_observatory_worker_enabled is True + assert plan.to_dict()["changes"]["local_observatory_worker_enabled"] is True + + +def test_launch_agent_plan_rejects_missing_local_observatory_roots( + tmp_path: Path, +) -> None: + repository = tmp_path / "repo" + repository.mkdir() + (repository / ".runtime" / "mission-core").mkdir(parents=True, mode=0o700) + uv_entrypoint = tmp_path / "uv" + uv_entrypoint.write_text("#!/bin/sh\n") + uv_entrypoint.chmod(0o700) + agent = tmp_path / "agent.plist" + _write_agent(path=agent, repository=repository, uv_entrypoint=uv_entrypoint) + + with pytest.raises( + MissionCoreLaunchAgentError, + match="local Observatory artifact store is unavailable", + ): + plan_mission_core_launch_agent( + repository_root=repository, + agent_path=agent, + enable_local_observatory_worker=True, + ) def test_launch_agent_plan_rejects_cross_repository_migration_by_default( diff --git a/tests/test_m49_portable_executor_promotion.py b/tests/test_m49_portable_executor_promotion.py new file mode 100644 index 0000000..5302f97 --- /dev/null +++ b/tests/test_m49_portable_executor_promotion.py @@ -0,0 +1,160 @@ +from __future__ import annotations + +import hashlib +import json +import re +from pathlib import Path +from typing import cast + +from k1link.observatory.m49_portable_executor import ( + M49_COMPILED_RUNNER_BUILD_SCHEMA, + M49_PORTABLE_RUNTIME_PHASES, +) +from k1link.observatory.portable_result_contract import canonical_json +from k1link.observatory.portable_run_definitions import PortableRunDefinitionRegistry +from k1link.observatory.portable_worker_runtime import PortableWorkerRuntimeRegistry + +REPOSITORY_ROOT = Path(__file__).resolve().parents[1] +PROMOTION_SCRIPT = ( + REPOSITORY_ROOT + / "experiments" + / "perception" + / "worker" + / "observatory_portable" + / "Invoke-M49PortableExecutorPromotion.ps1" +) +DEFINITION_REGISTRY = REPOSITORY_ROOT / "config" / "observatory-portable-run-definitions.json" +RUNTIME_REGISTRY = REPOSITORY_ROOT / "config" / "observatory-worker-runtime-candidates.json" + +DEFINITION_SHA256 = "f56d6321bd794ccdfb7d2e3b05d044b11f616ffb81ee29517386cc253046d4eb" +CANDIDATE_SHA256 = "65cd2063146a1dd320e30d5f4e21e4bf0aab1ff683e846926cbbfbe25a9f8a5e" +RUNNER_SHA256 = "7be449392ef161fb8713b4c984705d2373bc3cd05645332d92b88a8bff0c7db3" +BUILD_SEAL_SHA256 = "e3bb2e91c70712eff74e8e69718075a407e042fb494616d65984500da42b21a9" +RELEASE_SHA256 = "c5b0670d943fe0452ef4bbfbc144ab2439a1a674f9ef164798ad9f8b1ecc29fa" +RECEIPT_SHA256 = "b560ff9e02746cbb760502f2a3b4b7bd564f95ab52d1149d189927a326b1645a" +EXECUTOR_IMAGE_SHA256 = "f9278ab21aa65045be993dd19bffc25f49955e19598893ac78cc4761ca63ecf3" + + +def _script() -> str: + return PROMOTION_SCRIPT.read_text(encoding="utf-8") + + +def _payload(name: str) -> bytes: + matched = re.search( + rf"\${name}Payload = @'\n(.*?)\n'@\.Trim\(\)", + _script(), + re.DOTALL, + ) + assert matched is not None + return matched.group(1).encode("utf-8") + + +def test_promotion_payloads_are_canonical_and_exact() -> None: + expected = { + "RuntimeBuildSeal": (1014, BUILD_SEAL_SHA256), + "InstalledRelease": (1693, RELEASE_SHA256), + "ReadyReceipt": (2672, RECEIPT_SHA256), + } + for name, (byte_length, sha256) in expected.items(): + payload = _payload(name) + document = json.loads(payload) + assert payload == canonical_json(document) + assert len(payload) == byte_length + assert hashlib.sha256(payload).hexdigest() == sha256 + + build_seal = cast(dict[str, object], json.loads(_payload("RuntimeBuildSeal"))) + assert build_seal["schema_version"] == M49_COMPILED_RUNNER_BUILD_SCHEMA + assert cast(dict[str, object], build_seal["binary"]) == { + "byte_length": 274168, + "file_name": "run_m49_tgs_portable", + "format": "elf", + "sha256": RUNNER_SHA256, + } + + +def test_ready_receipt_closes_the_candidate_without_self_hashing() -> None: + release = cast(dict[str, object], json.loads(_payload("InstalledRelease"))) + receipt = cast(dict[str, object], json.loads(_payload("ReadyReceipt"))) + + assert release["state"] == "ready" + assert release["blockers"] == [] + assert cast(dict[str, object], release["executor_image"])["sha256"] == (EXECUTOR_IMAGE_SHA256) + assert receipt["receipt_state"] == "installed-ready" + assert receipt["release_sha256"] == RELEASE_SHA256 + assert receipt["runtime_release_root"] == "/release" + assert receipt["blockers"] == [] + files = cast(list[dict[str, object]], receipt["files"]) + asset_ids = [cast(str, row["asset_id"]) for row in files] + assert asset_ids == sorted(asset_ids) + assert {row["asset_id"] for row in files} == { + "m49-portable-compiled-runner", + "m49-portable-compiled-runner-build-seal", + "m49-portable-executor-release", + "m49-portable-profile", + } + assert not any(row["asset_id"] == "m49-portable-worker-installation-receipt" for row in files) + + +def test_promotion_is_one_exact_additive_transition() -> None: + script = _script() + lowered = script.lower() + + assert "param()" in script + assert '$ExpectedComputer = "DESKTOP-OPJ8J04"' in script + assert ( + '$CandidateId = "beef090b79ab0e019cff976033d0446c45d4746ee921e02e8076a5a67f21817e"' + in script + ) + assert '$RuntimeReleaseRoot = "/release"' in script + assert "Assert-ProtectedRuntime $protectedBefore" in script + assert "Assert-ProtectedRuntime $protectedAfter" in script + assert "Assert-LegacyTaskReady" in script + assert "requires reconciliation" in script + assert "collided during promotion" in script + assert "Move-Item -LiteralPath $stagingRoot -Destination $readyRoot" in script + for forbidden in ( + "docker build", + "docker create", + "docker run", + "docker tag", + "docker restart", + "docker stop", + "invoke-webrequest", + "invoke-restmethod", + "start-process", + "invoke-expression", + "robocopy", + "scp ", + "ssh ", + ): + assert forbidden not in lowered + + +def test_promoted_v3_definition_and_runtime_registry_bind_exactly() -> None: + definitions = PortableRunDefinitionRegistry.from_file(DEFINITION_REGISTRY) + definition = definitions.resolve_setup("m49-tgs-portable-v2") + assert definition.version == 3 + assert definition.definition_sha256 == DEFINITION_SHA256 + assert definition.executor.ready is True + assert definition.executor.release_sha256 == RELEASE_SHA256 + assert definition.executor.image_sha256 == EXECUTOR_IMAGE_SHA256 + assert definition.components[-1].component_id == "m49-tgs-portable-runner-v1" + assert definition.components[-1].sha256 == BUILD_SEAL_SHA256 + + runtime = PortableWorkerRuntimeRegistry.from_file( + RUNTIME_REGISTRY, + definitions=definitions, + ) + candidate = runtime.resolve("m49-tgs-portable-v2", DEFINITION_SHA256) + assert candidate.ready is True + assert candidate.candidate_sha256 == CANDIDATE_SHA256 + assert candidate.executor is not None + assert candidate.executor.release_sha256 == RELEASE_SHA256 + assert candidate.executor.image_sha256 == EXECUTOR_IMAGE_SHA256 + assert tuple(phase.phase_id for phase in candidate.phases) == M49_PORTABLE_RUNTIME_PHASES + assert all(phase.state == "implemented" for phase in candidate.phases) + assets = {asset.asset_id: asset for asset in candidate.reusable_assets} + assert assets["m49-portable-compiled-runner"].sha256 == RUNNER_SHA256 + assert assets["m49-portable-compiled-runner-build-seal"].sha256 == (BUILD_SEAL_SHA256) + assert assets["m49-portable-executor-release"].sha256 == RELEASE_SHA256 + assert assets["m49-portable-worker-installation-receipt"].sha256 == (RECEIPT_SHA256) diff --git a/tests/test_m49_portable_executor_release.py b/tests/test_m49_portable_executor_release.py index a7cf741..ae9087d 100644 --- a/tests/test_m49_portable_executor_release.py +++ b/tests/test_m49_portable_executor_release.py @@ -296,6 +296,11 @@ def _ready_m49_registry(tmp_path: Path) -> PortableRunDefinitionRegistry: "sha256": _sha256(RUNNER_PATH), } components = cast(list[object], selected["components"]) + components[:] = [ + component + for component in components + if cast(dict[str, object], component)["component_id"] != "m49-tgs-portable-runner-v1" + ] components.append(runner) components.sort(key=lambda value: cast(str, cast(dict[str, object], value)["component_id"])) executor = { @@ -818,7 +823,7 @@ def test_executor_release_candidate_is_deterministic_blocked_and_tamper_evident( assert "FROM ndc/mission-core-m49-t3-travel:20260826" in dockerfile assert builder.M49_TRAVEL_IMAGE_SHA256 in dockerfile assert "--network" not in dockerfile - assert "com.nodedc.authority=\"observation-only\"" in dockerfile + assert 'com.nodedc.authority="observation-only"' in dockerfile installer = INSTALLER_PATH.read_text(encoding="utf-8") assert "--pull=false --no-cache --network none" in installer assert "--network none --read-only" in installer diff --git a/tests/test_observatory_m49_worker_entrypoint.py b/tests/test_observatory_m49_worker_entrypoint.py new file mode 100644 index 0000000..f734403 --- /dev/null +++ b/tests/test_observatory_m49_worker_entrypoint.py @@ -0,0 +1,449 @@ +from __future__ import annotations + +import copy +import hashlib +import json +from pathlib import Path +from threading import Event +from typing import cast + +import httpx +import pytest + +import k1link.observatory.m49_worker_service as service_module +from k1link.observatory.m49_portable_executor import ( + M49_COMPILED_RUNNER_BUILD_SCHEMA, + M49_PORTABLE_COMPILED_RUNNER_ASSET_ID, + M49_PORTABLE_COMPILED_RUNNER_BUILD_SEAL_ASSET_ID, + M49_PORTABLE_COMPILER_CONTRACT, + M49_PORTABLE_PROFILE_ASSET_ID, + M49_PORTABLE_RUNNER_SOURCE_SHA256, + M49_PORTABLE_RUNNER_WRAPPER_SHA256, + M49_PORTABLE_RUNTIME_PHASES, + M49_PORTABLE_TRAVEL_BUILD_IMAGE_SHA256, + M49_PORTABLE_TRAVEL_IMAGE_ASSET_ID, + M49PortableSourceMaterializerAdapter, +) +from k1link.observatory.portable_result_contract import ( + OBSERVATION_ONLY_AUTHORITY, + canonical_json, +) +from k1link.observatory.portable_run_definitions import ( + PortableRunDefinitionRegistry, + canonical_sha256, +) +from k1link.observatory.portable_worker_runtime import PortableWorkerExecutorAdapter +from k1link.observatory.worker_http_transport import ObservatoryWorkerHttpGateway +from k1link.observatory.worker_service import ObservatoryWorkerServiceConfiguration + +REPOSITORY_ROOT = Path(__file__).resolve().parents[1] +DEFINITIONS_PATH = REPOSITORY_ROOT / "config" / "observatory-portable-run-definitions.json" +PROFILE_PATH = REPOSITORY_ROOT / "config" / "perception" / "m49-tgs-portable-v2.json" +SOURCE_REVISION = "f" * 40 +SOURCE_CANDIDATE_SHA256 = "a" * 64 +EXECUTOR_IMAGE_SHA256 = "b" * 64 + + +def _sha256(payload: bytes) -> str: + return hashlib.sha256(payload).hexdigest() + + +def _write(path: Path, payload: bytes, *, executable: bool = False) -> Path: + path.write_bytes(payload) + if executable: + path.chmod(0o755) + return path + + +def _release_files(release_root: Path) -> dict[str, tuple[Path, bytes]]: + binary_payload = b"\x7fELFfixed-m49-worker-test" + binary = _write( + release_root / "run_m49_tgs_portable", + binary_payload, + executable=True, + ) + profile_payload = PROFILE_PATH.read_bytes() + profile = _write(release_root / "m49-tgs-portable-v2.json", profile_payload) + build_seal_payload = canonical_json( + { + "schema_version": M49_COMPILED_RUNNER_BUILD_SCHEMA, + "source_revision": SOURCE_REVISION, + "source_state": "committed-snapshot", + "build_image_sha256": M49_PORTABLE_TRAVEL_BUILD_IMAGE_SHA256, + "profile_sha256": _sha256(profile_payload), + "runner_source_sha256": M49_PORTABLE_RUNNER_SOURCE_SHA256, + "runner_wrapper_sha256": M49_PORTABLE_RUNNER_WRAPPER_SHA256, + "compiler_contract": dict(M49_PORTABLE_COMPILER_CONTRACT), + "binary": { + "file_name": "run_m49_tgs_portable", + "format": "elf", + "byte_length": len(binary_payload), + "sha256": _sha256(binary_payload), + }, + "authority": dict(OBSERVATION_ONLY_AUTHORITY), + } + ) + build_seal = _write(release_root / "compiled-runner-build.json", build_seal_payload) + release_payload = canonical_json( + { + "schema_version": "missioncore.m49-tgs-portable-executor-installed/v1", + "release_id": service_module.M49_WORKER_RELEASE_ID, + "executor_image_sha256": EXECUTOR_IMAGE_SHA256, + "authority": dict(OBSERVATION_ONLY_AUTHORITY), + } + ) + release = _write(release_root / "executor-release.json", release_payload) + return { + M49_PORTABLE_COMPILED_RUNNER_ASSET_ID: (binary, binary_payload), + M49_PORTABLE_COMPILED_RUNNER_BUILD_SEAL_ASSET_ID: ( + build_seal, + build_seal_payload, + ), + M49_PORTABLE_PROFILE_ASSET_ID: (profile, profile_payload), + service_module.M49_EXECUTOR_RELEASE_ASSET_ID: (release, release_payload), + } + + +def _ready_definition_registry( + tmp_path: Path, + *, + release_sha256: str, +) -> tuple[Path, PortableRunDefinitionRegistry]: + root = cast( + dict[str, object], + json.loads(DEFINITIONS_PATH.read_text(encoding="utf-8")), + ) + selected = copy.deepcopy( + next( + cast(dict[str, object], value) + for value in cast(list[object], root["definitions"]) + if cast(dict[str, object], value)["setup_id"] == service_module.M49_WORKER_SETUP_ID + ) + ) + base = PortableRunDefinitionRegistry.from_file(DEFINITIONS_PATH).resolve_setup( + service_module.M49_WORKER_SETUP_ID + ) + components = cast(list[object], selected["components"]) + runner_component = { + "component_id": "m49-tgs-portable-runner-v1", + "kind": "runner", + "sha256": "c" * 64, + } + components[:] = [ + value + for value in components + if cast(dict[str, object], value)["component_id"] != runner_component["component_id"] + ] + components.append(runner_component) + components.sort(key=lambda value: cast(str, cast(dict[str, object], value)["component_id"])) + executor = { + "contour_id": "worker-006", + "state": "ready", + "release_id": service_module.M49_WORKER_RELEASE_ID, + "release_sha256": release_sha256, + "image_sha256": EXECUTOR_IMAGE_SHA256, + "reason_code": None, + "reason": None, + } + selected["executor"] = executor + identity = copy.deepcopy(base.identity_document()) + identity["components"] = copy.deepcopy(components) + identity["executor"] = { + key: executor[key] + for key in ("contour_id", "state", "release_id", "release_sha256", "image_sha256") + } + selected["definition_sha256"] = canonical_sha256(identity) + path = tmp_path / "definitions.json" + path.write_bytes( + canonical_json( + { + "schema_version": root["schema_version"], + "definitions": [selected], + } + ) + ) + return path, PortableRunDefinitionRegistry.from_file(path) + + +def _installation_receipt( + release_root: Path, + files: dict[str, tuple[Path, bytes]], +) -> tuple[Path, bytes]: + release_sha256 = _sha256(files[service_module.M49_EXECUTOR_RELEASE_ASSET_ID][1]) + rows = [ + { + "asset_id": asset_id, + "relative_path": path.name, + "byte_length": len(payload), + "sha256": _sha256(payload), + } + for asset_id, (path, payload) in sorted(files.items()) + ] + document = { + "schema_version": service_module.M49_WORKER_INSTALLATION_RECEIPT_SCHEMA, + "receipt_state": "installed-ready", + "worker_id": "worker-006", + "computer_name": "DESKTOP-OPJ8J04", + "source_revision": SOURCE_REVISION, + "candidate_sha256": SOURCE_CANDIDATE_SHA256, + "candidate_release_sha256": "d" * 64, + "candidate_worker_installation_receipt_sha256": "e" * 64, + "release_id": service_module.M49_WORKER_RELEASE_ID, + "release_sha256": release_sha256, + "release_root": ( + "D:\\NDC_MISSIONCORE\\runtime\\releases\\observatory-portable\\" + f"m49-tgs-portable-candidate-{SOURCE_CANDIDATE_SHA256}\\ready" + ), + "runtime_release_root": str(release_root), + "base_image_sha256": M49_PORTABLE_TRAVEL_BUILD_IMAGE_SHA256, + "executor_image_sha256": EXECUTOR_IMAGE_SHA256, + "files": rows, + "fixture_smoke": "passed", + "protected_runtime": [ + {"name": name, "container_id": f"{index:x}" * 64} + for index, name in enumerate( + service_module._PROTECTED_RUNTIME_NAMES, # noqa: SLF001 + start=1, + ) + ], + "legacy_m49_task_state": None, + "blockers": [], + "authority": dict(OBSERVATION_ONLY_AUTHORITY), + } + payload = canonical_json(document) + path = _write(release_root / "worker-installation-receipt.json", payload) + return path, payload + + +def _runtime_registry( + tmp_path: Path, + *, + definitions: PortableRunDefinitionRegistry, + files: dict[str, tuple[Path, bytes]], + receipt_payload: bytes, +) -> Path: + definition = definitions.resolve_setup(service_module.M49_WORKER_SETUP_ID) + assets = [ + { + "asset_id": M49_PORTABLE_COMPILED_RUNNER_ASSET_ID, + "kind": "local-file", + "sha256": _sha256(files[M49_PORTABLE_COMPILED_RUNNER_ASSET_ID][1]), + "byte_length": len(files[M49_PORTABLE_COMPILED_RUNNER_ASSET_ID][1]), + "component_id": None, + "model_release_id": None, + "model_artifact_role": None, + }, + { + "asset_id": M49_PORTABLE_COMPILED_RUNNER_BUILD_SEAL_ASSET_ID, + "kind": "local-file", + "sha256": _sha256(files[M49_PORTABLE_COMPILED_RUNNER_BUILD_SEAL_ASSET_ID][1]), + "byte_length": len(files[M49_PORTABLE_COMPILED_RUNNER_BUILD_SEAL_ASSET_ID][1]), + "component_id": None, + "model_release_id": None, + "model_artifact_role": None, + }, + { + "asset_id": service_module.M49_EXECUTOR_RELEASE_ASSET_ID, + "kind": "local-file", + "sha256": _sha256(files[service_module.M49_EXECUTOR_RELEASE_ASSET_ID][1]), + "byte_length": len(files[service_module.M49_EXECUTOR_RELEASE_ASSET_ID][1]), + "component_id": None, + "model_release_id": None, + "model_artifact_role": None, + }, + { + "asset_id": M49_PORTABLE_PROFILE_ASSET_ID, + "kind": "definition-component", + "sha256": _sha256(files[M49_PORTABLE_PROFILE_ASSET_ID][1]), + "byte_length": len(files[M49_PORTABLE_PROFILE_ASSET_ID][1]), + "component_id": "m49-tgs-portable-profile-v2", + "model_release_id": None, + "model_artifact_role": None, + }, + { + "asset_id": service_module.M49_WORKER_INSTALLATION_RECEIPT_ASSET_ID, + "kind": "local-file", + "sha256": _sha256(receipt_payload), + "byte_length": len(receipt_payload), + "component_id": None, + "model_release_id": None, + "model_artifact_role": None, + }, + { + "asset_id": M49_PORTABLE_TRAVEL_IMAGE_ASSET_ID, + "kind": "container-image", + "sha256": M49_PORTABLE_TRAVEL_BUILD_IMAGE_SHA256, + "byte_length": None, + "component_id": None, + "model_release_id": None, + "model_artifact_role": None, + }, + ] + assets.sort(key=lambda value: cast(str, value["asset_id"])) + identity = { + "schema_version": "missioncore.observatory-portable-worker-runtime-candidate/v1", + "adapter_id": service_module.M49_WORKER_ADAPTER_ID, + "setup_id": definition.setup_id, + "definition_id": definition.definition_id, + "definition_version": definition.version, + "definition_sha256": definition.definition_sha256, + "source_adapter_sha256": definition.source_adapter.contract_sha256, + "model_manifest_sha256": definition.model_manifest_sha256, + "resource_profile_sha256": definition.resource_profile.profile_sha256, + "result_contract_sha256": definition.result_contract.contract_sha256, + "state": "ready", + "executor": { + "release_id": service_module.M49_WORKER_RELEASE_ID, + "release_sha256": definition.executor.release_sha256, + "image_sha256": EXECUTOR_IMAGE_SHA256, + }, + "reusable_assets": assets, + "phases": [ + {"phase_id": phase_id, "state": "implemented"} + for phase_id in M49_PORTABLE_RUNTIME_PHASES + ], + "blockers": [], + "authority": dict(OBSERVATION_ONLY_AUTHORITY), + } + path = tmp_path / "runtime-registry.json" + path.write_bytes( + canonical_json( + { + "schema_version": ("missioncore.observatory-portable-worker-runtime-registry/v1"), + "candidates": [{**identity, "candidate_sha256": canonical_sha256(identity)}], + } + ) + ) + return path + + +def _fixture( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> tuple[service_module.M49WorkerEntrypointConfiguration, Path]: + release_root = tmp_path / "release" + release_root.mkdir() + monkeypatch.setattr(service_module, "_FIXED_RUNTIME_RELEASE_ROOT", release_root) + files = _release_files(release_root) + release_sha256 = _sha256(files[service_module.M49_EXECUTOR_RELEASE_ASSET_ID][1]) + definitions_path, definitions = _ready_definition_registry( + tmp_path, + release_sha256=release_sha256, + ) + receipt_path, receipt_payload = _installation_receipt(release_root, files) + runtime_path = _runtime_registry( + tmp_path, + definitions=definitions, + files=files, + receipt_payload=receipt_payload, + ) + token = _write(tmp_path / "worker.token", b"worker-006-test-bearer-token-000001") + token.chmod(0o600) + worker = ObservatoryWorkerServiceConfiguration( + base_url="http://127.0.0.1:18080", + bearer_token_file=token, + work_root=tmp_path / "work", + idle_poll_seconds=0.05, + transport_backoff_seconds=0.05, + max_consecutive_transport_failures=2, + ) + return ( + service_module.M49WorkerEntrypointConfiguration( + worker=worker, + definitions_file=definitions_path, + runtime_registry_file=runtime_path, + installation_receipt_file=receipt_path, + ), + receipt_path, + ) + + +def test_fixed_m49_composition_uses_one_gateway_for_agent_source_and_result( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + configuration, _receipt = _fixture(tmp_path, monkeypatch) + requests: list[httpx.Request] = [] + + def handle(request: httpx.Request) -> httpx.Response: + requests.append(request) + return httpx.Response(204) + + service = service_module.compose_installed_m49_worker_service( + configuration, + http_transport=httpx.MockTransport(handle), + ) + gateway = service.gateway + assert isinstance(gateway, ObservatoryWorkerHttpGateway) + assert service.agent._transport is gateway # noqa: SLF001 + registration = service.agent._executors.registrations[0] # noqa: SLF001 + adapter = cast(PortableWorkerExecutorAdapter, registration.adapter) + materializer = cast(M49PortableSourceMaterializerAdapter, adapter.source_materializer) + assert materializer.upstream is gateway + assert adapter.publisher is gateway + + service_module.run_installed_m49_worker(service, stop=Event(), once=True) + + assert [request.url.path for request in requests] == [ + "/api/v1/worker/observatory/recorded-jobs/claims" + ] + assert gateway._client.is_closed # noqa: SLF001 + + +def test_fixed_m49_composition_rejects_receipt_asset_drift_before_gateway( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + configuration, receipt_path = _fixture(tmp_path, monkeypatch) + receipt_path.write_bytes(receipt_path.read_bytes() + b"\n") + + with pytest.raises( + service_module.M49WorkerCompositionError, + match="asset admission is not ready", + ): + service_module.compose_installed_m49_worker_service( + configuration, + http_transport=httpx.MockTransport( + lambda _request: pytest.fail("gateway must not be reached") + ), + ) + + +def test_installation_receipt_rejects_nonfixed_runtime_root_and_links( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + configuration, receipt_path = _fixture(tmp_path, monkeypatch) + monkeypatch.setattr(service_module, "_FIXED_RUNTIME_RELEASE_ROOT", tmp_path / "other") + + with pytest.raises(service_module.M49WorkerCompositionError, match="root is not fixed"): + service_module.load_m49_worker_installation_receipt(receipt_path) + + monkeypatch.setattr(service_module, "_FIXED_RUNTIME_RELEASE_ROOT", receipt_path.parent) + link = tmp_path / "receipt-link.json" + link.symlink_to(receipt_path) + with pytest.raises(service_module.M49WorkerCompositionError, match="unsafe"): + service_module.load_m49_worker_installation_receipt(link) + + +def test_entrypoint_environment_requires_all_absolute_fixed_files(tmp_path: Path) -> None: + environment = { + "MISSIONCORE_OBSERVATORY_WORKER_TOKEN_FILE": str(tmp_path / "worker.token"), + "MISSIONCORE_OBSERVATORY_WORKER_WORK_ROOT": str(tmp_path / "work"), + service_module.M49_WORKER_DEFINITIONS_FILE_ENV: str(tmp_path / "definitions.json"), + service_module.M49_WORKER_RUNTIME_REGISTRY_FILE_ENV: str(tmp_path / "runtime.json"), + service_module.M49_WORKER_INSTALLATION_RECEIPT_FILE_ENV: str(tmp_path / "receipt.json"), + } + + configuration = service_module.M49WorkerEntrypointConfiguration.from_environment(environment) + + assert configuration.worker.base_url == "http://127.0.0.1:18080" + assert configuration.installation_receipt_file == tmp_path / "receipt.json" + with pytest.raises(service_module.M49WorkerCompositionError, match="is required"): + service_module.M49WorkerEntrypointConfiguration.from_environment( + { + key: value + for key, value in environment.items() + if key != service_module.M49_WORKER_RUNTIME_REGISTRY_FILE_ENV + } + ) diff --git a/tests/test_observatory_portable_run_definitions.py b/tests/test_observatory_portable_run_definitions.py index be8962a..90b0a43 100644 --- a/tests/test_observatory_portable_run_definitions.py +++ b/tests/test_observatory_portable_run_definitions.py @@ -21,7 +21,7 @@ REPOSITORY_ROOT = Path(__file__).resolve().parents[1] REGISTRY_PATH = REPOSITORY_ROOT / "config" / "observatory-portable-run-definitions.json" DEFINITION_SHA256 = "3692d41cec3949f348a36eb60a501fb2cd483fed1645679b0ec58061a2fc6dc2" MODEL_MANIFEST_SHA256 = "3fd2d43af73bd73f89d9ffae95d8770cfdeb46033ec967509124fac6ae4afe56" -M49_DEFINITION_SHA256 = "73611f24d70319ea1edca428726d6538a3cbad012a415cc0c1a7ecb7d9b4d910" +M49_DEFINITION_SHA256 = "f56d6321bd794ccdfb7d2e3b05d044b11f616ffb81ee29517386cc253046d4eb" M49_MODEL_MANIFEST_SHA256 = "489a43448f720a9b5c7993dc8279d167b77191a586f0d87b6d38b81cf728e2f1" @@ -99,7 +99,7 @@ def test_source_requirements_map_exactly_to_admission_contract() -> None: assert admission.adapter_sha256 == definition.source_adapter.contract_sha256 -def test_m49_portable_v2_is_model_free_and_contains_no_exact_source_binding() -> None: +def test_m49_portable_v3_is_ready_model_free_and_contains_no_exact_source_binding() -> None: definition = _registry().resolve_setup("m49-tgs-portable-v2") assert definition.definition_sha256 == M49_DEFINITION_SHA256 @@ -107,10 +107,15 @@ def test_m49_portable_v2_is_model_free_and_contains_no_exact_source_binding() -> assert definition.learned_models == () assert definition.model_manifest_sha256 == M49_MODEL_MANIFEST_SHA256 assert definition.resource_profile.accelerator_id == "cpu-only" - assert definition.executor.state == "not-installed" - assert definition.executor.release_id is None - assert definition.executor.release_sha256 is None - assert definition.executor.image_sha256 is None + assert definition.version == 3 + assert definition.executor.state == "ready" + assert definition.executor.release_id == "m49-tgs-portable-executor-v1" + assert definition.executor.release_sha256 == ( + "c5b0670d943fe0452ef4bbfbc144ab2439a1a674f9ef164798ad9f8b1ecc29fa" + ) + assert definition.executor.image_sha256 == ( + "f9278ab21aa65045be993dd19bffc25f49955e19598893ac78cc4761ca63ecf3" + ) identity = json.dumps(definition.identity_document(), sort_keys=True) assert "RAVNOVES00" not in identity assert "20260720T065719Z_viewer_live" not in identity @@ -123,6 +128,9 @@ def test_m49_portable_v2_is_model_free_and_contains_no_exact_source_binding() -> assert hashlib.sha256(profile_path.read_bytes()).hexdigest() == ( components["m49-tgs-portable-profile-v2"].sha256 ) + assert components["m49-tgs-portable-runner-v1"].sha256 == ( + "e3bb2e91c70712eff74e8e69718075a407e042fb494616d65984500da42b21a9" + ) assert profile["source_binding"] == { "mode": "admitted-k1-recording", "camera_timeline": "dynamic", @@ -132,18 +140,15 @@ def test_m49_portable_v2_is_model_free_and_contains_no_exact_source_binding() -> "filesystem_paths": "executor-resolved", } + without_runner = tuple( + component + for component in definition.components + if component.component_id != "m49-tgs-portable-runner-v1" + ) with pytest.raises(PortableRunDefinitionRegistryError, match="runner"): replace( definition, - executor=PortableExecutorAvailability( - contour_id="worker-006", - state="ready", - release_id="m49-tgs-portable-executor-v2", - release_sha256="1" * 64, - image_sha256="2" * 64, - reason_code=None, - reason=None, - ), + components=without_runner, ) @@ -280,7 +285,7 @@ def test_duplicate_definition_and_incomplete_ready_executor_are_rejected( PortableRunDefinitionRegistry.from_file(_write(tmp_path, incomplete)) -def test_production_definition_is_blocked_until_release_and_image_are_sealed() -> None: +def test_blocked_lab_definition_does_not_hide_ready_m49_definition() -> None: registry = _registry() definition = registry.definitions[0] @@ -293,8 +298,9 @@ def test_production_definition_is_blocked_until_release_and_image_are_sealed() - match="not sealed or installed", ): definition.to_recorded_run_definition() - with pytest.raises(PortableRunDefinitionUnavailableError): - registry.to_recorded_registry() + ready = registry.ready_recorded_definitions() + assert tuple(row.setup_id for row in ready) == ("m49-tgs-portable-v2",) + assert registry.to_recorded_registry().definitions == ready def test_conversion_to_recorded_definition_requires_and_preserves_sealed_identities() -> None: @@ -332,7 +338,7 @@ def test_conversion_to_recorded_definition_requires_and_preserves_sealed_identit def test_blocked_definition_does_not_hide_an_unrelated_ready_definition() -> None: registry = _registry() blocked_lab = registry.resolve_setup("lab-v1-eomt-ddrnet-portable-v1") - blocked_m49 = registry.resolve_setup("m49-tgs-portable-v2") + ready_m49 = registry.resolve_setup("m49-tgs-portable-v2") ready_executor = PortableExecutorAvailability( contour_id="worker-006", state="ready", @@ -349,11 +355,15 @@ def test_blocked_definition_does_not_hide_an_unrelated_ready_definition() -> Non executor=ready_executor, definition_sha256=canonical_sha256(identity), ) - mixed = PortableRunDefinitionRegistry((ready_lab, blocked_m49)) + mixed = PortableRunDefinitionRegistry((ready_lab, ready_m49)) - assert mixed.ready_recorded_definitions() == (ready_lab.to_recorded_run_definition(),) - assert mixed.to_recorded_registry().definitions == (ready_lab.to_recorded_run_definition(),) - assert mixed.resolve_setup("m49-tgs-portable-v2") is blocked_m49 + expected = ( + ready_lab.to_recorded_run_definition(), + ready_m49.to_recorded_run_definition(), + ) + assert mixed.ready_recorded_definitions() == expected + assert mixed.to_recorded_registry().definitions == expected + assert mixed.resolve_setup("m49-tgs-portable-v2") is ready_m49 def test_production_lab_v1_model_component_and_result_identities_are_exact() -> None: diff --git a/tests/test_observatory_portable_setup_api.py b/tests/test_observatory_portable_setup_api.py index 149fa5e..4938189 100644 --- a/tests/test_observatory_portable_setup_api.py +++ b/tests/test_observatory_portable_setup_api.py @@ -275,7 +275,9 @@ def test_portable_api_check_sha_fences_ready_submission(tmp_path: Path) -> None: assert binding.submit_count == 1 -def test_portable_api_rejects_blocked_executor_before_binding_submit(tmp_path: Path) -> None: +def test_portable_api_rejects_blocked_lab_executor_before_binding_submit( + tmp_path: Path, +) -> None: full_registry = PortableRunDefinitionRegistry.from_file(REGISTRY_PATH) ready_registry = _ready_lab_registry() queue = ObservatoryRecordedJobQueue( @@ -297,12 +299,12 @@ def test_portable_api_rejects_blocked_executor_before_binding_submit(tmp_path: P recorded_job_queue=queue, ) ) - definition = full_registry.resolve_setup("m49-tgs-portable-v2") + definition = full_registry.resolve_setup("lab-v1-eomt-ddrnet-portable-v1") response = TestClient(app).post( "/api/v1/observatory/runs", json={ "schema_version": "missioncore.observatory-recorded-run-submit/v1", - "idempotency_key": "portable:m49:blocked", + "idempotency_key": "portable:lab-v1:blocked", "source_session_id": SOURCE_SESSION_ID, "setup_id": definition.setup_id, "definition_sha256": definition.definition_sha256, diff --git a/tests/test_observatory_portable_setup_projection.py b/tests/test_observatory_portable_setup_projection.py index 09edd01..0823d9c 100644 --- a/tests/test_observatory_portable_setup_projection.py +++ b/tests/test_observatory_portable_setup_projection.py @@ -156,7 +156,8 @@ def test_generic_catalog_projects_lab_v1_and_model_free_m49_independently() -> N assert m49["display_name"] == PORTABLE_M49_DISPLAY_NAME assert m49["run_definition"]["models"] == [] assert m49["source_compatibility"]["outcome"] == "pass" - assert m49["executor"]["state"] == "not-installed" + assert m49["executor"]["state"] == "ready" + assert m49["executor"]["ready"] is True assert m49["preflight"]["submission_allowed"] is False @@ -165,14 +166,10 @@ def test_calculation_profile_policies_cover_both_exact_portable_definitions() -> profiles = portable_calculation_profile_registry(registry) resolved = { - definition.setup_id: profiles.resolve(definition) - for definition in registry.definitions + definition.setup_id: profiles.resolve(definition) for definition in registry.definitions } assert resolved["lab-v1-eomt-ddrnet-portable-v1"].lab_id == "LAB V1" - assert ( - resolved["lab-v1-eomt-ddrnet-portable-v1"].display_name - == PORTABLE_LAB_V1_DISPLAY_NAME - ) + assert resolved["lab-v1-eomt-ddrnet-portable-v1"].display_name == PORTABLE_LAB_V1_DISPLAY_NAME assert resolved["m49-tgs-portable-v2"].lab_id == "LAB M4.9T5" assert resolved["m49-tgs-portable-v2"].display_name == PORTABLE_M49_DISPLAY_NAME diff --git a/tests/test_observatory_portable_worker_integration.py b/tests/test_observatory_portable_worker_integration.py index 061e8d0..58eb95b 100644 --- a/tests/test_observatory_portable_worker_integration.py +++ b/tests/test_observatory_portable_worker_integration.py @@ -6,6 +6,7 @@ from typing import cast import pytest from k1link.artifact_gateway import CentralArtifactStore +from k1link.observatory import portable_worker_integration as integration_module from k1link.observatory.m49_portable_result import ( M49_PORTABLE_RESULT_CONTRACT_SHA256, validate_m49_portable_result, @@ -31,12 +32,14 @@ from k1link.observatory.portable_setup_projection import ( PORTABLE_M49_SETUP_ID, ) from k1link.observatory.portable_worker_integration import ( + OBSERVATORY_WORKER_LOCAL_ENABLED_ENV, OBSERVATORY_WORKER_RESULT_STAGING_ROOT_ENV, OBSERVATORY_WORKER_SOURCE_CAS_ROOT_ENV, PORTABLE_LAB_V1_RESULT_CONTRACT_SHA256, PortableWorkerIntegrationError, PortableWorkerStorageRoots, build_portable_observatory_worker_integration, + observatory_worker_local_enabled, portable_result_validator_registry, ) from k1link.observatory.recorded_jobs import ObservatoryRecordedJobQueue @@ -59,6 +62,23 @@ def test_exact_validator_registry_covers_both_portable_profiles() -> None: assert validators.resolve(M49_PORTABLE_RESULT_CONTRACT_SHA256) is validate_m49_portable_result +def test_local_worker_gate_is_fail_closed_and_accepts_only_exact_one() -> None: + assert observatory_worker_local_enabled({}) is False + assert observatory_worker_local_enabled({OBSERVATORY_WORKER_LOCAL_ENABLED_ENV: ""}) is False + assert observatory_worker_local_enabled( + {OBSERVATORY_WORKER_LOCAL_ENABLED_ENV: "1"} + ) is True + + for value in ("0", "true", " 1", "1 "): + with pytest.raises( + PortableWorkerIntegrationError, + match="must be exactly 1 when enabled", + ): + observatory_worker_local_enabled( + {OBSERVATORY_WORKER_LOCAL_ENABLED_ENV: value} + ) + + def test_validator_registry_fails_when_one_required_profile_is_absent() -> None: definitions = _definitions() only_lab_v1 = PortableRunDefinitionRegistry((definitions.definitions[0],)) @@ -173,6 +193,58 @@ def test_storage_roots_load_only_from_existing_disjoint_central_directories( ) +def test_storage_roots_accept_canonical_local_directories_under_data_dir( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + data_dir = tmp_path / "missioncore-data" + boundary = data_dir / "observatory-worker-local" + artifact_store = boundary / "artifact-store" + source_cas = boundary / "source-cas" + result_staging = boundary / "result-staging" + for path in (artifact_store, source_cas, result_staging): + path.mkdir(parents=True, exist_ok=True) + mount_checks: list[Path] = [] + monkeypatch.setattr( + integration_module.os.path, + "ismount", + lambda path: mount_checks.append(Path(path)) or False, + ) + + roots = PortableWorkerStorageRoots.from_paths( + artifact_store_root=artifact_store, + source_cas_root=source_cas, + result_staging_root=result_staging, + ) + + assert roots == PortableWorkerStorageRoots( + source_cas_root=source_cas, + result_staging_root=result_staging, + ) + assert mount_checks == [] + + +def test_explicit_volumes_storage_boundary_still_requires_a_mount( + monkeypatch: pytest.MonkeyPatch, +) -> None: + mount_checks: list[Path] = [] + monkeypatch.setattr( + integration_module.os.path, + "ismount", + lambda path: mount_checks.append(Path(path)) or False, + ) + + with pytest.raises( + PortableWorkerIntegrationError, + match="central artifact volume is not mounted: /Volumes/nodedc", + ): + integration_module._require_mounted_volume( + Path("/Volumes/nodedc/mission-core") + ) + + assert mount_checks == [Path("/Volumes/nodedc")] + + def test_server_integration_has_no_data_directory_storage_fallback( tmp_path: Path, ) -> None: diff --git a/tests/test_observatory_portable_worker_runtime.py b/tests/test_observatory_portable_worker_runtime.py index e08f2e1..f20b1e8 100644 --- a/tests/test_observatory_portable_worker_runtime.py +++ b/tests/test_observatory_portable_worker_runtime.py @@ -34,12 +34,8 @@ from k1link.observatory.worker_agent import ( ) REPOSITORY_ROOT = Path(__file__).resolve().parents[1] -DEFINITION_REGISTRY = ( - REPOSITORY_ROOT / "config" / "observatory-portable-run-definitions.json" -) -RUNTIME_REGISTRY = ( - REPOSITORY_ROOT / "config" / "observatory-worker-runtime-candidates.json" -) +DEFINITION_REGISTRY = REPOSITORY_ROOT / "config" / "observatory-portable-run-definitions.json" +RUNTIME_REGISTRY = REPOSITORY_ROOT / "config" / "observatory-worker-runtime-candidates.json" def _definitions() -> PortableRunDefinitionRegistry: @@ -55,33 +51,39 @@ def _runtime() -> PortableWorkerRuntimeRegistry: def _all_keys(value: object) -> set[str]: if isinstance(value, dict): - return set(value) | { - nested - for child in value.values() - for nested in _all_keys(child) - } + return set(value) | {nested for child in value.values() for nested in _all_keys(child)} if isinstance(value, list): return {nested for child in value for nested in _all_keys(child)} return set() -def test_production_candidates_bind_exact_definitions_but_remain_blocked() -> None: +def test_production_candidates_bind_exact_definitions_and_only_m49_is_ready() -> None: registry = _runtime() assert {candidate.setup_id for candidate in registry.candidates} == { "lab-v1-eomt-ddrnet-portable-v1", "m49-tgs-portable-v2", } - for candidate in registry.candidates: - assert candidate.ready is False - assert candidate.executor is None - assert "executor-release-unsealed" in candidate.blockers - assert any(phase.state == "missing" for phase in candidate.phases) - with pytest.raises( - PortableWorkerRuntimeUnavailableError, - match="no executor identity", - ): - candidate.executor_identity() + by_setup = {candidate.setup_id: candidate for candidate in registry.candidates} + lab_v1 = by_setup["lab-v1-eomt-ddrnet-portable-v1"] + assert lab_v1.ready is False + assert lab_v1.executor is None + assert "executor-release-unsealed" in lab_v1.blockers + assert any(phase.state == "missing" for phase in lab_v1.phases) + with pytest.raises( + PortableWorkerRuntimeUnavailableError, + match="no executor identity", + ): + lab_v1.executor_identity() + + m49 = by_setup["m49-tgs-portable-v2"] + assert m49.ready is True + assert m49.executor is not None + assert m49.blockers == () + assert all(phase.state == "implemented" for phase in m49.phases) + assert m49.executor_identity().release_sha256 == ( + "c5b0670d943fe0452ef4bbfbc144ab2439a1a674f9ef164798ad9f8b1ecc29fa" + ) def test_candidate_contract_is_source_independent_and_instruction_free() -> None: @@ -96,75 +98,34 @@ def test_candidate_contract_is_source_independent_and_instruction_free() -> None assert not any("priority" in key for key in keys) -def test_m49_reuses_portable_profile_generic_runner_and_travel_image() -> None: +def test_m49_ready_candidate_requires_promoted_runner_receipts_and_travel_image() -> None: candidate = _runtime().resolve( "m49-tgs-portable-v2", - "73611f24d70319ea1edca428726d6538a3cbad012a415cc0c1a7ecb7d9b4d910", + "f56d6321bd794ccdfb7d2e3b05d044b11f616ffb81ee29517386cc253046d4eb", ) bindings = { "m49-portable-profile": PortableWorkerLocalAssetBinding( asset_id="m49-portable-profile", file_path=REPOSITORY_ROOT / "config" / "perception" / "m49-tgs-portable-v2.json", ), - "m49-portable-runner-source": PortableWorkerLocalAssetBinding( - asset_id="m49-portable-runner-source", - file_path=( - REPOSITORY_ROOT - / "experiments" - / "perception" - / "worker" - / "observatory_portable" - / "run_m49_tgs_portable.cpp" - ), - ), - "m49-portable-runner-manifest": PortableWorkerLocalAssetBinding( - asset_id="m49-portable-runner-manifest", - file_path=( - REPOSITORY_ROOT - / "experiments" - / "perception" - / "worker" - / "observatory_portable" - / "m49-tgs-portable-runner-source.json" - ), - ), - "m49-portable-runner-wrapper": PortableWorkerLocalAssetBinding( - asset_id="m49-portable-runner-wrapper", - file_path=( - REPOSITORY_ROOT - / "experiments" - / "perception" - / "worker" - / "observatory_portable" - / "run_m49_tgs_portable.sh" - ), - ), - "m49-portable-smoke": PortableWorkerLocalAssetBinding( - asset_id="m49-portable-smoke", - file_path=( - REPOSITORY_ROOT - / "experiments" - / "perception" - / "worker" - / "observatory_portable" - / "smoke_m49_tgs_portable.sh" - ), - ), "travel-tgs-image": PortableWorkerLocalAssetBinding( asset_id="travel-tgs-image", - image_sha256=( - "7b412020f4d8392d1d1ed1b33beadc44140f0ea8f781e62dd69796042334300f" - ), + image_sha256=("7b412020f4d8392d1d1ed1b33beadc44140f0ea8f781e62dd69796042334300f"), ), } admission = inspect_runtime_candidate(candidate, bindings) - assert {item.state for item in admission.assets} == {"matched"} + states = {item.asset_id: item.state for item in admission.assets} + assert states["m49-portable-profile"] == "matched" + assert states["travel-tgs-image"] == "matched" + assert states["m49-portable-compiled-runner"] == "missing" + assert states["m49-portable-compiled-runner-build-seal"] == "missing" + assert states["m49-portable-executor-release"] == "missing" + assert states["m49-portable-worker-installation-receipt"] == "missing" assert admission.ready is False - assert "portable-tgs-runner-unsealed" in admission.blockers - assert "portable-camera-lidar-timeline-unimplemented" in admission.blockers - assert "portable-result-assembler-unimplemented" in admission.blockers + assert "asset-m49-portable-compiled-runner-missing" in admission.blockers + assert "asset-m49-portable-worker-installation-receipt-missing" in admission.blockers def test_m49_portable_runner_has_no_exact_source_or_frame_count_binding() -> None: @@ -195,13 +156,7 @@ def test_m49_portable_runner_has_no_exact_source_or_frame_count_binding() -> Non def test_m49_portable_runner_source_release_is_content_addressed() -> None: - root = ( - REPOSITORY_ROOT - / "experiments" - / "perception" - / "worker" - / "observatory_portable" - ) + root = REPOSITORY_ROOT / "experiments" / "perception" / "worker" / "observatory_portable" manifest = json.loads( (root / "m49-tgs-portable-runner-source.json").read_text(encoding="utf-8") ) @@ -233,9 +188,7 @@ def test_lab_candidate_verifies_reusable_repository_assets_without_claiming_exec bindings = { "ddrnet-goose-image": PortableWorkerLocalAssetBinding( asset_id="ddrnet-goose-image", - image_sha256=( - "591cb382c099eeb05e7ec16e2371e0b2da54d2bb5c49ec0f4ac88dbf72b0f0cd" - ), + image_sha256=("591cb382c099eeb05e7ec16e2371e0b2da54d2bb5c49ec0f4ac88dbf72b0f0cd"), ), "ddrnet-goose-runner": PortableWorkerLocalAssetBinding( asset_id="ddrnet-goose-runner", @@ -250,9 +203,7 @@ def test_lab_candidate_verifies_reusable_repository_assets_without_claiming_exec ), "eomt-image": PortableWorkerLocalAssetBinding( asset_id="eomt-image", - image_sha256=( - "58df7489c3f2276f9591d500a012dee03e23d35543ce3c390b4c001e6bf90794" - ), + image_sha256=("58df7489c3f2276f9591d500a012dee03e23d35543ce3c390b4c001e6bf90794"), ), "eomt-orchestrator": PortableWorkerLocalAssetBinding( asset_id="eomt-orchestrator", @@ -324,7 +275,7 @@ def test_lab_candidate_verifies_reusable_repository_assets_without_claiming_exec def test_local_asset_tampering_is_reported_without_execution(tmp_path: Path) -> None: candidate = _runtime().resolve( "m49-tgs-portable-v2", - "73611f24d70319ea1edca428726d6538a3cbad012a415cc0c1a7ecb7d9b4d910", + "f56d6321bd794ccdfb7d2e3b05d044b11f616ffb81ee29517386cc253046d4eb", ) tampered = tmp_path / "m49-profile.json" tampered.write_text("{}\n", encoding="utf-8") @@ -346,9 +297,10 @@ def test_local_asset_tampering_is_reported_without_execution(tmp_path: Path) -> states = {item.asset_id: item.state for item in admission.assets} assert states["m49-portable-profile"] == "mismatched" assert states["travel-tgs-image"] == "mismatched" - assert states["m49-portable-runner-source"] == "missing" - assert states["m49-portable-runner-manifest"] == "missing" - assert states["m49-portable-runner-wrapper"] == "missing" + assert states["m49-portable-compiled-runner"] == "missing" + assert states["m49-portable-compiled-runner-build-seal"] == "missing" + assert states["m49-portable-executor-release"] == "missing" + assert states["m49-portable-worker-installation-receipt"] == "missing" assert "asset-m49-portable-profile-mismatched" in admission.blockers assert "asset-travel-tgs-image-mismatched" in admission.blockers @@ -431,8 +383,7 @@ def test_ready_local_adapter_composes_only_local_ports_and_exact_job( image_sha256="2" * 64, ) phases = tuple( - PortableWorkerRuntimePhase(phase.phase_id, "implemented") - for phase in blocked.phases + PortableWorkerRuntimePhase(phase.phase_id, "implemented") for phase in blocked.phases ) candidate_identity = blocked.identity_document() candidate_identity["definition_sha256"] = definition.definition_sha256 @@ -470,9 +421,7 @@ def test_ready_local_adapter_composes_only_local_ports_and_exact_job( return PortableWorkerSourceStage( root=source_root, source_bundle_sha256=job.source_bundle_sha256, - source_capability_manifest_sha256=( - job.source_capability_manifest_sha256 - ), + source_capability_manifest_sha256=(job.source_capability_manifest_sha256), source_adapter_sha256=job.source_adapter_sha256, ) diff --git a/tests/test_observatory_worker_app_wiring.py b/tests/test_observatory_worker_app_wiring.py index 5d28ace..e881e65 100644 --- a/tests/test_observatory_worker_app_wiring.py +++ b/tests/test_observatory_worker_app_wiring.py @@ -1,5 +1,9 @@ from __future__ import annotations +import json +import os +import subprocess +import sys from pathlib import Path from k1link.observatory.m49_portable_result import ( @@ -8,6 +12,7 @@ from k1link.observatory.m49_portable_result import ( ) from k1link.observatory.portable_lab_v1_executor import validate_lab_v1_result_v2 from k1link.observatory.portable_worker_integration import ( + OBSERVATORY_WORKER_LOCAL_ENABLED_ENV, PORTABLE_LAB_V1_RESULT_CONTRACT_SHA256, ) from k1link.web import app as app_module @@ -15,7 +20,7 @@ from k1link.web import app as app_module WORKER_ROUTE_PREFIX = "/api/v1/worker/observatory" -def test_worker_router_is_hard_disabled_until_lease_and_publisher_exist() -> None: +def test_worker_router_is_fail_closed_without_explicit_local_gate() -> None: assert app_module.OBSERVATORY_RECORDED_JOB_QUEUE is not None assert app_module.OBSERVATORY_PORTABLE_RESULT_VALIDATORS is not None assert ( @@ -32,12 +37,15 @@ def test_worker_router_is_hard_disabled_until_lease_and_publisher_exist() -> Non ) assert app_module.OBSERVATORY_WORKER_CLAIM_LEASE_READY is False assert app_module.OBSERVATORY_WORKER_VERIFIED_RESULT_PUBLISHER_READY is False + assert app_module.OBSERVATORY_WORKER_LOCAL_ENABLED is False + assert app_module.OBSERVATORY_WORKER_API_GATE_ENABLED is False assert app_module.OBSERVATORY_WORKER_PRODUCTION_API_ENABLED is False assert app_module.OBSERVATORY_WORKER_DISPATCH_READY is False assert app_module.OBSERVATORY_WORKER_AUTHENTICATION is None assert app_module.OBSERVATORY_WORKER_AUTHENTICATION_ERROR is not None assert app_module.OBSERVATORY_WORKER_API_ERROR is not None - assert "hard-disabled" in app_module.OBSERVATORY_WORKER_API_ERROR + assert "local-only Worker API gate is disabled" in app_module.OBSERVATORY_WORKER_API_ERROR + assert OBSERVATORY_WORKER_LOCAL_ENABLED_ENV in app_module.OBSERVATORY_WORKER_API_ERROR if app_module.session_artifact_gateway is None: assert app_module.OBSERVATORY_PORTABLE_WORKER_INTEGRATION is None assert app_module.OBSERVATORY_PORTABLE_WORKER_INTEGRATION_ERROR is not None @@ -87,3 +95,74 @@ def test_valid_worker_credential_cannot_enable_production_router( getattr(route, "path", "").startswith(WORKER_ROUTE_PREFIX) for route in app_module.app.routes ) + + +def test_explicit_local_gate_enables_router_with_local_storage_and_authentication( + tmp_path: Path, +) -> None: + data_dir = tmp_path / "data" + evidence_dir = tmp_path / "evidence" + legacy_dir = tmp_path / "legacy" + boundary = data_dir / "observatory-worker-local" + artifact_store = boundary / "artifact-store" + source_cas = boundary / "source-cas" + result_staging = boundary / "result-staging" + token_path = data_dir / "worker-auth" / "observatory-worker.token" + for path in ( + data_dir, + evidence_dir, + legacy_dir, + artifact_store, + source_cas, + result_staging, + token_path.parent, + ): + path.mkdir(mode=0o700, parents=True, exist_ok=True) + token_path.write_text("worker-006-test-bearer-secret-32bytes", encoding="ascii") + token_path.chmod(0o600) + environment = os.environ.copy() + environment.update( + { + "MISSIONCORE_DATA_DIR": str(data_dir), + "MISSIONCORE_EVIDENCE_DIR": str(evidence_dir), + "MISSIONCORE_LEGACY_SESSIONS_DIR": str(legacy_dir), + "MISSIONCORE_ARTIFACT_STORE_ROOT": str(artifact_store), + "MISSIONCORE_OBSERVATORY_WORKER_SOURCE_CAS_ROOT": str(source_cas), + "MISSIONCORE_OBSERVATORY_WORKER_RESULT_STAGING_ROOT": str(result_staging), + OBSERVATORY_WORKER_LOCAL_ENABLED_ENV: "1", + } + ) + probe = """ +import json +from k1link.web import app as app_module +prefix = "/api/v1/worker/observatory" +paths = app_module.app.openapi().get("paths", {}) +print(json.dumps({ + "local_enabled": app_module.OBSERVATORY_WORKER_LOCAL_ENABLED, + "api_gate_enabled": app_module.OBSERVATORY_WORKER_API_GATE_ENABLED, + "claim_ready": app_module.OBSERVATORY_WORKER_CLAIM_LEASE_READY, + "publisher_ready": app_module.OBSERVATORY_WORKER_VERIFIED_RESULT_PUBLISHER_READY, + "dispatch_ready": app_module.OBSERVATORY_WORKER_DISPATCH_READY, + "api_error": app_module.OBSERVATORY_WORKER_API_ERROR, + "worker_route": any(path.startswith(prefix) for path in paths), +})) +""" + + completed = subprocess.run( + [sys.executable, "-c", probe], + check=True, + capture_output=True, + text=True, + env=environment, + ) + state = json.loads(completed.stdout) + + assert state == { + "local_enabled": True, + "api_gate_enabled": True, + "claim_ready": True, + "publisher_ready": True, + "dispatch_ready": True, + "api_error": None, + "worker_route": True, + }