diff --git a/solstone/think/providers/runtime_health.py b/solstone/think/providers/runtime_health.py index 86ed96fec..2ba6d14f4 100644 --- a/solstone/think/providers/runtime_health.py +++ b/solstone/think/providers/runtime_health.py @@ -228,6 +228,7 @@ REASON_CODE_GROUPS: dict[str, frozenset[str]] = { ), } REASON_CODES: frozenset[str] = frozenset().union(*REASON_CODE_GROUPS.values()) +ADMISSION_ONLY_REASON_CODES: frozenset[str] = frozenset({"ram-insufficient"}) if RUNTIME_PHASES & REASON_CODES: raise RuntimeError("runtime health phase and reason-code vocabularies overlap") @@ -1021,6 +1022,7 @@ def _validate_record_kind(value: object) -> RecordKind: __all__ = [ + "ADMISSION_ONLY_REASON_CODES", "InspectionStatus", "REASON_CODE_GROUPS", "REASON_CODES", diff --git a/solstone/think/supervisor.py b/solstone/think/supervisor.py index 6d32ab029..dbdc63581 100644 --- a/solstone/think/supervisor.py +++ b/solstone/think/supervisor.py @@ -77,6 +77,7 @@ from solstone.think.providers.install_state import ( from solstone.think.providers.memory import read_available_bytes from solstone.think.providers.mlx_server import MLX_SERVER_PROCESS_NAME from solstone.think.providers.runtime_health import ( + ADMISSION_ONLY_REASON_CODES, ReasonCode, RuntimeHealthConflictError, RuntimeHealthMalformedError, @@ -5143,6 +5144,28 @@ def _handle_provider_truth_result(state: ProviderRuntimeState) -> bool: ) _finish_provider_startup_condition(state, observation.phase) return True + if ( + not fingerprint_changed + and observation.phase == "host-blocked" + and observation.reason_code in ADMISSION_ONLY_REASON_CODES + and state.latest_phase in {"ready", "ready-proof-unavailable"} + ): + # Available-RAM headroom is an admission gate, not a liveness signal: the + # running provider's own resident footprint has already been subtracted + # from the reading the floor is compared against. Re-applying it here would + # make a successful model load the thing that evicts the model. + _write_provider_runtime( + state, + phase=state.latest_phase, + reason_code="stale-result-ignored", + detail={ + "slot": "truth", + "latched_phase": observation.phase, + "latched_reason_code": observation.reason_code, + }, + process=_current_provider_process_record(state.provider), + ) + return True if observation.phase in {"not-desired", "host-blocked"} and state.latest_phase in { "ready", "ready-proof-unavailable", diff --git a/tests/test_provider_runtime_health.py b/tests/test_provider_runtime_health.py index 0cbc026d0..920edec86 100644 --- a/tests/test_provider_runtime_health.py +++ b/tests/test_provider_runtime_health.py @@ -98,6 +98,26 @@ def test_phase_and_reason_code_vocabularies_are_disjoint() -> None: ) +def test_admission_only_reason_codes_are_intentional() -> None: + assert runtime_health.ADMISSION_ONLY_REASON_CODES == frozenset({"ram-insufficient"}) + assert not ( + runtime_health.ADMISSION_ONLY_REASON_CODES + & { + "platform-unsupported", + "package-unavailable", + "openmp-runtime-unavailable", + "gpu-probe-failed", + "gpu-unavailable", + "host-admission-blocked", + } + ) + + +def test_admission_only_reason_codes_are_known_reason_codes() -> None: + # Guard snake pre-map spelling against kebab post-map spelling mismatches. + assert runtime_health.ADMISSION_ONLY_REASON_CODES <= runtime_health.REASON_CODES + + def test_invalid_provider_rejected_on_entry_points(tmp_path: Path) -> None: for call in ( lambda: runtime_health_path("mlx", journal_path=tmp_path), diff --git a/tests/test_supervisor_parakeet.py b/tests/test_supervisor_parakeet.py index 1df6d4df7..25f557588 100644 --- a/tests/test_supervisor_parakeet.py +++ b/tests/test_supervisor_parakeet.py @@ -1245,3 +1245,45 @@ def test_parakeet_bootstrap_launches_only_via_reconciliation(monkeypatch) -> Non supervisor._submit_provider_start_if_needed(state, []) assert launches == [plan] + + +def test_ready_parakeet_host_admission_blocked_still_defers_admission_exclusive_stop( + monkeypatch, +) -> None: + """Pin deliberately-unchanged parakeet host-blocked stop deferral.""" + state = supervisor.ProviderRuntimeState("parakeet") + plan = _parakeet_plan("cpu") + state.latest_phase = "ready" + state.latest_plan = plan + state.desired_fingerprint = plan.desired_fingerprint_sha256 + state.retry.desired_fingerprint = plan.desired_fingerprint_sha256 + monkeypatch.setattr( + supervisor, + "_provider_runtime_states", + { + "local": supervisor.ProviderRuntimeState("local"), + "parakeet": state, + }, + ) + monkeypatch.setattr( + supervisor, "_write_provider_runtime", lambda *_args, **_kwargs: None + ) + state.truth_fence = supervisor._provider_fence(state, 0) + state.truth_future = _InlineExecutor().submit( + lambda: supervisor.ProviderTruthObservation( + provider="parakeet", + phase="host-blocked", + reason_code="host-admission-blocked", + detail={"host": {"reason": "stt admission pressure"}}, + desired_fingerprint_json=plan.desired_fingerprint_json, + desired_fingerprint_sha256=plan.desired_fingerprint_sha256, + boot_required=True, + ) + ) + + assert supervisor._handle_provider_truth_result(state) is True + + assert state.latest_phase == "stop-deferred" + assert state.pending_stop_admission_exclusive is True + assert state.pending_stop_target_phase == "host-blocked" + assert state.pending_stop_target_reason_code == "host-admission-blocked" diff --git a/tests/test_supervisor_provider_runtime.py b/tests/test_supervisor_provider_runtime.py index a401fd894..02d542791 100644 --- a/tests/test_supervisor_provider_runtime.py +++ b/tests/test_supervisor_provider_runtime.py @@ -32,6 +32,7 @@ from solstone.think.providers.runtime_health import ( RUNTIME_PHASES, ReasonCode, RuntimeHealthRecord, + RuntimeHealthUnavailableError, RuntimePhase, read_retry_token, read_runtime_health, @@ -1951,6 +1952,189 @@ def test_owned_local_host_blocked_defers_admission_exclusive_stop() -> None: assert state.pending_stop_target_phase == "host-blocked" +def _local_host_blocked_observation( + plan: supervisor.LocalServerLaunchPlan, + reason_code: ReasonCode, +) -> supervisor.ProviderTruthObservation: + return supervisor.ProviderTruthObservation( + provider="local", + phase="host-blocked", + reason_code=reason_code, + detail={"host": {"reason": reason_code}}, + desired_fingerprint_sha256=plan.desired_fingerprint_sha256, + boot_required=True, + ) + + +def test_ready_local_ram_insufficient_host_blocked_does_not_defer_stop() -> None: + plan = _local_plan() + state = supervisor._provider_runtime_states["local"] + _set_provider_ready("local", state, plan) + state.truth_fence = supervisor._provider_fence(state, state.retry.attempt_count) + state.truth_future = _future_with( + _local_host_blocked_observation(plan, "ram-insufficient") + ) + + assert supervisor._handle_provider_truth_result(state) is True + + assert state.latest_phase == "ready" + assert state.pending_stop_admission_exclusive is False + assert state.pending_stop_target_phase == "stopped" + assert state.pending_stop_request is None + + +@pytest.mark.parametrize( + "reason_code", + [ + "platform-unsupported", + "package-unavailable", + "openmp-runtime-unavailable", + "gpu-probe-failed", + "gpu-unavailable", + "host-admission-blocked", + ], +) +def test_ready_local_liveness_host_blocked_still_defers_admission_exclusive_stop( + reason_code: ReasonCode, +) -> None: + plan = _local_plan() + state = supervisor._provider_runtime_states["local"] + _set_provider_ready("local", state, plan) + state.truth_fence = supervisor._provider_fence(state, state.retry.attempt_count) + state.truth_future = _future_with( + _local_host_blocked_observation(plan, reason_code) + ) + + assert supervisor._handle_provider_truth_result(state) is True + + assert state.latest_phase == "stop-deferred" + assert state.pending_stop_admission_exclusive is True + assert state.pending_stop_target_phase == "host-blocked" + assert state.pending_stop_target_reason_code == reason_code + + +def test_ready_local_ram_insufficient_decline_preserves_plan_and_process() -> None: + plan = _local_plan() + managed = _FakeManaged() + state = supervisor._provider_runtime_states["local"] + process = { + "name": managed.name, + "pid": managed.process.pid, + "ref": managed.ref, + "port": 45678, + } + _set_provider_ready("local", state, plan) + write_runtime_health( + _runtime_record( + "local", + phase="ready", + fingerprint=plan.desired_fingerprint_sha256, + generation=state.generation, + attempt=state.retry.attempt_count, + process=process, + ) + ) + state.truth_fence = supervisor._provider_fence(state, state.retry.attempt_count) + state.truth_future = _future_with( + _local_host_blocked_observation(plan, "ram-insufficient") + ) + + assert supervisor._handle_provider_truth_result(state) is True + + health = read_runtime_health("local") + assert state.latest_phase == "ready" + assert state.latest_plan is plan + assert health["process"] == process + assert health["reason_code"] == "stale-result-ignored" + assert health["detail"]["slot"] == "truth" + assert health["detail"]["latched_phase"] == "host-blocked" + assert health["detail"]["latched_reason_code"] == "ram-insufficient" + + +def test_ready_local_ram_insufficient_decline_still_submits_probe(monkeypatch) -> None: + class _RecordingExecutor: + def __init__(self) -> None: + self.calls: list[tuple[Any, tuple[Any, ...]]] = [] + self.future: concurrent.futures.Future = concurrent.futures.Future() + + def submit(self, fn, *args): + self.calls.append((fn, args)) + return self.future + + plan = _local_plan() + state = supervisor._provider_runtime_states["local"] + _set_provider_ready("local", state, plan) + state.truth_fence = supervisor._provider_fence(state, state.retry.attempt_count) + state.truth_future = _future_with( + _local_host_blocked_observation(plan, "ram-insufficient") + ) + + assert supervisor._handle_provider_truth_result(state) is True + assert state.latest_phase == "ready" + + executor = _RecordingExecutor() + monkeypatch.setattr(supervisor, "_provider_executor", lambda: executor) + monkeypatch.setattr(supervisor, "read_service_port", lambda _service: 45678) + + supervisor._submit_provider_probe_if_needed(state) + + assert state.probe_future is executor.future + assert state.probe_fence is not None + assert len(executor.calls) == 1 + fn, args = executor.calls[0] + assert fn is supervisor._provider_probe_worker + assert args == ("local", 45678, state.probe_fence) + + +def test_ready_local_ram_insufficient_decline_write_unavailable_does_not_defer_stop( + monkeypatch, +) -> None: + """The error path must not re-enter the FALL-THROUGH trap. + + The plan is preserved, the process record is preserved, no stop is deferred, + and the phase is the honest pre-existing write-failure latch rather than a + false host-blocked. + """ + plan = _local_plan() + managed = _FakeManaged() + state = supervisor._provider_runtime_states["local"] + process = { + "name": managed.name, + "pid": managed.process.pid, + "ref": managed.ref, + "port": 45678, + } + _set_provider_ready("local", state, plan) + write_runtime_health( + _runtime_record( + "local", + phase="ready", + fingerprint=plan.desired_fingerprint_sha256, + generation=state.generation, + attempt=state.retry.attempt_count, + process=process, + ) + ) + + def fail_write(*_args, **_kwargs): + raise RuntimeHealthUnavailableError("runtime health write unavailable") + + monkeypatch.setattr(supervisor, "write_runtime_health", fail_write) + state.truth_fence = supervisor._provider_fence(state, state.retry.attempt_count) + state.truth_future = _future_with( + _local_host_blocked_observation(plan, "ram-insufficient") + ) + + assert supervisor._handle_provider_truth_result(state) is True + + assert state.pending_stop_request is None + assert state.pending_stop_admission_exclusive is False + assert state.pending_stop_target_phase == "stopped" + assert state.latest_plan is plan + assert read_runtime_health("local")["process"] == process + assert state.latest_phase == "state-unavailable" + + def test_parakeet_stt_admission_latch_survives_ready_probe_and_restart( monkeypatch, ) -> None: