atlas-iac/testing/tests/test_hermes_cli_recovery_edges.py

445 lines
16 KiB
Python
Raw Permalink Normal View History

2026-08-17 08:16:35 -03:00
"""Recovery convergence and hidden-authority edge coverage."""
from __future__ import annotations
import hashlib
import os
import sys
from contextlib import nullcontext
2026-08-17 08:16:35 -03:00
from pathlib import Path
from types import SimpleNamespace
import pytest
from testing.tests.test_hermes_cli_support import (
_pending_terminal_record,
lanes,
)
class _Snapshot:
def __init__(self, document=None):
self.document = document
self.closed = 0
def close(self):
self.closed += 1
def test_terminal_absence_distinguishes_present_missing_and_inaccessible(
tmp_path: Path,
monkeypatch,
):
path = tmp_path / "journal"
path.write_text("data", encoding="utf-8")
assert lanes._terminal_entry_absent(path) is False
path.unlink()
assert lanes._terminal_entry_absent(path) is True
monkeypatch.setattr(
lanes.os,
"open",
lambda *_args, **_kwargs: (_ for _ in ()).throw(PermissionError("denied")),
)
assert lanes._terminal_entry_absent(path) is False
def _staged_path(root: Path, identity, document: dict) -> tuple[Path, Path]:
if identity.state in {"pending", "committed"}:
canonical = lanes._terminal_path(
lanes.state_path(identity.board, identity.task_id),
identity.run_id,
identity.state,
)
else:
canonical = lanes._terminal_evidence_path(identity, identity.state, document)
digest = hashlib.sha256(canonical.name.encode("utf-8")).hexdigest()[:16]
return root / identity.board / f".retire.{digest}.{'a' * 16}.0", canonical
@pytest.mark.parametrize("state", ["prepared", "conflict"])
def test_staged_authority_accepts_valid_evidence_states(
tmp_path: Path,
monkeypatch,
state: str,
):
monkeypatch.setattr(lanes, "STATE_ROOT", tmp_path / "lanes")
identity = lanes.TerminalIdentity("board", "task", 5, state)
document = _pending_terminal_record("board", "task", 5, "accepted")
document["kanban_state"] = state
staged, canonical = _staged_path(lanes.STATE_ROOT, identity, document)
staged.parent.mkdir(parents=True)
authority = lanes._staged_terminal_authority(staged, _Snapshot(document))
assert authority == (identity, canonical)
def test_staged_authority_rejects_outside_nested_and_unknown_state(
tmp_path: Path,
monkeypatch,
):
monkeypatch.setattr(lanes, "STATE_ROOT", tmp_path / "lanes")
document = _pending_terminal_record("board", "task", 5, "accepted")
outside = tmp_path / ".retire.x"
assert lanes._staged_terminal_authority(outside, _Snapshot(document)) is None
nested = lanes.STATE_ROOT / "board" / "nested" / ".retire.name"
assert lanes._staged_terminal_authority(nested, _Snapshot(document)) is None
board = lanes.STATE_ROOT / "board"
board.mkdir(parents=True)
staged = board / f".retire.{'a' * 16}.{'b' * 16}.0"
document["kanban_state"] = "unknown"
assert lanes._staged_terminal_authority(staged, _Snapshot(document)) is None
def test_restore_staged_authority_handles_every_first_writer_state(
tmp_path: Path,
monkeypatch,
):
document = _pending_terminal_record("board", "task", 6, "accepted")
canonical = tmp_path / "canonical"
pending = lanes.TerminalIdentity("board", "task", 6, "pending")
calls = []
monkeypatch.setattr(
lanes,
"_persist_prepared_evidence",
lambda identity, value: calls.append(("prepared", identity, value)),
)
lanes._restore_staged_terminal_authority(pending, canonical, document)
assert calls[0][0] == "prepared"
monkeypatch.setattr(lanes, "_write_json_noreplace", lambda *_args: True)
committed = lanes.TerminalIdentity("board", "task", 6, "committed")
lanes._restore_staged_terminal_authority(committed, canonical, document)
monkeypatch.setattr(lanes, "_write_json_noreplace", lambda *_args: False)
monkeypatch.setattr(lanes, "_load_small_json", lambda _path: {})
monkeypatch.setattr(
lanes,
"_persist_conflict_evidence",
lambda *args: calls.append(("conflict", *args)),
)
lanes._restore_staged_terminal_authority(committed, canonical, document)
assert calls[-1][0] == "conflict"
prepared = lanes.TerminalIdentity("board", "task", 6, "prepared")
monkeypatch.setattr(lanes, "_terminal_evidence_valid", lambda *_args: False)
with pytest.raises(OSError, match="conflicting result"):
lanes._restore_staged_terminal_authority(prepared, canonical, document)
def test_restore_staged_authority_accepts_identical_existing_evidence(
tmp_path: Path,
monkeypatch,
):
document = _pending_terminal_record("board", "task", 7, "accepted")
document["kanban_state"] = "prepared"
identity = lanes.TerminalIdentity("board", "task", 7, "prepared")
monkeypatch.setattr(lanes, "_write_json_noreplace", lambda *_args: False)
monkeypatch.setattr(lanes, "_load_small_json", lambda _path: document)
monkeypatch.setattr(lanes, "_terminal_evidence_valid", lambda *_args: True)
lanes._restore_staged_terminal_authority(identity, tmp_path / "existing", document)
def test_staging_recovery_bounds_missing_snapshots_and_restore_errors(
monkeypatch,
):
path = Path("/tmp/staged")
monkeypatch.setattr(lanes, "_retirement_staging_paths", lambda: [path])
monkeypatch.setattr(lanes, "_open_terminal_recovery_snapshot", lambda _path: None)
assert lanes._recover_retirement_staging() == 0
identity = lanes.TerminalIdentity("board", "task", 8, "pending")
snapshot = _Snapshot(_pending_terminal_record("board", "task", 8, "accepted"))
monkeypatch.setattr(lanes, "_open_terminal_recovery_snapshot", lambda _path: snapshot)
monkeypatch.setattr(
lanes,
"_staged_terminal_authority",
lambda *_args: (identity, Path("/tmp/canonical")),
)
monkeypatch.setattr(
lanes,
"_restore_staged_terminal_authority",
lambda *_args: (_ for _ in ()).throw(OSError("restore")),
)
errors = []
monkeypatch.setattr(
lanes,
"_record_board_access_error",
lambda board, error: errors.append((board, str(error))),
)
assert lanes._recover_retirement_staging() == 0
assert errors == [("board", "restore")]
assert snapshot.closed == 1
def test_drain_staging_ignores_disappearing_entries(monkeypatch):
class Vanished:
def stat(self, **_kwargs):
raise FileNotFoundError
monkeypatch.setattr(lanes, "_retirement_staging_paths", lambda: [Vanished()])
assert lanes._drain_retirement_staging() == 0
def _prepared_path(tmp_path: Path, monkeypatch) -> tuple[Path, dict]:
monkeypatch.setattr(lanes, "STATE_ROOT", tmp_path / "lanes")
pending = _pending_terminal_record("board", "task", 9, "accepted")
identity = lanes.TerminalIdentity("board", "task", 9, "pending")
path = lanes._persist_prepared_evidence(identity, pending)
return path, pending
def test_prepared_recovery_defers_without_discarding_evidence(
tmp_path: Path,
monkeypatch,
):
path, _document = _prepared_path(tmp_path, monkeypatch)
monkeypatch.setattr(lanes, "_finalize_document_db", lambda *_args: "deferred")
assert lanes._recover_prepared_finalizations(object()) == 0
assert path.exists()
def test_prepared_recovery_bounds_db_and_publication_exceptions(
tmp_path: Path,
monkeypatch,
):
path, _document = _prepared_path(tmp_path, monkeypatch)
errors = []
monkeypatch.setattr(
lanes,
"_record_board_access_error",
lambda board, error: errors.append((board, str(error))),
)
monkeypatch.setattr(
lanes,
"_finalize_document_db",
lambda *_args: (_ for _ in ()).throw(OSError("database")),
)
assert lanes._recover_prepared_finalizations(object()) == 0
assert errors[-1] == ("board", "database")
assert path.exists()
monkeypatch.setattr(lanes, "_finalize_document_db", lambda *_args: "committed")
monkeypatch.setattr(
lanes,
"_promote_prepared_evidence",
lambda *_args, **_kwargs: (_ for _ in ()).throw(OSError("publish")),
)
assert lanes._recover_prepared_finalizations(object()) == 0
assert errors[-1] == ("board", "publish")
def test_pending_recovery_bounds_finalizer_errors_and_replacement_churn(
tmp_path: Path,
monkeypatch,
capsys,
):
monkeypatch.setattr(lanes, "STATE_ROOT", tmp_path / "lanes")
path = lanes._terminal_path(lanes.state_path("board", "task"), 10)
lanes.atomic_json(path, _pending_terminal_record("board", "task", 10, "accepted"))
monkeypatch.setitem(sys.modules, "hermes_cli", SimpleNamespace(kanban_db=object()))
errors = []
monkeypatch.setattr(
lanes,
"_record_board_access_error",
lambda board, error: errors.append((board, str(error))),
)
monkeypatch.setattr(
lanes,
"_finalize_terminal_record",
lambda *_args, **_kwargs: (_ for _ in ()).throw(OSError("finalize")),
)
assert lanes.recover_pending_finalizations() == 0
assert errors[-1] == ("board", "finalize")
def invalid(*_args, snapshot=None, **_kwargs):
snapshot.close()
return "invalid"
monkeypatch.setattr(lanes, "_finalize_terminal_record", invalid)
assert lanes.recover_pending_finalizations() == 0
assert "replacement churn deferred" in capsys.readouterr().err
def test_pending_recovery_classifies_foreign_identity(
tmp_path: Path,
monkeypatch,
):
monkeypatch.setattr(lanes, "STATE_ROOT", tmp_path / "lanes")
path = lanes._terminal_path(lanes.state_path("board", "task"), 11)
lanes.atomic_json(path, _pending_terminal_record("other", "task", 11, "accepted"))
monkeypatch.setitem(sys.modules, "hermes_cli", SimpleNamespace(kanban_db=object()))
reasons = []
def quarantine(_path, _identity, reason, *, snapshot):
reasons.append(reason)
os.unlink(_path)
monkeypatch.setattr(lanes, "_quarantine_terminal", quarantine)
monkeypatch.setattr(lanes, "_recover_exact_run", lambda *_args: False)
assert lanes.recover_pending_finalizations() == 0
assert reasons == ["foreign-identity"]
def test_pending_guard_rejects_invalid_run_and_inaccessible_or_bad_artifacts(
tmp_path: Path,
monkeypatch,
):
monkeypatch.setattr(lanes, "STATE_ROOT", tmp_path / "lanes")
assert lanes._has_pending_finalization("board", "task", True) is False
path = lanes._terminal_path(lanes.state_path("board", "task"), 12)
path.parent.mkdir(parents=True)
path.write_text("bad", encoding="utf-8")
real_stat = Path.stat
def denied(candidate, *args, **kwargs):
if candidate == path:
raise PermissionError("denied")
return real_stat(candidate, *args, **kwargs)
monkeypatch.setattr(Path, "stat", denied)
assert lanes._has_pending_finalization("board", "task", 12) is False
def test_staged_authority_rejects_a_mismatched_pathname_digest(
tmp_path: Path,
monkeypatch,
):
"""A staged entry whose name does not bind its payload has no authority."""
monkeypatch.setattr(lanes, "STATE_ROOT", tmp_path / "lanes")
document = _pending_terminal_record("board", "task", 5, "accepted")
board = lanes.STATE_ROOT / "board"
board.mkdir(parents=True)
forged = board / f".retire.{'f' * 16}.{'a' * 16}.0"
assert lanes._staged_terminal_authority(forged, _Snapshot(document)) is None
def test_drain_staging_defers_unbounded_replacement_churn(
tmp_path: Path,
monkeypatch,
):
"""A staging surface that keeps changing stops after bounded attempts."""
monkeypatch.setattr(lanes, "STATE_ROOT", tmp_path / "lanes")
board = lanes.STATE_ROOT / "board"
board.mkdir(parents=True)
churning = board / ".retire.churning"
churning.write_text("0", encoding="utf-8")
def churn():
churning.write_text(churning.read_text(encoding="utf-8") + "x")
return 1
monkeypatch.setattr(lanes, "_recover_retirement_staging", churn)
assert lanes._drain_retirement_staging() == 16
def test_prepared_recovery_orders_and_bounds_disappearing_journals(
tmp_path: Path,
monkeypatch,
):
"""A journal whose stat fails is ordered last instead of aborting replay."""
path, _document = _prepared_path(tmp_path, monkeypatch)
real_stat = Path.stat
calls = []
def vanishing(candidate, *args, **kwargs):
if candidate == path and not calls:
calls.append("ordered")
raise OSError("stat raced")
return real_stat(candidate, *args, **kwargs)
monkeypatch.setattr(Path, "stat", vanishing)
monkeypatch.setattr(lanes, "_finalize_document_db", lambda *_args: "deferred")
assert lanes._recover_prepared_finalizations(object()) == 0
assert calls == ["ordered"]
assert path.exists()
def test_prepared_recovery_defers_a_journal_that_cannot_be_quarantined(
tmp_path: Path,
monkeypatch,
):
"""An unquarantinable malformed name never spins recovery forever."""
monkeypatch.setattr(lanes, "STATE_ROOT", tmp_path / "lanes")
board = lanes.STATE_ROOT / "board"
board.mkdir(parents=True)
misnamed = board / "task.run-9.terminal.prepared-zz.json"
misnamed.write_text("{}", encoding="utf-8")
attempts = []
monkeypatch.setattr(
lanes,
"_quarantine_terminal",
lambda *_args, **_kwargs: attempts.append("quarantine"),
)
assert lanes._recover_prepared_finalizations(object()) == 0
assert len(attempts) == 16
assert misnamed.exists()
def test_prepared_recovery_counts_one_recovery_per_exact_run(
tmp_path: Path,
monkeypatch,
):
"""Duplicate prepared evidence for one run converges to a single winner."""
monkeypatch.setattr(lanes, "STATE_ROOT", tmp_path / "lanes")
identity = lanes.TerminalIdentity("board", "task", 9, "pending")
first = _pending_terminal_record("board", "task", 9, "first accepted")
second = _pending_terminal_record("board", "task", 9, "second accepted")
lanes._persist_prepared_evidence(identity, first)
lanes._persist_prepared_evidence(identity, second)
task = SimpleNamespace(
id="task",
status="running",
current_run_id=9,
completed_run_id=None,
result=None,
assignee="cli-auto",
)
db = SimpleNamespace(
scoped_current_board=lambda _board: nullcontext(),
connect=lambda board: SimpleNamespace(close=lambda: None),
get_task=lambda _conn, _task_id: task,
complete_task=lambda *_args, **_kwargs: True,
)
assert lanes._recover_prepared_finalizations(db) == 1
conflicts = list((lanes.STATE_ROOT / "board").glob("*.terminal.conflict-*.json"))
assert len(conflicts) == 1
def test_pending_guard_ignores_symlinked_boards_and_invalid_artifacts(
tmp_path: Path,
monkeypatch,
):
"""Invalid journals, evidence, and staging never hold a claim open."""
monkeypatch.setattr(lanes, "STATE_ROOT", tmp_path / "lanes")
board = lanes.STATE_ROOT / "board"
board.mkdir(parents=True)
pending = lanes._terminal_path(lanes.state_path("board", "task"), 13)
pending.write_text("{}", encoding="utf-8")
misnamed = board / "task.run-13.terminal.prepared-zz.json"
misnamed.write_text("{}", encoding="utf-8")
invalid = board / f"task.run-13.terminal.prepared-{'a' * 32}.json"
invalid.write_text("{}", encoding="utf-8")
(board / ".retire.dangling").symlink_to(board / "missing-target")
foreign = _pending_terminal_record("board", "other", 13, "accepted")
staged, _canonical = _staged_path(
lanes.STATE_ROOT,
lanes.TerminalIdentity("board", "other", 13, "pending"),
foreign,
)
lanes.atomic_json(staged, foreign)
assert lanes._has_pending_finalization("board", "task", 13) is False
assert lanes._has_pending_finalization("board", "other", 13) is True
aliased_root = tmp_path / "aliased"
aliased_root.mkdir()
(aliased_root / "linked").symlink_to(board, target_is_directory=True)
monkeypatch.setattr(lanes, "STATE_ROOT", aliased_root)
linked_pending = aliased_root / "linked" / "task.run-13.terminal.pending.json"
assert linked_pending.exists()
assert lanes._has_pending_finalization("linked", "task", 13) is False