"""Focused contracts for the typed Kubernetes mutation boundary.""" from __future__ import annotations import inspect from types import SimpleNamespace import pytest from src.models.incident import Incident, IncidentStatus, Severity, Signal from src.models.playbook import ActionType, RepairStep, RiskLevel from src.plugins.mcp.interfaces import MCPToolResult from src.plugins.mcp.mcp_bridge import MCPBridge, MCPServer, MCPTransport from src.plugins.mcp.providers.k8s_provider import K8sProvider from src.plugins.mcp.registry import AuditedMCPToolProvider from src.services.approval_execution import ApprovalExecutionService from src.services.auto_repair_service import AutoRepairService from src.services.callback_dispatcher import dispatch_action, get_action_spec from src.services.executor import ( ActionExecutor, DryRunResult, ExecutionResult, OperationType, ) from src.services.kubernetes_controlled_executor import ( build_kubernetes_callback_execution_claim, build_kubernetes_claim_for_incident, build_kubernetes_controlled_execution_claim, execute_kubernetes_controlled_claim, execute_kubernetes_readonly_action, issue_kubernetes_callback_execution_capability, validate_kubernetes_callback_execution_claim, ) from src.utils.timezone import now_taipei def _route() -> dict: return { "schema_version": "typed_domain_target_route_v2", "resolution_status": "resolved", "route_id": "typed:kubernetes_workload:service-awoooi-api", "target_kind": "kubernetes_workload", "target_resource": "awoooi-api", "source_namespace": "awoooi-prod", "canonical_asset_id": "service:awoooi-api", "executor": "kubernetes_controlled_executor", "verifier": "kubernetes_rollout_verifier", "controlled_apply_allowed": True, "cross_domain_fallback_allowed": False, "identity_evidence": { "status": "verified", "namespace": "awoooi-prod", "kind": "deployment", "name": "awoooi-api", }, } def _action() -> str: return ( "kubectl rollout restart deployment/awoooi-api " "-n awoooi-prod" ) def _incident() -> Incident: now = now_taipei() return Incident( incident_id="INC-K8S-001", status=IncidentStatus.INVESTIGATING, severity=Severity.P2, affected_services=["awoooi-api"], alert_category="kubernetes", signals=[ Signal( alert_name="DeploymentUnavailable", severity=Severity.P2, source="alertmanager", fired_at=now, labels={ "alertname": "DeploymentUnavailable", "namespace": "awoooi-prod", "deployment": "awoooi-api", "workload_kind": "deployment", "workload_name": "awoooi-api", }, ) ], ) def test_exact_rollout_restart_builds_bound_controlled_claim() -> None: claim = build_kubernetes_controlled_execution_claim( action=_action(), typed_target_route=_route(), incident_id="INC-K8S-001", ) assert claim["status"] == "ready_for_controlled_executor" assert claim["ready"] is True assert claim["operation_type"] == "RESTART_DEPLOYMENT" assert claim["namespace"] == "awoooi-prod" assert claim["workload_name"] == "awoooi-api" assert len(claim["action_sha256"]) == 64 assert len(claim["claim_binding_sha256"]) == 64 assert "kubectl" not in str(claim) @pytest.mark.parametrize( ("route_patch", "action", "reason"), [ ({"resolution_status": "asset_identity_unresolved"}, _action(), "typed_route_unresolved"), ({"target_kind": "host_systemd"}, _action(), "typed_route_wrong_domain"), ({"executor": "host_ansible_executor"}, _action(), "typed_route_wrong_executor"), ({"verifier": "host_runtime_independent_verifier"}, _action(), "typed_route_wrong_verifier"), ({"controlled_apply_allowed": False}, _action(), "typed_route_apply_not_allowed"), ({"cross_domain_fallback_allowed": True}, _action(), "cross_domain_fallback_not_fail_closed"), ({"identity_evidence": {"status": "unverified"}}, _action(), "kubernetes_identity_evidence_unverified"), ({"source_namespace": None}, "kubectl rollout restart deployment/awoooi-api", "typed_route_namespace_missing"), ({}, "kubectl rollout restart deployment/other -n awoooi-prod", "action_workload_name_route_mismatch"), ({}, "kubectl rollout restart deployment/awoooi-api -n other", "action_namespace_route_mismatch"), ({}, "kubectl rollout restart deployment/awoooi-api -n awoooi-prod; id", "action_parser_rejected:forbidden_shell_metachar"), ({}, "kubectl scale deployment/awoooi-api --replicas=2 -n awoooi-prod", "action_not_allowlisted_rollout_restart"), ], ) def test_claim_builder_fails_closed_on_route_or_action_drift( route_patch: dict, action: str, reason: str, ) -> None: route = _route() route.update(route_patch) claim = build_kubernetes_controlled_execution_claim( action=action, typed_target_route=route, incident_id="INC-K8S-001", ) assert claim["status"] == "blocked_no_write" assert claim["ready"] is False assert claim["reason"] == reason assert claim["runtime_write_performed"] is False def test_incident_builder_consumes_canonical_router_identity() -> None: claim = build_kubernetes_claim_for_incident( incident=_incident(), action=_action(), ) assert claim["ready"] is True assert claim["route_id"] == "typed:kubernetes_workload:service-awoooi-api" assert claim["canonical_asset_id"] == "service:awoooi-api" def test_verified_route_supplies_namespace_for_legacy_restart_action() -> None: claim = build_kubernetes_controlled_execution_claim( action="kubectl rollout restart deployment/awoooi-api", typed_target_route=_route(), incident_id="INC-K8S-001", ) assert claim["ready"] is True assert claim["namespace"] == "awoooi-prod" assert claim["namespace_source"] == "typed_route" class _TypedExecutor: def __init__(self) -> None: self.calls: list[tuple] = [] async def validate_action(self, operation_type, resource_name, namespace): self.calls.append(("validate", operation_type, resource_name, namespace)) return DryRunResult( passed=True, message="exact workload exists", resource_exists=True, ) async def restart_workload(self, workload_kind, workload_name, namespace): self.calls.append(("restart", workload_kind, workload_name, namespace)) return ExecutionResult( success=True, message="rollout restart triggered", operation_type=OperationType.RESTART_DEPLOYMENT, target_resource=f"{workload_kind}/{workload_name}", namespace=namespace, duration_ms=4, k8s_response={"generation": 7, "uid": "raw-workload-uid"}, ) async def validate_deployment_exists(self, deployment, namespace): self.calls.append(("validate_deployment", deployment, namespace)) return DryRunResult( passed=True, message="exact deployment exists", resource_exists=True, ) async def restart_deployment(self, deployment, namespace): self.calls.append(("restart_deployment", deployment, namespace)) return ExecutionResult( success=True, message="deployment restart triggered", operation_type=OperationType.RESTART_DEPLOYMENT, target_resource=f"deployment/{deployment}", namespace=namespace, duration_ms=6, ) async def execute_with_audit( self, *, approval, operation_type, resource_name, namespace, ): self.calls.append( ("execute_with_audit", approval, operation_type, resource_name, namespace) ) return ExecutionResult( success=True, message="approved rollout restart triggered", operation_type=operation_type, target_resource=f"deployment/{resource_name}", namespace=namespace, duration_ms=5, ) async def write_controlled_execution_audit( self, *, approval, result, dry_run_passed, dry_run_message, ): self.calls.append( ( "controlled_audit", approval, result.success, dict(result.k8s_response or {}), dry_run_passed, dry_run_message, ) ) return True class _VerifiedRollout: async def prepare(self, claim, *, deadline_monotonic=None): assert deadline_monotonic is not None return object(), {"claim_id": claim["claim_id"], "phase": "pre"} async def record_execution_outcome( self, *, claim, execution_context, runtime_write_outcome, deadline_monotonic=None, ): assert execution_context is not None assert deadline_monotonic is not None return { "status": "applied_pending_verification", "claim_id": claim["claim_id"], "runtime_write_outcome": runtime_write_outcome, "durable_writeback_ack": True, "stores_raw_uid": False, "stores_raw_output": False, } async def record_terminal_outcome( self, *, claim, execution_context, status, runtime_write_outcome, blockers, deadline_monotonic=None, ): assert execution_context is not None assert deadline_monotonic is not None return { "status": status, "claim_id": claim["claim_id"], "runtime_write_outcome": runtime_write_outcome, "blockers": list(blockers), "terminal": True, "durable_writeback_ack": True, "stores_raw_uid": False, "stores_raw_output": False, } async def verify( self, *, claim, execution_context, pre_state, deadline_monotonic=None, ): assert execution_context is not None assert deadline_monotonic is not None assert pre_state["claim_id"] == claim["claim_id"] return { "schema_version": "kubernetes_rollout_verifier_receipt_v1", "status": "verified_success", "verified": True, "verifier": "kubernetes_rollout_verifier", "claim_id": claim["claim_id"], "incident_id": claim["incident_id"], "blockers": [], "durable_writeback_ack": True, "stores_raw_uid": False, "stores_raw_output": False, } class _FailedRollout(_VerifiedRollout): async def verify( self, *, claim, execution_context, pre_state, deadline_monotonic=None, ): receipt = await super().verify( claim=claim, execution_context=execution_context, pre_state=pre_state, deadline_monotonic=deadline_monotonic, ) return { **receipt, "status": "failed_closed", "verified": False, "blockers": ["rollout_replicas_not_converged"], } class _TimedOutRollout(_VerifiedRollout): async def verify( self, *, claim, execution_context, pre_state, deadline_monotonic=None, ): raise TimeoutError("independent readback timed out") class _NoPollRollout(_VerifiedRollout): async def verify(self, **_kwargs): raise AssertionError("no-write outcome must not poll rollout state") @pytest.mark.asyncio async def test_controlled_executor_uses_typed_api_and_preserves_receipt_binding() -> None: executor = _TypedExecutor() claim = build_kubernetes_controlled_execution_claim( action=_action(), typed_target_route=_route(), incident_id="INC-K8S-001", ) result = await execute_kubernetes_controlled_claim( claim, executor=executor, verifier=_VerifiedRollout(), ) assert result.success is True assert [call[0] for call in executor.calls] == ["validate", "restart"] assert result.k8s_response["generation"] == 7 assert result.k8s_response["controlled_claim_id"] == claim["claim_id"] assert result.k8s_response["route_id"] == claim["route_id"] assert result.k8s_response["verifier"] == "kubernetes_rollout_verifier" assert result.k8s_response["verifier_receipt"]["verified"] is True assert result.k8s_response["runtime_write_performed"] is True assert result.k8s_response["stores_raw_uid"] is False assert "uid" not in result.k8s_response @pytest.mark.asyncio async def test_controlled_executor_does_not_close_when_verifier_fails() -> None: executor = _TypedExecutor() claim = build_kubernetes_controlled_execution_claim( action=_action(), typed_target_route=_route(), incident_id="INC-K8S-001", ) result = await execute_kubernetes_controlled_claim( claim, executor=executor, verifier=_FailedRollout(), ) assert [call[0] for call in executor.calls] == ["validate", "restart"] assert result.success is False assert result.error == ( "kubernetes_rollout_verifier:rollout_replicas_not_converged" ) assert result.k8s_response["verifier_receipt"]["verified"] is False assert result.k8s_response["runtime_write_performed"] is True @pytest.mark.asyncio async def test_post_write_verifier_timeout_never_retries_mutation() -> None: executor = _TypedExecutor() claim = build_kubernetes_controlled_execution_claim( action=_action(), typed_target_route=_route(), incident_id="INC-K8S-001", ) result = await execute_kubernetes_controlled_claim( claim, executor=executor, verifier=_TimedOutRollout(), ) assert [call[0] for call in executor.calls] == ["validate", "restart"] assert result.success is False assert result.k8s_response["runtime_write_performed"] is True assert ApprovalExecutionService._is_retryable_execution_result(result) is False @pytest.mark.asyncio async def test_shadow_mode_closes_as_durable_no_write_without_polling() -> None: class ShadowExecutor(_TypedExecutor): async def restart_workload(self, workload_kind, workload_name, namespace): self.calls.append(("restart", workload_kind, workload_name, namespace)) return ExecutionResult( success=True, message="shadow simulation", operation_type=OperationType.RESTART_DEPLOYMENT, target_resource=f"{workload_kind}/{workload_name}", namespace=namespace, duration_ms=1, k8s_response={"shadow_mode": True, "dry_run": True}, ) result = await execute_kubernetes_controlled_claim( build_kubernetes_controlled_execution_claim( action=_action(), typed_target_route=_route(), incident_id="INC-K8S-SHADOW", ), executor=ShadowExecutor(), verifier=_NoPollRollout(), ) assert result.success is False assert result.error == "kubernetes_controlled_executor:shadow_mode_no_write" assert result.k8s_response["runtime_write_performed"] is False assert result.k8s_response["lifecycle_receipt"]["status"] == ( "terminal_shadow_no_write" ) assert result.k8s_response["verifier_receipt"] is None @pytest.mark.asyncio async def test_preflight_failure_closes_as_durable_no_write() -> None: class MissingExecutor(_TypedExecutor): async def validate_action(self, operation_type, resource_name, namespace): self.calls.append(("validate", operation_type, resource_name, namespace)) return DryRunResult( passed=False, message="workload missing", resource_exists=False, ) result = await execute_kubernetes_controlled_claim( build_kubernetes_controlled_execution_claim( action=_action(), typed_target_route=_route(), incident_id="INC-K8S-MISSING", ), executor=MissingExecutor(), verifier=_NoPollRollout(), ) assert result.success is False assert result.k8s_response["runtime_write_outcome"] == "confirmed_no_write" assert result.k8s_response["lifecycle_receipt"]["status"] == ( "terminal_preflight_no_write" ) assert result.k8s_response["verifier_receipt"] is None @pytest.mark.asyncio async def test_tampered_controlled_claim_performs_zero_writes() -> None: executor = _TypedExecutor() claim = build_kubernetes_controlled_execution_claim( action=_action(), typed_target_route=_route(), incident_id="INC-K8S-001", ) claim["workload_name"] = "other" result = await execute_kubernetes_controlled_claim(claim, executor=executor) assert result.success is False assert result.error == ( "kubernetes_controlled_executor:" "controlled_execution_claim_binding_invalid" ) assert executor.calls == [] @pytest.mark.asyncio async def test_approved_execution_uses_same_controlled_adapter() -> None: executor = _TypedExecutor() approval = object() claim = build_kubernetes_controlled_execution_claim( action=_action(), typed_target_route=_route(), incident_id="INC-K8S-001", ) result = await execute_kubernetes_controlled_claim( claim, approval=approval, executor=executor, verifier=_VerifiedRollout(), ) assert result.success is True assert [call[0] for call in executor.calls] == [ "validate", "restart", "controlled_audit", ] audit_response = executor.calls[-1][3] assert executor.calls[-1][2] is True assert audit_response["verifier_receipt"]["verified"] is True assert "uid" not in audit_response @pytest.mark.asyncio async def test_approval_audit_is_written_after_failed_verifier_without_raw_uid() -> None: executor = _TypedExecutor() approval = object() claim = build_kubernetes_controlled_execution_claim( action=_action(), typed_target_route=_route(), incident_id="INC-K8S-001", ) result = await execute_kubernetes_controlled_claim( claim, approval=approval, executor=executor, verifier=_FailedRollout(), ) assert result.success is False assert [call[0] for call in executor.calls] == [ "validate", "restart", "controlled_audit", ] audit_response = executor.calls[-1][3] assert executor.calls[-1][2] is False assert audit_response["runtime_write_performed"] is True assert "uid" not in audit_response @pytest.mark.asyncio async def test_real_controlled_audit_writer_defensively_drops_raw_uid( monkeypatch, ) -> None: executor = ActionExecutor() captured: dict = {} async def write_audit(**kwargs): captured.update(kwargs) return True monkeypatch.setattr(executor, "_write_audit_log", write_audit) persisted = await executor.write_controlled_execution_audit( approval=SimpleNamespace(id="approval-1", requested_by="operator"), result=ExecutionResult( success=False, message="verifier failed", operation_type=OperationType.RESTART_DEPLOYMENT, target_resource="deployment/awoooi-api", namespace="awoooi-prod", duration_ms=5, k8s_response={ "uid": "raw-workload-uid", "controlled_claim_id": "claim-1", "runtime_write_performed": True, "verifier_receipt": {"verified": False}, }, error="kubernetes_rollout_verifier:failed_closed", ), dry_run_passed=True, dry_run_message="exact workload exists", ) assert persisted is True assert captured["success"] is False assert "uid" not in captured["k8s_response"] assert captured["k8s_response"]["runtime_write_performed"] is True @pytest.mark.asyncio async def test_ambiguous_dispatch_audit_failure_is_visible_and_never_retried() -> None: class AmbiguousExecutor(_TypedExecutor): async def restart_workload(self, workload_kind, workload_name, namespace): self.calls.append(("restart", workload_kind, workload_name, namespace)) return ExecutionResult( success=False, message="operation timed out", operation_type=OperationType.RESTART_DEPLOYMENT, target_resource=f"{workload_kind}/{workload_name}", namespace=namespace, duration_ms=30_000, k8s_response={ "runtime_write_outcome": "unknown_after_dispatch", }, error="Operation timed out after 30s", ) async def write_controlled_execution_audit(self, **kwargs): await super().write_controlled_execution_audit(**kwargs) return False executor = AmbiguousExecutor() claim = build_kubernetes_controlled_execution_claim( action=_action(), typed_target_route=_route(), incident_id="INC-K8S-001", ) result = await execute_kubernetes_controlled_claim( claim, approval=object(), executor=executor, verifier=_FailedRollout(), ) assert [call[0] for call in executor.calls] == [ "validate", "restart", "controlled_audit", ] assert result.success is False assert result.k8s_response["runtime_write_outcome"] == "unknown_after_dispatch" assert ApprovalExecutionService._is_retryable_execution_result(result) is False assert executor.calls[-1][2] is False assert result.k8s_response["audit_writeback_ack"] is False assert "kubernetes_controlled_audit:writeback_missing" in result.error def test_callback_claim_allows_only_exact_typed_restart() -> None: parameters = {"namespace": "awoooi-prod", "deployment": "awoooi-api"} claim = build_kubernetes_callback_execution_claim( tool_name="kubectl_restart", parameters=parameters, typed_target_route=_route(), incident_id="INC-K8S-CALLBACK", risk="medium", ) assert claim["ready"] is True assert validate_kubernetes_callback_execution_claim( tool_name="kubectl_restart", parameters={**parameters, "_mcp_audit": {"trace_id": "trace"}}, claim=claim, ) == (True, "exact_typed_callback_route_match") @pytest.mark.parametrize( ("tool_name", "risk", "reason"), [ ("kubectl_scale", "medium", "callback_mutation_requires_typed_executor_support"), ("kubectl_rollout_undo", "high", "callback_mutation_requires_typed_executor_support"), ("kubectl_delete", "critical", "critical_break_glass_required"), ], ) def test_callback_non_rollout_mutations_fail_closed( tool_name: str, risk: str, reason: str, ) -> None: claim = build_kubernetes_callback_execution_claim( tool_name=tool_name, parameters={"namespace": "awoooi-prod", "deployment": "awoooi-api"}, typed_target_route=_route(), incident_id="INC-K8S-CALLBACK", risk=risk, ) assert claim["ready"] is False assert claim["reason"] == reason assert claim["runtime_write_performed"] is False @pytest.mark.asyncio async def test_kubernetes_provider_refuses_mutation_without_typed_claim() -> None: provider = K8sProvider() executor = _TypedExecutor() provider._executor = executor provider._rollout_verifier = _VerifiedRollout() result = await provider.execute( "kubectl_restart", {"namespace": "awoooi-prod", "deployment": "awoooi-api"}, ) assert result.success is False assert result.error == ( "kubernetes_controlled_executor:controlled_execution_capability_missing" ) assert executor.calls == [] @pytest.mark.asyncio async def test_kubernetes_provider_consumes_exact_callback_claim() -> None: provider = K8sProvider() executor = _TypedExecutor() provider._executor = executor provider._rollout_verifier = _VerifiedRollout() parameters = {"namespace": "awoooi-prod", "deployment": "awoooi-api"} claim = build_kubernetes_callback_execution_claim( tool_name="kubectl_restart", parameters=parameters, typed_target_route=_route(), incident_id="INC-K8S-CALLBACK", risk="medium", ) capability = issue_kubernetes_callback_execution_capability( tool_name="kubectl_restart", parameters=parameters, claim=claim, ) result = await provider.execute( "kubectl_restart", { **parameters, "_controlled_execution_claim": claim, "_controlled_execution_capability": capability, }, ) assert result.success is True assert result.output["restarted"] is True assert executor.calls == [ ( "validate", OperationType.RESTART_DEPLOYMENT, "awoooi-api", "awoooi-prod", ), ("restart", "deployment", "awoooi-api", "awoooi-prod"), ] @pytest.mark.asyncio async def test_kubernetes_provider_rejects_replayed_callback_capability() -> None: provider = K8sProvider() provider._executor = _TypedExecutor() provider._rollout_verifier = _VerifiedRollout() parameters = {"namespace": "awoooi-prod", "deployment": "awoooi-api"} claim = build_kubernetes_callback_execution_claim( tool_name="kubectl_restart", parameters=parameters, typed_target_route=_route(), incident_id="INC-K8S-CALLBACK", risk="medium", ) capability = issue_kubernetes_callback_execution_capability( tool_name="kubectl_restart", parameters=parameters, claim=claim, ) envelope = { **parameters, "_controlled_execution_claim": claim, "_controlled_execution_capability": capability, } first = await provider.execute("kubectl_restart", envelope) replay = await provider.execute("kubectl_restart", envelope) assert first.success is True assert replay.success is False assert replay.error == ( "kubernetes_controlled_executor:" "controlled_execution_capability_already_consumed" ) @pytest.mark.asyncio async def test_callback_capability_nonce_is_never_written_to_mcp_audit( monkeypatch, ) -> None: provider = K8sProvider() provider._executor = _TypedExecutor() provider._rollout_verifier = _VerifiedRollout() audited_provider = AuditedMCPToolProvider(provider) parameters = {"namespace": "awoooi-prod", "deployment": "awoooi-api"} claim = build_kubernetes_callback_execution_claim( tool_name="kubectl_restart", parameters=parameters, typed_target_route=_route(), incident_id="INC-K8S-CALLBACK", risk="medium", ) capability = issue_kubernetes_callback_execution_capability( tool_name="kubectl_restart", parameters=parameters, claim=claim, ) captured: dict = {} async def record_mcp_call(**kwargs): captured.update(kwargs) monkeypatch.setattr( "src.services.mcp_audit_service.record_mcp_call", record_mcp_call, ) result = await audited_provider.execute( "kubectl_restart", { **parameters, "_controlled_execution_claim": claim, "_controlled_execution_capability": capability, "_mcp_audit": {"incident_id": "INC-K8S-CALLBACK"}, }, ) assert result.success is True assert "_controlled_execution_capability" not in captured["input_params"] assert "nonce" not in repr(captured["input_params"]) @pytest.mark.asyncio async def test_plain_fabricated_callback_claim_cannot_authorize_provider() -> None: provider = K8sProvider() executor = _TypedExecutor() provider._executor = executor parameters = {"namespace": "awoooi-prod", "deployment": "awoooi-api"} claim = build_kubernetes_callback_execution_claim( tool_name="kubectl_restart", parameters=parameters, typed_target_route=_route(), incident_id="INC-K8S-CALLBACK", risk="medium", ) result = await provider.execute( "kubectl_restart", {**parameters, "_controlled_execution_claim": claim}, ) assert result.success is False assert result.error == ( "kubernetes_controlled_executor:" "controlled_execution_capability_missing" ) assert executor.calls == [] @pytest.mark.asyncio @pytest.mark.parametrize("failure_stage", ["dry_run", "restart"]) async def test_kubernetes_provider_propagates_inner_failure( failure_stage: str, ) -> None: class FailingExecutor(_TypedExecutor): async def validate_action(self, operation_type, deployment, namespace): if failure_stage == "dry_run": return DryRunResult( passed=False, message="deployment missing", resource_exists=False, ) return await super().validate_action( operation_type, deployment, namespace, ) async def restart_workload(self, workload_kind, deployment, namespace): if failure_stage == "restart": return ExecutionResult( success=False, message="restart rejected", operation_type=OperationType.RESTART_DEPLOYMENT, target_resource=f"deployment/{deployment}", namespace=namespace, duration_ms=2, error="api rejected patch", ) return await super().restart_workload( workload_kind, deployment, namespace, ) provider = K8sProvider() provider._executor = FailingExecutor() provider._rollout_verifier = _VerifiedRollout() parameters = {"namespace": "awoooi-prod", "deployment": "awoooi-api"} claim = build_kubernetes_callback_execution_claim( tool_name="kubectl_restart", parameters=parameters, typed_target_route=_route(), incident_id="INC-K8S-CALLBACK", risk="medium", ) capability = issue_kubernetes_callback_execution_capability( tool_name="kubectl_restart", parameters=parameters, claim=claim, ) result = await provider.execute( "kubectl_restart", { **parameters, "_controlled_execution_claim": claim, "_controlled_execution_capability": capability, }, ) assert result.success is False assert result.error in { "kubernetes_controlled_executor:deployment missing", "api rejected patch", } @pytest.mark.asyncio @pytest.mark.parametrize( "tool_name", ["kubectl_delete", "kubectl_scale", "kubectl_restart", "kubectl_rollout_undo"], ) async def test_legacy_mcp_bridge_has_no_raw_mutation_fallback( monkeypatch, tool_name: str, ) -> None: monkeypatch.setattr("src.plugins.mcp.registry.get_provider", lambda _name: None) bridge = MCPBridge() server = MCPServer( name="kubernetes", transport=MCPTransport.HTTP, endpoint="internal://kubernetes", ) try: result = await bridge._execute_http( server, tool_name, {"namespace": "awoooi-prod", "deployment": "awoooi-api"}, ) finally: await bridge.close() assert result == { "error": ( "kubernetes_controlled_executor:" "provider_required_no_raw_mutation_fallback" ), "runtime_write_performed": False, } @pytest.mark.asyncio async def test_callback_dispatcher_blocks_identity_drift_before_provider( monkeypatch, ) -> None: def fail_if_called(_name: str): raise AssertionError("provider must not be resolved for an unverified target") monkeypatch.setattr("src.plugins.mcp.registry.get_provider", fail_if_called) result = await dispatch_action( action_name="k8s_restart", incident_id="INC-K8S-CALLBACK", labels={ "alertname": "DeploymentUnavailable", "namespace": "awoooi-prod", "deployment": "awoooi-api", }, ) assert result.success is False assert result.error and result.error.startswith("kubernetes_controlled_executor:") @pytest.mark.asyncio async def test_callback_dispatcher_passes_durable_claim_to_provider( monkeypatch, ) -> None: captured: dict = {} class Provider: async def execute(self, tool_name: str, parameters: dict) -> MCPToolResult: captured["tool_name"] = tool_name captured["parameters"] = parameters return MCPToolResult(success=True, output={"restarted": True}) monkeypatch.setattr( "src.plugins.mcp.registry.get_provider", lambda _name: Provider(), ) result = await dispatch_action( action_name="k8s_restart", incident_id="INC-K8S-CALLBACK", labels={ "alertname": "DeploymentUnavailable", "namespace": "awoooi-prod", "deployment": "awoooi-api", "workload_kind": "deployment", "workload_name": "awoooi-api", }, ) assert result.success is True claim = captured["parameters"]["_controlled_execution_claim"] assert captured["tool_name"] == "kubectl_restart" assert claim["ready"] is True assert claim["incident_id"] == "INC-K8S-CALLBACK" assert claim["route_id"] == "typed:kubernetes_workload:service-awoooi-api" @pytest.mark.asyncio async def test_auto_repair_step_uses_controlled_claim_adapter(monkeypatch) -> None: captured: dict = {} async def execute_claim(claim): captured["claim"] = claim return ExecutionResult( success=True, message="restart triggered", operation_type=OperationType.RESTART_DEPLOYMENT, target_resource="deployment/awoooi-api", namespace="awoooi-prod", duration_ms=3, ) monkeypatch.setattr( "src.services.auto_repair_service.execute_kubernetes_controlled_claim", execute_claim, ) service = AutoRepairService() step = RepairStep( step_number=1, action_type=ActionType.KUBECTL, command=_action(), risk_level=RiskLevel.MEDIUM, ) result = await service._execute_step(_incident(), step) assert result.startswith("SUCCESS: kubernetes_controlled_executor:k8s-claim:") assert captured["claim"]["canonical_asset_id"] == "service:awoooi-api" @pytest.mark.asyncio async def test_auto_repair_readonly_kubectl_remains_compatible(monkeypatch) -> None: calls: list[str] = [] async def execute_readonly(action: str): calls.append(action) return ExecutionResult( success=True, message="read-only complete", operation_type=OperationType.INVESTIGATE, target_resource="pods", namespace="awoooi-prod", duration_ms=1, ) monkeypatch.setattr( "src.services.auto_repair_service.execute_kubernetes_readonly_action", execute_readonly, ) service = AutoRepairService() step = RepairStep( step_number=1, action_type=ActionType.KUBECTL, command="kubectl get pods -n awoooi-prod", risk_level=RiskLevel.LOW, ) result = await service._execute_step(_incident(), step) assert result == "SUCCESS: kubernetes_readonly_executor" assert calls == ["kubectl get pods -n awoooi-prod"] @pytest.mark.asyncio async def test_readonly_adapter_blocks_appended_mutation_before_executor() -> None: executor = _TypedExecutor() result = await execute_kubernetes_readonly_action( ( "kubectl get pods -n awoooi-prod; " "kubectl rollout restart deployment/other -n awoooi-prod" ), executor=executor, ) assert result.success is False assert result.error == ( "kubernetes_readonly_executor:forbidden_shell_metachar" ) assert executor.calls == [] def test_callback_restart_timeout_covers_mutation_and_bounded_verifier() -> None: spec = get_action_spec("k8s_restart") assert spec is not None assert spec.timeout_sec == 120 @pytest.mark.asyncio async def test_approval_claim_loads_incident_before_execution(monkeypatch) -> None: incident = _incident() class IncidentService: async def get_from_working_memory(self, incident_id: str): assert incident_id == incident.incident_id return incident async def get_from_episodic_memory(self, _incident_id: str): raise AssertionError("working memory hit should stop the lookup") monkeypatch.setattr( "src.services.incident_service.get_incident_service", lambda: IncidentService(), ) approval = SimpleNamespace( incident_id=incident.incident_id, action=_action(), ) claim = await ApprovalExecutionService._kubernetes_claim_for_approval(approval) assert claim["ready"] is True assert claim["canonical_asset_id"] == "service:awoooi-api" @pytest.mark.asyncio @pytest.mark.parametrize( ("failure_mode", "reason"), [ ("lookup_exception", "incident_lookup_failed"), ("not_found", "incident_not_found"), ("identity_mismatch", "incident_identity_mismatch"), ("router_exception", "typed_route_resolution_failed"), ], ) async def test_approval_incident_dependency_failures_are_durable_blocked_claims( monkeypatch, failure_mode: str, reason: str, ) -> None: class IncidentService: async def get_from_working_memory(self, _incident_id: str): if failure_mode == "lookup_exception": raise RuntimeError("working memory unavailable") if failure_mode == "identity_mismatch": return {"incident_id": "INC-OTHER"} if failure_mode == "router_exception": return _incident() return None async def get_from_episodic_memory(self, _incident_id: str): return None monkeypatch.setattr( "src.services.incident_service.get_incident_service", lambda: IncidentService(), ) if failure_mode == "router_exception": def raise_router_error(**_kwargs): raise RuntimeError("typed router unavailable") monkeypatch.setattr( "src.services.approval_execution.build_kubernetes_claim_for_incident", raise_router_error, ) approval = SimpleNamespace(incident_id="INC-K8S-001", action=_action()) claim = await ApprovalExecutionService._kubernetes_claim_for_approval(approval) assert claim["status"] == "blocked_no_write" assert claim["reason"] == reason assert claim["runtime_write_performed"] is False def test_all_active_mutation_entries_reference_controlled_adapter() -> None: auto_repair_source = inspect.getsource(AutoRepairService._execute_step) approval_source = inspect.getsource(ApprovalExecutionService.execute_approved_action) assert "execute_kubernetes_controlled_claim" in auto_repair_source assert "execute_kubernetes_controlled_claim" in approval_source assert "execute_kubernetes_readonly_action" in approval_source assert "execute_kubectl_command(command)" not in auto_repair_source