"""Fault isolation contracts for the execution-pool maintenance passes. No single poisoned row, task, or board may abort the work queued behind it, and a coordinator-side fault must never be converted into a Kanban mutation. Startup maintenance is held to exactly the same standard as the steady-state loop, because the durable store lives on a PVC and would otherwise crash-loop. """ from __future__ import annotations import sqlite3 import sys from pathlib import Path import pytest ROOT = Path(__file__).parents[2] SCRIPTS = ROOT / "services/hermes/scripts" SCM_SCRIPTS = ROOT / "services/hermes/scm-common/scripts" sys.path[:0] = [str(SCRIPTS), str(SCM_SCRIPTS)] import execution_pool_coordinator as coordinator # noqa: E402 import execution_pool_maintenance as maintenance # noqa: E402 import execution_pool_protocol as protocol # noqa: E402 import execution_pool_server as server # noqa: E402 import execution_pool_store as pool_store # noqa: E402 from testing.tests.test_hermes_execution_pool_coordinator_v2 import ( # noqa: E402 MASTER, assignment_payload, binding, install_kanban, task, ) from testing.tests.test_hermes_execution_pool_support import ( # noqa: E402 LOCKED, clear_pool_deferrals, # noqa: F401 - autouse fixture, resolved by name exhausted_store, states, ) def test_one_poisoned_board_and_task_never_abort_the_work_behind_them( tmp_path, monkeypatch, capsys ): first = task(id="t_first", current_run_id=41) second = task(id="t_second", current_run_id=42) kanban = install_kanban( monkeypatch, [first, second], tmp_path, boards=["broken", "metis"] ) real_list_tasks = kanban.list_tasks payloads = {"calls": 0} def flaky_list_tasks(connection): if kanban.current == "broken": raise LOCKED return real_list_tasks(connection) def flaky_payload(_db, _connection, item, _board): payloads["calls"] += 1 if item.id == "t_first": raise LOCKED return assignment_payload() def scoped(board): kanban.current = board return kanban.scope(board) kanban.current = "" kanban.scope = kanban.scoped_current_board kanban.scoped_current_board = scoped kanban.list_tasks = flaky_list_tasks monkeypatch.setattr(coordinator, "assignment_payload", flaky_payload) store = pool_store.PoolStore(tmp_path / "poison.db") coordinator.Coordinator(MASTER, store).reconcile() assert payloads["calls"] == 2 assert states(store) == [("42", 0, 1, "assigned")] assert kanban.blocked == [] captured = capsys.readouterr().err assert "broken adoption" in captured assert "t_first recovery payload" in captured def test_adoption_conflict_leaves_the_ordinal_free_and_kanban_untouched( tmp_path, monkeypatch, capsys ): live = task() kanban = install_kanban(monkeypatch, [live], tmp_path) monkeypatch.setattr( coordinator, "assignment_payload", lambda _db, _connection, _task, _board: assignment_payload(), ) store = pool_store.PoolStore(tmp_path / "conflict.db") pool = coordinator.Coordinator(MASTER, store) def conflict(_binding, _payload): raise protocol.ProtocolError("conflicting duplicate assignment") monkeypatch.setattr(store, "add", conflict) pool.reconcile() assert store.available_ordinals() == [0, 1, 2] assert kanban.blocked == [] assert "t_deadbeef adoption" in capsys.readouterr().err def test_dispatch_defers_storage_and_identity_faults_without_blocking_the_run( tmp_path, monkeypatch, capsys ): live = task() kanban = install_kanban(monkeypatch, [live], tmp_path) monkeypatch.setattr( maintenance.cli_lane_dispatch, "claim_ready", lambda *_args, **_kwargs: [("metis", live.id)], ) store = pool_store.PoolStore(tmp_path / "dispatch.db") pool = coordinator.Coordinator(MASTER, store) monkeypatch.setattr( coordinator, "assignment_payload", lambda *_args: (_ for _ in ()).throw(LOCKED), ) pool.dispatch() assert kanban.blocked == [] assert store.active_assignments() == [] assert "metis/t_deadbeef dispatch" in capsys.readouterr().err monkeypatch.setattr( coordinator, "assignment_payload", lambda _db, _connection, _task, _board: assignment_payload(), ) monkeypatch.setattr( store, "add", lambda *_args: (_ for _ in ()).throw( protocol.ProtocolError("worker ordinal already has a live assignment") ), ) pool.dispatch() assert kanban.blocked == [] def test_dispatch_still_parks_the_exact_run_on_a_real_capability_failure( tmp_path, monkeypatch ): live = task() kanban = install_kanban(monkeypatch, [live], tmp_path) monkeypatch.setattr( maintenance.cli_lane_dispatch, "claim_ready", lambda *_args, **_kwargs: [("metis", live.id)], ) monkeypatch.setattr( coordinator, "assignment_payload", lambda *_args: (_ for _ in ()).throw(RuntimeError("context exceeds 32KiB")), ) store = pool_store.PoolStore(tmp_path / "capability.db") coordinator.Coordinator(MASTER, store).dispatch() assert kanban.blocked[-1][0] == "t_deadbeef" assert kanban.blocked[-1][1]["kind"] == "capability" assert kanban.blocked[-1][1]["expected_run_id"] == 23 assert "32KiB" in kanban.blocked[-1][1]["reason"] def test_a_failing_park_write_is_absorbed_rather_than_wedging_the_pass( tmp_path, monkeypatch, capsys ): live = task() kanban = install_kanban(monkeypatch, [live], tmp_path) kanban.block_task = lambda *_args, **_kwargs: (_ for _ in ()).throw(LOCKED) monkeypatch.setattr( maintenance.cli_lane_dispatch, "claim_ready", lambda *_args, **_kwargs: [("metis", live.id)], ) monkeypatch.setattr( coordinator, "assignment_payload", lambda *_args: (_ for _ in ()).throw(RuntimeError("bad")), ) store = pool_store.PoolStore(tmp_path / "park.db") coordinator.Coordinator(MASTER, store).dispatch() assert "t_deadbeef block" in capsys.readouterr().err def test_one_unprocessable_result_never_blocks_the_results_behind_it( tmp_path, monkeypatch, capsys ): install_kanban(monkeypatch, [task(), task(id="t_second", current_run_id=24)], tmp_path) store = pool_store.PoolStore(tmp_path / "results.db") pool = coordinator.Coordinator(MASTER, store) handled = [] def finalize(record): if record["task_id"] == "t_deadbeef": raise protocol.ProtocolError("result payload must be an object") handled.append(record["task_id"]) monkeypatch.setattr( store, "pending_results", lambda: [{**binding(), "result": []}, {**binding(task_id="t_second"), "result": {}}], ) monkeypatch.setattr(pool, "finalize", finalize) pool.recover_results() assert handled == ["t_second"] assert "result recovery deferred" in capsys.readouterr().err def test_startup_maintenance_is_as_safe_as_the_steady_state_loop( tmp_path, monkeypatch, capsys ): """A poisoned store must not stop the coordinator from binding its port.""" live = task() kanban = install_kanban(monkeypatch, [live], tmp_path) kanban.block_task = lambda *_args, **_kwargs: (_ for _ in ()).throw(LOCKED) kanban.list_boards = lambda include_archived=False: (_ for _ in ()).throw( RuntimeError("board registry unavailable") ) key = tmp_path / "key" key.write_bytes(b"m" * 48) key.chmod(0o600) monkeypatch.setattr(server, "STATE_ROOT", tmp_path / "state") monkeypatch.setattr(server, "KEY_PATH", key) monkeypatch.setattr(sys, "argv", ["execution_pool_server", "--once"]) exhausted_store(tmp_path / "state/assignments.db") assert server.run(coordinator.Coordinator) == 0 captured = capsys.readouterr().err assert "pool maintenance deferred" in captured row = sqlite3.connect(tmp_path / "state/assignments.db").execute( "SELECT state FROM assignments" ).fetchone() assert row[0] == "lease_failed" def test_maintenance_cycle_counts_deferrals_and_keeps_running(capsys): calls = [] count = server.maintenance_cycle( ( lambda: calls.append("first"), lambda: (_ for _ in ()).throw(RuntimeError("boom")), lambda: calls.append("last"), ) ) assert count == 1 assert calls == ["first", "last"] assert "pool maintenance deferred: RuntimeError: boom" in capsys.readouterr().err def test_the_poison_row_survives_restart_on_the_same_pvc_and_then_resolves( tmp_path, monkeypatch, capsys ): """The store lives on hermes-agent-home, so recovery must survive restarts.""" live = task() kanban = install_kanban(monkeypatch, [live], tmp_path) kanban.block_task = lambda *_args, **_kwargs: (_ for _ in ()).throw(LOCKED) database = tmp_path / "pvc/assignments.db" first = exhausted_store(database) coordinator.Coordinator(MASTER, first).expire_leases() assert states(first) == [("23", 0, 3, "lease_failed")] del first capsys.readouterr() # A fresh process attaches the same PVC path and finds the row waiting. restarted = pool_store.PoolStore(database) assert restarted.failed_leases()[0]["state"] == "lease_failed" assert restarted.available_ordinals() == [0, 1, 2] kanban.block_task = lambda _connection, task_id, **values: ( kanban.blocked.append((task_id, values)) or True ) coordinator.Coordinator(MASTER, restarted).expire_leases() assert states(restarted) == [("23", 0, 3, "finalized")] assert kanban.blocked[-1][1]["expected_run_id"] == 23 def test_a_store_write_fault_while_recording_leaves_the_row_retryable( tmp_path, monkeypatch, capsys ): live = task() install_kanban(monkeypatch, [live], tmp_path) store = exhausted_store(tmp_path / "store-fault.db") pool = coordinator.Coordinator(MASTER, store) real_finalize = store.finalize monkeypatch.setattr( store, "finalize", lambda *_args: (_ for _ in ()).throw(LOCKED) ) pool.expire_leases() assert states(store) == [("23", 0, 3, "lease_failed")] assert "lease park" in capsys.readouterr().err monkeypatch.setattr(store, "finalize", real_finalize) pool.expire_leases() assert states(store) == [("23", 0, 3, "finalized")] def test_finalize_reports_whether_the_exact_row_changed(tmp_path): store = pool_store.PoolStore(tmp_path / "finalize.db") store.add(binding(), assignment_payload()) assert store.finalize(binding(attempt=9), "stale") is False assert store.finalize(binding(), "stale") is True with pytest.raises(protocol.ProtocolError, match="terminal"): store.finalize(binding(), "lease_failed") def test_deferral_reporting_is_deduplicated_and_bounded(capsys): maintenance._defer("ctx", LOCKED) maintenance._defer("ctx", LOCKED) assert capsys.readouterr().err.count("pool work deferred") == 1 maintenance._defer("ctx", RuntimeError("changed")) assert "changed" in capsys.readouterr().err for index in range(maintenance.MAX_DEFERRALS + 1): maintenance._defer(f"ctx-{index}", LOCKED) assert len(maintenance.DEFERRALS) <= maintenance.MAX_DEFERRALS capsys.readouterr() def test_protocol_exposes_only_the_store_on_its_compatibility_surface(): assert protocol.PoolStore is pool_store.PoolStore missing = "SomethingElse" with pytest.raises(AttributeError, match="no attribute"): getattr(protocol, missing) def test_adoption_stops_as_soon_as_the_last_ordinal_is_taken(tmp_path, monkeypatch): items = [task(id=f"t_{index}", current_run_id=50 + index) for index in range(4)] kanban = install_kanban(monkeypatch, items, tmp_path) monkeypatch.setattr( coordinator, "assignment_payload", lambda _db, _connection, _task, _board: assignment_payload(), ) store = pool_store.PoolStore(tmp_path / "full.db") coordinator.Coordinator(MASTER, store).reconcile() assert len(store.active_assignments()) == 3 assert store.available_ordinals() == [] assert kanban.blocked == [] def test_an_unreadable_board_registry_defers_instead_of_escaping( tmp_path, monkeypatch, capsys ): kanban = install_kanban(monkeypatch, [task()], tmp_path) kanban.list_boards = lambda include_archived=False: (_ for _ in ()).throw(LOCKED) store = pool_store.PoolStore(tmp_path / "registry.db") coordinator.Coordinator(MASTER, store).reconcile() assert store.active_assignments() == [] assert "board-registry" in capsys.readouterr().err def test_an_unreadable_claim_path_defers_instead_of_escaping( tmp_path, monkeypatch, capsys ): install_kanban(monkeypatch, [task()], tmp_path) monkeypatch.setattr( maintenance.cli_lane_dispatch, "claim_ready", lambda *_args, **_kwargs: (_ for _ in ()).throw(LOCKED), ) store = pool_store.PoolStore(tmp_path / "claim.db") coordinator.Coordinator(MASTER, store).dispatch() assert store.active_assignments() == [] assert "claim-ready" in capsys.readouterr().err def test_an_owned_workspace_without_a_canonical_run_is_skipped_silently( tmp_path, monkeypatch ): owned = task(workspace_path="/owned", current_run_id=None) kanban = install_kanban(monkeypatch, [owned], tmp_path) monkeypatch.setattr( maintenance.cli_lane_dispatch, "claim_ready", lambda *_args, **_kwargs: [("metis", owned.id)], ) store = pool_store.PoolStore(tmp_path / "owned.db") coordinator.Coordinator(MASTER, store).dispatch() assert kanban.blocked == [] assert store.active_assignments() == [] def test_a_result_cannot_be_accepted_onto_an_already_terminal_row(tmp_path): store = pool_store.PoolStore(tmp_path / "terminal.db") store.add(binding(), assignment_payload()) store.finalize(binding(), "finalized") envelope = {**binding(), "payload_digest": "d" * 64, "payload": {}} with pytest.raises(protocol.ProtocolError, match="cannot accept a result"): store.accept_result(envelope) def test_a_nonstring_message_kind_is_rejected_before_any_lookup(): envelope = { "version": protocol.PROTOCOL_VERSION, "kind": 7, "board": "metis", "task_id": "t_deadbeef", "run_id": "23", "worker_ordinal": 0, "attempt": 1, "delivery_id": "d1", "issued_at": 0, "expires_at": 1, "payload_digest": "x", "payload": {}, "signature": "s", } with pytest.raises(protocol.ProtocolError, match="unexpected message kind"): protocol.verify_envelope(b"k" * 32, envelope)