"""Tests for the MCP restart coordinator and impact analysis (#658). Multi-session fixtures exercise every verdict branch: safe, unsafe (live work), override, and the fail-closed deny on incomplete inventory. Also covers the critical-section deny path and the new ``ControlPlaneDB.list_sessions``. """ from __future__ import annotations import os import tempfile import unittest from datetime import datetime, timedelta, timezone import restart_coordinator as rc from control_plane_db import ControlPlaneDB NOW = datetime(2026, 7, 24, 6, 0, 0, tzinfo=timezone.utc) def _ts(dt: datetime) -> str: return dt.isoformat() def _live_pid() -> int: return os.getpid() def _dead_pid() -> int: # A pid that is essentially never alive. os.kill(0) on it raises # ProcessLookupError → is_process_alive False. return 2_000_000_000 def _session(session_id, *, pid, status="active", heartbeat=None, role="author"): return { "session_id": session_id, "role": role, "profile": "prgs-author", "pid": pid, "status": status, "last_heartbeat_at": _ts(heartbeat or NOW), } def _lease( lease_id, *, session_id, freshness, kind="issue", number=658, phase="allocated", worktree=None, role="author", ): return { "lease_id": lease_id, "session_id": session_id, "role": role, "phase": phase, "work_kind": kind, "work_number": number, "worktree_path": worktree, "freshness": {"freshness": freshness}, } class EvaluateRestartImpactTest(unittest.TestCase): def test_incomplete_inventory_denies_fail_closed(self) -> None: report = rc.evaluate_restart_impact( {"inventory_complete": False, "incomplete_reasons": ["db down"]}, now=NOW, ) self.assertEqual(report.verdict, rc.VERDICT_UNSAFE) self.assertFalse(report.allow_restart) self.assertFalse(report.restart_performed) self.assertIn("db down", report.incomplete_reasons) self.assertTrue( any("fail closed" in reasoning for reasoning in report.reasons) ) def test_missing_completeness_flag_denies(self) -> None: # No inventory_complete key at all → treated as incomplete. report = rc.evaluate_restart_impact({}, now=NOW) self.assertEqual(report.verdict, rc.VERDICT_UNSAFE) self.assertFalse(report.allow_restart) def test_no_other_work_is_safe(self) -> None: report = rc.evaluate_restart_impact( { "inventory_complete": True, "sessions": [_session("requester", pid=_live_pid())], "leases": [], }, now=NOW, requesting_session_id="requester", ) self.assertEqual(report.verdict, rc.VERDICT_SAFE) self.assertTrue(report.allow_restart) self.assertEqual(report.blast_radius, rc.BLAST_NONE) self.assertEqual(report.affected_issues, []) def test_dead_foreign_session_and_lease_are_not_disruptive(self) -> None: report = rc.evaluate_restart_impact( { "inventory_complete": True, "sessions": [ _session("requester", pid=_live_pid()), _session("dead", pid=_dead_pid()), ], "leases": [ _lease("l-dead", session_id="dead", freshness="stale_dead_process") ], }, now=NOW, requesting_session_id="requester", ) self.assertEqual(report.verdict, rc.VERDICT_SAFE) self.assertTrue(report.allow_restart) self.assertEqual(report.counts["leases_disruptive"], 0) self.assertEqual(report.counts["sessions_live_other"], 0) def test_live_foreign_lease_denies_without_override(self) -> None: report = rc.evaluate_restart_impact( { "inventory_complete": True, "sessions": [ _session("requester", pid=_live_pid()), _session("worker", pid=_live_pid()), ], "leases": [ _lease( "l1", session_id="worker", freshness="active", worktree="/tmp/wt-658", phase="implementing", ) ], }, now=NOW, requesting_session_id="requester", ) self.assertEqual(report.verdict, rc.VERDICT_UNSAFE) self.assertFalse(report.allow_restart) # Critical section detected: active lease with a live owner. self.assertEqual(len(report.critical_sections), 1) self.assertEqual(report.affected_issues, [658]) self.assertEqual(report.counts["mutations"], 1) self.assertTrue(report.override_would_allow) self.assertEqual(report.blast_radius, rc.BLAST_HIGH) # Placeholder ack state for the affected session. self.assertEqual(report.ack_state.get("worker"), "pending") def test_operator_override_allows_despite_live_work(self) -> None: inv = { "inventory_complete": True, "sessions": [ _session("requester", pid=_live_pid()), _session("worker", pid=_live_pid()), ], "leases": [_lease("l1", session_id="worker", freshness="active")], } report = rc.evaluate_restart_impact( inv, now=NOW, requesting_session_id="requester", operator_override=True, ) self.assertEqual(report.verdict, rc.VERDICT_OVERRIDE) self.assertTrue(report.allow_restart) self.assertFalse(report.restart_performed) def test_deny_when_critical_section_open(self) -> None: # A single live author lease in a mutating phase is a critical section # that must deny an un-overridden restart. report = rc.evaluate_restart_impact( { "inventory_complete": True, "sessions": [_session("worker", pid=_live_pid())], "leases": [ _lease( "l1", session_id="worker", freshness="active", phase="merging", kind="pr", number=900, ) ], }, now=NOW, requesting_session_id="requester", ) self.assertEqual(report.verdict, rc.VERDICT_UNSAFE) self.assertFalse(report.allow_restart) self.assertEqual(report.affected_prs, [900]) self.assertEqual(len(report.critical_sections), 1) def test_terminal_lock_makes_restart_unsafe(self) -> None: report = rc.evaluate_restart_impact( { "inventory_complete": True, "sessions": [_session("requester", pid=_live_pid())], "leases": [], "terminal_lock": {"terminal_pr": 812}, }, now=NOW, requesting_session_id="requester", ) self.assertEqual(report.verdict, rc.VERDICT_UNSAFE) self.assertFalse(report.allow_restart) self.assertIsNotNone(report.terminal_lock) self.assertTrue( any("terminal" in reasoning for reasoning in report.reasons) ) def test_other_live_session_without_lease_is_disruptive(self) -> None: report = rc.evaluate_restart_impact( { "inventory_complete": True, "sessions": [ _session("requester", pid=_live_pid()), _session("idle-but-live", pid=_live_pid()), ], "leases": [], }, now=NOW, requesting_session_id="requester", ) self.assertEqual(report.verdict, rc.VERDICT_UNSAFE) self.assertEqual(report.counts["sessions_live_other"], 1) def test_stale_heartbeat_session_not_counted_live(self) -> None: stale = NOW - timedelta(hours=2) report = rc.evaluate_restart_impact( { "inventory_complete": True, "sessions": [ _session("requester", pid=_live_pid()), _session("stale", pid=_live_pid(), heartbeat=stale), ], "leases": [], }, now=NOW, requesting_session_id="requester", ) self.assertEqual(report.verdict, rc.VERDICT_SAFE) self.assertEqual(report.counts["sessions_live_other"], 0) def test_prior_recovery_attempts_echoed(self) -> None: report = rc.evaluate_restart_impact( { "inventory_complete": True, "sessions": [_session("requester", pid=_live_pid())], "leases": [], "prior_recovery_attempts": [ {"kind": "client_reconnect", "at": _ts(NOW)} ], }, now=NOW, requesting_session_id="requester", ) self.assertEqual(len(report.prior_recovery_attempts), 1) self.assertEqual(report.counts["prior_recovery_attempts"], 1) def test_bare_string_freshness_accepted(self) -> None: lease = _lease("l1", session_id="worker", freshness="active") lease["freshness"] = "active" # bare string, not a dict report = rc.evaluate_restart_impact( { "inventory_complete": True, "sessions": [_session("worker", pid=_live_pid())], "leases": [lease], }, now=NOW, requesting_session_id="requester", ) self.assertEqual(report.counts["leases_disruptive"], 1) def test_as_dict_is_serializable_dto(self) -> None: import json report = rc.evaluate_restart_impact( { "inventory_complete": True, "sessions": [_session("requester", pid=_live_pid())], "leases": [], }, now=NOW, requesting_session_id="requester", ) payload = report.as_dict() # Round-trips through JSON — safe for the console DTO. encoded = json.dumps(payload) decoded = json.loads(encoded) self.assertEqual(decoded["verdict"], rc.VERDICT_SAFE) self.assertIn("audit_record", decoded) self.assertEqual(decoded["audit_record"]["event"], "restart_impact_evaluated") self.assertFalse(decoded["restart_performed"]) self.assertIn("coordinator_version", decoded) class ListSessionsTest(unittest.TestCase): def setUp(self) -> None: self._tmp = tempfile.TemporaryDirectory() self.db = ControlPlaneDB(os.path.join(self._tmp.name, "cp.sqlite3")) def tearDown(self) -> None: self._tmp.cleanup() def test_list_sessions_filters_by_status(self) -> None: self.db.upsert_session(session_id="a", role="author", pid=1, status="active") self.db.upsert_session(session_id="b", role="author", pid=2, status="ended") active = self.db.list_sessions(statuses=("active",)) ids = {row["session_id"] for row in active} self.assertEqual(ids, {"a"}) every = self.db.list_sessions() self.assertEqual({row["session_id"] for row in every}, {"a", "b"}) def test_list_sessions_feeds_coordinator(self) -> None: self.db.upsert_session( session_id="requester", role="author", pid=os.getpid(), status="active" ) report = rc.evaluate_restart_impact( { "inventory_complete": True, "sessions": self.db.list_sessions(statuses=("active",)), "leases": [], }, now=NOW, requesting_session_id="requester", ) self.assertEqual(report.counts["sessions_total"], 1) if __name__ == "__main__": # pragma: no cover unittest.main()