atlas-iac/testing/tests/test_hermes_execution_pool_dispatch_v2.py
Hermes Agent 5d963c9354 hermes: make pool lease recovery and release isolation safe
Independent review t_5975c06a blocked this branch on a P1: a Kanban write that
failed while a lease expired left a `lease_failed` row that was invisible to
every pass, immortal to garbage collection, and fatal to the coordinator. It
poisoned `reconcile()` forever with a conflicting-duplicate primary key,
produced a spurious capability `block_task` from `dispatch()`, and -- because
startup maintenance ran unguarded before the port bound, against a store on a
PVC -- crash-looped the coordinator with no automatic recovery.

`lease_failed` is now a retryable state that every maintenance pass drains, and
a row only reaches a terminal state on authoritative evidence about its exact
Kanban run, so nothing is collected before its outcome is known and nothing is
silently dropped. Each row, task, and board is processed in isolation, and a
coordinator-side fault is never converted into a Kanban mutation. Startup runs
through the same guarded cycle as the steady-state loop.

The wire protocol and the durable store are now separate modules, and the
maintenance passes moved out of the coordinator, so each file stays under the
managed line ceiling with room for the recovery logic.

Also closes three consequential handoff risks the same review raised:

* mediator-N pinned itself hard to worker-N while sharing a ReadWriteOnce
  claim, so a drain or preemption that moved only the lower-priority worker
  deadlocked the ordinal on Multi-Attach until an operator deleted a Pod. The
  shared workspace is now ReadWriteMany (as the hermes-chat tenant workspaces
  already are on the same class), colocation is a preference, and the mediator
  shares the worker's preemption priority, so each Pod reschedules on its own.
* the broker permits only branch creation, so a retry that added commits could
  never submit and the run's work was discarded with the failure. Submission
  now targets a fresh attempt- or content-scoped ref in the same reviewed
  namespace -- never an update -- and is idempotent under replay. A refused
  submission downgrades the result and says why instead of unwinding the run.
* the provider CLIs were reinstalled into an emptyDir on every Pod start inside
  the 10m Flux health window for the whole hermes app. They now install once
  per pinned version onto a durable volume, re-verified against the real
  binaries and time-bounded, and the best-effort pool no longer gates the
  health of the app its dependents wait on.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-17 16:31:15 +00:00

412 lines
14 KiB
Python

"""Coordinator dispatch migration and versioned server contracts."""
from __future__ import annotations
import sqlite3
import sys
from pathlib import Path
from types import SimpleNamespace
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
from testing.tests.test_hermes_execution_pool_coordinator_v2 import ( # noqa: E402
MASTER,
assignment_payload,
binding,
install_kanban,
store_with_assignment,
task,
)
from testing.tests.test_hermes_execution_pool_mediator import http_request # noqa: E402
def test_reconcile_fences_old_run_then_recovers_current_pathless_run(
tmp_path, monkeypatch
):
item = task(current_run_id=24)
install_kanban(monkeypatch, [item], tmp_path)
store = store_with_assignment(tmp_path)
pool = coordinator.Coordinator(MASTER, store)
monkeypatch.setattr(
coordinator, "assignment_payload",
lambda _db, _connection, _task, _board: assignment_payload(),
)
pool.reconcile()
rows = store._connect().execute(
"SELECT run_id,state,worker_ordinal FROM assignments ORDER BY run_id"
).fetchall()
assert [tuple(row) for row in rows] == [
("23", "stale", 0), ("24", "assigned", 0)
]
def test_reconcile_preserves_owned_workspace_and_exactly_blocks_registry_failure(
tmp_path, monkeypatch
):
owned = task(id="t_owned", current_run_id=30, workspace_path="/owned")
broken = task(id="t_broken", current_run_id=31)
kanban = install_kanban(monkeypatch, [owned, broken], tmp_path)
store = protocol.PoolStore(tmp_path / "pool.db")
pool = coordinator.Coordinator(MASTER, store)
def prepare(_db, _connection, item, _board):
if item.id == "t_broken":
raise RuntimeError("registry unavailable")
return assignment_payload()
monkeypatch.setattr(coordinator, "assignment_payload", prepare)
pool.reconcile()
assert store.active_assignments() == []
assert kanban.blocked[0][0] == "t_broken"
assert kanban.blocked[0][1]["expected_run_id"] == 31
assert all(call[0] != "t_owned" for call in kanban.blocked)
def test_dispatch_claims_only_pathless_task_with_safe_assignment_branch(
tmp_path, monkeypatch
):
item = task(status="running")
install_kanban(monkeypatch, [item], tmp_path)
store = protocol.PoolStore(tmp_path / "pool.db")
pool = coordinator.Coordinator(MASTER, store)
observed = []
def claim(active, limit, eligible):
observed.append((active, limit, eligible("metis", item)))
return [("metis", item.id)]
monkeypatch.setattr(maintenance.cli_lane_dispatch, "claim_ready", claim)
monkeypatch.setattr(
coordinator, "assignment_payload",
lambda *_a: assignment_payload(branch="wt/t_deadbeef"),
)
pool.dispatch()
assert observed == [(set(), 3, True)]
record = store.active_assignments()[0]
assert record["run_id"] == "23"
assert record["payload"]["branch"] == "wt/t_deadbeef"
def test_dispatch_incompatible_unversioned_claim_api_fails_closed(
tmp_path, monkeypatch
):
owned = task(workspace_path="/opt/data/workspace/live")
kanban = install_kanban(monkeypatch, [owned], tmp_path)
store = protocol.PoolStore(tmp_path / "pool.db")
pool = coordinator.Coordinator(MASTER, store)
calls = []
def old_claim(active, limit):
calls.append((active, limit))
return [("metis", owned.id)]
monkeypatch.setattr(maintenance.cli_lane_dispatch, "claim_ready", old_claim)
with pytest.raises(TypeError):
pool.dispatch()
assert calls == []
assert store.active_assignments() == []
assert kanban.reclaimed == []
def test_dispatch_workspace_ownership_race_is_exactly_fenced(
tmp_path, monkeypatch
):
owned = task(workspace_path="/opt/data/workspace/live")
kanban = install_kanban(monkeypatch, [owned], tmp_path)
store = protocol.PoolStore(tmp_path / "pool.db")
pool = coordinator.Coordinator(MASTER, store)
monkeypatch.setattr(
maintenance.cli_lane_dispatch,
"claim_ready",
lambda *_a: [("metis", owned.id)],
)
pool.dispatch()
assert store.active_assignments() == []
assert kanban.blocked == [
(
owned.id,
{
"reason": (
"Distributed claim fenced because an existing workspace "
"is owned by the local lane"
),
"kind": "capability",
"expected_run_id": 23,
},
)
]
def test_dispatch_preparation_failure_is_surfaced_and_full_pool_does_not_claim(
tmp_path, monkeypatch
):
broken = task(current_run_id=None)
kanban = install_kanban(monkeypatch, [broken], tmp_path)
store = protocol.PoolStore(tmp_path / "pool.db")
pool = coordinator.Coordinator(MASTER, store)
monkeypatch.setattr(
maintenance.cli_lane_dispatch,
"claim_ready",
lambda *_a: [("metis", broken.id)],
)
pool.dispatch()
assert kanban.blocked[0][1]["expected_run_id"] is None
assert "canonical run ID" in kanban.blocked[0][1]["reason"]
full = protocol.PoolStore(tmp_path / "full.db")
for ordinal in range(3):
full.add(
binding(task_id=f"t_{ordinal}", run_id=str(ordinal + 1), worker_ordinal=ordinal),
assignment_payload(),
)
pool = coordinator.Coordinator(MASTER, full)
monkeypatch.setattr(
maintenance.cli_lane_dispatch,
"claim_ready",
lambda *_a: (_ for _ in ()).throw(AssertionError("must not claim")),
)
pool.dispatch()
def test_lease_expiry_with_noncanonical_run_is_released_as_stale(tmp_path, monkeypatch):
install_kanban(monkeypatch, [task()], tmp_path)
store = protocol.PoolStore(tmp_path / "pool.db")
exact = binding(run_id="bad-run", attempt=3)
store.add(exact, assignment_payload())
store.offer(0)
with store._connect() as connection:
connection.execute("UPDATE assignments SET lease_until=1")
coordinator.Coordinator(MASTER, store).expire_leases()
assert store._connect().execute("SELECT state FROM assignments").fetchone()[0] == "stale"
def test_terminal_activity_marker_prevents_duplicate_append(tmp_path):
class Kanban:
worker_log_path = staticmethod(
lambda task_id, board=None: str(tmp_path / f"{board}-{task_id}.log")
)
record = {
**binding(),
"result_digest": "a" * 64,
"result": {"final_activity": "terminal evidence\n"},
}
coordinator._append_terminal_activity(Kanban, record)
coordinator._append_terminal_activity(Kanban, record)
text = (tmp_path / "metis-t_deadbeef.log").read_text()
assert text.count("execution-pool-result") == 1
empty = {**record, "result": {"final_activity": ""}}
coordinator._append_terminal_activity(Kanban, empty)
def test_activity_and_terminal_log_reject_invalid_or_symlink_targets(tmp_path):
class Kanban:
worker_log_path = staticmethod(lambda *_a, **_k: str(tmp_path / "worker.log"))
with pytest.raises(protocol.ProtocolError, match="payload"):
coordinator._append_activity(Kanban, {**binding(), "payload": []})
coordinator._append_activity(Kanban, {**binding(), "payload": {"activity": ""}})
(tmp_path / "worker.log").symlink_to(tmp_path / "target")
with pytest.raises(protocol.ProtocolError, match="symlink"):
coordinator._append_terminal_activity(
Kanban,
{
**binding(), "result_digest": "a" * 64,
"result": {"final_activity": "terminal"},
},
)
def test_heartbeat_bad_run_duplicate_and_unstructured_result_paths(
tmp_path, monkeypatch
):
item = task()
kanban = install_kanban(monkeypatch, [item], tmp_path)
bad_binding = binding(run_id="bad-run")
bad_store = protocol.PoolStore(tmp_path / "bad.db")
bad_store.add(bad_binding, assignment_payload())
bad_store.offer(0)
pool = coordinator.Coordinator(MASTER, bad_store)
key = protocol.derive_ordinal_key(MASTER, 0)
with pytest.raises(protocol.ProtocolError, match="canonical"):
pool.heartbeat(protocol.sign_envelope(key, "heartbeat", bad_binding, {}))
store = store_with_assignment(tmp_path / "duplicate")
store.offer(0)
pool = coordinator.Coordinator(MASTER, store)
request = protocol.sign_envelope(
key, "heartbeat", binding(), {"activity": "once"},
delivery_id="same-heartbeat",
)
pool.heartbeat(request)
pool.heartbeat(request)
assert (tmp_path / "metis/t_deadbeef.log").read_text() == "once"
second = binding(run_id="24", worker_ordinal=1)
item.current_run_id = 24
store.add(second, assignment_payload())
pool.result(
protocol.sign_envelope(
protocol.derive_ordinal_key(MASTER, 1), "result", second,
{"structured": [], "returncode": 1},
)
)
assert kanban.blocked[-1][1]["expected_run_id"] == 24
def test_reconcile_and_dispatch_cover_empty_error_and_missing_task_paths(
tmp_path, monkeypatch
):
item = task()
kanban = install_kanban(monkeypatch, [item], tmp_path, boards=["", "metis"])
store = store_with_assignment(tmp_path)
pool = coordinator.Coordinator(MASTER, store)
original_connect = kanban.connect
kanban.list_boards = lambda include_archived=False: []
kanban.connect = lambda board=None: (_ for _ in ()).throw(OSError("busy"))
pool.reconcile()
kanban.connect = original_connect
kanban.list_boards = lambda include_archived=False: [""]
pool.reconcile()
full = SimpleNamespace(
active_assignments=lambda: [], available_ordinals=lambda: []
)
coordinator.Coordinator(MASTER, full).reconcile()
empty = protocol.PoolStore(tmp_path / "missing.db")
monkeypatch.setattr(
maintenance.cli_lane_dispatch, "claim_ready", lambda *_a: [("metis", "missing")]
)
coordinator.Coordinator(MASTER, empty).dispatch()
def test_nonterminal_lease_expiry_is_reoffered_without_kanban_block(tmp_path, monkeypatch):
kanban = install_kanban(monkeypatch, [task()], tmp_path)
store = store_with_assignment(tmp_path)
store.offer(0)
with store._connect() as connection:
connection.execute("UPDATE assignments SET lease_until=1")
coordinator.Coordinator(MASTER, store).expire_leases()
assert kanban.blocked == []
assert store.active_assignments()[0]["attempt"] == 2
class HTTPStore:
def __init__(self, fail=False):
self.fail = fail
def available_ordinals(self):
if self.fail:
raise sqlite3.Error("busy")
return [0, 1, 2]
class HTTPCoordinator:
def __init__(self, fail_ready=False):
self.store = HTTPStore(fail_ready)
@staticmethod
def poll(value):
return {"route": "poll", "value": value}
@staticmethod
def heartbeat(_value):
raise protocol.ProtocolError("stale")
@staticmethod
def result(_value):
raise RuntimeError("storage")
def test_versioned_server_routes_readiness_and_errors():
handler = server.handler_factory(HTTPCoordinator())
assert http_request(handler, "/ready") == (200, {"ready": True, "version": 2})
assert http_request(handler, "/missing")[0] == 404
body = protocol.canonical_json({"safe": True})
assert http_request(handler, "/v1/poll", body=body) == (
200, {"route": "poll", "value": {"safe": True}}
)
assert http_request(handler, "/v1/heartbeat", body=body)[0] == 409
assert http_request(handler, "/v1/result", body=body)[0] == 503
assert http_request(handler, "/unknown", body=body)[0] == 404
assert http_request(handler, "/v1/poll", body=b"")[0] == 409
assert http_request(server.handler_factory(HTTPCoordinator(True)), "/ready")[0] == 503
class RunCoordinator:
instances = []
def __init__(self, key, store):
self.key = key
self.store = store
self.calls = []
self.dispatch_count = 0
self.__class__.instances.append(self)
def expire_leases(self):
self.calls.append("expire")
def recover_results(self):
self.calls.append("recover")
def reconcile(self):
self.calls.append("reconcile")
def dispatch(self):
self.calls.append("dispatch")
self.dispatch_count += 1
if self.dispatch_count > 1:
raise RuntimeError("deferred")
def test_server_once_and_maintenance_loop_are_versioned_and_resilient(
tmp_path, monkeypatch, capsys
):
monkeypatch.setattr(server, "STATE_ROOT", tmp_path)
monkeypatch.setattr(server, "read_key", lambda _path: MASTER)
monkeypatch.setattr(sys, "argv", ["pool", "--once"])
assert server.run(RunCoordinator) == 0
assert RunCoordinator.instances[-1].calls == [
"expire", "recover", "reconcile", "dispatch"
]
started = []
class FakeServer:
def __init__(self, address, _handler, max_workers):
started.append((address, max_workers))
def serve_forever(self):
return
monkeypatch.setattr(sys, "argv", ["pool"])
monkeypatch.setattr(server, "BoundedHTTPServer", FakeServer)
monkeypatch.setattr(
server.time,
"sleep",
lambda _seconds: (_ for _ in ()).throw(StopIteration()),
)
with pytest.raises(StopIteration):
server.run(RunCoordinator)
assert started == [(("0.0.0.0", server.PORT), 8)]
assert "maintenance deferred" in capsys.readouterr().err
def test_coordinator_compatibility_exports_delegate_to_server(monkeypatch):
marker = object()
monkeypatch.setattr(server, "handler_factory", lambda value: (marker, value))
assert coordinator.handler_factory(marker) == (marker, marker)
monkeypatch.setattr(server, "run", lambda value: 17 if value is coordinator.Coordinator else 0)
assert coordinator.main() == 17