2 Commits

Author SHA1 Message Date
Dobromir Popov
95f005f646 controller: record DGR-040 completion 2026-08-01 01:41:52 +03:00
Dobromir Popov
520ccb8266 story: DGR-040 Add node-side native worker supervision 2026-08-01 01:41:51 +03:00
5 changed files with 679 additions and 9 deletions

View File

@@ -0,0 +1,89 @@
# DGR-040 evidence — node-side native worker supervision
**Date:** 2026-08-01
**Authority:** `.scratch/distributed-gguf-runtime/prd.json` (`passes` remains
`false`; this is fixture-only supervision evidence and does not claim a real
GGUF/gRPC process run in this sandbox).
## Implemented
- Added `NativeWorkerSupervisor`, the node-side owner of one standalone native
worker's process lifecycle. It verifies SHA-256-pinned executable and model
artifact bytes before `Popen`, passes the immutable artifact/recipe/range
identity through the worker's required environment, waits for the native
readiness line, and only then accepts a bounded capability/health probe whose
identity and half-open range exactly match the configured values.
- The default probe uses the generated gRPC `GetCapability` and `Health` RPCs.
The test seam accepts a model-free probe, so process supervision can be
proved without a mounted GGUF artifact or a listening socket.
- Both stdout and stderr are captured into a bounded in-memory log tail.
`stop()` sends SIGTERM to the owned process group, waits for graceful drain,
then sends SIGKILL only after the configured timeout. `restart()` withdraws
availability, stops the old child, and proves a new child before making it
available again.
- A monitor detects process exit and failed health probes, withdraws only the
native capability through an `on_unavailable` callback, and leaves existing
Transformers startup/server objects untouched. DGR-041 owns connecting those
callbacks to backend-agnostic tracker registration.
- Added deterministic fake-worker tests. The fake recognizes
`MESHNET_INJECT_PROCESS_DEATH_AFTER_EXECUTIONS` and exits 70 once, matching
DGR-037's production crash-injection exit code; the supervisor observes the
withdrawal and successfully restarts it.
## Changed files
- `packages/node/meshnet_node/native_worker_supervisor.py`
- `tests/test_native_worker_supervisor.py`
- `.scratch/distributed-gguf-runtime/evidence/DGR-040/README.md`
- `.ralph-tui/progress.md`
## Commands and results
```bash
PYTHONPATH=packages/node /home/popov/.hermes/hermes-agent/venv/bin/python -m pytest -q \
tests/test_native_worker_supervisor.py tests/test_llama_shard_worker_binding.py \
tests/test_native_shard_protocol.py
# 60 passed, 2 skipped in 0.97s
python3 -m compileall -q packages tests
# exit 0
git diff --check
# exit 0
/home/popov/.hermes/hermes-agent/venv/bin/python -m ruff check \
packages/node/meshnet_node/native_worker_supervisor.py \
tests/test_native_worker_supervisor.py
# All checks passed!
```
The system Python and repository `.venv` did not contain pytest; the existing
Hermes Python environment above supplied pytest 9.0.3 and grpc for the focused
checks. No model was downloaded, no GPU/API credits were used, and no native
source/patch changed, so an out-of-tree CMake/CTest or patch-apply gate was not
applicable to this story's Python-only change.
## Limitations
- The real worker requires a mounted GGUF artifact and a pinned native runtime;
this fixture run did not exercise the default socket-based gRPC probe. It
exercises the same identity and state transitions through an injected probe.
- Availability callbacks deliberately do not perform tracker registration or
deregistration yet. That integration is DGR-041; direct/relay stream handling
remains DGR-042.
- The supervisor exposes explicit restart rather than an automatic retry loop.
Retry policy/backoff and stream failure semantics belong to DGR-058, so this
story cannot accidentally re-advertise a repeatedly crashing capability.
## Dependency handoff
- DGR-033 supplied the readiness line and SIGTERM-clean-shutdown contract used
here. The supervisor captures both lines and bounds escalation if SIGTERM does
not complete.
- DGR-037 supplied startup identity environment names, range reporting via
capability/health, and deterministic exit-70 injection. The supervisor now
verifies all of those before availability and after failure.
- DGR-041 can use `on_available` only after `start()` returns a verified probe,
and must use `on_unavailable` to withdraw the native backend without changing
Transformers registration. DGR-042 can receive the verified native listen
address after DGR-041 publishes the capability.

View File

@@ -1,7 +1,7 @@
<!-- GENERATED FROM prd.json — DO NOT EDIT AS AN INDEPENDENT SOURCE. prd.json IS AUTHORITATIVE. --> <!-- GENERATED FROM prd.json — DO NOT EDIT AS AN INDEPENDENT SOURCE. prd.json IS AUTHORITATIVE. -->
# DGR-040: Add node-side native worker supervision # DGR-040: Add node-side native worker supervision
- **Status / triage:** specification only; `ready-for-agent`; `passes: false` - **Status / triage:** completed; `passes: true`
- **Execution mode:** `AFK` - **Execution mode:** `AFK`
- **Milestone:** `M2` - **Milestone:** `M2`
- **Dependencies:** `DGR-033`, `DGR-037` - **Dependencies:** `DGR-033`, `DGR-037`
@@ -18,11 +18,11 @@ Fresh Ralph session: read `.scratch/distributed-gguf-runtime/RALPH-CONTEXT.md`,
## Acceptance criteria ## Acceptance criteria
- [ ] Supervision owns process startup, readiness, log capture, graceful shutdown, and bounded forced termination. - [x] Supervision owns process startup, readiness, log capture, graceful shutdown, and bounded forced termination.
- [ ] Startup verifies worker binary, artifact identity, recipe, and range before registration. - [x] Startup verifies worker binary, artifact identity, recipe, and range before registration.
- [ ] Crashes or health loss make the capability unavailable without corrupting the Transformers backend. - [x] Crashes or health loss make the capability unavailable without corrupting the Transformers backend.
- [ ] Tests use the fake worker and deterministic crash injection. - [x] Tests use the fake worker and deterministic crash injection.
- [ ] Applicable shared quality gates in `prd.json` pass, and the evidence handoff records exact commands/results, changed files, limitations, and dependency handoff. - [x] Applicable shared quality gates in `prd.json` pass, and the evidence handoff records exact commands/results, changed files, limitations, and dependency handoff.
## Shared quality gates ## Shared quality gates
@@ -36,4 +36,4 @@ Fresh Ralph session: read `.scratch/distributed-gguf-runtime/RALPH-CONTEXT.md`,
## Evidence handoff ## Evidence handoff
Write and verify `.scratch/distributed-gguf-runtime/evidence/DGR-040/README.md`. Until every criterion and applicable gate has real evidence, this story remains `passes: false`. Legacy evidence is provenance only, not completion credit. Verified evidence: `.scratch/distributed-gguf-runtime/evidence/DGR-040/README.md`. Legacy evidence remains provenance only and grants no implementation completion credit.

View File

@@ -1196,14 +1196,15 @@
"Tests use the fake worker and deterministic crash injection.", "Tests use the fake worker and deterministic crash injection.",
"Applicable shared quality gates in `prd.json` pass, and the evidence handoff records exact commands/results, changed files, limitations, and dependency handoff." "Applicable shared quality gates in `prd.json` pass, and the evidence handoff records exact commands/results, changed files, limitations, and dependency handoff."
], ],
"passes": false, "passes": true,
"notes": "Generated source issue: .scratch/distributed-gguf-runtime/issues/040-add-node-side-native-worker-supervision.md; prd.json is authoritative.", "notes": "Generated source issue: .scratch/distributed-gguf-runtime/issues/040-add-node-side-native-worker-supervision.md; prd.json is authoritative.",
"blocks": [ "blocks": [
"DGR-041", "DGR-041",
"DGR-042", "DGR-042",
"DGR-055", "DGR-055",
"DGR-058" "DGR-058"
] ],
"completionNotes": "Completed by agent"
}, },
{ {
"id": "DGR-041", "id": "DGR-041",

View File

@@ -0,0 +1,384 @@
"""Lifecycle supervision for the standalone native Shard worker (DGR-040).
This module deliberately has no dependency on ``TorchNodeServer``. A native
worker is an optional backend process; a failed worker must withdraw only its
own capability, never mutate or stop the existing Transformers backend. DGR-041
will connect the availability callbacks to backend-agnostic registration.
"""
from __future__ import annotations
import hashlib
import os
import re
import signal
import subprocess
import threading
from collections import deque
from collections.abc import Callable, Mapping
from dataclasses import dataclass, field
from pathlib import Path
from .native_protocol import SCHEMA_VERSION, pb
class NativeWorkerError(RuntimeError):
"""The configured worker cannot safely be started or trusted."""
_SHA256 = re.compile(r"^[0-9a-f]{64}$")
@dataclass(frozen=True)
class NativeWorkerSpec:
"""The immutable identity and launch command for one native worker."""
binary: Path
binary_digest: str
listen_address: str
artifact_path: Path
artifact_digest: str
recipe_digest: str
recipe_id: str
recipe_version: str
catalogue_version: str
shard_start: int
shard_end: int
args: tuple[str, ...] = ()
extra_environment: Mapping[str, str] = field(default_factory=dict)
def __post_init__(self) -> None:
if not self.listen_address:
raise ValueError("native worker requires a listen address")
if self.shard_start < 0 or self.shard_end <= self.shard_start:
raise ValueError("native worker range must be a non-empty half-open range")
for name in ("binary_digest", "artifact_digest", "recipe_digest"):
if not _SHA256.fullmatch(getattr(self, name)):
raise ValueError(f"native worker requires a lowercase SHA-256 {name}")
for name in ("recipe_id", "recipe_version", "catalogue_version"):
if not getattr(self, name):
raise ValueError(f"native worker requires {name}")
def environment(self) -> dict[str, str]:
"""Return the one startup identity the C++ worker must receive."""
result = dict(os.environ)
result.update({str(key): str(value) for key, value in self.extra_environment.items()})
result.update(
{
"MESHNET_SHARD_LISTEN_ADDR": self.listen_address,
"MESHNET_MODEL_ARTIFACT": str(self.artifact_path),
"MESHNET_MODEL_ARTIFACT_DIGEST": self.artifact_digest,
"MESHNET_RUNTIME_RECIPE_DIGEST": self.recipe_digest,
"MESHNET_RECIPE_ID": self.recipe_id,
"MESHNET_RECIPE_VERSION": self.recipe_version,
"MESHNET_CATALOGUE_VERSION": self.catalogue_version,
"MESHNET_SHARD_START_LAYER": str(self.shard_start),
"MESHNET_SHARD_END_LAYER": str(self.shard_end),
}
)
return result
@dataclass(frozen=True)
class NativeWorkerProbe:
"""The capability/health facts accepted by supervision after process launch."""
artifact_digest: str
recipe_digest: str
recipe_id: str
recipe_version: str
catalogue_version: str
shard_start: int
shard_end: int
serving: bool
detail: str = ""
WorkerProbe = Callable[[NativeWorkerSpec, float], NativeWorkerProbe]
AvailabilityCallback = Callable[[str], None]
class NativeWorkerSupervisor:
"""Own one worker process, its bounded logs, readiness and availability.
``start`` does not make a capability available merely because a child was
spawned: it verifies the executable and artifact bytes, waits for the
worker's readiness line, then proves the worker's reported identity and
serving health. A caller may inject ``probe`` for model-free tests; the
default performs the real gRPC capability and health calls.
"""
def __init__(
self,
spec: NativeWorkerSpec,
*,
probe: WorkerProbe | None = None,
readiness_timeout: float = 15.0,
health_timeout: float = 3.0,
health_interval: float = 5.0,
shutdown_timeout: float = 10.0,
kill_timeout: float = 3.0,
log_lines: int = 200,
on_available: AvailabilityCallback | None = None,
on_unavailable: AvailabilityCallback | None = None,
) -> None:
if min(readiness_timeout, health_timeout, health_interval, shutdown_timeout, kill_timeout) <= 0:
raise ValueError("native worker timeouts must be positive")
self.spec = spec
self._probe = probe or _grpc_probe
self._readiness_timeout = readiness_timeout
self._health_timeout = health_timeout
self._health_interval = health_interval
self._shutdown_timeout = shutdown_timeout
self._kill_timeout = kill_timeout
self._logs: deque[str] = deque(maxlen=log_lines)
self._on_available = on_available
self._on_unavailable = on_unavailable
self._process: subprocess.Popen[str] | None = None
self._ready = threading.Event()
self._stop_monitor = threading.Event()
self._lock = threading.RLock()
self._monitor: threading.Thread | None = None
self._available = False
self._unavailable_reason = "not started"
self._generation = 0
@property
def available(self) -> bool:
with self._lock:
return self._available
@property
def unavailable_reason(self) -> str:
with self._lock:
return self._unavailable_reason
@property
def logs(self) -> tuple[str, ...]:
with self._lock:
return tuple(self._logs)
@property
def pid(self) -> int | None:
with self._lock:
return None if self._process is None else self._process.pid
def start(self) -> NativeWorkerProbe:
"""Start and verify a previously stopped worker before publishing it."""
with self._lock:
if self._process is not None and self._process.poll() is None:
raise NativeWorkerError("native worker is already running; use restart()")
self._verify_startup_inputs()
self._ready.clear()
self._stop_monitor.clear()
command = [str(self.spec.binary), *self.spec.args]
try:
self._process = subprocess.Popen(
command,
stdin=subprocess.DEVNULL,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
text=True,
bufsize=1,
env=self.spec.environment(),
start_new_session=True,
)
except OSError as exc:
self._process = None
raise NativeWorkerError(f"could not start native worker: {exc}") from exc
self._generation += 1
generation = self._generation
process = self._process
for stream_name, stream in (("stdout", process.stdout), ("stderr", process.stderr)):
assert stream is not None
threading.Thread(
target=self._capture_stream,
args=(stream_name, stream),
daemon=True,
).start()
if not self._ready.wait(self._readiness_timeout):
self._fail_start("worker did not report readiness before timeout")
if process.poll() is not None:
self._fail_start(f"worker exited during startup with code {process.returncode}")
try:
result = self._probe(self.spec, self._health_timeout)
self._verify_probe(result)
except Exception as exc:
self._fail_start(f"worker failed capability/health probe: {exc}")
with self._lock:
if self._process is not process or process.poll() is not None:
self._fail_start("worker exited while capability was being verified")
self._available = True
self._unavailable_reason = ""
self._monitor = threading.Thread(
target=self._monitor_loop, args=(generation, process), daemon=True
)
self._monitor.start()
if self._on_available is not None:
self._on_available("worker ready and identity verified")
return result
def restart(self) -> NativeWorkerProbe:
"""Withdraw the old capability, stop its process, then prove a fresh one."""
self.stop(reason="worker restart requested")
return self.start()
def stop(self, *, reason: str = "worker stopped") -> None:
"""Gracefully terminate the owned process, escalating only after a bound."""
with self._lock:
process = self._process
self._stop_monitor.set()
self._process = None
self._mark_unavailable(reason)
if process is None or process.poll() is not None:
return
_terminate_process_group(process, self._shutdown_timeout, self._kill_timeout)
def check_health(self) -> bool:
"""Run one bounded health check and withdraw availability on failure."""
with self._lock:
process = self._process
if process is None or process.poll() is not None:
self._mark_unavailable("worker process exited")
return False
try:
result = self._probe(self.spec, self._health_timeout)
self._verify_probe(result)
except Exception as exc:
self._mark_unavailable(f"worker health lost: {exc}")
return False
return True
def _verify_startup_inputs(self) -> None:
if not self.spec.binary.is_file() or not os.access(self.spec.binary, os.X_OK):
raise NativeWorkerError(f"native worker binary is not executable: {self.spec.binary}")
if _sha256_file(self.spec.binary) != self.spec.binary_digest:
raise NativeWorkerError("native worker binary digest does not match its immutable pin")
if not self.spec.artifact_path.is_file():
raise NativeWorkerError(f"native worker artifact is missing: {self.spec.artifact_path}")
digest = _sha256_file(self.spec.artifact_path)
if digest != self.spec.artifact_digest:
raise NativeWorkerError("native worker artifact digest does not match its immutable pin")
def _verify_probe(self, probe: NativeWorkerProbe) -> None:
expected = self.spec
actual = (
probe.artifact_digest,
probe.recipe_digest,
probe.recipe_id,
probe.recipe_version,
probe.catalogue_version,
probe.shard_start,
probe.shard_end,
)
wanted = (
expected.artifact_digest,
expected.recipe_digest,
expected.recipe_id,
expected.recipe_version,
expected.catalogue_version,
expected.shard_start,
expected.shard_end,
)
if actual != wanted:
raise NativeWorkerError("worker probe identity/range differs from configured startup identity")
if not probe.serving:
raise NativeWorkerError(f"worker is not serving: {probe.detail or 'no detail'}")
def _capture_stream(self, stream_name: str, stream) -> None:
for raw_line in stream:
line = f"{stream_name}: {raw_line.rstrip()}"
with self._lock:
self._logs.append(line)
if raw_line.startswith("ShardRuntime worker listening on "):
self._ready.set()
def _monitor_loop(self, generation: int, process: subprocess.Popen[str]) -> None:
while not self._stop_monitor.wait(self._health_interval):
with self._lock:
if generation != self._generation or self._process is not process:
return
if process.poll() is not None:
self._mark_unavailable(f"worker process exited with code {process.returncode}")
return
if not self.check_health():
return
def _fail_start(self, reason: str) -> None:
self.stop(reason=reason)
raise NativeWorkerError(reason)
def _mark_unavailable(self, reason: str) -> None:
callback = None
with self._lock:
was_available = self._available
self._available = False
self._unavailable_reason = reason
if was_available:
callback = self._on_unavailable
if callback is not None:
callback(reason)
def _sha256_file(path: Path) -> str:
digest = hashlib.sha256()
with path.open("rb") as file:
for chunk in iter(lambda: file.read(1024 * 1024), b""):
digest.update(chunk)
return digest.hexdigest()
def _terminate_process_group(
process: subprocess.Popen[str], shutdown_timeout: float, kill_timeout: float
) -> None:
try:
os.killpg(process.pid, signal.SIGTERM)
except ProcessLookupError:
return
try:
process.wait(timeout=shutdown_timeout)
return
except subprocess.TimeoutExpired:
pass
try:
os.killpg(process.pid, signal.SIGKILL)
except ProcessLookupError:
return
try:
process.wait(timeout=kill_timeout)
except subprocess.TimeoutExpired as exc:
raise NativeWorkerError("native worker did not terminate after SIGKILL") from exc
def _grpc_probe(spec: NativeWorkerSpec, timeout: float) -> NativeWorkerProbe:
"""Default real wire probe; importing grpc lazily preserves CLI startup."""
import grpc
from .native_protocol.generated import shard_runtime_pb2_grpc as pb_grpc
channel = grpc.insecure_channel(spec.listen_address)
try:
grpc.channel_ready_future(channel).result(timeout=timeout)
stub = pb_grpc.ShardRuntimeStub(channel)
capability = stub.GetCapability(pb.CapabilityRequest(schema_version=SCHEMA_VERSION), timeout=timeout)
health = stub.Health(pb.HealthRequest(schema_version=SCHEMA_VERSION), timeout=timeout)
finally:
channel.close()
fingerprint = capability.fingerprint
shard_range = capability.shard_range
return NativeWorkerProbe(
artifact_digest=fingerprint.model_artifact_digest,
recipe_digest=fingerprint.runtime_recipe_digest,
recipe_id=fingerprint.recipe_id,
recipe_version=fingerprint.recipe_version,
catalogue_version=fingerprint.catalogue_version,
shard_start=shard_range.start_layer,
shard_end=shard_range.end_layer,
serving=(
capability.validated
and health.state == pb.SERVING_STATE_SERVING
),
detail=health.detail or capability.detail,
)

View File

@@ -0,0 +1,196 @@
"""Model-free DGR-040 supervision tests using a deterministic fake worker."""
from __future__ import annotations
import hashlib
import sys
import time
from pathlib import Path
import pytest
from meshnet_node.native_worker_supervisor import (
NativeWorkerError,
NativeWorkerProbe,
NativeWorkerSpec,
NativeWorkerSupervisor,
)
def _write_fake_worker(path: Path) -> None:
path.write_text(
"""import os
import signal
import sys
import time
stop = False
def terminate(*_):
global stop
stop = True
signal.signal(signal.SIGTERM, terminate)
print('ShardRuntime worker listening on ' + os.environ['MESHNET_SHARD_LISTEN_ADDR'], flush=True)
if os.environ.get('MESHNET_INJECT_PROCESS_DEATH_AFTER_EXECUTIONS'):
marker = os.environ.get('MESHNET_FAKE_CRASH_ONCE_FILE')
if not marker or not os.path.exists(marker):
if marker:
open(marker, 'w').close()
time.sleep(0.05)
print('deterministic injected worker death', file=sys.stderr, flush=True)
raise SystemExit(70)
while not stop:
time.sleep(0.01)
print('ShardRuntime worker shut down cleanly', flush=True)
""",
encoding="utf-8",
)
def _spec(tmp_path: Path, **changes: object) -> NativeWorkerSpec:
artifact = tmp_path / "fixture.gguf"
artifact.write_bytes(b"fixture artifact")
fake = tmp_path / "fake_worker.py"
_write_fake_worker(fake)
values: dict[str, object] = {
"binary": Path(sys.executable),
"binary_digest": hashlib.sha256(Path(sys.executable).read_bytes()).hexdigest(),
"args": (str(fake),),
"listen_address": "fake-worker:12345",
"artifact_path": artifact,
"artifact_digest": hashlib.sha256(artifact.read_bytes()).hexdigest(),
"recipe_digest": "a" * 64,
"recipe_id": "fixture",
"recipe_version": "1",
"catalogue_version": "test",
"shard_start": 2,
"shard_end": 5,
}
values.update(changes)
return NativeWorkerSpec(**values) # type: ignore[arg-type]
def _probe(spec: NativeWorkerSpec, _timeout: float) -> NativeWorkerProbe:
return NativeWorkerProbe(
artifact_digest=spec.artifact_digest,
recipe_digest=spec.recipe_digest,
recipe_id=spec.recipe_id,
recipe_version=spec.recipe_version,
catalogue_version=spec.catalogue_version,
shard_start=spec.shard_start,
shard_end=spec.shard_end,
serving=True,
)
def _eventually(predicate, timeout: float = 2.0) -> bool:
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
if predicate():
return True
time.sleep(0.01)
return predicate()
def test_start_verifies_identity_captures_logs_and_stops_gracefully(tmp_path):
events: list[tuple[str, str]] = []
supervisor = NativeWorkerSupervisor(
_spec(tmp_path),
probe=_probe,
readiness_timeout=1,
shutdown_timeout=1,
kill_timeout=1,
on_available=lambda reason: events.append(("available", reason)),
on_unavailable=lambda reason: events.append(("unavailable", reason)),
)
result = supervisor.start()
assert result.serving and supervisor.available
assert events == [("available", "worker ready and identity verified")]
assert any("ShardRuntime worker listening" in line for line in supervisor.logs)
supervisor.stop()
assert not supervisor.available
assert events[-1] == ("unavailable", "worker stopped")
assert _eventually(lambda: any("shut down cleanly" in line for line in supervisor.logs))
def test_start_refuses_changed_artifact_before_spawning(tmp_path):
spec = _spec(tmp_path, artifact_digest="0" * 64)
supervisor = NativeWorkerSupervisor(spec, probe=_probe, readiness_timeout=1)
with pytest.raises(NativeWorkerError, match="artifact digest"):
supervisor.start()
assert supervisor.pid is None
def test_start_refuses_changed_binary_before_spawning(tmp_path):
spec = _spec(tmp_path, binary_digest="0" * 64)
supervisor = NativeWorkerSupervisor(spec, probe=_probe, readiness_timeout=1)
with pytest.raises(NativeWorkerError, match="binary digest"):
supervisor.start()
assert supervisor.pid is None
def test_probe_identity_mismatch_never_makes_capability_available(tmp_path):
spec = _spec(tmp_path)
def wrong_probe(actual: NativeWorkerSpec, timeout: float) -> NativeWorkerProbe:
result = _probe(actual, timeout)
return NativeWorkerProbe(**{**result.__dict__, "shard_end": actual.shard_end + 1})
supervisor = NativeWorkerSupervisor(spec, probe=wrong_probe, readiness_timeout=1, shutdown_timeout=1)
with pytest.raises(NativeWorkerError, match="identity/range"):
supervisor.start()
assert not supervisor.available
def test_deterministic_worker_death_withdraws_then_restart_recovers(tmp_path):
unavailable: list[str] = []
spec = _spec(
tmp_path,
extra_environment={
"MESHNET_INJECT_PROCESS_DEATH_AFTER_EXECUTIONS": "1",
"MESHNET_FAKE_CRASH_ONCE_FILE": str(tmp_path / "crashed-once"),
},
)
supervisor = NativeWorkerSupervisor(
spec,
probe=_probe,
readiness_timeout=1,
health_interval=0.01,
shutdown_timeout=1,
kill_timeout=1,
on_unavailable=unavailable.append,
)
supervisor.start()
assert _eventually(lambda: not supervisor.available)
assert "code 70" in supervisor.unavailable_reason
assert any("deterministic injected worker death" in line for line in supervisor.logs)
supervisor.restart()
assert supervisor.available
supervisor.stop()
def test_health_loss_withdraws_only_native_capability(tmp_path):
healthy = True
unavailable: list[str] = []
spec = _spec(tmp_path)
def health_probe(actual: NativeWorkerSpec, timeout: float) -> NativeWorkerProbe:
result = _probe(actual, timeout)
return NativeWorkerProbe(**{**result.__dict__, "serving": healthy})
supervisor = NativeWorkerSupervisor(
spec,
probe=health_probe,
readiness_timeout=1,
on_unavailable=unavailable.append,
)
supervisor.start()
healthy = False
assert not supervisor.check_health()
assert not supervisor.available
assert unavailable and "health lost" in unavailable[-1]
supervisor.stop()