diff --git a/apps/api/src/core/config.py b/apps/api/src/core/config.py
index 650758ae3..7e694413f 100644
--- a/apps/api/src/core/config.py
+++ b/apps/api/src/core/config.py
@@ -981,6 +981,31 @@ class Settings(BaseSettings):
"then enqueue them for the existing check-mode worker."
),
)
+ ENABLE_KUBERNETES_CONTROLLED_CANDIDATE_WORKER: bool = Field(
+ default=False,
+ description=(
+ "True only in the dedicated Kubernetes executor broker; API, signal "
+ "worker, and Ansible broker must keep this disabled."
+ ),
+ )
+ KUBERNETES_CONTROLLED_CANDIDATE_INTERVAL_SECONDS: int = Field(
+ default=30,
+ ge=5,
+ le=300,
+ description="Dedicated Kubernetes candidate claim interval.",
+ )
+ KUBERNETES_CONTROLLED_CANDIDATE_BATCH_LIMIT: int = Field(
+ default=1,
+ ge=1,
+ le=1,
+ description="Single-writer Kubernetes candidate batch size.",
+ )
+ KUBERNETES_CONTROLLED_CANDIDATE_STARTUP_SLEEP_SECONDS: int = Field(
+ default=15,
+ ge=0,
+ le=300,
+ description="Bounded startup delay for the dedicated executor broker.",
+ )
ENABLE_SECURITY_CONTROL_PLANE_MAINTENANCE_WORKER: bool = Field(
default=False,
description=(
diff --git a/apps/api/src/jobs/kubernetes_controlled_candidate_job.py b/apps/api/src/jobs/kubernetes_controlled_candidate_job.py
new file mode 100644
index 000000000..f8dcd25d6
--- /dev/null
+++ b/apps/api/src/jobs/kubernetes_controlled_candidate_job.py
@@ -0,0 +1,131 @@
+"""Bounded owner loop for durable Kubernetes controlled candidates.
+
+This module is a standalone executor entry point. It is deliberately not
+started by the API or signal worker: runtime activation requires a dedicated
+Kubernetes executor identity and remains a release/runtime gate.
+"""
+
+from __future__ import annotations
+
+import asyncio
+from collections.abc import Awaitable, Callable, Mapping
+from typing import Any
+
+import structlog
+
+from src.core.config import settings
+from src.services.kubernetes_controlled_candidate_service import (
+ claim_next_kubernetes_controlled_candidate,
+ run_claimed_kubernetes_controlled_candidate,
+)
+
+logger = structlog.get_logger(__name__)
+
+ClaimNext = Callable[..., Awaitable[dict[str, Any] | None]]
+RunCandidate = Callable[..., Awaitable[dict[str, Any]]]
+
+
+async def run_kubernetes_controlled_candidate_once(
+ *,
+ project_id: str = "awoooi",
+ batch_limit: int | None = None,
+ claim_next: ClaimNext | None = None,
+ run_candidate: RunCandidate | None = None,
+) -> dict[str, Any]:
+ """Claim a bounded batch; every claim crosses its own dispatch barrier."""
+
+ limit = min(
+ max(
+ 1,
+ int(
+ batch_limit
+ or settings.KUBERNETES_CONTROLLED_CANDIDATE_BATCH_LIMIT
+ ),
+ ),
+ 1,
+ )
+ claim_one = claim_next or claim_next_kubernetes_controlled_candidate
+ execute_one = run_candidate or run_claimed_kubernetes_controlled_candidate
+ stats: dict[str, Any] = {
+ "claimed": 0,
+ "verified_closed": 0,
+ "pending_closure": 0,
+ "duplicate_apply_suppressed": 0,
+ "errors": 0,
+ "runtime_activation_enabled": bool(
+ settings.ENABLE_KUBERNETES_CONTROLLED_CANDIDATE_WORKER
+ ),
+ "single_writer_executor": "awoooi-kubernetes-executor-broker",
+ }
+ for _ in range(limit):
+ candidate = await claim_one(project_id=project_id)
+ if candidate is None:
+ break
+ stats["claimed"] += 1
+ try:
+ result = await execute_one(candidate, project_id=project_id)
+ except asyncio.CancelledError:
+ raise
+ except Exception as exc: # noqa: BLE001 - dispatch state remains fail closed
+ stats["errors"] += 1
+ logger.error(
+ "kubernetes_controlled_candidate_worker_failed_closed",
+ event_id=str(candidate.get("event_id") or ""),
+ claim_id=str(
+ (candidate.get("claim") or {}).get("claim_id")
+ if isinstance(candidate.get("claim"), Mapping)
+ else ""
+ ),
+ error_type=type(exc).__name__,
+ replay_apply_allowed=False,
+ )
+ continue
+ status = str(result.get("status") or "")
+ if result.get("closed") is True:
+ stats["verified_closed"] += 1
+ elif status == "duplicate_apply_suppressed":
+ stats["duplicate_apply_suppressed"] += 1
+ else:
+ stats["pending_closure"] += 1
+ return stats
+
+
+async def run_kubernetes_controlled_candidate_loop() -> None:
+ """Run only under the dedicated, explicitly enabled executor process."""
+
+ if not settings.ENABLE_KUBERNETES_CONTROLLED_CANDIDATE_WORKER:
+ logger.info("kubernetes_controlled_candidate_worker_disabled")
+ return
+ await asyncio.sleep(
+ settings.KUBERNETES_CONTROLLED_CANDIDATE_STARTUP_SLEEP_SECONDS
+ )
+ logger.info(
+ "kubernetes_controlled_candidate_worker_started",
+ interval_seconds=(
+ settings.KUBERNETES_CONTROLLED_CANDIDATE_INTERVAL_SECONDS
+ ),
+ batch_limit=settings.KUBERNETES_CONTROLLED_CANDIDATE_BATCH_LIMIT,
+ single_writer_executor="awoooi-kubernetes-executor-broker",
+ )
+ while True:
+ try:
+ stats = await run_kubernetes_controlled_candidate_once()
+ logger.info("kubernetes_controlled_candidate_worker_tick", **stats)
+ except asyncio.CancelledError:
+ raise
+ except Exception as exc: # noqa: BLE001 - loop stays fail visible
+ logger.error(
+ "kubernetes_controlled_candidate_worker_tick_failed",
+ error_type=type(exc).__name__,
+ )
+ await asyncio.sleep(
+ settings.KUBERNETES_CONTROLLED_CANDIDATE_INTERVAL_SECONDS
+ )
+
+
+def main() -> None:
+ asyncio.run(run_kubernetes_controlled_candidate_loop())
+
+
+if __name__ == "__main__":
+ main()
diff --git a/apps/api/src/services/controlled_alert_target_router.py b/apps/api/src/services/controlled_alert_target_router.py
index 8bca0b8c2..8fb42e847 100644
--- a/apps/api/src/services/controlled_alert_target_router.py
+++ b/apps/api/src/services/controlled_alert_target_router.py
@@ -128,7 +128,14 @@ def _kubernetes_identity_evidence(
) -> dict[str, Any]:
"""Verify a Kubernetes identity without trusting a workload-kind label alone."""
- kind = _text_token(
+ registry_namespace = _text_token(
+ str(identity.get("kubernetes_namespace") or "")
+ )
+ registry_kind = _text_token(str(identity.get("kubernetes_kind") or ""))
+ registry_workload_name = _text_token(
+ str(identity.get("kubernetes_name") or "")
+ )
+ label_kind = _text_token(
str(
labels.get("workload_kind")
or labels.get("kubernetes_kind")
@@ -136,6 +143,7 @@ def _kubernetes_identity_evidence(
or ""
)
)
+ kind = label_kind or registry_kind
name = _text_token(
str(
labels.get("workload_name")
@@ -153,6 +161,19 @@ def _kubernetes_identity_evidence(
requested_namespace = _text_token(namespace)
target_name = _text_token(target_resource)
registry_name = _text_token(str(identity.get("service_name") or ""))
+ registry_tuple_present = bool(
+ registry_namespace and registry_kind and registry_workload_name
+ )
+ registry_scope_exact = bool(
+ registry_tuple_present
+ and registry_namespace == requested_namespace
+ and registry_kind == kind
+ and registry_workload_name == name
+ )
+ explicit_label_scope_exact = bool(
+ not any((registry_namespace, registry_kind, registry_workload_name))
+ and label_kind
+ )
registry_exact = bool(
identity.get("resolution_status") == "resolved"
and identity.get("asset_domain") == "kubernetes_workload"
@@ -162,6 +183,7 @@ def _kubernetes_identity_evidence(
and name
and name == target_name
and (not registry_name or registry_name == target_name)
+ and (registry_scope_exact or explicit_label_scope_exact)
)
runtime_receipt = (
diff --git a/apps/api/src/services/decision_manager.py b/apps/api/src/services/decision_manager.py
index 2259e6129..ea69acafd 100644
--- a/apps/api/src/services/decision_manager.py
+++ b/apps/api/src/services/decision_manager.py
@@ -2070,16 +2070,64 @@ class DecisionManager:
await self._save_token(token)
return
+ kubernetes_action = action.casefold().startswith("kubectl ")
try:
- from src.jobs.awooop_ansible_candidate_backfill_job import (
- enqueue_ai_decision_ansible_candidate,
+ incident_dump = getattr(incident, "model_dump", None)
+ if callable(incident_dump):
+ typed_incident = incident_dump(mode="json")
+ else:
+ typed_incident = {
+ "incident_id": incident.incident_id,
+ "affected_services": list(incident.affected_services or []),
+ "signals": [
+ {
+ "alert_name": getattr(signal, "alert_name", ""),
+ "labels": dict(getattr(signal, "labels", {}) or {}),
+ }
+ for signal in incident.signals or []
+ ],
+ }
+ from src.services.controlled_alert_target_router import (
+ resolve_typed_incident_target,
)
- handoff = await enqueue_ai_decision_ansible_candidate(
- incident=incident,
- proposal_data=proposal_data,
- project_id=str(getattr(incident, "project_id", None) or "awoooi"),
+ typed_target_route = resolve_typed_incident_target(typed_incident)
+ except Exception as exc:
+ logger.warning(
+ "ai_decision_typed_target_route_failed_closed",
+ incident_id=incident.incident_id,
+ error_type=type(exc).__name__,
)
+ typed_target_route = {}
+ kubernetes_lane = bool(
+ kubernetes_action
+ or typed_target_route.get("target_kind") == "kubernetes_workload"
+ )
+ try:
+ if kubernetes_lane:
+ from src.services.kubernetes_controlled_candidate_service import (
+ enqueue_ai_decision_kubernetes_candidate,
+ )
+
+ handoff = await enqueue_ai_decision_kubernetes_candidate(
+ incident=incident,
+ proposal_data=proposal_data,
+ project_id=str(
+ getattr(incident, "project_id", None) or "awoooi"
+ ),
+ )
+ else:
+ from src.jobs.awooop_ansible_candidate_backfill_job import (
+ enqueue_ai_decision_ansible_candidate,
+ )
+
+ handoff = await enqueue_ai_decision_ansible_candidate(
+ incident=incident,
+ proposal_data=proposal_data,
+ project_id=str(
+ getattr(incident, "project_id", None) or "awoooi"
+ ),
+ )
except Exception as exc:
logger.warning(
"ai_decision_controlled_executor_queue_failed",
@@ -2091,28 +2139,60 @@ class DecisionManager:
"status": "controlled_executor_queue_failed",
"queued": False,
"side_effect_performed": False,
- "single_writer_executor": "awoooi-ansible-executor-broker",
+ "single_writer_executor": (
+ "awoooi-kubernetes-executor-broker"
+ if kubernetes_lane
+ else "awoooi-ansible-executor-broker"
+ ),
+ "cross_domain_fallback_allowed": False,
"active_blockers": ["controlled_executor_queue_failed"],
}
queued = handoff.get("queued") is True
+ single_writer_executor = str(
+ handoff.get("single_writer_executor")
+ or (
+ "awoooi-kubernetes-executor-broker"
+ if kubernetes_lane
+ else "awoooi-ansible-executor-broker"
+ )
+ )
proposal_data.update({
"auto_executed": False,
"direct_runtime_execution_performed": False,
"approval_execution_invoked": False,
- "check_mode_required_before_apply": True,
- "single_writer_executor": "awoooi-ansible-executor-broker",
+ "check_mode_required_before_apply": not kubernetes_lane,
+ "typed_dry_run_required_before_apply": kubernetes_lane,
+ "single_writer_executor": single_writer_executor,
"controlled_executor_handoff": handoff,
"automation_run_id": str(handoff.get("automation_run_id") or ""),
"automation_state": (
- "controlled_check_mode_queued"
+ "controlled_kubernetes_candidate_queued"
+ if queued and kubernetes_lane
+ else "controlled_check_mode_queued"
if queued
- else "blocked_with_safe_next_action"
+ else str(
+ handoff.get("status") or "blocked_with_safe_next_action"
+ )
),
"safe_next_action": (
- "awooop_ansible_check_mode_worker_claims_candidate"
+ str(
+ handoff.get("next_safe_action")
+ or (
+ "kubernetes_controlled_candidate_worker_claims_pending_candidate"
+ if kubernetes_lane
+ else "awooop_ansible_check_mode_worker_claims_candidate"
+ )
+ )
if queued
- else "repair_candidate_backfill_or_playbook_coverage_gap"
+ else str(
+ handoff.get("next_safe_action")
+ or (
+ "repair_canonical_kubernetes_identity_before_new_candidate"
+ if kubernetes_lane
+ else "repair_candidate_backfill_or_playbook_coverage_gap"
+ )
+ )
),
})
if queued:
diff --git a/apps/api/src/services/kubernetes_controlled_candidate_service.py b/apps/api/src/services/kubernetes_controlled_candidate_service.py
new file mode 100644
index 000000000..5cfb6f83b
--- /dev/null
+++ b/apps/api/src/services/kubernetes_controlled_candidate_service.py
@@ -0,0 +1,661 @@
+"""Durable queue and single-dispatch closure for Kubernetes AUTO candidates.
+
+The decision layer may enqueue one exact typed Kubernetes claim, but only this
+service may claim it for the existing controlled executor. The durable event
+row is also the idempotency boundary: once dispatch starts it is never eligible
+for another apply, even if closure writeback is interrupted.
+"""
+
+from __future__ import annotations
+
+import hashlib
+import json
+from collections.abc import Awaitable, Callable, Mapping
+from html import escape
+from typing import Any
+from uuid import NAMESPACE_URL, UUID, uuid4, uuid5
+
+import structlog
+from sqlalchemy import text
+
+from src.services.audit_sink import sanitize
+from src.services.kubernetes_controlled_executor import (
+ build_kubernetes_claim_for_incident,
+ execute_kubernetes_controlled_claim,
+)
+
+logger = structlog.get_logger(__name__)
+
+CANDIDATE_SCHEMA_VERSION = "kubernetes_controlled_candidate_v1"
+HANDOFF_SCHEMA_VERSION = "ai_decision_controlled_executor_handoff_v1"
+_PROVIDER_EVENT_PREFIX = "k8s-controlled-candidate:"
+_CLAIM_LEASE_SECONDS = 60
+
+PersistCandidate = Callable[[dict[str, Any], str], Awaitable[dict[str, Any]]]
+BeginDispatch = Callable[[Mapping[str, Any], str], Awaitable[bool]]
+FinalizeCandidate = Callable[
+ [Mapping[str, Any], str, Mapping[str, Any]], Awaitable[bool]
+]
+ExecuteClaim = Callable[[Mapping[str, Any]], Awaitable[Any]]
+ReceiptWriter = Callable[[Mapping[str, Any], Any, str], Awaitable[dict[str, Any]]]
+
+
+def _incident_project_id(incident: Any, project_id: str | None) -> str:
+ return str(project_id or getattr(incident, "project_id", None) or "awoooi")
+
+
+def _action_from_proposal(proposal_data: Mapping[str, Any]) -> str:
+ return str(
+ proposal_data.get("kubectl_command") or proposal_data.get("action") or ""
+ ).strip()
+
+
+def _provider_event_id(claim_id: str) -> str:
+ return f"{_PROVIDER_EVENT_PREFIX}{claim_id}"
+
+
+def _candidate_run_id(claim_id: str) -> str:
+ return str(uuid5(NAMESPACE_URL, f"awoooi:kubernetes-candidate:{claim_id}"))
+
+
+def _safe_receipt_id(prefix: str, *values: str) -> str:
+ digest = hashlib.sha256("|".join(values).encode("utf-8")).hexdigest()[:24]
+ return f"{prefix}:{digest}"
+
+
+async def persist_kubernetes_controlled_candidate(
+ record: dict[str, Any],
+ project_id: str,
+) -> dict[str, Any]:
+ """Insert one candidate or return the exact existing durable event."""
+
+ from src.db.base import get_db_context
+
+ claim = record["claim"]
+ claim_id = str(claim["claim_id"])
+ provider_event_id = _provider_event_id(claim_id)
+ run_id = UUID(str(record["automation_run_id"]))
+ content_hash = str(claim["claim_binding_sha256"])
+ envelope = sanitize(record)
+ source_envelope = json.dumps(envelope, ensure_ascii=False, default=str)
+ preview = f"kubernetes_candidate:pending:{claim_id}"[:256]
+
+ async with get_db_context(project_id) as db:
+ inserted = await db.execute(
+ text(
+ """
+ INSERT INTO awooop_conversation_event (
+ project_id, channel_type, provider_event_id,
+ run_id, content_type, content_hash, content_preview,
+ content_redacted, redaction_version, source_envelope,
+ is_duplicate, received_at
+ ) VALUES (
+ :project_id, 'internal', :provider_event_id,
+ :run_id, 'command', :content_hash, :preview,
+ :preview, 'audit_sink_v1', CAST(:source_envelope AS jsonb),
+ FALSE, NOW()
+ )
+ ON CONFLICT (project_id, channel_type, provider_event_id)
+ DO NOTHING
+ RETURNING event_id, source_envelope
+ """
+ ),
+ {
+ "project_id": project_id,
+ "provider_event_id": provider_event_id,
+ "run_id": run_id,
+ "content_hash": content_hash,
+ "preview": preview,
+ "source_envelope": source_envelope,
+ },
+ )
+ inserted_row = inserted.fetchone()
+ if inserted_row is not None:
+ return {
+ "event_id": str(inserted_row[0]),
+ "created": True,
+ "record": envelope,
+ }
+
+ existing = await db.execute(
+ text(
+ """
+ SELECT event_id, source_envelope
+ FROM awooop_conversation_event
+ WHERE project_id = :project_id
+ AND channel_type = 'internal'
+ AND provider_event_id = :provider_event_id
+ LIMIT 1
+ """
+ ),
+ {
+ "project_id": project_id,
+ "provider_event_id": provider_event_id,
+ },
+ )
+ existing_row = existing.fetchone()
+ if existing_row is None:
+ raise RuntimeError("kubernetes_candidate_dedupe_receipt_missing")
+ stored = existing_row[1]
+ if isinstance(stored, str):
+ stored = json.loads(stored)
+ if not isinstance(stored, Mapping):
+ raise RuntimeError("kubernetes_candidate_dedupe_receipt_invalid")
+ stored_claim = stored.get("claim")
+ if (
+ not isinstance(stored_claim, Mapping)
+ or stored_claim.get("claim_id") != claim_id
+ or stored_claim.get("claim_binding_sha256") != content_hash
+ ):
+ raise RuntimeError("kubernetes_candidate_dedupe_receipt_mismatch")
+ return {
+ "event_id": str(existing_row[0]),
+ "created": False,
+ "record": dict(stored),
+ }
+
+
+async def enqueue_ai_decision_kubernetes_candidate(
+ *,
+ incident: Any,
+ proposal_data: Mapping[str, Any],
+ project_id: str | None = None,
+ persist_candidate: PersistCandidate | None = None,
+) -> dict[str, Any]:
+ """Queue one exact K8s claim; unresolved identity never falls back."""
+
+ action = _action_from_proposal(proposal_data)
+ claim = build_kubernetes_claim_for_incident(incident=incident, action=action)
+ project = _incident_project_id(incident, project_id)
+ if claim.get("ready") is not True:
+ reason = str(claim.get("reason") or "kubernetes_claim_not_ready")
+ return {
+ "schema_version": HANDOFF_SCHEMA_VERSION,
+ "status": "kubernetes_candidate_blocked_no_write",
+ "queued": False,
+ "deduplicated": False,
+ "side_effect_performed": False,
+ "runtime_mutation_performed": False,
+ "single_writer_executor": "awoooi-kubernetes-executor-broker",
+ "typed_domain": claim.get("typed_domain") or "unknown",
+ "canonical_asset_id": claim.get("canonical_asset_id"),
+ "cross_domain_fallback_allowed": False,
+ "active_blockers": [reason],
+ "next_safe_action": (
+ "repair_canonical_kubernetes_identity_before_new_candidate"
+ ),
+ }
+
+ claim_id = str(claim["claim_id"])
+ automation_run_id = _candidate_run_id(claim_id)
+ record = {
+ "schema_version": CANDIDATE_SCHEMA_VERSION,
+ "status": "pending",
+ "terminal_execution": False,
+ "claim": claim,
+ "automation_run_id": automation_run_id,
+ "project_id": project,
+ "incident_id": str(claim.get("incident_id") or ""),
+ "matched_playbook_id": str(
+ proposal_data.get("matched_playbook_id")
+ or proposal_data.get("playbook_id")
+ or ""
+ ),
+ "source_commitment": "AIA-CONV-033",
+ "legacy_ansible_candidate_consulted": False,
+ "cross_domain_fallback_allowed": False,
+ "runtime_mutation_performed": False,
+ }
+ persistence = await (persist_candidate or persist_kubernetes_controlled_candidate)(
+ record,
+ project,
+ )
+ stored = persistence.get("record")
+ stored_status = str(
+ stored.get("status") if isinstance(stored, Mapping) else "pending"
+ )
+ already_terminal = bool(
+ isinstance(stored, Mapping) and stored.get("terminal_execution") is True
+ )
+ created = persistence.get("created") is True
+ queued = not already_terminal
+ return {
+ "schema_version": HANDOFF_SCHEMA_VERSION,
+ "status": (
+ "controlled_kubernetes_candidate_queued"
+ if created
+ else "controlled_kubernetes_candidate_exists"
+ if queued
+ else "controlled_kubernetes_apply_duplicate_suppressed"
+ ),
+ "queued": queued,
+ "deduplicated": not created,
+ "side_effect_performed": False,
+ "runtime_mutation_performed": False,
+ "candidate_event_id": str(persistence.get("event_id") or ""),
+ "automation_run_id": automation_run_id,
+ "claim_id": claim_id,
+ "candidate_status": stored_status,
+ "single_writer_executor": "awoooi-kubernetes-executor-broker",
+ "typed_domain": "kubernetes_workload",
+ "canonical_asset_id": claim.get("canonical_asset_id"),
+ "cross_domain_fallback_allowed": False,
+ "active_blockers": (
+ [] if queued else ["exact_kubernetes_apply_already_terminal"]
+ ),
+ "next_safe_action": (
+ "kubernetes_controlled_candidate_worker_claims_pending_candidate"
+ if queued
+ else "read_existing_kubernetes_closure_receipt"
+ ),
+ }
+
+
+async def claim_next_kubernetes_controlled_candidate(
+ *,
+ project_id: str = "awoooi",
+) -> dict[str, Any] | None:
+ """Claim one pending or pre-dispatch lease-expired candidate."""
+
+ from src.db.base import get_db_context
+
+ claim_token = str(uuid4())
+ async with get_db_context(project_id) as db:
+ claimed = await db.execute(
+ text(
+ """
+ WITH next_candidate AS (
+ SELECT event_id
+ FROM awooop_conversation_event
+ WHERE project_id = :project_id
+ AND channel_type = 'internal'
+ AND provider_event_id LIKE :provider_prefix
+ AND source_envelope ->> 'schema_version' = :schema_version
+ AND (
+ source_envelope ->> 'status' = 'pending'
+ OR (
+ source_envelope ->> 'status' = 'claimed'
+ AND (
+ source_envelope ->> 'claim_lease_expires_at'
+ )::timestamptz <= NOW()
+ )
+ )
+ ORDER BY received_at ASC
+ FOR UPDATE SKIP LOCKED
+ LIMIT 1
+ )
+ UPDATE awooop_conversation_event AS event
+ SET source_envelope = event.source_envelope || jsonb_build_object(
+ 'status', 'claimed',
+ 'claim_token', :claim_token,
+ 'claimed_at', NOW(),
+ 'claim_lease_expires_at',
+ NOW() + make_interval(secs => :lease_seconds),
+ 'claim_attempt_count',
+ COALESCE(
+ (event.source_envelope ->> 'claim_attempt_count')::int,
+ 0
+ ) + 1
+ )
+ FROM next_candidate
+ WHERE event.event_id = next_candidate.event_id
+ RETURNING event.event_id, event.run_id, event.source_envelope
+ """
+ ),
+ {
+ "project_id": project_id,
+ "provider_prefix": f"{_PROVIDER_EVENT_PREFIX}%",
+ "schema_version": CANDIDATE_SCHEMA_VERSION,
+ "claim_token": claim_token,
+ "lease_seconds": _CLAIM_LEASE_SECONDS,
+ },
+ )
+ row = claimed.fetchone()
+ if row is None:
+ return None
+ envelope = row[2]
+ if isinstance(envelope, str):
+ envelope = json.loads(envelope)
+ if not isinstance(envelope, Mapping):
+ raise RuntimeError("claimed_kubernetes_candidate_invalid")
+ return {
+ **dict(envelope),
+ "event_id": str(row[0]),
+ "automation_run_id": str(row[1]),
+ "claim_token": claim_token,
+ }
+
+
+async def begin_kubernetes_candidate_dispatch(
+ candidate: Mapping[str, Any],
+ project_id: str,
+) -> bool:
+ """Cross the irreversible single-dispatch barrier exactly once."""
+
+ from src.db.base import get_db_context
+
+ async with get_db_context(project_id) as db:
+ updated = await db.execute(
+ text(
+ """
+ UPDATE awooop_conversation_event
+ SET source_envelope = source_envelope || jsonb_build_object(
+ 'status', 'dispatching',
+ 'dispatch_started_at', NOW(),
+ 'terminal_execution', FALSE,
+ 'replay_apply_allowed', FALSE
+ )
+ WHERE project_id = :project_id
+ AND event_id = CAST(:event_id AS uuid)
+ AND source_envelope ->> 'schema_version' = :schema_version
+ AND source_envelope ->> 'status' = 'claimed'
+ AND source_envelope ->> 'claim_token' = :claim_token
+ RETURNING event_id
+ """
+ ),
+ {
+ "project_id": project_id,
+ "event_id": str(candidate.get("event_id") or ""),
+ "schema_version": CANDIDATE_SCHEMA_VERSION,
+ "claim_token": str(candidate.get("claim_token") or ""),
+ },
+ )
+ return updated.fetchone() is not None
+
+
+async def finalize_kubernetes_controlled_candidate(
+ candidate: Mapping[str, Any],
+ project_id: str,
+ closure: Mapping[str, Any],
+) -> bool:
+ """Persist public-safe closure receipts without reopening apply."""
+
+ from src.db.base import get_db_context
+
+ safe_closure = json.dumps(sanitize(dict(closure)), ensure_ascii=False, default=str)
+ status = str(closure.get("status") or "applied_pending_closure")
+ async with get_db_context(project_id) as db:
+ updated = await db.execute(
+ text(
+ """
+ UPDATE awooop_conversation_event
+ SET source_envelope = source_envelope || jsonb_build_object(
+ 'status', :status,
+ 'terminal_execution', TRUE,
+ 'replay_apply_allowed', FALSE,
+ 'closure', CAST(:closure AS jsonb),
+ 'terminal_at', NOW()
+ )
+ WHERE project_id = :project_id
+ AND event_id = CAST(:event_id AS uuid)
+ AND source_envelope ->> 'schema_version' = :schema_version
+ AND source_envelope ->> 'status' = 'dispatching'
+ AND source_envelope ->> 'claim_token' = :claim_token
+ RETURNING event_id
+ """
+ ),
+ {
+ "project_id": project_id,
+ "event_id": str(candidate.get("event_id") or ""),
+ "schema_version": CANDIDATE_SCHEMA_VERSION,
+ "claim_token": str(candidate.get("claim_token") or ""),
+ "status": status,
+ "closure": safe_closure,
+ },
+ )
+ return updated.fetchone() is not None
+
+
+async def _write_learning_receipt(
+ candidate: Mapping[str, Any],
+ execution: Any,
+ execution_run_id: str,
+) -> dict[str, Any]:
+ from src.models.approval import ApprovalRequest, ApprovalStatus, RiskLevel
+ from src.services.learning_service import (
+ ExecutionResult as LearningExecutionResult,
+ )
+ from src.services.learning_service import get_learning_service
+
+ claim = candidate.get("claim")
+ assert isinstance(claim, Mapping)
+ claim_id = str(claim.get("claim_id") or "")
+ matched_playbook_id = str(candidate.get("matched_playbook_id") or "")
+ if not matched_playbook_id:
+ return {
+ "schema_version": "kubernetes_candidate_learning_receipt_v1",
+ "receipt_id": _safe_receipt_id(
+ "k8s-learning", claim_id, execution_run_id
+ ),
+ "run_id": execution_run_id,
+ "claim_id": claim_id,
+ "incident_id": str(claim.get("incident_id") or ""),
+ "durable_writeback_ack": False,
+ "blocker": "matched_playbook_identity_missing",
+ "provider_call_performed": False,
+ }
+ action = (
+ f"kubectl rollout restart {claim.get('workload_kind')}/"
+ f"{claim.get('workload_name')} -n {claim.get('namespace')}"
+ )
+ approval_id = uuid5(NAMESPACE_URL, f"{claim_id}:{execution_run_id}:learning")
+ approval = ApprovalRequest(
+ id=approval_id,
+ action=action,
+ description="Typed Kubernetes controlled candidate learning writeback",
+ risk_level=RiskLevel.MEDIUM,
+ requested_by="auto_approve",
+ required_signatures=0,
+ status=ApprovalStatus.APPROVED,
+ incident_id=str(claim.get("incident_id") or ""),
+ matched_playbook_id=matched_playbook_id,
+ metadata={
+ "claim_id": claim_id,
+ "execution_run_id": execution_run_id,
+ "source_commitment": "AIA-CONV-033",
+ },
+ )
+ record = await get_learning_service().process_execution_result(
+ approval=approval,
+ result=LearningExecutionResult(
+ approval_id=str(approval_id),
+ incident_id=str(claim.get("incident_id") or ""),
+ action=action,
+ success=bool(getattr(execution, "success", False)),
+ error_message=str(getattr(execution, "error", "") or "") or None,
+ ),
+ )
+ return {
+ "schema_version": "kubernetes_candidate_learning_receipt_v1",
+ "receipt_id": _safe_receipt_id("k8s-learning", claim_id, execution_run_id),
+ "run_id": execution_run_id,
+ "claim_id": claim_id,
+ "incident_id": str(claim.get("incident_id") or ""),
+ "feedback_type": record.feedback_type.value,
+ "durable_writeback_ack": True,
+ }
+
+
+async def _send_telegram_receipt(
+ candidate: Mapping[str, Any],
+ execution: Any,
+ execution_run_id: str,
+) -> dict[str, Any]:
+ from src.services.telegram_gateway import (
+ _telegram_send_delivery_succeeded,
+ get_telegram_gateway,
+ )
+
+ claim = candidate.get("claim")
+ assert isinstance(claim, Mapping)
+ claim_id = str(claim.get("claim_id") or "")
+ success = bool(getattr(execution, "success", False))
+ gateway = get_telegram_gateway()
+ delivery = await gateway.send_canonical_message(
+ product_id="awoooi",
+ signal_family="incident_lifecycle",
+ severity="P1",
+ text=(
+ "✅ Kubernetes 受控修復已完成獨立驗證\n"
+ if success
+ else "⛔ Kubernetes 受控修復未完成 closure\n"
+ )
+ + f"incident={escape(str(claim.get('incident_id') or ''))}\n"
+ + f"claim={escape(claim_id)}\n"
+ + f"run={escape(execution_run_id)}\n"
+ + "asset="
+ + escape(str(claim.get("canonical_asset_id") or ""))
+ + "\n"
+ + "executor=kubernetes_controlled_executor | "
+ + "verifier=kubernetes_rollout_verifier | cross_domain_fallback=false",
+ parse_mode="HTML",
+ )
+ acknowledged = _telegram_send_delivery_succeeded(delivery)
+ return {
+ "schema_version": "kubernetes_candidate_telegram_receipt_v1",
+ "receipt_id": _safe_receipt_id("k8s-telegram", claim_id, execution_run_id),
+ "run_id": execution_run_id,
+ "claim_id": claim_id,
+ "provider_delivery_acknowledged": acknowledged,
+ "durable_writeback_ack": acknowledged,
+ }
+
+
+def _execution_receipts(execution: Any) -> tuple[dict[str, Any], dict[str, Any]]:
+ response = getattr(execution, "k8s_response", None)
+ if not isinstance(response, Mapping):
+ return {}, {}
+ lifecycle = response.get("lifecycle_receipt")
+ verifier = response.get("verifier_receipt")
+ return (
+ dict(lifecycle) if isinstance(lifecycle, Mapping) else {},
+ dict(verifier) if isinstance(verifier, Mapping) else {},
+ )
+
+
+async def run_claimed_kubernetes_controlled_candidate(
+ candidate: Mapping[str, Any],
+ *,
+ project_id: str = "awoooi",
+ begin_dispatch: BeginDispatch | None = None,
+ execute_claim: ExecuteClaim | None = None,
+ write_learning_receipt: ReceiptWriter | None = None,
+ send_telegram_receipt: ReceiptWriter | None = None,
+ finalize_candidate: FinalizeCandidate | None = None,
+) -> dict[str, Any]:
+ """Dispatch once and close only when all receipts share one run id."""
+
+ claim = candidate.get("claim")
+ if not isinstance(claim, Mapping) or not str(candidate.get("claim_token") or ""):
+ return {
+ "status": "candidate_claim_invalid_no_write",
+ "closed": False,
+ "runtime_mutation_performed": False,
+ }
+
+ started = await (begin_dispatch or begin_kubernetes_candidate_dispatch)(
+ candidate,
+ project_id,
+ )
+ if not started:
+ return {
+ "status": "duplicate_apply_suppressed",
+ "closed": False,
+ "runtime_mutation_performed": False,
+ "claim_id": str(claim.get("claim_id") or ""),
+ }
+
+ execution = await (execute_claim or execute_kubernetes_controlled_claim)(claim)
+ lifecycle, verifier = _execution_receipts(execution)
+ lifecycle_run_id = str(lifecycle.get("run_id") or "")
+ verifier_run_id = str(verifier.get("run_id") or "")
+ execution_run_id = (
+ lifecycle_run_id
+ if lifecycle_run_id and lifecycle_run_id == verifier_run_id
+ else ""
+ )
+ response = getattr(execution, "k8s_response", None)
+ response = response if isinstance(response, Mapping) else {}
+ runtime_mutation = response.get("runtime_write_performed") is True
+ verifier_closed = bool(
+ getattr(execution, "success", False)
+ and runtime_mutation
+ and lifecycle.get("durable_writeback_ack") is True
+ and verifier.get("verified") is True
+ and verifier.get("durable_writeback_ack") is True
+ and execution_run_id
+ )
+
+ if execution_run_id:
+ try:
+ learning = await (write_learning_receipt or _write_learning_receipt)(
+ candidate, execution, execution_run_id
+ )
+ except Exception as exc: # noqa: BLE001 - closure must fail closed
+ learning = {
+ "run_id": execution_run_id,
+ "durable_writeback_ack": False,
+ "blocker": f"learning_receipt_failed:{type(exc).__name__}",
+ }
+ try:
+ telegram = await (send_telegram_receipt or _send_telegram_receipt)(
+ candidate, execution, execution_run_id
+ )
+ except Exception as exc: # noqa: BLE001 - closure must fail closed
+ telegram = {
+ "run_id": execution_run_id,
+ "durable_writeback_ack": False,
+ "blocker": f"telegram_receipt_failed:{type(exc).__name__}",
+ }
+ else:
+ learning = {"run_id": None, "durable_writeback_ack": False}
+ telegram = {"run_id": None, "durable_writeback_ack": False}
+
+ same_run_receipts = bool(
+ execution_run_id
+ and learning.get("run_id") == execution_run_id
+ and telegram.get("run_id") == execution_run_id
+ )
+ learning_ack = learning.get("durable_writeback_ack") is True
+ telegram_ack = bool(
+ telegram.get("durable_writeback_ack") is True
+ and telegram.get("provider_delivery_acknowledged") is True
+ )
+ closed = bool(
+ verifier_closed and same_run_receipts and learning_ack and telegram_ack
+ )
+ status = (
+ "verified_closed"
+ if closed
+ else "terminal_no_write"
+ if not runtime_mutation
+ else "applied_pending_closure"
+ )
+ closure = {
+ "schema_version": "kubernetes_controlled_candidate_closure_v1",
+ "status": status,
+ "closed": closed,
+ "claim_id": str(claim.get("claim_id") or ""),
+ "incident_id": str(claim.get("incident_id") or ""),
+ "execution_run_id": execution_run_id or None,
+ "runtime_mutation_performed": runtime_mutation,
+ "runtime_write_outcome": response.get("runtime_write_outcome"),
+ "lifecycle_receipt": lifecycle,
+ "verifier_receipt": verifier,
+ "learning_receipt": learning,
+ "telegram_receipt": telegram,
+ "same_run_receipts": same_run_receipts,
+ "replay_apply_allowed": False,
+ "cross_domain_fallback_allowed": False,
+ }
+ finalized = await (finalize_candidate or finalize_kubernetes_controlled_candidate)(
+ candidate, project_id, closure
+ )
+ if not finalized:
+ return {
+ **closure,
+ "status": "durable_closure_writeback_failed",
+ "closed": False,
+ "durable_closure_ack": False,
+ }
+ return {**closure, "durable_closure_ack": True}
diff --git a/apps/api/src/services/service_registry.py b/apps/api/src/services/service_registry.py
index 8c3c56a6c..eff56fffe 100644
--- a/apps/api/src/services/service_registry.py
+++ b/apps/api/src/services/service_registry.py
@@ -81,6 +81,9 @@ class ServiceInfo:
self.containers: list[str] = data.get("containers", [])
self.aliases: list[str] = data.get("aliases", [])
self.allowed_catalog_ids: list[str] = data.get("allowed_catalog_ids", [])
+ self.kubernetes_namespace: str = data.get("kubernetes_namespace", "")
+ self.kubernetes_kind: str = data.get("kubernetes_kind", "")
+ self.kubernetes_name: str = data.get("kubernetes_name", "")
class ServiceRegistryClient:
@@ -189,6 +192,9 @@ class ServiceRegistryClient:
"executor": info.executor,
"verifier": info.verifier,
"allowed_catalog_ids": list(info.allowed_catalog_ids),
+ "kubernetes_namespace": info.kubernetes_namespace or None,
+ "kubernetes_kind": info.kubernetes_kind or None,
+ "kubernetes_name": info.kubernetes_name or None,
"controlled_apply_allowed": info.stateful_level == StatefulLevel.AUTO,
"drift_work_item_id": None,
"registry_error": self._load_error,
diff --git a/apps/api/tests/test_decision_manager_bare_metal_kubectl_guard.py b/apps/api/tests/test_decision_manager_bare_metal_kubectl_guard.py
index d7cc82904..cbfcf2c82 100644
--- a/apps/api/tests/test_decision_manager_bare_metal_kubectl_guard.py
+++ b/apps/api/tests/test_decision_manager_bare_metal_kubectl_guard.py
@@ -75,6 +75,26 @@ def manager(monkeypatch):
"src.jobs.awooop_ansible_candidate_backfill_job.enqueue_ai_decision_ansible_candidate",
_queue_candidate,
)
+ async def _queue_kubernetes_candidate(**_kwargs):
+ return {
+ "schema_version": "ai_decision_controlled_executor_handoff_v1",
+ "status": "controlled_kubernetes_candidate_queued",
+ "automation_run_id": "00000000-0000-0000-0000-000000000102",
+ "queued": True,
+ "side_effect_performed": False,
+ "single_writer_executor": "awoooi-kubernetes-executor-broker",
+ "typed_domain": "kubernetes_workload",
+ "cross_domain_fallback_allowed": False,
+ "active_blockers": [],
+ "next_safe_action": (
+ "kubernetes_controlled_candidate_worker_claims_pending_candidate"
+ ),
+ }
+
+ monkeypatch.setattr(
+ "src.services.kubernetes_controlled_candidate_service.enqueue_ai_decision_kubernetes_candidate",
+ _queue_kubernetes_candidate,
+ )
return mgr
@@ -175,12 +195,48 @@ class TestBareMetalKubectlGuard:
assert token.proposal_data["blocked_reason"] == "critical_break_glass_required"
assert token.proposal_data["direct_runtime_execution_performed"] is False
+ @pytest.mark.asyncio
+ async def test_d037_kubernetes_auto_queues_without_ansible_candidate(self, manager):
+ incident = _fake_incident(
+ host_type="kubernetes",
+ alertname="AwoooPAutoRepairCanaryT16",
+ )
+ incident.incident_id = "INC-20260711-D037E5"
+ incident.affected_services = ["awoooi-auto-repair-canary"]
+ incident.signals[0].labels.update({
+ "namespace": "awoooi-prod",
+ "deployment": "awoooi-auto-repair-canary",
+ "workload_kind": "deployment",
+ })
+ token = _fake_token(
+ "kubectl rollout restart "
+ "deployment/awoooi-auto-repair-canary -n awoooi-prod"
+ )
+
+ with patch(
+ "src.jobs.awooop_ansible_candidate_backfill_job."
+ "enqueue_ai_decision_ansible_candidate",
+ new=AsyncMock(side_effect=AssertionError("wrong executor domain")),
+ ) as ansible_queue:
+ await manager._auto_execute(incident, token)
+
+ assert token.state == DecisionState.EXECUTING
+ assert token.proposal_data["single_writer_executor"] == (
+ "awoooi-kubernetes-executor-broker"
+ )
+ assert token.proposal_data["check_mode_required_before_apply"] is False
+ assert token.proposal_data["typed_dry_run_required_before_apply"] is True
+ assert token.proposal_data["direct_runtime_execution_performed"] is False
+ ansible_queue.assert_not_awaited()
+
def test_active_auto_execute_only_routes_to_single_writer_queue() -> None:
source = inspect.getsource(DecisionManager._auto_execute)
legacy_source = inspect.getsource(DecisionManager._legacy_auto_execute_disabled)
+ assert "enqueue_ai_decision_kubernetes_candidate" in source
assert "enqueue_ai_decision_ansible_candidate" in source
+ assert "kubernetes_action" in source
assert "ApprovalExecutionService" not in source
assert "execute_approved_action" not in source
assert "_ssh_execute" not in source
diff --git a/apps/api/tests/test_kubernetes_controlled_candidate_service.py b/apps/api/tests/test_kubernetes_controlled_candidate_service.py
new file mode 100644
index 000000000..6a516585e
--- /dev/null
+++ b/apps/api/tests/test_kubernetes_controlled_candidate_service.py
@@ -0,0 +1,345 @@
+from __future__ import annotations
+
+from collections.abc import Mapping
+from unittest.mock import AsyncMock
+
+import pytest
+
+from src.jobs.kubernetes_controlled_candidate_job import (
+ run_kubernetes_controlled_candidate_once,
+)
+from src.services.executor import ExecutionResult, OperationType
+from src.services.kubernetes_controlled_candidate_service import (
+ enqueue_ai_decision_kubernetes_candidate,
+ run_claimed_kubernetes_controlled_candidate,
+)
+
+
+class _Incident:
+ project_id = "awoooi"
+ incident_id = "INC-20260711-D037E5"
+
+ def __init__(self, target: str = "awoooi-auto-repair-canary") -> None:
+ self.target = target
+
+ def model_dump(self, *, mode: str = "json") -> dict[str, object]:
+ assert mode == "json"
+ return {
+ "incident_id": self.incident_id,
+ "project_id": self.project_id,
+ "affected_services": [self.target],
+ "signals": [
+ {
+ "alert_name": "AwoooPAutoRepairCanaryT16",
+ "labels": {
+ "alertname": "AwoooPAutoRepairCanaryT16",
+ "namespace": "awoooi-prod",
+ "deployment": self.target,
+ "host_type": "kubernetes",
+ },
+ }
+ ],
+ }
+
+
+def _proposal() -> dict[str, str]:
+ return {
+ "kubectl_command": (
+ "kubectl rollout restart "
+ "deployment/awoooi-auto-repair-canary -n awoooi-prod"
+ ),
+ "risk_level": "medium",
+ "matched_playbook_id": "pb:auto-repair-canary",
+ }
+
+
+@pytest.mark.asyncio
+async def test_d037_queues_typed_kubernetes_candidate_without_ansible_fallback() -> (
+ None
+):
+ persisted: list[dict[str, object]] = []
+
+ async def persist(record: dict[str, object], project_id: str) -> dict[str, object]:
+ assert project_id == "awoooi"
+ persisted.append(record)
+ return {
+ "event_id": "00000000-0000-0000-0000-000000000033",
+ "created": True,
+ "record": record,
+ }
+
+ handoff = await enqueue_ai_decision_kubernetes_candidate(
+ incident=_Incident(),
+ proposal_data=_proposal(),
+ persist_candidate=persist,
+ )
+
+ assert handoff["status"] == "controlled_kubernetes_candidate_queued"
+ assert handoff["queued"] is True
+ assert handoff["single_writer_executor"] == (
+ "awoooi-kubernetes-executor-broker"
+ )
+ assert handoff["typed_domain"] == "kubernetes_workload"
+ assert handoff["canonical_asset_id"] == ("service:awoooi-auto-repair-canary")
+ assert handoff["cross_domain_fallback_allowed"] is False
+ assert len(persisted) == 1
+ record = persisted[0]
+ assert record["legacy_ansible_candidate_consulted"] is False
+ assert record["runtime_mutation_performed"] is False
+ claim = record["claim"]
+ assert isinstance(claim, Mapping)
+ assert claim["executor"] == "kubernetes_controlled_executor"
+ assert claim["verifier"] == "kubernetes_rollout_verifier"
+
+
+@pytest.mark.asyncio
+async def test_exact_duplicate_reuses_durable_candidate_instead_of_rebuilding() -> None:
+ first_record: dict[str, object] = {}
+
+ async def persist(record: dict[str, object], _project_id: str) -> dict[str, object]:
+ nonlocal first_record
+ if not first_record:
+ first_record = record
+ created = True
+ else:
+ created = False
+ return {
+ "event_id": "00000000-0000-0000-0000-000000000033",
+ "created": created,
+ "record": first_record,
+ }
+
+ first = await enqueue_ai_decision_kubernetes_candidate(
+ incident=_Incident(),
+ proposal_data=_proposal(),
+ persist_candidate=persist,
+ )
+ duplicate = await enqueue_ai_decision_kubernetes_candidate(
+ incident=_Incident(),
+ proposal_data=_proposal(),
+ persist_candidate=persist,
+ )
+
+ assert first["deduplicated"] is False
+ assert duplicate["deduplicated"] is True
+ assert duplicate["status"] == "controlled_kubernetes_candidate_exists"
+ assert duplicate["claim_id"] == first["claim_id"]
+ assert duplicate["automation_run_id"] == first["automation_run_id"]
+
+
+@pytest.mark.asyncio
+async def test_unknown_asset_fails_closed_without_candidate_or_runtime_mutation() -> (
+ None
+):
+ persist = AsyncMock()
+
+ handoff = await enqueue_ai_decision_kubernetes_candidate(
+ incident=_Incident(target="unregistered-canary"),
+ proposal_data={
+ "kubectl_command": (
+ "kubectl rollout restart deployment/unregistered-canary "
+ "-n awoooi-prod"
+ )
+ },
+ persist_candidate=persist,
+ )
+
+ assert handoff["status"] == "kubernetes_candidate_blocked_no_write"
+ assert handoff["queued"] is False
+ assert handoff["runtime_mutation_performed"] is False
+ assert handoff["cross_domain_fallback_allowed"] is False
+ assert handoff["active_blockers"]
+ persist.assert_not_awaited()
+
+
+def _claimed_candidate() -> dict[str, object]:
+ return {
+ "schema_version": "kubernetes_controlled_candidate_v1",
+ "status": "claimed",
+ "event_id": "00000000-0000-0000-0000-000000000033",
+ "claim_token": "00000000-0000-0000-0000-000000000034",
+ "automation_run_id": "00000000-0000-0000-0000-000000000035",
+ "matched_playbook_id": "pb:auto-repair-canary",
+ "claim": {
+ "claim_id": "k8s-claim:d037e5",
+ "incident_id": "INC-20260711-D037E5",
+ "canonical_asset_id": "service:awoooi-auto-repair-canary",
+ "workload_kind": "deployment",
+ "workload_name": "awoooi-auto-repair-canary",
+ "namespace": "awoooi-prod",
+ },
+ }
+
+
+def _verified_execution(run_id: str = "run:k8s:d037e5") -> ExecutionResult:
+ return ExecutionResult(
+ success=True,
+ message="verified",
+ operation_type=OperationType.RESTART_DEPLOYMENT,
+ target_resource="deployment/awoooi-auto-repair-canary",
+ namespace="awoooi-prod",
+ duration_ms=25,
+ k8s_response={
+ "runtime_write_performed": True,
+ "runtime_write_outcome": "confirmed_applied",
+ "lifecycle_receipt": {
+ "receipt_id": "lifecycle:d037e5",
+ "run_id": run_id,
+ "durable_writeback_ack": True,
+ },
+ "verifier_receipt": {
+ "receipt_id": "verifier:d037e5",
+ "run_id": run_id,
+ "verified": True,
+ "durable_writeback_ack": True,
+ },
+ },
+ )
+
+
+@pytest.mark.asyncio
+async def test_single_dispatch_closes_only_with_same_run_receipts() -> None:
+ dispatch_started = False
+ execute = AsyncMock(return_value=_verified_execution())
+ finalized: list[Mapping[str, object]] = []
+
+ async def begin(_candidate: Mapping[str, object], _project: str) -> bool:
+ nonlocal dispatch_started
+ if dispatch_started:
+ return False
+ dispatch_started = True
+ return True
+
+ async def learning(
+ _candidate: Mapping[str, object],
+ _execution: object,
+ run_id: str,
+ ) -> dict[str, object]:
+ return {
+ "receipt_id": "learning:d037e5",
+ "run_id": run_id,
+ "durable_writeback_ack": True,
+ }
+
+ async def telegram(
+ _candidate: Mapping[str, object],
+ _execution: object,
+ run_id: str,
+ ) -> dict[str, object]:
+ return {
+ "receipt_id": "telegram:d037e5",
+ "run_id": run_id,
+ "provider_delivery_acknowledged": True,
+ "durable_writeback_ack": True,
+ }
+
+ async def finalize(
+ _candidate: Mapping[str, object],
+ _project: str,
+ closure: Mapping[str, object],
+ ) -> bool:
+ finalized.append(closure)
+ return True
+
+ first = await run_claimed_kubernetes_controlled_candidate(
+ _claimed_candidate(),
+ begin_dispatch=begin,
+ execute_claim=execute,
+ write_learning_receipt=learning,
+ send_telegram_receipt=telegram,
+ finalize_candidate=finalize,
+ )
+ duplicate = await run_claimed_kubernetes_controlled_candidate(
+ _claimed_candidate(),
+ begin_dispatch=begin,
+ execute_claim=execute,
+ write_learning_receipt=learning,
+ send_telegram_receipt=telegram,
+ finalize_candidate=finalize,
+ )
+
+ assert first["status"] == "verified_closed"
+ assert first["closed"] is True
+ assert first["same_run_receipts"] is True
+ assert first["durable_closure_ack"] is True
+ assert first["replay_apply_allowed"] is False
+ assert duplicate["status"] == "duplicate_apply_suppressed"
+ assert duplicate["runtime_mutation_performed"] is False
+ execute.assert_awaited_once()
+ assert len(finalized) == 1
+
+
+@pytest.mark.asyncio
+async def test_run_mismatch_is_durable_pending_closure_and_never_reapply() -> None:
+ execution = _verified_execution()
+ execute = AsyncMock(return_value=execution)
+ finalized: list[Mapping[str, object]] = []
+
+ async def begin(_candidate: Mapping[str, object], _project: str) -> bool:
+ return True
+
+ async def learning(
+ _candidate: Mapping[str, object],
+ _execution: object,
+ _run_id: str,
+ ) -> dict[str, object]:
+ return {
+ "run_id": "run:k8s:wrong",
+ "durable_writeback_ack": True,
+ }
+
+ async def telegram(
+ _candidate: Mapping[str, object],
+ _execution: object,
+ run_id: str,
+ ) -> dict[str, object]:
+ return {
+ "run_id": run_id,
+ "provider_delivery_acknowledged": True,
+ "durable_writeback_ack": True,
+ }
+
+ async def finalize(
+ _candidate: Mapping[str, object],
+ _project: str,
+ closure: Mapping[str, object],
+ ) -> bool:
+ finalized.append(closure)
+ return True
+
+ result = await run_claimed_kubernetes_controlled_candidate(
+ _claimed_candidate(),
+ begin_dispatch=begin,
+ execute_claim=execute,
+ write_learning_receipt=learning,
+ send_telegram_receipt=telegram,
+ finalize_candidate=finalize,
+ )
+
+ assert result["status"] == "applied_pending_closure"
+ assert result["closed"] is False
+ assert result["same_run_receipts"] is False
+ assert result["replay_apply_allowed"] is False
+ assert finalized[0]["runtime_mutation_performed"] is True
+
+
+@pytest.mark.asyncio
+async def test_dedicated_worker_claims_only_one_candidate_per_tick() -> None:
+ claim_next = AsyncMock(return_value=_claimed_candidate())
+ run_candidate = AsyncMock(
+ return_value={"status": "verified_closed", "closed": True}
+ )
+
+ stats = await run_kubernetes_controlled_candidate_once(
+ batch_limit=99,
+ claim_next=claim_next,
+ run_candidate=run_candidate,
+ )
+
+ assert stats["claimed"] == 1
+ assert stats["verified_closed"] == 1
+ assert stats["single_writer_executor"] == (
+ "awoooi-kubernetes-executor-broker"
+ )
+ claim_next.assert_awaited_once_with(project_id="awoooi")
+ run_candidate.assert_awaited_once()
diff --git a/apps/api/tests/test_sre_k3s_controlled_automation_work_items_api.py b/apps/api/tests/test_sre_k3s_controlled_automation_work_items_api.py
index 231c99880..156e490d6 100644
--- a/apps/api/tests/test_sre_k3s_controlled_automation_work_items_api.py
+++ b/apps/api/tests/test_sre_k3s_controlled_automation_work_items_api.py
@@ -155,8 +155,8 @@ def test_loader_returns_fixed_architecture_provider_order_and_agent99_bridge() -
"active_or_completed_commitments": 68,
"by_status": {
"analysis_or_governance_complete": 6,
- "source_implemented_runtime_pending": 30,
- "in_progress": 30,
+ "source_implemented_runtime_pending": 31,
+ "in_progress": 29,
"not_started_or_no_current_evidence": 2,
"superseded": 2,
},
@@ -188,6 +188,12 @@ def test_loader_returns_fixed_architecture_provider_order_and_agent99_bridge() -
assert "blocks Recover before single-flight claim" in " ".join(
commitments["AIA-CONV-032"]["source_evidence"]
)
+ assert commitments["AIA-CONV-033"]["status"] == (
+ "source_implemented_runtime_pending"
+ )
+ assert "single-dispatch barrier" in " ".join(
+ commitments["AIA-CONV-033"]["source_evidence"]
+ )
assert "Host112" in commitments["AIA-CONV-049"]["title"]
assert commitments["AIA-CONV-049"]["status"] == (
"source_implemented_runtime_pending"
diff --git a/docs/operations/sre-ai-agent-conversation-commitments.snapshot.json b/docs/operations/sre-ai-agent-conversation-commitments.snapshot.json
index f7ab15a50..10a92f129 100644
--- a/docs/operations/sre-ai-agent-conversation-commitments.snapshot.json
+++ b/docs/operations/sre-ai-agent-conversation-commitments.snapshot.json
@@ -50,7 +50,7 @@
{"id":"AIA-CONV-031","category":"telegram","status":"source_implemented_runtime_pending","title":"同 fingerprint 去重、聚合、抑噪及 resolved/recovered 更新","linked_work_items":["AIA-SRE-017"],"terminal_condition":"同事件只維護 canonical lifecycle,恢復卡有 destination-bound provider acknowledgement。","source_evidence":["existing firing fingerprint convergence and parent-child alert storm aggregation","single-row durable resolved lifecycle reservation with bounded retry","duplicate resolved webhook suppression before incident or Telegram side effects","recovery card accepted only with destination-bound provider acknowledgement"]},
{"id":"AIA-CONV-032","category":"named_incident","status":"source_implemented_runtime_pending","title":"INC-20260711-11C751 cold-start-gate 補 PlayBook、transport、rollback 與 verifier","linked_work_items":["AIA-SRE-004","AIA-SRE-006","AIA-SRE-015"],"terminal_condition":"namespace identity 修正後完成一次 bounded same-run repair 或明確 no-write terminal。","source_evidence":["source namespace separated from normalized host_recovery execution identity","exact inventory scope absence blocks Recover before single-flight claim or Agent99 transport","durable alert operation receipt records the explicit no-write terminal","consumer-validated source no-write verifier receipt with runtime closure false"]},
- {"id":"AIA-CONV-033","category":"named_incident","status":"in_progress","title":"INC-20260711-D037E5 修復 stale candidate 阻擋、自動排隊與 durable closure","linked_work_items":["AIA-SRE-007","AIA-SRE-015"],"terminal_condition":"新候選可 claim,apply 不重複,verifier、Telegram 與 learning receipts 同 run 完成。"},
+ {"id":"AIA-CONV-033","category":"named_incident","status":"source_implemented_runtime_pending","title":"INC-20260711-D037E5 修復 stale candidate 阻擋、自動排隊與 durable closure","linked_work_items":["AIA-SRE-007","AIA-SRE-015"],"terminal_condition":"新候選可 claim,apply 不重複,verifier、Telegram 與 learning receipts 同 run 完成。","source_evidence":["D037E5 canonical workload identity resolves to the Kubernetes executor/verifier without consulting legacy Ansible candidates","fingerprint-stable durable candidate event supports pending and pre-dispatch lease-expired atomic claim","dedicated feature-gated Kubernetes executor job owns a bounded one-candidate claim tick and is not started by API, signal worker or Ansible broker","single-dispatch barrier suppresses duplicate apply after dispatch starts","durable closure requires lifecycle, independent verifier, Telegram provider acknowledgement and learning receipts on one execution run"]},
{"id":"AIA-CONV-034","category":"named_incident","status":"in_progress","title":"AWOOOI CPU 高負載與 P99 上升完成 RCA、bounded repair 與 verifier","linked_work_items":["AIA-SRE-004","AIA-SRE-006","AIA-SRE-010"],"terminal_condition":"相關 metrics/logs/changes 對齊 canonical asset,修復後 latency 與資源 verifier 關閉。"},
{"id":"AIA-CONV-035","category":"named_incident","status":"in_progress","title":"Backup/restore escrow 補 freshness、escrow metadata、restore drill 與 DR scorecard","linked_work_items":["AIA-SRE-011","AIA-SRE-016"],"terminal_condition":"read-only evidence 齊全;缺 escrow/restore receipt 時保持 blocked 而非假綠。"},
{"id":"AIA-CONV-036","category":"named_incident","status":"source_implemented_runtime_pending","title":"Sentry Snuba DockerContainerUnhealthy 反覆告警自動修復","linked_work_items":["AIA-SRE-005"],"terminal_condition":"exact-container replay 完成 bounded recovery、健康 verifier 與 recurrence fence。"},
@@ -98,8 +98,8 @@
"active_or_completed_commitments": 68,
"by_status": {
"analysis_or_governance_complete": 6,
- "source_implemented_runtime_pending": 30,
- "in_progress": 30,
+ "source_implemented_runtime_pending": 31,
+ "in_progress": 29,
"not_started_or_no_current_evidence": 2,
"superseded": 2
},
diff --git a/k8s/awoooi-prod/15-service-registry-configmap.yaml b/k8s/awoooi-prod/15-service-registry-configmap.yaml
index 8e9f31666..1f40fa2dc 100644
--- a/k8s/awoooi-prod/15-service-registry-configmap.yaml
+++ b/k8s/awoooi-prod/15-service-registry-configmap.yaml
@@ -9,7 +9,7 @@ metadata:
app: awoooi
component: service-registry
annotations:
- awoooi.wooo.work/source-sha256: "3977ad4862d761e93fdc13562d5a60708631b26d8cf93fbc2f5c50da3c071db0"
+ awoooi.wooo.work/source-sha256: "1a2777467c6d7b7912617f2fa6d22f590fa95a469aab0a8ff477d9725b425854"
data:
service-registry.yaml: |
# ops/config/service-registry.yaml
@@ -199,6 +199,20 @@ data:
asset_domain: kubernetes_workload
containers: []
+ - name: awoooi-auto-repair-canary
+ canonical_id: "service:awoooi-auto-repair-canary"
+ display_name: "AWOOOI Auto Repair Canary (K3s)"
+ host: "k3s"
+ stateful_level: AUTO
+ asset_domain: kubernetes_workload
+ executor: kubernetes_controlled_executor
+ verifier: kubernetes_rollout_verifier
+ kubernetes_namespace: awoooi-prod
+ kubernetes_kind: deployment
+ kubernetes_name: awoooi-auto-repair-canary
+ aliases: ["AwoooPAutoRepairCanaryT16"]
+ containers: []
+
- name: blackbox-exporter
display_name: "Blackbox Exporter"
host: "192.168.0.110"
diff --git a/ops/config/service-registry.yaml b/ops/config/service-registry.yaml
index c4b725e90..c08392609 100644
--- a/ops/config/service-registry.yaml
+++ b/ops/config/service-registry.yaml
@@ -185,6 +185,20 @@ services:
asset_domain: kubernetes_workload
containers: []
+ - name: awoooi-auto-repair-canary
+ canonical_id: "service:awoooi-auto-repair-canary"
+ display_name: "AWOOOI Auto Repair Canary (K3s)"
+ host: "k3s"
+ stateful_level: AUTO
+ asset_domain: kubernetes_workload
+ executor: kubernetes_controlled_executor
+ verifier: kubernetes_rollout_verifier
+ kubernetes_namespace: awoooi-prod
+ kubernetes_kind: deployment
+ kubernetes_name: awoooi-auto-repair-canary
+ aliases: ["AwoooPAutoRepairCanaryT16"]
+ containers: []
+
- name: blackbox-exporter
display_name: "Blackbox Exporter"
host: "192.168.0.110"