392 lines
14 KiB
Python
392 lines
14 KiB
Python
|
|
"""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)
|