Files
awoooi/apps/api/tests/test_kubernetes_controlled_executor.py
2026-07-18 21:55:52 +08:00

1217 lines
40 KiB
Python

"""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