feat: enforce MCP restart class permissions (#663)
This commit is contained in:
+370
-14
@@ -28,11 +28,12 @@ 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.0.0-issue-658"
|
||||
COORDINATOR_VERSION = "1.1.0-issue-663"
|
||||
|
||||
# Restart verdicts. Exactly the three the acceptance criteria name.
|
||||
VERDICT_SAFE = "safe"
|
||||
@@ -54,6 +55,194 @@ LEASE_FRESHNESS_LIVE = "active"
|
||||
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)
|
||||
|
||||
@@ -75,6 +264,7 @@ class SessionImpact:
|
||||
heartbeat_stale: bool
|
||||
is_requester: bool
|
||||
live: bool
|
||||
connector: str | None = None
|
||||
|
||||
def as_dict(self) -> dict[str, Any]:
|
||||
return {
|
||||
@@ -87,6 +277,7 @@ class SessionImpact:
|
||||
"heartbeat_stale": self.heartbeat_stale,
|
||||
"is_requester": self.is_requester,
|
||||
"live": self.live,
|
||||
"connector": self.connector,
|
||||
}
|
||||
|
||||
|
||||
@@ -105,6 +296,7 @@ class LeaseImpact:
|
||||
disruptive: bool
|
||||
is_mutation: bool
|
||||
is_critical_section: bool
|
||||
connector: str | None = None
|
||||
|
||||
def as_dict(self) -> dict[str, Any]:
|
||||
return {
|
||||
@@ -119,6 +311,7 @@ class LeaseImpact:
|
||||
"disruptive": self.disruptive,
|
||||
"is_mutation": self.is_mutation,
|
||||
"is_critical_section": self.is_critical_section,
|
||||
"connector": self.connector,
|
||||
}
|
||||
|
||||
|
||||
@@ -127,6 +320,13 @@ 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
|
||||
@@ -153,6 +353,13 @@ class RestartImpactReport:
|
||||
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,
|
||||
@@ -206,6 +413,7 @@ def _classify_session(
|
||||
requesting_session_id and session_id == requesting_session_id
|
||||
),
|
||||
live=live,
|
||||
connector=(str(row.get("connector") or "").strip() or None),
|
||||
)
|
||||
|
||||
|
||||
@@ -258,6 +466,7 @@ def _classify_lease(row: Mapping[str, Any]) -> LeaseImpact:
|
||||
disruptive=disruptive,
|
||||
is_mutation=is_mutation,
|
||||
is_critical_section=disruptive,
|
||||
connector=(str(row.get("connector") or "").strip() or None),
|
||||
)
|
||||
|
||||
|
||||
@@ -279,6 +488,14 @@ def evaluate_restart_impact(
|
||||
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.
|
||||
|
||||
@@ -301,6 +518,61 @@ def evaluate_restart_impact(
|
||||
"""
|
||||
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 [])]
|
||||
@@ -323,15 +595,67 @@ def evaluate_restart_impact(
|
||||
]
|
||||
lease_impacts = [_classify_lease(l) for l in leases_raw]
|
||||
|
||||
# Only *other* live sessions and live leases constitute blast radius: a
|
||||
# restart that would kill only the requesting session with no other work in
|
||||
# flight is safe.
|
||||
# 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 session_impacts if s.live and not s.is_requester
|
||||
s for s in scoped_sessions if s.live and not s.is_requester
|
||||
]
|
||||
disruptive_leases = [l for l in lease_impacts if l.disruptive]
|
||||
critical_sections = [l for l in lease_impacts if l.is_critical_section]
|
||||
mutations = [l for l in lease_impacts if l.is_mutation]
|
||||
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(
|
||||
{
|
||||
@@ -348,9 +672,24 @@ def evaluate_restart_impact(
|
||||
}
|
||||
)
|
||||
|
||||
disruptive = bool(disruptive_leases or other_live_sessions or terminal_lock)
|
||||
disruptive = bool(
|
||||
disruptive_leases or other_live_sessions or terminal_lock_in_scope
|
||||
)
|
||||
|
||||
if not inventory_complete:
|
||||
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(
|
||||
@@ -381,7 +720,7 @@ def evaluate_restart_impact(
|
||||
f"{len(critical_sections)} critical section(s) in flight "
|
||||
"(active lease with a live owner)"
|
||||
)
|
||||
if terminal_lock:
|
||||
if terminal_lock_in_scope:
|
||||
reasons.append("active terminal (merge) lock present")
|
||||
|
||||
override_would_allow = bool(inventory_complete and disruptive)
|
||||
@@ -411,6 +750,12 @@ def evaluate_restart_impact(
|
||||
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),
|
||||
@@ -424,6 +769,15 @@ def evaluate_restart_impact(
|
||||
|
||||
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,
|
||||
@@ -440,9 +794,11 @@ def evaluate_restart_impact(
|
||||
affected_issues=affected_issues,
|
||||
affected_prs=affected_prs,
|
||||
mutations=mutations,
|
||||
terminal_lock=dict(terminal_lock)
|
||||
if isinstance(terminal_lock, Mapping)
|
||||
else terminal_lock,
|
||||
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,
|
||||
|
||||
Reference in New Issue
Block a user