"""MCP restart coordinator and impact analysis (#658). Before any sanctioned MCP restart, a central coordinator must evaluate the live control-plane state — active sessions, leases/locks, in-flight issue/PR work, mutations, worktrees, and recovery history — and produce an *impact preview* so operators (and the web console, #642/#652) can see the blast radius **before** concurrent LLM work is disrupted. Design rules (mirrors the read-only posture of ``workflow_dashboard`` / ``lease_lifecycle``): * **Pure classification.** :func:`evaluate_restart_impact` takes an already gathered inventory and returns a structured report. It never touches the network, the filesystem, or a live process, so multi-session fixtures can drive every branch in unit tests. The coordinator *never restarts anything*; a mutative apply path is a later child gated by a drain proof (non-goal here). * **Fail closed.** If the inventory is not explicitly complete, the verdict is ``unsafe`` / deny — an incomplete evaluation must never green-light a restart. * **No secrets.** Session ids, pids, and profiles are operational metadata, not credentials; nothing secret flows through this module. The single sanctioned entry point post-#657 is the MCP tool ``gitea_request_mcp_restart`` (dry-run by default), which gathers the inventory from the #613 control-plane DB and calls :func:`evaluate_restart_impact`. """ from __future__ import annotations from dataclasses import dataclass, field from datetime import datetime, timezone from enum import Enum from typing import Any, Mapping, Sequence import lease_lifecycle COORDINATOR_VERSION = "1.1.0-issue-663" # Restart verdicts. Exactly the three the acceptance criteria name. VERDICT_SAFE = "safe" VERDICT_UNSAFE = "unsafe" VERDICT_OVERRIDE = "override" # Blast-radius severity bands. BLAST_NONE = "none" BLAST_LOW = "low" BLAST_MEDIUM = "medium" BLAST_HIGH = "high" # A live lease with a live owner process is treated as active in-flight work. LEASE_FRESHNESS_LIVE = "active" # Default staleness window for a session heartbeat (seconds). A session whose # last heartbeat is older than this is not counted as live even if its row is # still marked ``active`` — it is assumed dead/detached. DEFAULT_SESSION_HEARTBEAT_STALE_SECONDS = 900 class RestartClass(str, Enum): """The only restart/recovery classes accepted by the coordinator.""" CLIENT_RECONNECT = "client_reconnect" SESSION_RECONNECT = "session_reconnect" WORKER_RESTART = "worker_restart" ROLE_RUNTIME_RESTART = "role_runtime_restart" CONNECTOR_RESTART = "connector_restart" CONFIGURATION_RELOAD = "configuration_reload" ROLLING_MCP_RESTART = "rolling_mcp_restart" FULL_MCP_RESTART = "full_mcp_restart" HOST_RESTART = "host_restart" @dataclass(frozen=True) class RestartClassPolicy: """Least-privilege policy for one :class:`RestartClass`.""" restart_class: RestartClass required_permission: str expected_blast_radius: str drain_requirement: str full_drain_required: bool approval_requirement: str audit_requirement: str recovery_behavior: str request_roles: tuple[str, ...] execution_roles: tuple[str, ...] def as_dict(self) -> dict[str, Any]: return { "restart_class": self.restart_class.value, "required_permission": self.required_permission, "expected_blast_radius": self.expected_blast_radius, "drain_requirement": self.drain_requirement, "full_drain_required": self.full_drain_required, "approval_requirement": self.approval_requirement, "audit_requirement": self.audit_requirement, "recovery_behavior": self.recovery_behavior, "request_roles": list(self.request_roles), "execution_roles": list(self.execution_roles), } WORKER_ROLES = ("author", "reviewer", "merger", "reconciler") CONTROL_ROLES = ("controller", "operator", "admin") ALL_REQUEST_ROLES = WORKER_ROLES + CONTROL_ROLES RESTART_CLASS_POLICIES: dict[RestartClass, RestartClassPolicy] = { RestartClass.CLIENT_RECONNECT: RestartClassPolicy( RestartClass.CLIENT_RECONNECT, "mcp.reconnect.client", BLAST_NONE, "none", False, "self_service", "record class, actor, client namespace, reason, and outcome", "Reconnect only the caller's client transport; no daemon or peer session changes.", ALL_REQUEST_ROLES, ALL_REQUEST_ROLES, ), RestartClass.SESSION_RECONNECT: RestartClassPolicy( RestartClass.SESSION_RECONNECT, "mcp.reconnect.session", BLAST_LOW, "requesting_session_safe_point", False, "self_service", "record class, actor, session id, reason, and outcome", "Rebind identity, capability, and workspace state for one session.", ALL_REQUEST_ROLES, ALL_REQUEST_ROLES, ), RestartClass.WORKER_RESTART: RestartClassPolicy( RestartClass.WORKER_RESTART, "mcp.restart.worker.request", BLAST_LOW, "target_worker", False, "controller_approval_and_automated_gates", "record class, actor, target worker, approval, drain proof, and outcome", "Restart one worker after its own lease and mutation scope is drained.", ALL_REQUEST_ROLES, ("operator", "admin"), ), RestartClass.ROLE_RUNTIME_RESTART: RestartClassPolicy( RestartClass.ROLE_RUNTIME_RESTART, "mcp.restart.role_runtime.request", BLAST_MEDIUM, "target_role_runtime", False, "controller_approval_and_automated_gates", "record class, actor, role namespace, approval, drain proof, and outcome", "Restart only the selected role runtime and then re-probe that namespace.", ALL_REQUEST_ROLES, ("operator", "admin"), ), RestartClass.CONNECTOR_RESTART: RestartClassPolicy( RestartClass.CONNECTOR_RESTART, "mcp.restart.connector.request", BLAST_MEDIUM, "target_connector", False, "controller_approval_and_automated_gates", "record class, actor, connector id, approval, drain proof, and outcome", "Restart one connector while unrelated role runtimes remain available.", ALL_REQUEST_ROLES, ("operator", "admin"), ), RestartClass.CONFIGURATION_RELOAD: RestartClassPolicy( RestartClass.CONFIGURATION_RELOAD, "mcp.reload.configuration.request", BLAST_LOW, "mutation_quiesce", False, "controller_approval_and_automated_gates", "record class, actor, configuration revision, approval, and outcome", "Gracefully reload configuration without replacing the daemon process.", ALL_REQUEST_ROLES, ("operator", "admin"), ), RestartClass.ROLLING_MCP_RESTART: RestartClassPolicy( RestartClass.ROLLING_MCP_RESTART, "mcp.restart.rolling.request", BLAST_MEDIUM, "one_instance_at_a_time", False, "controller_approval_and_automated_gates", "record class, actor, instance order, approval, per-instance drains, and outcome", "Drain, restart, verify, and restore one instance before advancing to the next.", CONTROL_ROLES, ("operator", "admin"), ), RestartClass.FULL_MCP_RESTART: RestartClassPolicy( RestartClass.FULL_MCP_RESTART, "mcp.restart.full.request", BLAST_HIGH, "all_sessions_and_mutations", True, "controller_approval_and_automated_gates", "record class, actor, full impact report, approval, drain proof, and outcome", "Stop and restore the complete MCP runtime only after a verified full drain.", CONTROL_ROLES, ("operator", "admin"), ), RestartClass.HOST_RESTART: RestartClassPolicy( RestartClass.HOST_RESTART, "mcp.restart.host.request", BLAST_HIGH, "all_host_work", True, "controller_approval_plus_infrastructure_operator", "record class, actor, host, incident or change id, approval, drain proof, and outcome", "Hand off to infrastructure ownership; reconcile every runtime after the host returns.", ("controller", "operator", "admin"), ("operator", "admin"), ), } def resolve_restart_class(value: RestartClass | str) -> RestartClass: """Resolve a restart class or fail closed for an unknown value.""" if isinstance(value, RestartClass): return value try: return RestartClass(str(value).strip()) except ValueError as exc: raise ValueError(f"unknown restart class {value!r}; deny (fail closed)") from exc def restart_class_policy(value: RestartClass | str) -> RestartClassPolicy: """Return the canonical policy for *value*.""" return RESTART_CLASS_POLICIES[resolve_restart_class(value)] def permissions_for_role(role: str | None) -> tuple[str, ...]: """Return request permissions granted to a workflow role by this policy.""" normalized = str(role or "").strip().lower() return tuple( policy.required_permission for policy in RESTART_CLASS_POLICIES.values() if normalized in policy.request_roles ) def _utc_now() -> datetime: return datetime.now(timezone.utc) def _parse_ts(value: str | None) -> datetime | None: return lease_lifecycle._parse_ts(value) @dataclass(frozen=True) class SessionImpact: """One MCP session a restart would terminate.""" session_id: str role: str | None profile: str | None pid: int | None status: str | None alive: bool | None heartbeat_stale: bool is_requester: bool live: bool connector: str | None = None def as_dict(self) -> dict[str, Any]: return { "session_id": self.session_id, "role": self.role, "profile": self.profile, "pid": self.pid, "status": self.status, "alive": self.alive, "heartbeat_stale": self.heartbeat_stale, "is_requester": self.is_requester, "live": self.live, "connector": self.connector, } @dataclass(frozen=True) class LeaseImpact: """One control-plane lease a restart would disrupt.""" lease_id: str | None session_id: str | None role: str | None phase: str | None freshness: str | None work_kind: str | None work_number: int | None worktree_path: str | None disruptive: bool is_mutation: bool is_critical_section: bool connector: str | None = None def as_dict(self) -> dict[str, Any]: return { "lease_id": self.lease_id, "session_id": self.session_id, "role": self.role, "phase": self.phase, "freshness": self.freshness, "work_kind": self.work_kind, "work_number": self.work_number, "worktree_path": self.worktree_path, "disruptive": self.disruptive, "is_mutation": self.is_mutation, "is_critical_section": self.is_critical_section, "connector": self.connector, } @dataclass(frozen=True) class RestartImpactReport: """Impact preview DTO returned to the console / operator (#642/#652).""" coordinator_version: str restart_class: str restart_policy: dict[str, Any] policy_enforced: bool permission_authorized: bool role_authorized: bool approval_satisfied: bool authorization_reasons: list[str] evaluated_at: str dry_run: bool restart_performed: bool inventory_complete: bool verdict: str allow_restart: bool override_would_allow: bool operator_override: bool blast_radius: str reasons: list[str] affected_sessions: list[SessionImpact] affected_leases: list[LeaseImpact] critical_sections: list[LeaseImpact] affected_issues: list[int] affected_prs: list[int] mutations: list[LeaseImpact] terminal_lock: dict[str, Any] | None ack_state: dict[str, str] prior_recovery_attempts: list[dict[str, Any]] counts: dict[str, int] audit_record: dict[str, Any] incomplete_reasons: list[str] = field(default_factory=list) def as_dict(self) -> dict[str, Any]: return { "coordinator_version": self.coordinator_version, "restart_class": self.restart_class, "restart_policy": dict(self.restart_policy), "policy_enforced": self.policy_enforced, "permission_authorized": self.permission_authorized, "role_authorized": self.role_authorized, "approval_satisfied": self.approval_satisfied, "authorization_reasons": list(self.authorization_reasons), "evaluated_at": self.evaluated_at, "dry_run": self.dry_run, "restart_performed": self.restart_performed, "inventory_complete": self.inventory_complete, "incomplete_reasons": list(self.incomplete_reasons), "verdict": self.verdict, "allow_restart": self.allow_restart, "override_would_allow": self.override_would_allow, "operator_override": self.operator_override, "blast_radius": self.blast_radius, "reasons": list(self.reasons), "affected_sessions": [s.as_dict() for s in self.affected_sessions], "affected_leases": [l.as_dict() for l in self.affected_leases], "critical_sections": [l.as_dict() for l in self.critical_sections], "affected_issues": list(self.affected_issues), "affected_prs": list(self.affected_prs), "mutations": [l.as_dict() for l in self.mutations], "terminal_lock": self.terminal_lock, "ack_state": dict(self.ack_state), "prior_recovery_attempts": list(self.prior_recovery_attempts), "counts": dict(self.counts), "audit_record": dict(self.audit_record), } def _classify_session( row: Mapping[str, Any], *, now: datetime, requesting_session_id: str | None, heartbeat_stale_seconds: int, ) -> SessionImpact: session_id = str(row.get("session_id") or "") pid = row.get("pid") status = (row.get("status") or "").strip().lower() or None alive = lease_lifecycle.is_process_alive(pid) if pid is not None else None hb = _parse_ts(row.get("last_heartbeat_at")) heartbeat_stale = bool( hb is not None and (now - hb).total_seconds() > heartbeat_stale_seconds ) live = bool(status == "active" and alive is not False and not heartbeat_stale) return SessionImpact( session_id=session_id, role=row.get("role"), profile=row.get("profile"), pid=pid, status=status, alive=alive, heartbeat_stale=heartbeat_stale, is_requester=bool( requesting_session_id and session_id == requesting_session_id ), live=live, connector=(str(row.get("connector") or "").strip() or None), ) # Lease phases that represent an active mutation in flight (as opposed to a # mere allocation/claim with no work committed yet). An active lease in any of # these phases is a critical section a restart must not sever. _MUTATING_PHASES = frozenset( { "implementing", "publishing", "pushing", "committing", "reviewing", "merging", "reconciling", "conflict_fix", } ) def _classify_lease(row: Mapping[str, Any]) -> LeaseImpact: freshness_obj = row.get("freshness") if isinstance(freshness_obj, Mapping): freshness = str(freshness_obj.get("freshness") or "").strip().lower() or None else: freshness = str(freshness_obj or "").strip().lower() or None phase = (row.get("phase") or "").strip().lower() or None worktree = row.get("worktree_path") disruptive = freshness == LEASE_FRESHNESS_LIVE # A live lease is a mutation-in-flight if it carries an author worktree or # its phase names a mutating step. All disruptive leases are critical # sections a restart would sever regardless. is_mutation = bool( disruptive and (bool(worktree) or (phase in _MUTATING_PHASES)) ) number = row.get("work_number") try: number = int(number) if number is not None else None except (TypeError, ValueError): number = None return LeaseImpact( lease_id=row.get("lease_id"), session_id=row.get("session_id"), role=row.get("role"), phase=phase, freshness=freshness, work_kind=(str(row.get("work_kind") or "").strip().lower() or None), work_number=number, worktree_path=worktree, disruptive=disruptive, is_mutation=is_mutation, is_critical_section=disruptive, connector=(str(row.get("connector") or "").strip() or None), ) def _blast_radius(*, session_count: int, work_count: int, mutation_count: int) -> str: if mutation_count > 0 or work_count >= 3 or session_count >= 3: return BLAST_HIGH if work_count > 0 or session_count == 2: return BLAST_MEDIUM if session_count == 1: return BLAST_LOW return BLAST_NONE def evaluate_restart_impact( inventory: Mapping[str, Any], *, now: datetime | None = None, operator_override: bool = False, requesting_session_id: str | None = None, dry_run: bool = True, session_heartbeat_stale_seconds: int = DEFAULT_SESSION_HEARTBEAT_STALE_SECONDS, restart_class: RestartClass | str | None = None, requester_role: str | None = None, requester_permissions: Sequence[str] | None = None, controller_approved: bool = False, operator_authorized: bool = False, target_session_id: str | None = None, target_role: str | None = None, target_connector: str | None = None, ) -> RestartImpactReport: """Evaluate a proposed MCP restart and return an impact preview. ``inventory`` is a mapping with: * ``sessions`` — session rows (session_id, role, profile, pid, status, last_heartbeat_at). * ``leases`` — control-plane lease rows, each ideally carrying an enriched ``freshness`` dict (as :func:`lease_lifecycle.list_active_leases` returns); a bare string freshness is also accepted. * ``terminal_lock`` — the active terminal (merge) lock, if any. * ``prior_recovery_attempts`` — narrower recovery attempts already tried (e.g. sanctioned reconnects) so the operator sees escalation history. * ``inventory_complete`` — bool. **Must** be explicitly True; a missing or falsy value forces a deny (fail closed). * ``incomplete_reasons`` — optional reasons the inventory is incomplete. The coordinator never restarts anything: ``restart_performed`` is always False and the mutative apply path is a later drain-gated child. """ moment = now or _utc_now() reasons: list[str] = [] authorization_reasons: list[str] = [] # ``None`` preserves the pre-#663 impact-only API for callers that have not # yet been migrated. All MCP requests pass an explicit class and therefore # take the fail-closed policy path. policy_enforced = restart_class is not None try: resolved_class = resolve_restart_class( restart_class or RestartClass.FULL_MCP_RESTART ) policy = RESTART_CLASS_POLICIES[resolved_class] unknown_class = False except ValueError as exc: resolved_class = None policy = None unknown_class = True authorization_reasons.append(str(exc)) normalized_role = str(requester_role or "").strip().lower() granted = {str(p).strip() for p in (requester_permissions or ())} if policy_enforced and policy is not None: permission_authorized = policy.required_permission in granted role_authorized = normalized_role in policy.request_roles if not permission_authorized: authorization_reasons.append( f"missing required permission {policy.required_permission!r}" ) if not role_authorized: authorization_reasons.append( f"role {normalized_role or 'unknown'!r} may not request " f"{policy.restart_class.value}" ) elif unknown_class: permission_authorized = False role_authorized = False else: permission_authorized = True role_authorized = True if policy_enforced and policy is not None: approval = policy.approval_requirement if approval == "self_service": approval_satisfied = True elif approval == "controller_approval_plus_infrastructure_operator": approval_satisfied = bool(controller_approved and operator_authorized) else: approval_satisfied = bool(controller_approved) if not approval_satisfied: authorization_reasons.append( f"approval requirement not satisfied: {approval}" ) elif unknown_class: approval_satisfied = False else: approval_satisfied = True inventory_complete = bool(inventory.get("inventory_complete", False)) incomplete_reasons = [str(r) for r in (inventory.get("incomplete_reasons") or [])] sessions_raw: Sequence[Mapping[str, Any]] = inventory.get("sessions") or [] leases_raw: Sequence[Mapping[str, Any]] = inventory.get("leases") or [] terminal_lock = inventory.get("terminal_lock") or None prior_recovery_attempts = [ dict(a) for a in (inventory.get("prior_recovery_attempts") or []) ] session_impacts = [ _classify_session( s, now=moment, requesting_session_id=requesting_session_id, heartbeat_stale_seconds=session_heartbeat_stale_seconds, ) for s in sessions_raw ] lease_impacts = [_classify_lease(l) for l in leases_raw] # Route impact through the selected class. Narrow classes never inherit a # full-runtime drain merely because unrelated work exists. target_complete = True if resolved_class in { RestartClass.CLIENT_RECONNECT, RestartClass.SESSION_RECONNECT, RestartClass.CONFIGURATION_RELOAD, }: scoped_sessions: list[SessionImpact] = [] scoped_leases: list[LeaseImpact] = [] elif resolved_class == RestartClass.WORKER_RESTART: selected_session = (target_session_id or "").strip() target_complete = bool(selected_session) scoped_sessions = [ s for s in session_impacts if s.session_id == selected_session ] scoped_leases = [ l for l in lease_impacts if l.session_id == selected_session ] elif resolved_class == RestartClass.ROLE_RUNTIME_RESTART: selected_role = (target_role or "").strip().lower() target_complete = bool(selected_role) scoped_sessions = [ s for s in session_impacts if str(s.role or "").lower() == selected_role ] scoped_leases = [ l for l in lease_impacts if str(l.role or "").lower() == selected_role ] elif resolved_class == RestartClass.CONNECTOR_RESTART: selected_connector = (target_connector or "").strip() target_complete = bool(selected_connector) scoped_sessions = [ s for s in session_impacts if s.connector == selected_connector ] scoped_leases = [ l for l in lease_impacts if l.connector == selected_connector ] else: scoped_sessions = list(session_impacts) scoped_leases = list(lease_impacts) if policy_enforced and not target_complete: authorization_reasons.append( f"target required for {resolved_class.value if resolved_class else 'unknown class'}" ) other_live_sessions = [ s for s in scoped_sessions if s.live and not s.is_requester ] disruptive_leases = [l for l in scoped_leases if l.disruptive] critical_sections = [l for l in scoped_leases if l.is_critical_section] mutations = [l for l in scoped_leases if l.is_mutation] terminal_lock_in_scope = ( terminal_lock if resolved_class not in { RestartClass.CLIENT_RECONNECT, RestartClass.SESSION_RECONNECT, } else None ) affected_issues = sorted( { l.work_number for l in disruptive_leases if l.work_kind == "issue" and l.work_number is not None } ) affected_prs = sorted( { l.work_number for l in disruptive_leases if l.work_kind == "pr" and l.work_number is not None } ) disruptive = bool( disruptive_leases or other_live_sessions or terminal_lock_in_scope ) authorization_ok = bool( not unknown_class and permission_authorized and role_authorized and approval_satisfied and target_complete ) if policy_enforced and not authorization_ok: verdict = VERDICT_UNSAFE allow_restart = False reasons.append("restart class authorization denied (fail closed)") reasons.extend(authorization_reasons) elif not inventory_complete: verdict = VERDICT_UNSAFE allow_restart = False reasons.append( "inventory incomplete: restart evaluation cannot confirm blast " "radius — deny (fail closed, #658)" ) reasons.extend(incomplete_reasons) elif not disruptive: verdict = VERDICT_SAFE allow_restart = True reasons.append("no other live sessions, live leases, or terminal lock") elif operator_override: verdict = VERDICT_OVERRIDE allow_restart = True reasons.append( "live work present; operator override accepts the blast radius" ) else: verdict = VERDICT_UNSAFE allow_restart = False reasons.append( "live work would be disrupted; restart denied without operator " "override" ) if critical_sections and inventory_complete: reasons.append( f"{len(critical_sections)} critical section(s) in flight " "(active lease with a live owner)" ) if terminal_lock_in_scope: reasons.append("active terminal (merge) lock present") override_would_allow = bool(inventory_complete and disruptive) blast_radius = _blast_radius( session_count=len(other_live_sessions), work_count=len(affected_issues) + len(affected_prs), mutation_count=len(mutations), ) # Acknowledgement is a later child (drain protocol); expose per-session # placeholders so the console can render the ack column now. ack_state = {s.session_id: "pending" for s in other_live_sessions} counts = { "sessions_total": len(session_impacts), "sessions_live_other": len(other_live_sessions), "leases_total": len(lease_impacts), "leases_disruptive": len(disruptive_leases), "critical_sections": len(critical_sections), "mutations": len(mutations), "affected_issues": len(affected_issues), "affected_prs": len(affected_prs), "prior_recovery_attempts": len(prior_recovery_attempts), } audit_record = { "event": "restart_impact_evaluated", "coordinator_version": COORDINATOR_VERSION, "restart_class": ( resolved_class.value if resolved_class else str(restart_class or "") ), "required_permission": ( policy.required_permission if policy is not None else None ), "evaluated_at": moment.isoformat(), "dry_run": dry_run, "operator_override": bool(operator_override), "requesting_session_id": requesting_session_id, "inventory_complete": inventory_complete, "verdict": verdict, "allow_restart": allow_restart, "blast_radius": blast_radius, "counts": counts, } return RestartImpactReport( coordinator_version=COORDINATOR_VERSION, restart_class=( resolved_class.value if resolved_class else str(restart_class or "") ), restart_policy=policy.as_dict() if policy is not None else {}, policy_enforced=policy_enforced, permission_authorized=permission_authorized, role_authorized=role_authorized, approval_satisfied=approval_satisfied, authorization_reasons=authorization_reasons, evaluated_at=moment.isoformat(), dry_run=dry_run, restart_performed=False, inventory_complete=inventory_complete, verdict=verdict, allow_restart=allow_restart, override_would_allow=override_would_allow, operator_override=bool(operator_override), blast_radius=blast_radius, reasons=reasons, affected_sessions=session_impacts, affected_leases=lease_impacts, critical_sections=critical_sections, affected_issues=affected_issues, affected_prs=affected_prs, mutations=mutations, terminal_lock=( dict(terminal_lock_in_scope) if isinstance(terminal_lock_in_scope, Mapping) else terminal_lock_in_scope ), ack_state=ack_state, prior_recovery_attempts=prior_recovery_attempts, counts=counts, audit_record=audit_record, incomplete_reasons=incomplete_reasons, )