diff --git a/experiments/perception/worker/streaming_profile_stage1/collect_worker_host.ps1 b/experiments/perception/worker/streaming_profile_stage1/collect_worker_host.ps1 new file mode 100644 index 0000000..f168e8f --- /dev/null +++ b/experiments/perception/worker/streaming_profile_stage1/collect_worker_host.ps1 @@ -0,0 +1,137 @@ +# Trusted Windows Worker adapter. Read-only Docker/NVIDIA commands only. +# Run on the host, NEVER in the model container. Requests and responses are +# separate directories; mount only responses read-only into the runtime. +# Docker GPU access inventory is NOT a native host/WSL process audit. +param( + [Parameter(Mandatory=$true)][string]$Requests, + [Parameter(Mandatory=$true)][string]$Responses, + [Parameter(Mandatory=$true)][string]$Container, + [string]$WorkerId = 'worker-006', + [ValidateRange(1,300)][int]$Seconds = 240, + # Explicit bounded lab fault, off by default. Never supplied by input data. + [string]$DiagnosticDelayAfterProgressFile = '', + [ValidateRange(0,5000)][int]$DiagnosticDelayMilliseconds = 0 +) +$ErrorActionPreference = 'Stop' +if ($Container -notmatch '^[a-zA-Z0-9][a-zA-Z0-9_.-]{0,127}$') { throw 'Invalid launcher target' } +if ([IO.Path]::GetFullPath($Requests) -eq [IO.Path]::GetFullPath($Responses)) { throw 'Separate control mounts required' } +$utf8 = New-Object System.Text.UTF8Encoding($false) +$deadline = [Diagnostics.Stopwatch]::StartNew() +$targetId = $null +$lastNonce = $null +$samples = $errors = 0 +$lastErrorPhase = $null +$phase = 'request' +$delayEvidence = $null + +function Read-Command([string]$File, [string]$Arguments, [int]$Limit = 65536) { + $info = New-Object Diagnostics.ProcessStartInfo + $info.FileName = $File + $info.Arguments = $Arguments + $info.UseShellExecute = $false + $info.CreateNoWindow = $true + $info.RedirectStandardOutput = $true + $info.RedirectStandardError = $true + $process = New-Object Diagnostics.Process + $process.StartInfo = $info + try { + if (-not $process.Start()) { throw 'Command did not start' } + # Every command below has a fixed projection and <=64 inspected objects. + # stderr is bounded by the known CLI, never copied into the channel/log. + $output = $process.StandardOutput.ReadToEndAsync() + $errorOutput = $process.StandardError.ReadToEndAsync() + if (-not $process.WaitForExit(750)) { + $process.Kill() + $process.WaitForExit() + throw 'Read-only command timeout' + } + $value = $output.GetAwaiter().GetResult() + if ($process.ExitCode -ne 0 -or $value.Length -gt $Limit) { throw 'Read-only command failed or exceeded bound' } + return $value.Trim() + } finally { $process.Dispose() } +} + +# Never serialize Config.Env, mounts, credentials or full Docker inspect. +$template = '{"id":{{json .Id}},"image":{{json .Image}},"running":{{json .State.Running}},"runtime":{{json .HostConfig.Runtime}},"privileged":{{json .HostConfig.Privileged}},"devices":{{if .HostConfig.Devices}}true{{else}}false{{end}},"gpu":{{if .HostConfig.DeviceRequests}}true{{else}}false{{end}},"nano":{{json .HostConfig.NanoCpus}},"quota":{{json .HostConfig.CpuQuota}},"period":{{json .HostConfig.CpuPeriod}},"cpuset":{{json .HostConfig.CpusetCpus}},"memory":{{json .HostConfig.Memory}}}' +$escapedTemplate = $template.Replace('"', '\"') +try { + while ($deadline.Elapsed.TotalSeconds -lt $Seconds -and -not (Test-Path (Join-Path $Requests 'stop'))) { + Start-Sleep -Milliseconds 25 + try { + $phase = 'request' + $path = Join-Path $Requests 'request.json' + if (-not (Test-Path $path) -or (Get-Item $path).Length -gt 1024) { continue } + $request = [IO.File]::ReadAllText($path) | ConvertFrom-Json + if (@($request.PSObject.Properties.Name).Count -ne 4 -or + $request.schema_version -ne 'missioncore.worker-host-observation/v1' -or + $request.nonce -cnotmatch '^[a-f0-9]{64}$' -or + $request.activation_sha256 -cnotmatch '^[a-f0-9]{64}$' -or + $request.sequence -lt 1 -or $request.nonce -eq $lastNonce) { continue } + $facts = [ordered]@{ + worker_id=$WorkerId; container_id=$null; image_sha256=$null + cpu_limit_millicores=$null; memory_limit_mib=$null + gpu_name=$null; driver_version=$null; sm_clock_mhz=$null; memory_clock_mhz=$null + gpu_telemetry_available=$false; inventory_complete=$false + competing_gpu_clients=@(); inventory_scope='docker-gpu-access'; host_process_inventory='unproved' + } + try { + $phase = 'gpu' + $gpu = (Read-Command 'nvidia-smi.exe' '--query-gpu=name,driver_version,clocks.sm,clocks.mem --format=csv,noheader,nounits -i 0' 1024).Split(',') + if ($gpu.Count -ne 4) { throw 'GPU schema' } + $facts.gpu_name=$gpu[0].Trim(); $facts.driver_version=$gpu[1].Trim() + $facts.sm_clock_mhz=[int]$gpu[2].Trim(); $facts.memory_clock_mhz=[int]$gpu[3].Trim() + $facts.gpu_telemetry_available=$true + } catch { $errors++; $lastErrorPhase=$phase } + try { + $phase = 'docker' + if (-not $targetId) { + $targetId = Read-Command 'docker.exe' "inspect --format {{.Id}} $Container" 128 + if ($targetId -cnotmatch '^[a-f0-9]{64}$') { $targetId=$null; throw 'Target not ready' } + } + $ids = @((Read-Command 'docker.exe' 'ps --no-trunc -q' 8192) -split '\r?\n' | Where-Object { $_ }) + if ($ids.Count -gt 64 -or @($ids | Where-Object { $_ -cnotmatch '^[a-f0-9]{64}$' }).Count) { throw 'Inventory bound' } + $inspectIds = @($ids + $targetId | Sort-Object -Unique) + $rows = (Read-Command 'docker.exe' ('inspect --format "' + $escapedTemplate + '" ' + ($inspectIds -join ' '))) -split '\r?\n' | ForEach-Object { $_ | ConvertFrom-Json } + $target = @($rows | Where-Object { $_.id -eq $targetId }) + if ($target.Count -eq 1 -and $target[0].running) { + $t = $target[0] + $facts.container_id=$t.id; $facts.image_sha256=$t.image.Substring(7) + if ($t.nano -gt 0) { $facts.cpu_limit_millicores=[long]($t.nano / 1000000) } + elseif ($t.quota -gt 0 -and $t.period -gt 0) { $facts.cpu_limit_millicores=[long](1000 * $t.quota / $t.period) } + elseif (-not $t.cpuset) { $facts.cpu_limit_millicores=0 } + $facts.memory_limit_mib=[long][Math]::Ceiling($t.memory / 1048576) + } + $facts.competing_gpu_clients=@($rows | Where-Object { + $_.id -ne $targetId -and $_.id -in $ids -and + ($_.gpu -or $_.devices -or $_.privileged -or $_.runtime -eq 'nvidia') + } | ForEach-Object { 'docker:' + $_.id }) + $after = @((Read-Command 'docker.exe' 'ps --no-trunc -q' 8192) -split '\r?\n' | Where-Object { $_ }) + $facts.inventory_complete = (($ids | Sort-Object) -join ',') -eq (($after | Sort-Object) -join ',') + } catch { $errors++; $lastErrorPhase=$phase } + $phase = 'publish' + if ($DiagnosticDelayMilliseconds -gt 0 -and -not $delayEvidence -and + (Test-Path $DiagnosticDelayAfterProgressFile) -and + (Get-Item $DiagnosticDelayAfterProgressFile).Length -le 65536 -and + (Get-Content -Raw $DiagnosticDelayAfterProgressFile) -match '"completed": 16') { + $delayStarted = [Diagnostics.Stopwatch]::GetTimestamp() + Start-Sleep -Milliseconds $DiagnosticDelayMilliseconds + $delayEvidence = @{request_sequence=$request.sequence;requested_ms=$DiagnosticDelayMilliseconds;started_ticks=$delayStarted.ToString();ended_ticks=[Diagnostics.Stopwatch]::GetTimestamp().ToString();clock_domain='windows-stopwatch'} + } + $reply = [ordered]@{ + schema_version=$request.schema_version; activation_sha256=$request.activation_sha256 + nonce=$request.nonce; sequence=$request.sequence; facts=$facts + } | ConvertTo-Json -Depth 5 -Compress + if ($utf8.GetByteCount($reply) -gt 16384) { throw 'Reply bound' } + $temporary = Join-Path $Responses 'response.tmp' + $destination = Join-Path $Responses 'response.json' + [IO.File]::WriteAllText($temporary, $reply, $utf8) + # Windows PowerShell 5 coerces $null to an empty string for this .NET + # overload (invalid backup path). NullString preserves a true null. + if (Test-Path $destination) { [IO.File]::Replace($temporary, $destination, [NullString]::Value) } + else { [IO.File]::Move($temporary, $destination) } + $lastNonce=$request.nonce; $samples++ + } catch { $errors++; $lastErrorPhase=$phase } + } +} finally { + [ordered]@{schema_version='missioncore.worker-host-collector-exit/v1'; utc=[DateTime]::UtcNow.ToString('o'); monotonic_ticks=[Diagnostics.Stopwatch]::GetTimestamp().ToString(); clock_domain='windows-stopwatch'; elapsed_ms=$deadline.Elapsed.TotalMilliseconds; samples=$samples; errors=$errors; last_error_phase=$lastErrorPhase; container_id=$targetId; host_mutations=0; host_process_inventory='unproved';diagnostic_delay=$delayEvidence} | ConvertTo-Json -Depth 4 -Compress +} diff --git a/experiments/perception/worker/streaming_profile_stage1/pilot_lifecycle.py b/experiments/perception/worker/streaming_profile_stage1/pilot_lifecycle.py index a8280ee..80a2283 100644 --- a/experiments/perception/worker/streaming_profile_stage1/pilot_lifecycle.py +++ b/experiments/perception/worker/streaming_profile_stage1/pilot_lifecycle.py @@ -16,8 +16,12 @@ from pathlib import Path from pilot_freshness import CLOCK_DOMAIN from k1link.perception.realtime_contract import StreamStart +from k1link.perception.streaming_continuity import StreamSuspended from k1link.perception.streaming_lifecycle import StreamingLifecycle +from k1link.perception.worker_control import WorkerControlChannel +from k1link.perception.worker_control_pump import WorkerControlPump from k1link.perception.worker_lease import WorkerLeaseError +from k1link.perception.worker_operating_envelope import WorkerOperatingEnvelope def bounded_digest(path, limit=65536): @@ -59,6 +63,21 @@ class FencedIngress: class PilotController: def __init__(self, args, report, mailbox, stop): + control_requests = getattr(args, "worker_control_requests", None) + control_config = None + if control_requests: + if not os.statvfs(args.worker_control_responses).f_flag & os.ST_RDONLY: + raise ValueError("host responses must be a read-only controller mount") + profile = json.loads(Path("/out/candidate-profile.json").read_text()) + control_config = { + "envelope": { + **profile["operating_envelope"]["reference_conditions"], + "inventory_scope": "docker-gpu-access", + }, + "mode": args.worker_readiness_mode, + "channel": "missioncore.worker-host-observation/v1", + "host_process_inventory": "unproved", + } config = { "pilot_options": { key: value @@ -79,6 +98,8 @@ class PilotController: for key in ("OPENBLAS_NUM_THREADS", "OMP_NUM_THREADS", "MKL_NUM_THREADS") }, } + if control_config: + config["worker_control"] = control_config self.start = StreamStart( run_id=args.run_id, source_id="RAVNOVES00-20260720T065719Z-viewer-live", @@ -110,6 +131,25 @@ class PilotController: target=self._renew, name="pilot-controller-heartbeat", daemon=True ) self.thread.start() + self.control = None + self.report = report + try: + if control_config: + self.control = WorkerControlPump( + self.runtime, + WorkerControlChannel( + self.start, + Path(control_requests), + Path(args.worker_control_responses), + container_id=os.environ["HOSTNAME"], + clock_domain_id="worker-linux-monotonic", + ), + WorkerOperatingEnvelope(**control_config["envelope"]), + mode=control_config["mode"], + ) + except BaseException: + self.close("worker-control-bootstrap-failed") + raise report["runtime_binding"] = self.start.to_dict() report["runtime_effective_config"] = config report["lease_ttl_ms"] = 2000 @@ -133,6 +173,26 @@ class PilotController: self.heartbeat_stop.set() self.thread.join(timeout=1) + def ready(self): + if self.control: + self.control.ready() + else: + self.runtime.ready() + + def spawn(self, factory): + # A bounded diagnostic still has its outer wall-clock watchdog. While + # warmup facts lag, preserve already loaded children and the local lease. + while True: + try: + return self.runtime.spawn(factory) + except StreamSuspended: + if not self.control: + raise + self.heartbeat_stop.wait(0.025) + def close(self, reason): + if self.control: + self.report["worker_control_thread_released"] = self.control.close() + self.report["worker_control"] = self.control.snapshot() self.stop_renewals() return self.runtime.close(reason) diff --git a/experiments/perception/worker/streaming_profile_stage1/run_joint_pilot.py b/experiments/perception/worker/streaming_profile_stage1/run_joint_pilot.py index 2922b65..32447b4 100644 --- a/experiments/perception/worker/streaming_profile_stage1/run_joint_pilot.py +++ b/experiments/perception/worker/streaming_profile_stage1/run_joint_pilot.py @@ -303,7 +303,7 @@ def run(args): start_new_session=True, ) - process = controller.runtime.spawn(spawn) if controller else spawn() + process = controller.spawn(spawn) if controller else spawn() children.append(process) return process @@ -523,7 +523,7 @@ def run(args): return compute_gpu_impl(bundle) if controller: - controller.runtime.ready() + controller.ready() if args.telemetry_mode == "nvml": monitor = threading.Thread( @@ -919,6 +919,13 @@ if __name__ == "__main__": parser.add_argument("--worker-lease-root") parser.add_argument("--lease-generation", type=int, default=1) parser.add_argument("--controller-image-sha256") + parser.add_argument("--worker-control-requests") + parser.add_argument("--worker-control-responses") + parser.add_argument( + "--worker-readiness-mode", + choices=("strict-envelope", "labelled-experiment"), + default="labelled-experiment", + ) parser.add_argument("--stop-renew-after-sequence", type=int, default=-1) parser.add_argument("--recover-input", action="store_true") parser.add_argument("--input-gap", action="append", default=[], metavar="SEQUENCE:MILLISECONDS") @@ -943,6 +950,10 @@ if __name__ == "__main__": parser.error("invalid bounded input gap plan") if args.recover_input and args.input_transport != "binary-ipc": parser.error("resumable full profile requires binary input") + if bool(args.worker_control_requests) != bool(args.worker_control_responses): + parser.error("host control needs separate request/response mounts") + if args.worker_control_requests and not args.recover_input: + parser.error("host control requires resumable input and common lifecycle") if gaps and (not args.recover_input or any(seq >= args.frames for seq, _ in gaps)): parser.error("input gaps require recovery and must be inside the bounded source") raise SystemExit(run(args)) diff --git a/src/k1link/perception/streaming_lifecycle.py b/src/k1link/perception/streaming_lifecycle.py index 890e2c7..9c1c4ec 100644 --- a/src/k1link/perception/streaming_lifecycle.py +++ b/src/k1link/perception/streaming_lifecycle.py @@ -113,8 +113,9 @@ class StreamingLifecycle: if self.continuity is None: self.request_stop("worker-not-ready") raise - self.continuity.pause("worker-telemetry") - self.mailbox.pause() + if self.state == GraphState.RUNNING: + self.continuity.pause("worker-telemetry") + self.mailbox.pause() if not allow_unavailable: raise StreamSuspended(str(exc)) from exc except WorkerReadinessError: @@ -142,12 +143,29 @@ class StreamingLifecycle: if self.continuity is None: self.request_stop("worker-not-ready") raise - self.continuity.pause("worker-telemetry") - self.mailbox.pause() + if self.state == GraphState.RUNNING: + self.continuity.pause("worker-telemetry") + self.mailbox.pause() except WorkerReadinessError: self.request_stop("worker-not-ready") raise + def attach_readiness(self, monitor: WorkerReadinessMonitor) -> None: + """Trusted bootstrap after acquiring the local lease, BEFORE any child. + + Allows a real collector to report the acquired owner, not a fabricated + pre-lease claim. No GPU work may be spawned by this bootstrap adapter + before this call. Installation is one-shot and cannot weaken a monitor. + """ + with self._lock: + self._check(self.start, starting=True) + if self.state != GraphState.STARTING or self._children or self._readiness is not None: + raise WorkerReadinessError("readiness must be attached once before model startup") + if self.continuity is not None and not monitor.recoverable: + raise WorkerReadinessError("recoverable input requires recoverable readiness") + monitor.check(self.start, now_monotonic_ns=self._clock_ns(), require_warmup=False) + self._readiness = monitor + def renew(self, start: StreamStart) -> None: with self._lock: self._check(start, starting=True, allow_unavailable=True) diff --git a/src/k1link/perception/worker_control.py b/src/k1link/perception/worker_control.py new file mode 100644 index 0000000..ad92569 --- /dev/null +++ b/src/k1link/perception/worker_control.py @@ -0,0 +1,206 @@ +"""Bounded Worker-local host observation channel, separate from sensor ingress. + +The launcher gives the controller a private request directory and a SEPARATE +read-only response mount. Only the trusted host collector writes responses; +neither a Docker socket nor host commands are exposed to the AI container. +This filesystem boundary is not authentication for a network/GCS connection. + +One outstanding nonce requires a new host read after each request. Observation +time is the request's LOCAL monotonic start (a conservative lower bound), never +the reply's arrival or the foreign Windows clock. Delays cannot rejuvenate data. +""" + +from __future__ import annotations + +import hashlib +import json +import os +import secrets +import time +from collections.abc import Callable +from pathlib import Path +from typing import Any + +from .realtime_contract import RealtimeContractError, StreamStart, _digest, _integer +from .worker_operating_envelope import WorkerSnapshot + +CONTROL_SCHEMA = "missioncore.worker-host-observation/v1" +MAX_CONTROL_BYTES = 16384 +INVENTORY_SCOPE = "docker-gpu-access" +_FACT_FIELDS = { + "worker_id", + "container_id", + "image_sha256", + "cpu_limit_millicores", + "memory_limit_mib", + "gpu_name", + "driver_version", + "sm_clock_mhz", + "memory_clock_mhz", + "gpu_telemetry_available", + "inventory_complete", + "competing_gpu_clients", + "inventory_scope", + "host_process_inventory", +} + + +def activation_digest(start: StreamStart) -> str: + return hashlib.sha256( + json.dumps(start.to_dict(), sort_keys=True, separators=(",", ":")).encode() + ).hexdigest() + + +class WorkerControlChannel: + """Single-reader/single-writer adapter owned by the local controller thread. + + Paths and target container prefix come from the launcher, never a source + packet. Container identity is pinned on the first response and then exact. + Only 12..64 lowercase hex Docker IDs are accepted (not names). + """ + + def __init__( + self, + start: StreamStart, + requests: Path, + responses: Path, + *, + container_id: str, + clock_domain_id: str, + maximum_age_ms: int = 1000, + clock_ns: Callable[[], int] = time.monotonic_ns, + ) -> None: + if not 12 <= len(container_id) <= 64 or any( + c not in "0123456789abcdef" for c in container_id + ): + raise ValueError("launcher must select a Docker ID, not a name") + if requests.resolve() == responses.resolve(): + raise ValueError("request and read-only response mounts must be separate") + _integer(maximum_age_ms, "control maximum age", minimum=1) + self.start, self.activation = start, activation_digest(start) + self.requests, self.responses = requests, responses + self.container_id, self.clock_domain_id = container_id, clock_domain_id + self.clock_ns, self.maximum_age_ns = clock_ns, maximum_age_ms * 1_000_000 + self.pending: dict[str, Any] | None = None + self.requested_ns = 0 + self.sequence = self.accepted = self.expired = 0 + self.last_error: str | None = None + self.last_roundtrip_ms: float | None = None + + def request(self) -> None: + if self.pending is not None: + return + self.sequence += 1 + self.requested_ns = self.clock_ns() + self.pending = { + "schema_version": CONTROL_SCHEMA, + "activation_sha256": self.activation, + "nonce": secrets.token_hex(32), + "sequence": self.sequence, + } + # Two fixed files; no queue or per-sample history. Requests contain no + # executable command, secret, source data or caller-selected target. + path = self.requests / "request.json" + temporary = self.requests / "request.tmp" + temporary.write_text(json.dumps(self.pending), encoding="utf-8") + os.replace(temporary, path) + + def poll(self, *, owner: tuple[str, int], warmup_complete: bool) -> WorkerSnapshot | None: + if self.pending is None: + return None + now = self.clock_ns() + if now < self.requested_ns: + raise RealtimeContractError("controller clock moved backwards") + try: + with (self.responses / "response.json").open("rb") as stream: + raw = stream.read(MAX_CONTROL_BYTES + 1) + if len(raw) > MAX_CONTROL_BYTES: + raise ValueError("size") + reply = json.loads(raw.decode("utf-8-sig")) + if not isinstance(reply, dict) or set(reply) != set(self.pending) | {"facts"}: + raise ValueError("schema") + if any(reply[key] != value for key, value in self.pending.items()): + # Old/different activation cannot poison or refresh this owner. + raise ValueError("binding") + snapshot = self._snapshot(reply["facts"], owner, warmup_complete) + except (OSError, ValueError, TypeError, KeyError): + self.last_error = "response-unavailable-or-invalid" + if now - self.requested_ns > self.maximum_age_ns: + self.expired += 1 + self.pending = None + return None + self.pending = None + self.last_roundtrip_ms = (now - self.requested_ns) / 1_000_000 + if now - self.requested_ns > self.maximum_age_ns: + self.expired += 1 + self.last_error = "response-expired" + # Preserve known conflicts even in a delayed reply; the monitor + # separately rejects its OLD age. It must never become fresh. + else: + self.last_error = None + self.accepted += 1 + return snapshot + + def _snapshot(self, facts: object, owner: tuple[str, int], warmup: bool) -> WorkerSnapshot: + if not isinstance(facts, dict) or set(facts) != _FACT_FIELDS: + raise ValueError("facts schema") + target = facts["container_id"] + if target is not None: + _digest(target, "container ID") + if not target.startswith(self.container_id): + # Do not disguise a changed container as a missing observation. + raise ContainerIdentityConflict("host observed another container") + if ( + facts["inventory_scope"] != INVENTORY_SCOPE + or facts["host_process_inventory"] != "unproved" + ): + raise ValueError("unsupported inventory claim") + if type(facts["inventory_complete"]) is not bool: + raise ValueError("inventory completeness") + clients = facts["competing_gpu_clients"] + if not isinstance(clients, list) or len(clients) > 64: + raise ValueError("inventory bound") + # Partial inventory may prove a conflict, but cannot prove its absence. + competitors = tuple(clients) if clients or facts["inventory_complete"] else None + snapshot = WorkerSnapshot( + worker_id=facts["worker_id"], + clock_domain_id=self.clock_domain_id, + observed_monotonic_ns=self.requested_ns, + gpu_name=facts["gpu_name"], + driver_version=facts["driver_version"], + image_sha256=facts["image_sha256"] if target else None, + effective_config_sha256=self.start.effective_config_sha256, + cpu_limit_millicores=facts["cpu_limit_millicores"], + memory_limit_mib=facts["memory_limit_mib"], + sm_clock_mhz=facts["sm_clock_mhz"], + memory_clock_mhz=facts["memory_clock_mhz"], + gpu_owner_run_id=owner[0], + lease_generation=owner[1], + competing_gpu_clients=competitors, + warmup_complete=warmup, + inventory_scope=INVENTORY_SCOPE, + gpu_telemetry_available=facts["gpu_telemetry_available"], + ) + if target is not None: + self.container_id = target + return snapshot + + def snapshot(self) -> dict[str, object]: + return { + "schema_version": CONTROL_SCHEMA, + "requests": self.sequence, + "accepted": self.accepted, + "expired": self.expired, + "pending": self.pending is not None, + "last_error": self.last_error, + "last_roundtrip_ms": self.last_roundtrip_ms, + "container_id": self.container_id, + "inventory_scope": INVENTORY_SCOPE, + "host_process_inventory": "unproved", + "network_authenticated": False, + "realtime_qualified": False, + } + + +class ContainerIdentityConflict(RuntimeError): + """Trusted host returned a known different container; terminal, not a lag.""" diff --git a/src/k1link/perception/worker_control_pump.py b/src/k1link/perception/worker_control_pump.py new file mode 100644 index 0000000..4e97fa2 --- /dev/null +++ b/src/k1link/perception/worker_control_pump.py @@ -0,0 +1,107 @@ +"""Local controller adapter: host I/O never holds lifecycle locks or renews lease.""" + +from __future__ import annotations + +import threading +import time + +from .streaming_continuity import StreamSuspended +from .streaming_lifecycle import StreamingLifecycle +from .worker_control import ContainerIdentityConflict, WorkerControlChannel +from .worker_operating_envelope import WorkerOperatingEnvelope, WorkerSnapshot +from .worker_readiness import ReadinessMode, WorkerReadinessMonitor, WorkerTelemetryUnavailable + + +class WorkerControlPump: + def __init__( + self, + runtime: StreamingLifecycle, + channel: WorkerControlChannel, + envelope: WorkerOperatingEnvelope, + *, + mode: ReadinessMode, + bootstrap_seconds: float = 5.0, + ) -> None: + self.runtime, self.channel = runtime, channel + self.stop_event, self.warmed, self.warm_observed = ( + threading.Event(), + threading.Event(), + threading.Event(), + ) + self.io_failures = 0 + self.error: str | None = None + deadline = time.monotonic() + bootstrap_seconds + initial = None + while initial is None and time.monotonic() < deadline: + initial = self._poll() + if initial is None: + self.stop_event.wait(0.025) + if initial is None: + raise WorkerTelemetryUnavailable( + "host collector bootstrap timed out; no models started" + ) + runtime.attach_readiness( + WorkerReadinessMonitor( + runtime.start, + envelope, + initial, + mode=mode, + clock_domain_id=channel.clock_domain_id, + now_monotonic_ns=channel.clock_ns(), + recoverable=True, + ) + ) + self.thread = threading.Thread(target=self._run, name="worker-host-observer", daemon=True) + runtime.track_thread(self.thread) + self.thread.start() + + def _poll(self) -> WorkerSnapshot | None: + # Owner/config/warmup come from THIS controller, not a host JSON reply + # or a source heartbeat. The lease is checked again on observe_worker. + self.runtime.check_current(self.runtime.start, starting=True) + try: + self.channel.request() + return self.channel.poll( + owner=(self.runtime.start.run_id, self.runtime.start.lease_generation), + warmup_complete=self.warmed.is_set(), + ) + except OSError: + self.io_failures += 1 + return None # The unchanged observation expires independently. + + def _run(self) -> None: + try: + while not self.stop_event.wait(0.025) and not self.runtime.stop_event.is_set(): + observed = self._poll() + if observed is not None: + self.runtime.observe_worker(self.runtime.start, observed) + if observed.warmup_complete: + self.warm_observed.set() + self.stop_event.wait(0.2) + except Exception as exc: + self.error = type(exc).__name__ + self.runtime.request_stop( + "worker-container-conflict" + if isinstance(exc, ContainerIdentityConflict) + else "worker-control-failed" + ) + + def ready(self, seconds: float = 5.0) -> None: + self.warmed.set() + deadline = time.monotonic() + seconds + while time.monotonic() < deadline: + if self.warm_observed.wait(0.025): + try: + self.runtime.ready() + return + except StreamSuspended: + self.stop_event.wait(0.025) + raise WorkerTelemetryUnavailable("post-warmup host readiness not established") + + def close(self) -> bool: + self.stop_event.set() + self.thread.join(timeout=1) + return not self.thread.is_alive() + + def snapshot(self) -> dict[str, object]: + return {**self.channel.snapshot(), "io_failures": self.io_failures, "error": self.error} diff --git a/src/k1link/perception/worker_operating_envelope.py b/src/k1link/perception/worker_operating_envelope.py index 85bf459..2e17a59 100644 --- a/src/k1link/perception/worker_operating_envelope.py +++ b/src/k1link/perception/worker_operating_envelope.py @@ -31,9 +31,13 @@ class WorkerOperatingEnvelope: minimum_sm_clock_mhz: int minimum_memory_clock_mhz: int maximum_snapshot_age_ms: int = 1000 + # A Docker inventory is NOT an audit of native host/WSL GPU processes. + inventory_scope: str = "host-compute" def __post_init__(self) -> None: _identifier(self.envelope_id, "envelope_id") + if self.inventory_scope not in ("host-compute", "docker-gpu-access"): + raise RealtimeContractError("unsupported inventory scope") for field in ("gpu_name", "driver_version"): value = getattr(self, field) if not isinstance(value, str) or not value.strip() or len(value) > 160: @@ -66,9 +70,15 @@ class WorkerSnapshot: lease_generation: int | None competing_gpu_clients: tuple[str, ...] | None warmup_complete: bool | None + inventory_scope: str = "host-compute" + gpu_telemetry_available: bool = True def __post_init__(self) -> None: _identifier(self.worker_id, "worker_id") + if self.inventory_scope not in ("host-compute", "docker-gpu-access"): + raise RealtimeContractError("unsupported inventory scope") + if type(self.gpu_telemetry_available) is not bool: + raise RealtimeContractError("GPU telemetry availability must be boolean") _identifier(self.clock_domain_id, "clock_domain_id") _integer(self.observed_monotonic_ns, "observed_monotonic_ns") for field in ("image_sha256", "effective_config_sha256"): @@ -129,6 +139,10 @@ def operating_envelope_failures( if age < 0: raise RealtimeContractError("worker snapshot is from the future") failures = [] + if observed.inventory_scope != expected.inventory_scope: + failures.append("inventory-scope-mismatch") + if not observed.gpu_telemetry_available: + failures.append("gpu-telemetry-unavailable") if age > expected.maximum_snapshot_age_ms * 1_000_000: failures.append("worker-snapshot-expired") for field in ("worker_id", "image_sha256", "effective_config_sha256"): diff --git a/src/k1link/perception/worker_readiness.py b/src/k1link/perception/worker_readiness.py index 1dc1e64..08f8d27 100644 --- a/src/k1link/perception/worker_readiness.py +++ b/src/k1link/perception/worker_readiness.py @@ -107,7 +107,11 @@ class WorkerReadinessMonitor: if self.recoverable: # Missing metrics are not a proof that another process owns GPU. # Known identity/conflict/clock failures still fence immediately. - uncertain = {"worker-snapshot-expired", "warmup-not-complete"} + uncertain = { + "worker-snapshot-expired", + "warmup-not-complete", + "gpu-telemetry-unavailable", + } for field in ("image_sha256", "effective_config_sha256"): if getattr(self._observed, field) is None: uncertain.add(f"{field}-mismatch-or-unknown") diff --git a/tests/test_perception_worker_control.py b/tests/test_perception_worker_control.py new file mode 100644 index 0000000..deb4c05 --- /dev/null +++ b/tests/test_perception_worker_control.py @@ -0,0 +1,331 @@ +"""Small synthetic controller tests. No GPU, Docker, host commands or network.""" + +import json +import os +import threading +import time +from dataclasses import replace + +import pytest + +from k1link.perception.realtime_contract import StreamStart +from k1link.perception.streaming_continuity import StreamSuspended +from k1link.perception.streaming_lifecycle import StreamingLifecycle +from k1link.perception.streaming_queue import StreamMailbox +from k1link.perception.worker_control import ContainerIdentityConflict, WorkerControlChannel +from k1link.perception.worker_control_pump import WorkerControlPump +from k1link.perception.worker_operating_envelope import WorkerOperatingEnvelope +from k1link.perception.worker_readiness import ( + WorkerReadinessError, + WorkerReadinessMonitor, + WorkerTelemetryUnavailable, +) + + +@pytest.fixture +def control(tmp_path): + start = StreamStart( + "run", + "source", + "worker", + "epoch", + 1, + *[c * 64 for c in "abcd"], + "source-clock", + "recorded-source-paced", + ) + requests, responses = tmp_path / "requests", tmp_path / "responses" + requests.mkdir() + responses.mkdir() + now = [1_000_000_000] + channel = WorkerControlChannel( + start, + requests, + responses, + container_id="e" * 12, + clock_domain_id="local-clock", + clock_ns=lambda: now[0], + ) + facts = { + "worker_id": "worker", + "container_id": "e" * 64, + "image_sha256": "b" * 64, + "cpu_limit_millicores": 8000, + "memory_limit_mib": 8192, + "gpu_name": "RTX4090", + "driver_version": "610.47", + "sm_clock_mhz": 2610, + "memory_clock_mhz": 10251, + "gpu_telemetry_available": True, + "inventory_complete": True, + "competing_gpu_clients": [], + "inventory_scope": "docker-gpu-access", + "host_process_inventory": "unproved", + } + + def reply(**changes): + channel.request() + payload = {**channel.pending, "facts": {**facts, **changes}} + (responses / "response.json").write_text(json.dumps(payload)) + return payload + + return channel, now, reply + + +def poll(channel): + return channel.poll(owner=("run", 1), warmup_complete=True) + + +def monitor(channel, initial, now): + return WorkerReadinessMonitor( + channel.start, + WorkerOperatingEnvelope( + "test/v1", + "RTX4090", + "610.47", + 8000, + 8192, + 2610, + 10251, + inventory_scope="docker-gpu-access", + ), + initial, + mode="labelled-experiment", + clock_domain_id="local-clock", + now_monotonic_ns=now[0], + recoverable=True, + ) + + +def test_host_reply_uses_request_time_and_local_authority_not_receipt_time(control): + channel, now, reply = control + payload = reply() + now[0] += 750_000_000 + observed = poll(channel) + assert observed.observed_monotonic_ns == 1_000_000_000 + assert observed.gpu_owner_run_id == "run" and observed.lease_generation == 1 + assert observed.effective_config_sha256 == channel.start.effective_config_sha256 + assert observed.competing_gpu_clients == () + assert channel.container_id == "e" * 64 + assert channel.snapshot()["last_roundtrip_ms"] == 750 + assert "observed_monotonic_ns" not in payload + assert poll(channel) is None # A response is consumable exactly once. + + +@pytest.mark.parametrize("change", ["nonce", "activation_sha256", "sequence", "extra", "oversize"]) +def test_stale_cross_activation_malformed_response_cannot_refresh_owner(control, change): + channel, now, reply = control + payload = reply() + if change == "extra": + payload["command"] = "must never run" + elif change == "sequence": + payload[change] += 1 + else: + payload[change] = "f" * (17000 if change == "oversize" else 64) + (channel.responses / "response.json").write_text(json.dumps(payload)) + assert poll(channel) is None + now[0] += 1_000_000_001 + assert poll(channel) is None and channel.expired == 1 + channel.request() + assert channel.sequence == 2 and channel.accepted == 0 + + +def test_delayed_reply_stays_expired_and_cannot_hide_a_known_conflict(control): + channel, now, reply = control + reply() + readiness = monitor(channel, poll(channel), now) + now[0] += 1 + reply() + now[0] += 1_200_000_000 + with pytest.raises(WorkerTelemetryUnavailable): + readiness.observe( + channel.start, poll(channel), now_monotonic_ns=now[0], require_warmup=True + ) + assert channel.expired == 1 + now[0] += 1 + reply(inventory_complete=False, competing_gpu_clients=["docker:" + "f" * 64]) + with pytest.raises(WorkerReadinessError, match="competing-gpu"): + readiness.observe( + channel.start, poll(channel), now_monotonic_ns=now[0], require_warmup=True + ) + + +@pytest.mark.parametrize( + "field,value", + [ + ("inventory_complete", False), + ("gpu_telemetry_available", False), + ("container_id", None), + ], +) +def test_missing_real_facts_cannot_be_waived_by_labelled_experiment(control, field, value): + channel, now, reply = control + reply() + readiness = monitor(channel, poll(channel), now) + now[0] += 1 + reply(**{field: value}) + with pytest.raises(WorkerTelemetryUnavailable): + readiness.observe( + channel.start, poll(channel), now_monotonic_ns=now[0], require_warmup=True + ) + assert not readiness.snapshot()["terminal_failures"] + now[0] += 1 + reply() + readiness.observe(channel.start, poll(channel), now_monotonic_ns=now[0], require_warmup=True) + + +def test_container_id_is_pinned_and_changed_target_is_terminal(control): + channel, _, reply = control + reply() + poll(channel) + reply(container_id="e" * 12 + "f" * 52) + with pytest.raises(ContainerIdentityConflict): + poll(channel) + + +def test_docker_scope_cannot_satisfy_host_wide_inventory_requirement(control): + channel, now, reply = control + reply() + observed = poll(channel) + with pytest.raises(WorkerReadinessError, match="inventory-scope-mismatch"): + monitor(channel, replace(observed, inventory_scope="host-compute"), now) + + +def test_real_observation_bridge_waits_resumes_but_does_not_renew_lease(control, tmp_path): + channel, now, reply = control + run = StreamingLifecycle( + channel.start, + tmp_path / "lease", + StreamMailbox(), + threading.Event(), + clock_ns=lambda: now[0], + recover_input=True, + source_clock_ns=lambda: now[0], + ) + try: + reply() + readiness = monitor(channel, poll(channel), now) + run.attach_readiness(readiness) + with pytest.raises(WorkerReadinessError, match="once"): + run.attach_readiness(readiness) + run.ready() + now[0] += 1_000_000_001 + run.renew(channel.start) + deadline = run.lease.deadline_ns + with pytest.raises(StreamSuspended): + run.check_input(channel.start) + assert not run.stop_event.is_set() + reply() + run.observe_worker(channel.start, poll(channel)) + assert run.lease.deadline_ns == deadline + assert run.continuity.phase == "waiting" # Fresh metrics aren't fresh sensor evidence. + epoch = run.begin_input(channel.start) + assert epoch.epoch_id != channel.start.epoch_id + finally: + assert run.close() + + +def test_telemetry_lag_during_warmup_blocks_spawn_without_poisoning_first_input(control, tmp_path): + channel, now, reply = control + run = StreamingLifecycle( + channel.start, + tmp_path / "lease", + StreamMailbox(), + threading.Event(), + clock_ns=lambda: now[0], + recover_input=True, + source_clock_ns=lambda: now[0], + ) + spawned = [] + try: + reply() + run.attach_readiness(monitor(channel, poll(channel), now)) + now[0] += 1_000_000_001 + run.renew(channel.start) + with pytest.raises(StreamSuspended): + run.spawn(lambda: spawned.append(True)) + assert not spawned and not run.stop_event.is_set() + assert run.continuity.phase == "active" # No input has started yet. + reply() + run.observe_worker(channel.start, poll(channel)) + run.ready() + run.check_input(channel.start) + finally: + assert run.close() + + +def test_controller_pump_bootstrap_refresh_ready_and_shutdown_are_independent_of_lease( + control, tmp_path +): + channel, _, reply = control + facts = reply()["facts"] + channel.pending = None + channel.clock_ns = time.monotonic_ns + stopped = threading.Event() + + def host_fixture(): + previous = None + while not stopped.wait(0.005): + request = json.loads((channel.requests / "request.json").read_text()) + if request["nonce"] == previous: + continue + path = channel.responses / "fixture.tmp" + path.write_text(json.dumps({**request, "facts": facts})) + os.replace(path, channel.responses / "response.json") + previous = request["nonce"] + + host = threading.Thread(target=host_fixture) + host.start() + run = StreamingLifecycle( + channel.start, + tmp_path / "lease", + StreamMailbox(), + threading.Event(), + recover_input=True, + source_clock_ns=time.monotonic_ns, + ) + pump = None + try: + pump = WorkerControlPump( + run, + channel, + WorkerOperatingEnvelope( + "test/v1", + "RTX4090", + "610.47", + 8000, + 8192, + 2610, + 10251, + inventory_scope="docker-gpu-access", + ), + mode="labelled-experiment", + ) + pump.ready(seconds=1) + run.check_input(channel.start) + assert channel.accepted >= 2 and run.lease.renewals == 0 + assert pump.snapshot()["host_process_inventory"] == "unproved" + finally: + if pump: + assert pump.close() + stopped.set() + host.join(timeout=1) + assert not host.is_alive() and run.close() + + +@pytest.mark.parametrize( + "changes", + [ + {"inventory_complete": "true"}, + {"competing_gpu_clients": ["x"] * 65}, + {"competing_gpu_clients": "empty"}, + {"memory_limit_mib": True}, + {"host_process_inventory": "complete"}, + {"gpu_telemetry_available": "yes"}, + ], +) +def test_schema_is_bounded_and_does_not_accept_false_host_claims(control, changes): + channel, _, reply = control + reply(**changes) + assert poll(channel) is None + assert channel.accepted == 0