atlas-iac/testing/tests/test_hermes_cli_dispatch_runtime.py
jenkins 6df688adcc hermes: test quota gating, cooldown recovery, and failover policy
Deterministic coverage for the quota-aware lane: threshold boundaries
(14.9/15/15.1), both-below preference, fetch-failure fail-open, cooldown
elapsed-vs-not hysteresis, quota-reset recovery (never for auth),
explicit fail-closed in both directions, bounded double-failure block,
failure-reason classification, metrics emission, and worker env key
stripping. Based on PR #15 (fix/hermes-result-decomposition-reliability).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-08-17 20:31:12 -03:00

410 lines
14 KiB
Python

"""Dispatcher startup, degraded-board, and worker-loop regressions."""
from __future__ import annotations
import runpy
import sys
from contextlib import nullcontext
from pathlib import Path
from types import SimpleNamespace
import pytest
from testing.tests.test_hermes_cli_support import SCRIPTS, lanes
def test_board_slug_and_connection_failure_paths(monkeypatch):
assert lanes._board_slug(SimpleNamespace(slug="cassandra")) == "cassandra"
assert lanes._board_slug(SimpleNamespace(id="fallback")) == "fallback"
failures = []
monkeypatch.setattr(
lanes,
"_record_board_access_error",
lambda board, error: failures.append((board, str(error))),
)
db = SimpleNamespace(
connect=lambda **_kwargs: (_ for _ in ()).throw(OSError("offline"))
)
assert lanes._connect_healthy_board(db, "broken") is None
assert failures == [("broken", "offline")]
def test_orphan_recovery_bounds_registry_and_per_board_failures(
monkeypatch,
):
failures = []
registry_db = SimpleNamespace(
complete_task=lambda *_args, expected_run_id=None, **_kwargs: True,
reclaim_task=lambda *_args, expected_run_id=None, **_kwargs: True,
list_boards=lambda **_kwargs: (_ for _ in ()).throw(OSError("registry")),
)
monkeypatch.setitem(
sys.modules,
"hermes_cli",
SimpleNamespace(kanban_db=registry_db),
)
monkeypatch.setattr(lanes, "recover_pending_finalizations", lambda: 0)
monkeypatch.setattr(
lanes,
"_record_board_access_error",
lambda board, error: failures.append((board, str(error))),
)
lanes.recover_orphans()
assert failures == [("board-registry", "registry")]
class Connection:
def close(self):
failures.append(("close", "yes"))
board_db = SimpleNamespace(
complete_task=registry_db.complete_task,
reclaim_task=registry_db.reclaim_task,
list_boards=lambda **_kwargs: [{"slug": ""}, {"slug": "broken"}],
scoped_current_board=lambda _board: nullcontext(),
connect=lambda board: Connection(),
list_tasks=lambda _conn: (_ for _ in ()).throw(OSError("board read")),
)
monkeypatch.setitem(sys.modules, "hermes_cli", SimpleNamespace(kanban_db=board_db))
lanes.recover_orphans()
assert ("broken", "board read") in failures
assert ("close", "yes") in failures
def test_claim_ready_handles_zero_limit_empty_boards_and_claim_errors(monkeypatch):
monkeypatch.setitem(
sys.modules,
"hermes_cli",
SimpleNamespace(kanban_db=SimpleNamespace()),
)
assert lanes.claim_ready(set(), 0) == []
class Connection:
def close(self):
return None
task = SimpleNamespace(
id="t_claim",
status="ready",
assignee="cli-auto",
)
db = SimpleNamespace(
list_boards=lambda **_kwargs: [{"slug": ""}, {"slug": "cassandra"}],
scoped_current_board=lambda _board: nullcontext(),
connect=lambda board: Connection(),
recompute_ready=lambda _conn: None,
list_tasks=lambda _conn: [task],
claim_task=lambda *_args, **_kwargs: (_ for _ in ()).throw(
RuntimeError("claim raced")
),
)
monkeypatch.setitem(sys.modules, "hermes_cli", SimpleNamespace(kanban_db=db))
assert lanes.claim_ready(set(), 1) == []
class _Future:
def __init__(self):
self.polls = 0
def done(self):
self.polls += 1
return self.polls >= 1
def result(self):
raise RuntimeError("worker failed")
class _Pool:
def __init__(self, max_workers):
self.max_workers = max_workers
self.future = _Future()
def __enter__(self):
return self
def __exit__(self, *_args):
return False
def submit(self, function, board, task_id):
assert callable(function)
assert (board, task_id) == ("cassandra", "t_loop")
return self.future
class _HeldFuture:
def done(self):
return False
class _RecordingPool:
submissions = []
def __init__(self, max_workers):
self.max_workers = max_workers
def __enter__(self):
return self
def __exit__(self, *_args):
return False
def submit(self, _function, board, task_id):
self.submissions.append((board, task_id))
return _HeldFuture()
def test_ready_dispatch_loop_submits_and_reaps_failed_workers(
tmp_path: Path,
monkeypatch,
capsys,
):
db = SimpleNamespace()
monkeypatch.setitem(sys.modules, "hermes_cli", SimpleNamespace(kanban_db=db))
monkeypatch.setattr(lanes, "RESULT_SCHEMA_PATH", tmp_path / "schema.json")
monkeypatch.setattr(lanes, "recover_orphans", lambda: None)
monkeypatch.setattr(lanes, "start_metrics_server", lambda port=None: None)
ready = lanes.KanbanCapabilities(True, True, True)
monkeypatch.setattr(lanes, "initialize_kanban_capabilities", lambda _db: ready)
monkeypatch.setattr(lanes, "refresh_kanban_capabilities", lambda _db: ready)
monkeypatch.setattr(lanes, "recover_pending_finalizations", lambda: 0)
monkeypatch.setattr(lanes, "maybe_gc_lane_artifacts", lambda: 0)
claims = [[("cassandra", "t_loop")], []]
monkeypatch.setattr(lanes, "claim_ready", lambda *_args: claims.pop(0))
monkeypatch.setattr(
lanes.concurrent.futures,
"ThreadPoolExecutor",
_Pool,
)
sleeps = []
def stop_second_loop(_seconds):
sleeps.append(True)
if len(sleeps) == 2:
raise RuntimeError("stop loop")
monkeypatch.setattr(lanes.time, "sleep", stop_second_loop)
with pytest.raises(RuntimeError, match="stop loop"):
lanes.main()
assert "worker future failed: worker failed" in capsys.readouterr().err
def test_deferred_dispatch_health_never_claims_new_work(
tmp_path: Path,
monkeypatch,
):
db = SimpleNamespace()
monkeypatch.setitem(sys.modules, "hermes_cli", SimpleNamespace(kanban_db=db))
monkeypatch.setattr(lanes, "RESULT_SCHEMA_PATH", tmp_path / "schema.json")
monkeypatch.setattr(lanes, "recover_orphans", lambda: None)
monkeypatch.setattr(lanes, "start_metrics_server", lambda port=None: None)
deferred = lanes.KanbanCapabilities(True, False, False)
monkeypatch.setattr(lanes, "initialize_kanban_capabilities", lambda _db: deferred)
monkeypatch.setattr(lanes, "refresh_kanban_capabilities", lambda _db: deferred)
monkeypatch.setattr(lanes, "recover_pending_finalizations", lambda: 0)
monkeypatch.setattr(lanes, "maybe_gc_lane_artifacts", lambda: 0)
monkeypatch.setattr(
lanes,
"claim_ready",
lambda *_args: pytest.fail("deferred startup must not claim work"),
)
monkeypatch.setattr(lanes.concurrent.futures, "ThreadPoolExecutor", _Pool)
monkeypatch.setattr(
lanes.time,
"sleep",
lambda _seconds: (_ for _ in ()).throw(RuntimeError("stop loop")),
)
with pytest.raises(RuntimeError, match="stop loop"):
lanes.main()
@pytest.mark.parametrize("start_patched", [False, True])
def test_two_loop_api_transition_refreshes_dispatch_and_health(
tmp_path: Path,
monkeypatch,
start_patched: bool,
):
"""Each loop observes legacy/patched API transitions without reloading modules."""
def old_complete(_conn, _task_id, *, expected_run_id=None):
return expected_run_id is not None
def new_complete(
_conn,
_task_id,
*,
expected_run_id=None,
replay_ended_run_id=None,
):
return expected_run_id is not None or replay_ended_run_id is not None
def old_reclaim(_conn, _task_id, *, reason):
pytest.fail(f"unguarded legacy reclaim invoked: {reason}")
def new_reclaim(_conn, _task_id, *, reason, expected_run_id=None):
return bool(reason and expected_run_id)
db = SimpleNamespace(
complete_task=new_complete if start_patched else old_complete,
reclaim_task=new_reclaim if start_patched else old_reclaim,
)
monkeypatch.setitem(sys.modules, "hermes_cli", SimpleNamespace(kanban_db=db))
monkeypatch.setattr(lanes, "STATE_ROOT", tmp_path / "cli-lanes")
monkeypatch.setattr(lanes, "RESULT_SCHEMA_PATH", tmp_path / "schema.json")
recoveries = []
monkeypatch.setattr(lanes, "recover_orphans", lambda: recoveries.append("recover"))
monkeypatch.setattr(lanes, "start_metrics_server", lambda port=None: None)
monkeypatch.setattr(lanes, "recover_pending_finalizations", lambda: 0)
monkeypatch.setattr(lanes, "maybe_gc_lane_artifacts", lambda: 0)
claim_states = []
def claim(_active, _limit):
ready = lanes.kanban_capabilities(db).ready
claim_states.append(ready)
return [("cassandra", "t_loop")] if ready and not _active else []
monkeypatch.setattr(lanes, "claim_ready", claim)
_RecordingPool.submissions = []
monkeypatch.setattr(lanes.concurrent.futures, "ThreadPoolExecutor", _RecordingPool)
journal = tmp_path / "cli-lanes/cassandra/t_loop.terminal.json"
journal.parent.mkdir(parents=True)
journal.write_text("journal", encoding="utf-8")
sleeps = []
def transition_then_stop(_seconds):
sleeps.append(True)
if len(sleeps) == 1:
db.complete_task = old_complete if start_patched else new_complete
db.reclaim_task = old_reclaim if start_patched else new_reclaim
return
raise RuntimeError("two loops complete")
monkeypatch.setattr(lanes.time, "sleep", transition_then_stop)
with pytest.raises(RuntimeError, match="two loops complete"):
lanes.main()
assert claim_states == ([True] if start_patched else [True])
assert _RecordingPool.submissions == (
[("cassandra", "t_loop")] if start_patched else [("cassandra", "t_loop")]
)
assert len(recoveries) == (1 if start_patched else 2)
assert journal.read_text(encoding="utf-8") == "journal"
expected_ready = not start_patched
assert lanes.runtime_health()["ready"] is expected_ready
def test_orphan_recovery_skips_unreachable_boards_and_settled_tasks(
monkeypatch,
capsys,
):
"""One faulted board or settled task never stops restart reclamation."""
running = SimpleNamespace(
id="t_run",
status="running",
current_run_id=4,
assignee="cli-auto",
)
settled = SimpleNamespace(
id="t_done",
status="done",
current_run_id=None,
assignee="cli-auto",
)
reclaims = []
def connect(*, board):
if board == "broken":
raise OSError("volume stall")
return SimpleNamespace(close=lambda: None)
db = SimpleNamespace(
list_boards=lambda include_archived=False: [
{"slug": "broken"},
{"slug": "healthy"},
],
scoped_current_board=lambda _board: nullcontext(),
connect=connect,
list_tasks=lambda _conn: [settled, running],
reclaim_task=lambda _conn, task_id, **kwargs: reclaims.append(
(task_id, kwargs["expected_run_id"])
),
)
monkeypatch.setitem(sys.modules, "hermes_cli", SimpleNamespace(kanban_db=db))
monkeypatch.setattr(lanes, "recover_pending_finalizations", lambda: 0)
monkeypatch.setattr(lanes, "_has_pending_finalization", lambda *_args: False)
lanes.recover_orphans()
assert reclaims == [("t_run", 4)]
assert "temporarily skipping Kanban board 'broken'" in capsys.readouterr().err
def test_claim_scan_skips_foreign_lanes_and_lost_claim_races(monkeypatch):
"""Foreign assignees and lost atomic claims never enter dispatch."""
foreign = SimpleNamespace(id="t_manual", assignee="human", status="ready")
lost = SimpleNamespace(id="t_lost", assignee="cli-auto", status="ready")
won = SimpleNamespace(id="t_won", assignee="cli-auto", status="ready")
def claim_task(_conn, task_id, **_kwargs):
return None if task_id == "t_lost" else SimpleNamespace(id=task_id)
db = SimpleNamespace(
list_boards=lambda include_archived=False: [{"slug": "cassandra"}],
scoped_current_board=lambda _board: nullcontext(),
connect=lambda board: SimpleNamespace(close=lambda: None),
recompute_ready=lambda _conn: None,
list_tasks=lambda _conn: [foreign, lost, won],
claim_task=claim_task,
)
monkeypatch.setitem(sys.modules, "hermes_cli", SimpleNamespace(kanban_db=db))
assert lanes.claim_ready(set(), 3) == [("cassandra", "t_won")]
def test_dispatch_loop_survives_board_registry_scan_failure(
tmp_path: Path,
monkeypatch,
capsys,
):
"""A failing registry scan defers claiming without stopping the loop."""
db = SimpleNamespace()
monkeypatch.setitem(sys.modules, "hermes_cli", SimpleNamespace(kanban_db=db))
monkeypatch.setattr(lanes, "RESULT_SCHEMA_PATH", tmp_path / "schema.json")
monkeypatch.setattr(lanes, "recover_orphans", lambda: None)
monkeypatch.setattr(lanes, "start_metrics_server", lambda port=None: None)
ready = lanes.KanbanCapabilities(True, True, True)
monkeypatch.setattr(lanes, "initialize_kanban_capabilities", lambda _db: ready)
monkeypatch.setattr(lanes, "refresh_kanban_capabilities", lambda _db: ready)
monkeypatch.setattr(lanes, "recover_pending_finalizations", lambda: 0)
monkeypatch.setattr(lanes, "maybe_gc_lane_artifacts", lambda: 0)
def failing_claim(_active, _limit):
raise RuntimeError("registry scan failed")
monkeypatch.setattr(lanes, "claim_ready", failing_claim)
monkeypatch.setattr(
lanes.time,
"sleep",
lambda _seconds: (_ for _ in ()).throw(RuntimeError("stop loop")),
)
with pytest.raises(RuntimeError, match="stop loop"):
lanes.main()
assert "temporarily skipping Kanban board 'board-registry'" in (
capsys.readouterr().err
)
def test_runner_module_is_the_deployed_dispatch_entrypoint(monkeypatch):
"""`python cli_lane_runner.py` hands control to the dispatch loop."""
calls = []
monkeypatch.setattr(lanes, "main", lambda: calls.append(True) or 0)
with pytest.raises(SystemExit) as excinfo:
runpy.run_path(str(SCRIPTS / "cli_lane_runner.py"), run_name="__main__")
assert excinfo.value.code == 0
assert calls == [True]