From 0aee24576155d8f829c4ca4cff424376e7f05af1 Mon Sep 17 00:00:00 2001 From: Your Name Date: Wed, 22 Jul 2026 16:00:05 +0800 Subject: [PATCH] feat(aiops): add typed CPU P99 recovery closure --- .../controlled_alert_target_router.py | 22 +- .../services/cpu_p99_controlled_recovery.py | 1423 +++++++++++++++++ .../services/independent_verifier_registry.py | 14 + .../tests/test_cpu_p99_controlled_recovery.py | 552 +++++++ ...3s_controlled_automation_work_items_api.py | 4 +- ...ent-conversation-commitments.snapshot.json | 8 +- .../15-service-registry-configmap.yaml | 7 +- ops/config/service-registry.yaml | 5 + ops/monitoring/alerts.yml | 2 +- ops/signoz/alerting/rules.yaml | 9 + 10 files changed, 2036 insertions(+), 10 deletions(-) create mode 100644 apps/api/src/services/cpu_p99_controlled_recovery.py create mode 100644 apps/api/tests/test_cpu_p99_controlled_recovery.py diff --git a/apps/api/src/services/controlled_alert_target_router.py b/apps/api/src/services/controlled_alert_target_router.py index 8fb42e847..733e5b2c0 100644 --- a/apps/api/src/services/controlled_alert_target_router.py +++ b/apps/api/src/services/controlled_alert_target_router.py @@ -67,6 +67,10 @@ _NODE_EXPORTER_ALERTS = { "nodeexporterscrapedown", "nodeexporterunhealthy", } +_HOST_CPU_PRESSURE_ALERTS = { + "hosthighcpuload", + "hostloadaveragesustainedhigh", +} _ALERTMANAGER_DELIVERY_ALERTS = {"alertchainbrokenalertmanager"} _PROVIDER_FRESHNESS_ALERTS = { "providerfreshnesssignal", @@ -704,6 +708,18 @@ def resolve_typed_alert_target( host110_pressure_alert = compact_alert.startswith( "host110sustainedmoderatepressure" ) + host_cpu_pressure_alert = ( + compact_alert in _HOST_CPU_PRESSURE_ALERTS or host110_pressure_alert + ) + pressure_target_identity = ( + registry.resolve_identity(target_resource) + if host_cpu_pressure_alert and target_resource + else {} + ) + pressure_target_is_exact_service = bool( + pressure_target_identity.get("resolution_status") == "resolved" + and pressure_target_identity.get("asset_domain") != "unknown" + ) allowed_catalog_ids: list[str] = [] if compact_alert in _DISK_ALERTS and host_scope is not None: if host_scope[0] in {"host_110", "host_188"}: @@ -714,9 +730,10 @@ def resolve_typed_alert_target( if host_scope[0] == "host_110": allowed_catalog_ids = ["ansible:110-devops"] elif ( - host110_pressure_alert + host_cpu_pressure_alert and host_scope is not None and host_scope[0] == "host_110" + and not pressure_target_is_exact_service ): allowed_catalog_ids = ["ansible:110-host-pressure-readonly"] elif ( @@ -729,7 +746,7 @@ def resolve_typed_alert_target( if host_scope is not None and ( compact_alert in _DISK_ALERTS or compact_alert in _NODE_EXPORTER_ALERTS - or host110_pressure_alert + or (host_cpu_pressure_alert and not pressure_target_is_exact_service) or "wazuh" in text ): return _typed_route( @@ -744,6 +761,7 @@ def resolve_typed_alert_target( risk_class="medium", allowed_catalog_ids=allowed_catalog_ids, allowed_inventory_hosts=[host_scope[0]], + controlled_apply_allowed=(False if host_cpu_pressure_alert else None), ) identity = registry.resolve_identity(target_resource) diff --git a/apps/api/src/services/cpu_p99_controlled_recovery.py b/apps/api/src/services/cpu_p99_controlled_recovery.py new file mode 100644 index 000000000..4b03c6f27 --- /dev/null +++ b/apps/api/src/services/cpu_p99_controlled_recovery.py @@ -0,0 +1,1423 @@ +"""Durable CPU/P99 correlation, typed candidate, and dual-verifier closure. + +This module is a control-plane contract. It accepts only receipt-backed, +public-safe evidence and never queries a provider or invokes an executor. A +runtime worker may consume the durable candidate later, but only the typed +domain executor named by the canonical asset route may apply a repair. +""" + +from __future__ import annotations + +import hashlib +import json +import re +from collections.abc import Awaitable, Callable, Mapping +from datetime import datetime +from typing import Any +from uuid import UUID + +import structlog +from sqlalchemy import text + +from src.services.audit_sink import sanitize +from src.services.controlled_alert_target_router import ( + resolve_typed_alert_target, +) +from src.services.service_registry import get_service_registry + +logger = structlog.get_logger(__name__) + +CANDIDATE_SCHEMA_VERSION = "cpu_p99_controlled_recovery_candidate_v1" +CLOSURE_SCHEMA_VERSION = "cpu_p99_controlled_recovery_closure_v1" +HANDOFF_SCHEMA_VERSION = "cpu_p99_controlled_recovery_handoff_v1" +API_CANONICAL_ASSET_ID = "service:awoooi-api" +RESOURCE_VERIFIER = "cpu_resource_independent_verifier" +LATENCY_VERIFIER = "signoz_api_latency_independent_verifier" +RESOURCE_READBACK_SCHEMA_VERSION = "prometheus_cpu_resource_post_readback_v1" +LATENCY_READBACK_SCHEMA_VERSION = "signoz_api_p99_post_readback_v1" +RESOURCE_VERIFIER_SCHEMA_VERSION = "cpu_resource_post_verifier_v1" +LATENCY_VERIFIER_SCHEMA_VERSION = "signoz_api_p99_post_verifier_v1" +RESOURCE_READBACK_SOURCE = "prometheus_independent_query" +LATENCY_READBACK_SOURCE = "signoz_independent_query" +MAX_EVIDENCE_SKEW_SECONDS = 15 * 60 + +_SAFE_ID = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._:@+-]{0,191}$") +_RESOURCE_SOURCES = { + "prometheus_readonly_query", + "host_sustained_load_controller", +} +_RESOURCE_METRICS = { + "cpu_usage_percent", + "container_cpu_cores", + "load5_per_core", +} + +PersistCandidate = Callable[[dict[str, Any], str], Awaitable[dict[str, Any]]] +PersistClosure = Callable[[Mapping[str, Any], Mapping[str, Any], str], Awaitable[bool]] + + +def _safe_text(value: Any) -> str: + return str(value or "").strip() + + +def _number(value: Any) -> float | None: + if isinstance(value, bool) or not isinstance(value, int | float): + return None + return float(value) + + +def _valid_public_receipt_id(value: Any) -> bool: + receipt_id = _safe_text(value) + return bool( + _SAFE_ID.fullmatch(receipt_id) + and re.fullmatch(r"[0-9a-f]{64}", receipt_id) is None + ) + + +def _timestamp(receipt: Mapping[str, Any], field: str) -> datetime | None: + raw = _safe_text(receipt.get(field)) + if not raw: + return None + try: + value = datetime.fromisoformat(raw.replace("Z", "+00:00")) + except ValueError: + return None + return value if value.tzinfo is not None else None + + +def _observed_at(receipt: Mapping[str, Any]) -> datetime | None: + return _timestamp(receipt, "observed_at") + + +def _common_receipt_fields( + receipts: tuple[Mapping[str, Any], ...], +) -> tuple[dict[str, str] | None, list[str]]: + blockers: list[str] = [] + first = receipts[0] + correlation = { + field: _safe_text(first.get(field)) + for field in ("trace_id", "run_id", "work_item_id") + } + for field, value in correlation.items(): + if not _SAFE_ID.fullmatch(value): + blockers.append(f"{field}_missing_or_invalid") + try: + UUID(correlation["run_id"]) + except (TypeError, ValueError): + blockers.append("run_id_not_uuid") + + observed: list[datetime] = [] + receipt_ids: list[str] = [] + for receipt in receipts: + receipt_id = _safe_text(receipt.get("receipt_id")) + if not _valid_public_receipt_id(receipt_id): + blockers.append("evidence_receipt_id_missing_or_invalid") + else: + receipt_ids.append(receipt_id) + max_age_seconds = _number(receipt.get("max_age_seconds")) + if ( + receipt.get("freshness_verified") is not True + or max_age_seconds is None + or max_age_seconds <= 0 + or max_age_seconds > MAX_EVIDENCE_SKEW_SECONDS + ): + blockers.append("evidence_freshness_unverified") + for field, expected in correlation.items(): + if _safe_text(receipt.get(field)) != expected: + blockers.append(f"same_run_{field}_mismatch") + timestamp = _observed_at(receipt) + if timestamp is None: + blockers.append("evidence_observed_at_missing_or_invalid") + else: + observed.append(timestamp) + + if len(receipt_ids) != len(set(receipt_ids)): + blockers.append("evidence_receipt_replayed") + + if observed: + try: + skew = (max(observed) - min(observed)).total_seconds() + except TypeError: + blockers.append("evidence_timezone_mismatch") + else: + if skew > MAX_EVIDENCE_SKEW_SECONDS: + blockers.append("evidence_window_not_correlated") + blockers = list(dict.fromkeys(blockers)) + return (None if blockers else correlation), blockers + + +def _candidate_fingerprint( + *, + correlation: Mapping[str, str], + canonical_asset_id: str, + receipts: tuple[Mapping[str, Any], ...], + route_binding: Mapping[str, Any], +) -> str: + material = { + "schema_version": CANDIDATE_SCHEMA_VERSION, + "trace_id": correlation["trace_id"], + "run_id": correlation["run_id"], + "work_item_id": correlation["work_item_id"], + "canonical_asset_id": canonical_asset_id, + "evidence_receipt_ids": [ + _safe_text(receipt.get("receipt_id")) for receipt in receipts + ], + "route_binding": dict(sorted(route_binding.items())), + } + canonical = json.dumps(material, sort_keys=True, separators=(",", ":")) + return hashlib.sha256(canonical.encode("utf-8")).hexdigest() + + +def _derived_receipt_id(prefix: str, *values: str) -> str: + digest = hashlib.sha256("|".join(values).encode("utf-8")).hexdigest()[:24] + return f"{prefix}:{digest}" + + +def _public_safe_record(record: Mapping[str, Any]) -> dict[str, Any]: + """Sanitize payload fields while preserving this module's computed digest IDs.""" + + safe = sanitize(dict(record)) + for field in ("fingerprint", "candidate_fingerprint"): + value = _safe_text(record.get(field)) + if re.fullmatch(r"[0-9a-f]{64}", value): + safe[field] = value + return safe + + +def _candidate_record_matches(expected: Mapping[str, Any], stored: Any) -> bool: + if not isinstance(stored, Mapping): + return False + immutable_fields = ( + "schema_version", + "kind", + "status", + "trace_id", + "run_id", + "work_item_id", + "fingerprint", + "canonical_asset_id", + "typed_domain", + "typed_route", + "evidence", + "repair_candidate", + "candidate_created", + "repair_dispatch_allowed", + "next_safe_action", + "policy", + "runtime_mutation_performed", + "cross_domain_fallback_allowed", + "source_commitment", + ) + public_expected = _public_safe_record(expected) + raw_match = all( + stored.get(field) == expected.get(field) for field in immutable_fields + ) + public_safe_match = all( + stored.get(field) == public_expected.get(field) for field in immutable_fields + ) + return raw_match or public_safe_match + + +async def persist_cpu_p99_candidate( + record: dict[str, Any], + project_id: str, +) -> dict[str, Any]: + """Insert one internal candidate event or return its exact duplicate.""" + + from src.db.base import get_db_context + + fingerprint = _safe_text(record["fingerprint"]) + provider_event_id = f"cpu-p99-controlled-candidate:{fingerprint}" + safe_record = _public_safe_record(record) + source_envelope = json.dumps(safe_record, ensure_ascii=False, default=str) + preview = f"cpu_p99:{record['status']}:{record['work_item_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 + """ + ), + { + "project_id": project_id, + "provider_event_id": provider_event_id, + "run_id": UUID(_safe_text(record["run_id"])), + "content_hash": fingerprint, + "preview": preview, + "source_envelope": source_envelope, + }, + ) + row = inserted.fetchone() + if row is not None: + return { + "event_id": str(row[0]), + "created": True, + "record": safe_record, + } + + 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("cpu_p99_candidate_dedupe_receipt_missing") + stored = existing_row[1] + if isinstance(stored, str): + stored = json.loads(stored) + if ( + not isinstance(stored, Mapping) + or stored.get("schema_version") != CANDIDATE_SCHEMA_VERSION + or stored.get("fingerprint") != fingerprint + or stored.get("work_item_id") != record["work_item_id"] + ): + raise RuntimeError("cpu_p99_candidate_dedupe_receipt_mismatch") + return { + "event_id": str(existing_row[0]), + "created": False, + "record": dict(stored), + } + + +def _blocked_result( + *, + status: str, + blockers: list[str], + next_safe_action: str, + correlation: Mapping[str, str] | None = None, +) -> dict[str, Any]: + correlation = correlation or {} + return { + "schema_version": HANDOFF_SCHEMA_VERSION, + "accepted": False, + "created": False, + "deduplicated": False, + "candidate_created": False, + "repair_dispatch_allowed": False, + "status": status, + "trace_id": correlation.get("trace_id"), + "run_id": correlation.get("run_id"), + "work_item_id": correlation.get("work_item_id"), + "active_blockers": blockers, + "next_safe_action": next_safe_action, + "provider_call_performed": False, + "paid_provider_call_performed": False, + "executor_invoked": False, + "agent99_dispatch_performed": False, + "runtime_mutation_performed": False, + "cross_domain_fallback_allowed": False, + } + + +def _repair_contract(route: Mapping[str, Any]) -> dict[str, Any]: + target_kind = _safe_text(route.get("target_kind")) + executor = _safe_text(route.get("executor")) + verifier = _safe_text(route.get("verifier")) + if ( + target_kind == "kubernetes_workload" + and executor == "kubernetes_controlled_executor" + and verifier == "kubernetes_rollout_verifier" + and route.get("controlled_apply_allowed") is True + ): + return { + "status": "bounded_repair_candidate_recorded", + "candidate_created": True, + "repair_dispatch_allowed": True, + "action_id": "kubernetes_rollout_restart_v1", + "single_writer_executor": "awoooi-kubernetes-executor-broker", + "next_safe_action": ( + "kubernetes_controlled_candidate_worker_runs_typed_dry_run_" + "then_single_apply" + ), + } + if target_kind == "host_systemd" and route.get("allowed_catalog_ids"): + return { + "status": "no_write_host_rca_candidate_recorded", + "candidate_created": True, + "repair_dispatch_allowed": False, + "action_id": _safe_text(route["allowed_catalog_ids"][0]), + "single_writer_executor": "awoooi-ansible-executor-broker", + "next_safe_action": ( + "host_ansible_worker_runs_exact_readonly_pressure_playbook_" + "then_emits_source_specific_repair_candidate" + ), + } + if target_kind == "database": + return { + "status": "no_write_database_rca_candidate_recorded", + "candidate_created": True, + "repair_dispatch_allowed": False, + "action_id": "database_readonly_cpu_p99_rca_v1", + "single_writer_executor": "awoooi-db-bounded-executor-broker", + "next_safe_action": ( + "db_independent_verifier_reads_query_lock_and_capacity_" + "evidence_before_any_bounded_db_candidate" + ), + } + return { + "status": "typed_domain_playbook_gap_recorded", + "candidate_created": False, + "repair_dispatch_allowed": False, + "action_id": "none", + "single_writer_executor": "none", + "next_safe_action": ( + "create_exact_same_domain_bounded_playbook_before_dispatch" + ), + } + + +async def build_cpu_p99_controlled_candidate( + *, + root_cause_target: str, + resource_evidence: Mapping[str, Any], + latency_evidence: Mapping[str, Any], + log_evidence: Mapping[str, Any], + recent_change_evidence: Mapping[str, Any], + project_id: str = "awoooi", + registry: Any | None = None, + persist_candidate: PersistCandidate | None = None, +) -> dict[str, Any]: + """Persist a typed candidate only after the four evidence lanes align.""" + + receipts = ( + resource_evidence, + latency_evidence, + log_evidence, + recent_change_evidence, + ) + correlation, blockers = _common_receipt_fields(receipts) + if correlation is None: + return _blocked_result( + status="correlation_evidence_unverified", + blockers=blockers, + next_safe_action=("recollect_resource_latency_logs_and_changes_in_one_run"), + ) + + evidence_labels = resource_evidence.get("labels") + labels = evidence_labels if isinstance(evidence_labels, Mapping) else {} + registry = registry or get_service_registry() + route = resolve_typed_alert_target( + alertname=_safe_text(resource_evidence.get("alertname") or "HostHighCpuLoad"), + target_resource=root_cause_target, + namespace=_safe_text(resource_evidence.get("namespace")), + labels=labels, + alert_category=_safe_text( + resource_evidence.get("alert_category") or "host_resource" + ), + registry=registry, + ) + if route.get("resolution_status") != "resolved": + canonical_asset_id = "" + fingerprint = _candidate_fingerprint( + correlation=correlation, + canonical_asset_id=( + _safe_text(route.get("drift_work_item_id")) or "unresolved" + ), + receipts=receipts, + route_binding={ + "typed_domain": "unknown", + "drift_work_item_id": _safe_text(route.get("drift_work_item_id")), + "repair_dispatch_allowed": False, + }, + ) + work_item_id = _safe_text(route.get("drift_work_item_id")) + record = { + "schema_version": CANDIDATE_SCHEMA_VERSION, + "kind": "asset_identity_drift_work_item", + "status": "asset_identity_unresolved", + **correlation, + "work_item_id": work_item_id or correlation["work_item_id"], + "fingerprint": fingerprint, + "root_cause_target_digest": hashlib.sha256( + root_cause_target.encode("utf-8") + ).hexdigest(), + "canonical_asset_id": canonical_asset_id, + "typed_domain": "unknown", + "typed_route": sanitize(route), + "candidate_created": False, + "repair_dispatch_allowed": False, + "runtime_mutation_performed": False, + "cross_domain_fallback_allowed": False, + "source_commitment": "AIA-CONV-034", + } + try: + persistence = await (persist_candidate or persist_cpu_p99_candidate)( + record, project_id + ) + except Exception as exc: # noqa: BLE001 - no receipt means no drift item + logger.warning( + "cpu_p99_drift_item_persistence_failed", + error_type=type(exc).__name__, + trace_id=correlation["trace_id"], + ) + return _blocked_result( + status="durable_drift_receipt_failed", + blockers=["durable_drift_receipt_failed"], + next_safe_action=( + "restore_durable_drift_store_before_any_auto_candidate" + ), + correlation=correlation, + ) + if not _candidate_record_matches(record, persistence.get("record")): + return _blocked_result( + status="durable_drift_receipt_mismatch", + blockers=["durable_drift_receipt_mismatch"], + next_safe_action=( + "repair_durable_drift_store_before_any_auto_candidate" + ), + correlation=correlation, + ) + return { + "schema_version": HANDOFF_SCHEMA_VERSION, + "accepted": True, + "created": persistence.get("created") is True, + "deduplicated": persistence.get("created") is not True, + "candidate_created": False, + "repair_dispatch_allowed": False, + "status": "asset_identity_unresolved", + **correlation, + "work_item_id": record["work_item_id"], + "candidate_event_id": _safe_text(persistence.get("event_id")), + "canonical_asset_id": "", + "typed_domain": "unknown", + "executor": None, + "verifier": None, + "active_blockers": ["canonical_asset_identity_unresolved"], + "next_safe_action": ( + "repair_canonical_asset_mapping_before_any_auto_candidate" + ), + "provider_call_performed": False, + "paid_provider_call_performed": False, + "executor_invoked": False, + "agent99_dispatch_performed": False, + "runtime_mutation_performed": False, + "cross_domain_fallback_allowed": False, + } + + canonical_asset_id = _safe_text(route.get("canonical_asset_id")) + typed_domain = _safe_text(route.get("target_kind")) + resource_value = _number(resource_evidence.get("observed_value")) + resource_threshold = _number(resource_evidence.get("threshold")) + latency_value = _number(latency_evidence.get("p99_seconds")) + latency_threshold = _number(latency_evidence.get("threshold_seconds")) + log_assets = { + _safe_text(value) + for value in log_evidence.get("canonical_asset_ids") or [] + if _safe_text(value) + } + change_assets = { + _safe_text(value) + for value in recent_change_evidence.get("canonical_asset_ids") or [] + if _safe_text(value) + } + evidence_blockers: list[str] = [] + if resource_evidence.get("schema_version") != "cpu_resource_evidence_v1": + evidence_blockers.append("resource_evidence_schema_invalid") + if resource_evidence.get("source") not in _RESOURCE_SOURCES: + evidence_blockers.append("resource_evidence_source_untrusted") + if resource_evidence.get("durable_readback_ack") is not True: + evidence_blockers.append("resource_evidence_not_durable") + if resource_evidence.get("metric_name") not in _RESOURCE_METRICS: + evidence_blockers.append("resource_metric_not_allowlisted") + if ( + resource_value is None + or resource_threshold is None + or resource_value <= resource_threshold + ): + evidence_blockers.append("resource_pressure_not_verified") + if _safe_text(resource_evidence.get("canonical_asset_id")) != canonical_asset_id: + evidence_blockers.append("resource_canonical_asset_mismatch") + if _safe_text(resource_evidence.get("asset_domain")) != typed_domain: + evidence_blockers.append("resource_typed_domain_mismatch") + + if latency_evidence.get("schema_version") != "signoz_api_p99_evidence_v1": + evidence_blockers.append("latency_evidence_schema_invalid") + if latency_evidence.get("source") != "signoz_readonly_query": + evidence_blockers.append("latency_evidence_source_untrusted") + if latency_evidence.get("durable_readback_ack") is not True: + evidence_blockers.append("latency_evidence_not_durable") + if ( + _safe_text(latency_evidence.get("canonical_asset_id")) != API_CANONICAL_ASSET_ID + or _safe_text(latency_evidence.get("service_name")) != "awoooi-api" + or _safe_text(latency_evidence.get("namespace")) != "awoooi-prod" + ): + evidence_blockers.append("latency_canonical_asset_mismatch") + if ( + latency_value is None + or latency_threshold is None + or latency_value <= latency_threshold + ): + evidence_blockers.append("p99_regression_not_verified") + + if ( + log_evidence.get("schema_version") != "sanitized_log_correlation_receipt_v1" + or log_evidence.get("durable_readback_ack") is not True + or log_evidence.get("sanitized") is not True + or log_evidence.get("untrusted_evidence") is not True + or log_evidence.get("raw_log_recorded") is not False + or not log_evidence.get("evidence_refs") + ): + evidence_blockers.append("sanitized_log_correlation_unverified") + if ( + canonical_asset_id not in log_assets + or API_CANONICAL_ASSET_ID not in log_assets + or _safe_text(log_evidence.get("root_cause_asset_id")) != canonical_asset_id + or log_evidence.get("rca_status") != "correlated" + ): + evidence_blockers.append("log_root_cause_asset_mismatch") + + if ( + recent_change_evidence.get("schema_version") + != "recent_change_correlation_receipt_v1" + or recent_change_evidence.get("durable_readback_ack") is not True + or recent_change_evidence.get("correlation_status") + not in {"matched", "no_change_in_window"} + or canonical_asset_id not in change_assets + or _safe_text(recent_change_evidence.get("root_cause_asset_id")) + != canonical_asset_id + ): + evidence_blockers.append("recent_change_correlation_unverified") + + evidence_blockers = list(dict.fromkeys(evidence_blockers)) + if evidence_blockers: + return _blocked_result( + status="canonical_evidence_alignment_blocked", + blockers=evidence_blockers, + next_safe_action=( + "recollect_canonical_asset_bound_metrics_logs_and_changes" + ), + correlation=correlation, + ) + + repair = _repair_contract(route) + fingerprint = _candidate_fingerprint( + correlation=correlation, + canonical_asset_id=canonical_asset_id, + receipts=receipts, + route_binding={ + "typed_domain": typed_domain, + "executor": _safe_text(route.get("executor")), + "verifier": _safe_text(route.get("verifier")), + "action_id": repair["action_id"], + "single_writer_executor": repair["single_writer_executor"], + "repair_dispatch_allowed": repair["repair_dispatch_allowed"], + }, + ) + record = { + "schema_version": CANDIDATE_SCHEMA_VERSION, + "kind": "typed_cpu_p99_controlled_candidate", + "status": repair["status"], + **correlation, + "fingerprint": fingerprint, + "canonical_asset_id": canonical_asset_id, + "typed_domain": typed_domain, + "typed_route": sanitize(route), + "evidence": { + "resource_receipt_id": _safe_text(resource_evidence.get("receipt_id")), + "resource_metric": resource_evidence.get("metric_name"), + "resource_before": resource_value, + "resource_threshold": resource_threshold, + "latency_receipt_id": _safe_text(latency_evidence.get("receipt_id")), + "latency_before_seconds": latency_value, + "latency_threshold_seconds": latency_threshold, + "log_receipt_id": _safe_text(log_evidence.get("receipt_id")), + "change_receipt_id": _safe_text(recent_change_evidence.get("receipt_id")), + "raw_logs_recorded": False, + }, + "repair_candidate": { + "action_id": repair["action_id"], + "executor": route.get("executor"), + "verifier": route.get("verifier"), + "single_writer_executor": repair["single_writer_executor"], + "candidate_created": repair["candidate_created"], + "repair_dispatch_allowed": repair["repair_dispatch_allowed"], + "apply_authorized_by_this_service": False, + }, + "next_safe_action": repair["next_safe_action"], + "policy": { + "provider_call_allowed": False, + "paid_provider_call_allowed": False, + "agent99_dispatch_allowed": False, + "direct_executor_invocation_allowed": False, + "runtime_mutation_allowed": False, + "cross_domain_fallback_allowed": False, + }, + "runtime_mutation_performed": False, + "source_commitment": "AIA-CONV-034", + } + try: + persistence = await (persist_candidate or persist_cpu_p99_candidate)( + record, project_id + ) + except Exception as exc: # noqa: BLE001 - no receipt means no candidate + logger.warning( + "cpu_p99_candidate_persistence_failed", + error_type=type(exc).__name__, + trace_id=correlation["trace_id"], + ) + return _blocked_result( + status="durable_candidate_receipt_failed", + blockers=["durable_candidate_receipt_failed"], + next_safe_action="restore_durable_candidate_store_before_dispatch", + correlation=correlation, + ) + if not _candidate_record_matches(record, persistence.get("record")): + return _blocked_result( + status="durable_candidate_receipt_mismatch", + blockers=["durable_candidate_receipt_mismatch"], + next_safe_action="repair_durable_candidate_store_before_dispatch", + correlation=correlation, + ) + + return { + "schema_version": HANDOFF_SCHEMA_VERSION, + "accepted": True, + "created": persistence.get("created") is True, + "deduplicated": persistence.get("created") is not True, + "candidate_created": repair["candidate_created"], + "repair_dispatch_allowed": repair["repair_dispatch_allowed"], + "status": repair["status"], + **correlation, + "candidate_event_id": _safe_text(persistence.get("event_id")), + "fingerprint": fingerprint, + "canonical_asset_id": canonical_asset_id, + "typed_domain": typed_domain, + "executor": route.get("executor"), + "verifier": route.get("verifier"), + "action_id": repair["action_id"], + "single_writer_executor": repair["single_writer_executor"], + "active_blockers": ( + [] + if repair["repair_dispatch_allowed"] + else ["typed_bounded_apply_not_yet_authorized"] + ), + "next_safe_action": repair["next_safe_action"], + "provider_call_performed": False, + "paid_provider_call_performed": False, + "executor_invoked": False, + "agent99_dispatch_performed": False, + "runtime_mutation_performed": False, + "cross_domain_fallback_allowed": False, + } + + +def _typed_execution_receipt_blockers( + candidate: Mapping[str, Any], + execution_receipt: Mapping[str, Any], +) -> list[str]: + blockers: list[str] = [] + correlation = { + field: _safe_text(candidate.get(field)) + for field in ("trace_id", "run_id", "work_item_id") + } + for field, value in correlation.items(): + if not _SAFE_ID.fullmatch(value): + blockers.append(f"candidate_{field}_missing_or_invalid") + try: + UUID(correlation["run_id"]) + except (TypeError, ValueError): + blockers.append("candidate_run_id_not_uuid") + + if candidate.get("schema_version") != CANDIDATE_SCHEMA_VERSION: + blockers.append("candidate_schema_invalid") + if candidate.get("kind") != "typed_cpu_p99_controlled_candidate": + blockers.append("candidate_kind_invalid") + candidate_fingerprint = _safe_text(candidate.get("fingerprint")) + if re.fullmatch(r"[0-9a-f]{64}", candidate_fingerprint) is None: + blockers.append("candidate_fingerprint_invalid") + canonical_asset_id = _safe_text(candidate.get("canonical_asset_id")) + if not _SAFE_ID.fullmatch(canonical_asset_id): + blockers.append("candidate_canonical_asset_invalid") + + repair = candidate.get("repair_candidate") + if not isinstance(repair, Mapping): + repair = {} + blockers.append("repair_candidate_missing") + if repair.get("repair_dispatch_allowed") is not True: + blockers.append("candidate_not_apply_authorized") + action_id = _safe_text(repair.get("action_id")) + executor = _safe_text(repair.get("executor")) + domain_verifier = _safe_text(repair.get("verifier")) + if not _SAFE_ID.fullmatch(action_id): + blockers.append("candidate_action_id_invalid") + if not _SAFE_ID.fullmatch(executor): + blockers.append("candidate_executor_invalid") + if not _SAFE_ID.fullmatch(domain_verifier) or domain_verifier == executor: + blockers.append("candidate_domain_verifier_invalid") + + execution_receipt_id = _safe_text(execution_receipt.get("receipt_id")) + domain_verifier_receipt_id = _safe_text( + execution_receipt.get("domain_verifier_receipt_id") + ) + if execution_receipt.get("schema_version") != "typed_bounded_execution_receipt_v1": + blockers.append("execution_receipt_schema_invalid") + if not _valid_public_receipt_id(execution_receipt_id): + blockers.append("execution_receipt_id_missing_or_invalid") + for field, expected in correlation.items(): + if _safe_text(execution_receipt.get(field)) != expected: + blockers.append(f"same_run_{field}_mismatch") + if execution_receipt.get("status") != "applied": + blockers.append("execution_not_applied") + if execution_receipt.get("durable_writeback_ack") is not True: + blockers.append("execution_writeback_missing") + if execution_receipt.get("runtime_mutation_performed") is not True: + blockers.append("execution_runtime_mutation_unverified") + if ( + _safe_text(execution_receipt.get("candidate_fingerprint")) + != candidate_fingerprint + ): + blockers.append("execution_candidate_fingerprint_mismatch") + if _safe_text(execution_receipt.get("canonical_asset_id")) != canonical_asset_id: + blockers.append("execution_canonical_asset_mismatch") + if _safe_text(execution_receipt.get("action_id")) != action_id: + blockers.append("execution_action_mismatch") + if _safe_text(execution_receipt.get("executor")) != executor: + blockers.append("execution_executor_mismatch") + if execution_receipt.get("cross_domain_fallback_performed") is not False: + blockers.append("execution_cross_domain_fallback_detected") + if _safe_text(execution_receipt.get("domain_verifier")) != domain_verifier: + blockers.append("domain_verifier_mismatch") + if execution_receipt.get("domain_verifier_status") != "verified": + blockers.append("domain_verifier_not_verified") + if execution_receipt.get("domain_verifier_independent") is not True: + blockers.append("domain_verifier_not_independent") + if execution_receipt.get("durable_domain_verifier_ack") is not True: + blockers.append("domain_verifier_writeback_missing") + if ( + not _valid_public_receipt_id(domain_verifier_receipt_id) + or domain_verifier_receipt_id == execution_receipt_id + ): + blockers.append("domain_verifier_receipt_missing_or_replayed") + return list(dict.fromkeys(blockers)) + + +def _post_readback_blockers( + *, + candidate: Mapping[str, Any], + execution_receipt: Mapping[str, Any], + readback: Mapping[str, Any], + expected_schema: str, + expected_source: str, + expected_asset_id: str, +) -> list[str]: + blockers = _typed_execution_receipt_blockers(candidate, execution_receipt) + correlation = { + field: _safe_text(candidate.get(field)) + for field in ("trace_id", "run_id", "work_item_id") + } + source_receipt_id = _safe_text(readback.get("receipt_id")) + if not _valid_public_receipt_id(source_receipt_id): + blockers.append("post_readback_receipt_id_missing_or_invalid") + if source_receipt_id in { + _safe_text(execution_receipt.get("receipt_id")), + _safe_text(execution_receipt.get("domain_verifier_receipt_id")), + }: + blockers.append("post_readback_receipt_replayed") + for field, expected in correlation.items(): + if _safe_text(readback.get(field)) != expected: + blockers.append(f"same_run_{field}_mismatch") + if readback.get("schema_version") != expected_schema: + blockers.append("post_readback_schema_invalid") + if _safe_text(readback.get("source")) != expected_source: + blockers.append("post_readback_source_untrusted") + if readback.get("trusted_server_readback") is not True: + blockers.append("post_readback_not_server_trusted") + if readback.get("independent") is not True: + blockers.append("post_readback_not_independent") + if readback.get("durable_readback_ack") is not True: + blockers.append("post_readback_not_durable") + max_age_seconds = _number(readback.get("max_age_seconds")) + if ( + readback.get("freshness_verified") is not True + or max_age_seconds is None + or max_age_seconds <= 0 + or max_age_seconds > MAX_EVIDENCE_SKEW_SECONDS + ): + blockers.append("post_readback_freshness_unverified") + if _safe_text(readback.get("candidate_fingerprint")) != _safe_text( + candidate.get("fingerprint") + ): + blockers.append("post_readback_candidate_fingerprint_mismatch") + if _safe_text(readback.get("execution_receipt_id")) != _safe_text( + execution_receipt.get("receipt_id") + ): + blockers.append("post_readback_execution_receipt_mismatch") + if _safe_text(readback.get("canonical_asset_id")) != expected_asset_id: + blockers.append("post_readback_canonical_asset_mismatch") + if readback.get("runtime_mutation_performed") is not False: + blockers.append("post_readback_runtime_mutation_detected") + if readback.get("cross_domain_fallback_performed") is not False: + blockers.append("post_readback_cross_domain_fallback_detected") + + applied_at = _timestamp(execution_receipt, "applied_at") + observed_at = _observed_at(readback) + if applied_at is None: + blockers.append("execution_applied_at_missing_or_invalid") + if observed_at is None: + blockers.append("post_readback_observed_at_missing_or_invalid") + if applied_at is not None and observed_at is not None and observed_at <= applied_at: + blockers.append("post_readback_not_after_execution") + return list(dict.fromkeys(blockers)) + + +def build_cpu_resource_post_verifier_receipt( + *, + candidate: Mapping[str, Any], + execution_receipt: Mapping[str, Any], + metric_readback: Mapping[str, Any], +) -> dict[str, Any]: + """Evaluate a public-safe Prometheus post-readback without performing I/O.""" + + canonical_asset_id = _safe_text(candidate.get("canonical_asset_id")) + blockers = _post_readback_blockers( + candidate=candidate, + execution_receipt=execution_receipt, + readback=metric_readback, + expected_schema=RESOURCE_READBACK_SCHEMA_VERSION, + expected_source=RESOURCE_READBACK_SOURCE, + expected_asset_id=canonical_asset_id, + ) + evidence = candidate.get("evidence") + if not isinstance(evidence, Mapping): + evidence = {} + blockers.append("candidate_evidence_missing") + metric_name = _safe_text(metric_readback.get("metric_name")) + before = _number(evidence.get("resource_before")) + threshold = _number(evidence.get("resource_threshold")) + observed = _number(metric_readback.get("observed_value")) + if ( + metric_name != _safe_text(evidence.get("resource_metric")) + or metric_name not in _RESOURCE_METRICS + ): + blockers.append("resource_post_metric_mismatch") + if ( + before is None + or threshold is None + or observed is None + or observed >= before + or observed > threshold + ): + blockers.append("resource_postcondition_not_met") + blockers = list(dict.fromkeys(blockers)) + source_receipt_id = _safe_text(metric_readback.get("receipt_id")) + execution_receipt_id = _safe_text(execution_receipt.get("receipt_id")) + candidate_fingerprint = _safe_text(candidate.get("fingerprint")) + return { + "schema_version": RESOURCE_VERIFIER_SCHEMA_VERSION, + "receipt_id": _derived_receipt_id( + "cpu-resource-verifier", + candidate_fingerprint, + execution_receipt_id, + source_receipt_id, + _safe_text(metric_readback.get("observed_at")), + metric_name, + _safe_text(observed), + ), + "status": "verified" if not blockers else "blocked", + "verifier": RESOURCE_VERIFIER, + "independent": True, + "read_only": metric_readback.get("runtime_mutation_performed") is False, + "trusted_server_readback": ( + metric_readback.get("trusted_server_readback") is True + ), + "durable_readback_ack": metric_readback.get("durable_readback_ack") is True, + "freshness_verified": metric_readback.get("freshness_verified") is True, + "max_age_seconds": metric_readback.get("max_age_seconds"), + "source": _safe_text(metric_readback.get("source")), + "source_receipt_id": source_receipt_id, + "observed_at": metric_readback.get("observed_at"), + "trace_id": candidate.get("trace_id"), + "run_id": candidate.get("run_id"), + "work_item_id": candidate.get("work_item_id"), + "candidate_fingerprint": candidate_fingerprint, + "execution_receipt_id": execution_receipt_id, + "canonical_asset_id": canonical_asset_id, + "metric_name": metric_name, + "before_value": before, + "threshold": threshold, + "observed_value": observed, + "active_blockers": blockers, + "runtime_mutation_performed": ( + metric_readback.get("runtime_mutation_performed") is True + ), + "cross_domain_fallback_performed": ( + metric_readback.get("cross_domain_fallback_performed") is True + ), + } + + +def build_signoz_api_p99_post_verifier_receipt( + *, + candidate: Mapping[str, Any], + execution_receipt: Mapping[str, Any], + signoz_readback: Mapping[str, Any], +) -> dict[str, Any]: + """Evaluate a public-safe SigNoz P99 post-readback without performing I/O.""" + + blockers = _post_readback_blockers( + candidate=candidate, + execution_receipt=execution_receipt, + readback=signoz_readback, + expected_schema=LATENCY_READBACK_SCHEMA_VERSION, + expected_source=LATENCY_READBACK_SOURCE, + expected_asset_id=API_CANONICAL_ASSET_ID, + ) + evidence = candidate.get("evidence") + if not isinstance(evidence, Mapping): + evidence = {} + blockers.append("candidate_evidence_missing") + before = _number(evidence.get("latency_before_seconds")) + threshold = _number(evidence.get("latency_threshold_seconds")) + observed = _number(signoz_readback.get("p99_seconds")) + if ( + _safe_text(signoz_readback.get("service_name")) != "awoooi-api" + or _safe_text(signoz_readback.get("namespace")) != "awoooi-prod" + ): + blockers.append("latency_post_service_identity_mismatch") + if ( + before is None + or threshold is None + or observed is None + or observed >= before + or observed > threshold + ): + blockers.append("latency_postcondition_not_met") + blockers = list(dict.fromkeys(blockers)) + source_receipt_id = _safe_text(signoz_readback.get("receipt_id")) + execution_receipt_id = _safe_text(execution_receipt.get("receipt_id")) + candidate_fingerprint = _safe_text(candidate.get("fingerprint")) + return { + "schema_version": LATENCY_VERIFIER_SCHEMA_VERSION, + "receipt_id": _derived_receipt_id( + "signoz-p99-verifier", + candidate_fingerprint, + execution_receipt_id, + source_receipt_id, + _safe_text(signoz_readback.get("observed_at")), + _safe_text(signoz_readback.get("service_name")), + _safe_text(signoz_readback.get("namespace")), + _safe_text(observed), + ), + "status": "verified" if not blockers else "blocked", + "verifier": LATENCY_VERIFIER, + "independent": True, + "read_only": signoz_readback.get("runtime_mutation_performed") is False, + "trusted_server_readback": ( + signoz_readback.get("trusted_server_readback") is True + ), + "durable_readback_ack": signoz_readback.get("durable_readback_ack") is True, + "freshness_verified": signoz_readback.get("freshness_verified") is True, + "max_age_seconds": signoz_readback.get("max_age_seconds"), + "source": _safe_text(signoz_readback.get("source")), + "source_receipt_id": source_receipt_id, + "observed_at": signoz_readback.get("observed_at"), + "trace_id": candidate.get("trace_id"), + "run_id": candidate.get("run_id"), + "work_item_id": candidate.get("work_item_id"), + "candidate_fingerprint": candidate_fingerprint, + "execution_receipt_id": execution_receipt_id, + "canonical_asset_id": API_CANONICAL_ASSET_ID, + "service_name": _safe_text(signoz_readback.get("service_name")), + "namespace": _safe_text(signoz_readback.get("namespace")), + "before_p99_seconds": before, + "threshold_seconds": threshold, + "p99_seconds": observed, + "active_blockers": blockers, + "runtime_mutation_performed": ( + signoz_readback.get("runtime_mutation_performed") is True + ), + "cross_domain_fallback_performed": ( + signoz_readback.get("cross_domain_fallback_performed") is True + ), + } + + +def build_cpu_p99_closure( + *, + candidate: Mapping[str, Any], + execution_receipt: Mapping[str, Any], + resource_verifier_receipt: Mapping[str, Any], + latency_verifier_receipt: Mapping[str, Any], +) -> dict[str, Any]: + """Require one run plus independent resource and latency postconditions.""" + + execution_blockers = _typed_execution_receipt_blockers(candidate, execution_receipt) + blockers = list(execution_blockers) + correlation = { + field: _safe_text(candidate.get(field)) + for field in ("trace_id", "run_id", "work_item_id") + } + + receipts = ( + execution_receipt, + resource_verifier_receipt, + latency_verifier_receipt, + ) + for receipt in receipts: + for field, expected in correlation.items(): + if _safe_text(receipt.get(field)) != expected: + blockers.append(f"same_run_{field}_mismatch") + if not _valid_public_receipt_id(receipt.get("receipt_id")): + blockers.append("closure_receipt_id_missing_or_invalid") + + canonical_asset_id = _safe_text(candidate.get("canonical_asset_id")) + candidate_fingerprint = _safe_text(candidate.get("fingerprint")) + execution_receipt_id = _safe_text(execution_receipt.get("receipt_id")) + execution_verified = not execution_blockers + if not execution_verified: + blockers.append("typed_execution_receipt_unverified") + + evidence = candidate.get("evidence") + if not isinstance(evidence, Mapping): + evidence = {} + blockers.append("candidate_evidence_missing") + resource_before = _number(evidence.get("resource_before")) + resource_threshold = _number(evidence.get("resource_threshold")) + resource_after = _number(resource_verifier_receipt.get("observed_value")) + resource_source_receipt_id = _safe_text( + resource_verifier_receipt.get("source_receipt_id") + ) + resource_max_age = _number(resource_verifier_receipt.get("max_age_seconds")) + resource_observed_at = _observed_at(resource_verifier_receipt) + execution_applied_at = _timestamp(execution_receipt, "applied_at") + resource_verified = bool( + execution_verified + and resource_verifier_receipt.get("schema_version") + == RESOURCE_VERIFIER_SCHEMA_VERSION + and resource_verifier_receipt.get("status") == "verified" + and resource_verifier_receipt.get("independent") is True + and resource_verifier_receipt.get("read_only") is True + and resource_verifier_receipt.get("trusted_server_readback") is True + and resource_verifier_receipt.get("durable_readback_ack") is True + and resource_verifier_receipt.get("freshness_verified") is True + and resource_max_age is not None + and 0 < resource_max_age <= MAX_EVIDENCE_SKEW_SECONDS + and resource_observed_at is not None + and execution_applied_at is not None + and resource_observed_at > execution_applied_at + and not resource_verifier_receipt.get("active_blockers") + and _safe_text(resource_verifier_receipt.get("candidate_fingerprint")) + == candidate_fingerprint + and _safe_text(resource_verifier_receipt.get("execution_receipt_id")) + == execution_receipt_id + and _safe_text(resource_verifier_receipt.get("canonical_asset_id")) + == canonical_asset_id + and _safe_text(resource_verifier_receipt.get("verifier")) == RESOURCE_VERIFIER + and _safe_text(resource_verifier_receipt.get("source")) + == RESOURCE_READBACK_SOURCE + and _valid_public_receipt_id(resource_source_receipt_id) + and _safe_text(resource_verifier_receipt.get("receipt_id")) + == _derived_receipt_id( + "cpu-resource-verifier", + candidate_fingerprint, + execution_receipt_id, + resource_source_receipt_id, + _safe_text(resource_verifier_receipt.get("observed_at")), + _safe_text(resource_verifier_receipt.get("metric_name")), + _safe_text(resource_after), + ) + and _safe_text(resource_verifier_receipt.get("metric_name")) + == _safe_text(evidence.get("resource_metric")) + and resource_verifier_receipt.get("runtime_mutation_performed") is False + and resource_verifier_receipt.get("cross_domain_fallback_performed") is False + and resource_before is not None + and resource_threshold is not None + and resource_after is not None + and resource_after < resource_before + and resource_after <= resource_threshold + ) + if not resource_verified: + blockers.append("resource_independent_verifier_not_closed") + + latency_before = _number(evidence.get("latency_before_seconds")) + latency_threshold = _number(evidence.get("latency_threshold_seconds")) + latency_after = _number(latency_verifier_receipt.get("p99_seconds")) + latency_source_receipt_id = _safe_text( + latency_verifier_receipt.get("source_receipt_id") + ) + latency_max_age = _number(latency_verifier_receipt.get("max_age_seconds")) + latency_observed_at = _observed_at(latency_verifier_receipt) + if ( + resource_source_receipt_id + and resource_source_receipt_id == latency_source_receipt_id + ): + blockers.append("independent_readback_receipt_replayed") + latency_verified = bool( + execution_verified + and latency_verifier_receipt.get("schema_version") + == LATENCY_VERIFIER_SCHEMA_VERSION + and latency_verifier_receipt.get("status") == "verified" + and latency_verifier_receipt.get("independent") is True + and latency_verifier_receipt.get("read_only") is True + and latency_verifier_receipt.get("trusted_server_readback") is True + and latency_verifier_receipt.get("durable_readback_ack") is True + and latency_verifier_receipt.get("freshness_verified") is True + and latency_max_age is not None + and 0 < latency_max_age <= MAX_EVIDENCE_SKEW_SECONDS + and latency_observed_at is not None + and execution_applied_at is not None + and latency_observed_at > execution_applied_at + and not latency_verifier_receipt.get("active_blockers") + and _safe_text(latency_verifier_receipt.get("candidate_fingerprint")) + == candidate_fingerprint + and _safe_text(latency_verifier_receipt.get("execution_receipt_id")) + == execution_receipt_id + and _safe_text(latency_verifier_receipt.get("canonical_asset_id")) + == API_CANONICAL_ASSET_ID + and _safe_text(latency_verifier_receipt.get("verifier")) == LATENCY_VERIFIER + and _safe_text(latency_verifier_receipt.get("source")) + == LATENCY_READBACK_SOURCE + and _valid_public_receipt_id(latency_source_receipt_id) + and _safe_text(latency_verifier_receipt.get("receipt_id")) + == _derived_receipt_id( + "signoz-p99-verifier", + candidate_fingerprint, + execution_receipt_id, + latency_source_receipt_id, + _safe_text(latency_verifier_receipt.get("observed_at")), + _safe_text(latency_verifier_receipt.get("service_name")), + _safe_text(latency_verifier_receipt.get("namespace")), + _safe_text(latency_after), + ) + and _safe_text(latency_verifier_receipt.get("service_name")) == "awoooi-api" + and _safe_text(latency_verifier_receipt.get("namespace")) == "awoooi-prod" + and latency_verifier_receipt.get("runtime_mutation_performed") is False + and latency_verifier_receipt.get("cross_domain_fallback_performed") is False + and latency_before is not None + and latency_threshold is not None + and latency_after is not None + and latency_after < latency_before + and latency_after <= latency_threshold + ) + if not latency_verified: + blockers.append("latency_independent_verifier_not_closed") + + blockers = list(dict.fromkeys(blockers)) + closed = not blockers + return { + "schema_version": CLOSURE_SCHEMA_VERSION, + "status": ( + "verified_closed" + if closed + else ( + "applied_pending_independent_verification" + if execution_verified + else "typed_execution_or_verification_blocked" + ) + ), + "closed": closed, + **correlation, + "candidate_fingerprint": candidate_fingerprint, + "canonical_asset_id": canonical_asset_id, + "typed_domain": candidate.get("typed_domain"), + "execution_receipt_id": execution_receipt_id, + "resource_verifier_receipt_id": resource_verifier_receipt.get("receipt_id"), + "latency_verifier_receipt_id": latency_verifier_receipt.get("receipt_id"), + "resource_verifier_closed": resource_verified, + "latency_verifier_closed": latency_verified, + "same_run_receipts": not any( + blocker.startswith("same_run_") for blocker in blockers + ), + "active_blockers": blockers, + "runtime_mutation_performed": execution_verified, + "cross_domain_fallback_performed": ( + execution_receipt.get("cross_domain_fallback_performed") is True + ), + "durable_closure_ack": False, + "next_safe_action": ( + "write_incident_lifecycle_and_km_learning_receipts" + if closed + else ( + "collect_missing_same_run_independent_verifier_receipts" + if execution_verified + else "obtain_verified_typed_execution_receipt" + ) + ), + } + + +async def persist_cpu_p99_closure( + candidate: Mapping[str, Any], + closure: Mapping[str, Any], + project_id: str, +) -> bool: + """Append a closure receipt without reopening or dispatching a repair.""" + + if not ( + candidate.get("schema_version") == CANDIDATE_SCHEMA_VERSION + and candidate.get("kind") == "typed_cpu_p99_controlled_candidate" + and closure.get("schema_version") == CLOSURE_SCHEMA_VERSION + and closure.get("status") == "verified_closed" + and closure.get("closed") is True + and closure.get("durable_closure_ack") is True + and not closure.get("active_blockers") + and closure.get("resource_verifier_closed") is True + and closure.get("latency_verifier_closed") is True + and closure.get("same_run_receipts") is True + and closure.get("runtime_mutation_performed") is True + and closure.get("cross_domain_fallback_performed") is False + and _safe_text(closure.get("candidate_fingerprint")) + == _safe_text(candidate.get("fingerprint")) + and all( + _safe_text(closure.get(field)) == _safe_text(candidate.get(field)) + for field in ("trace_id", "run_id", "work_item_id") + ) + and all( + _valid_public_receipt_id(closure.get(field)) + for field in ( + "execution_receipt_id", + "resource_verifier_receipt_id", + "latency_verifier_receipt_id", + ) + ) + ): + return False + + from src.db.base import get_db_context + + safe_closure = json.dumps( + _public_safe_record(closure), ensure_ascii=False, default=str + ) + 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( + 'closure', CAST(:closure AS jsonb), + 'closure_status', :closure_status, + 'closed', :closed, + 'closure_recorded_at', NOW() + ) + WHERE project_id = :project_id + AND channel_type = 'internal' + AND source_envelope ->> 'schema_version' = :schema_version + AND source_envelope ->> 'fingerprint' = :fingerprint + AND NOT (source_envelope ? 'closure') + RETURNING event_id + """ + ), + { + "project_id": project_id, + "schema_version": CANDIDATE_SCHEMA_VERSION, + "fingerprint": _safe_text(candidate.get("fingerprint")), + "closure": safe_closure, + "closure_status": _safe_text(closure.get("status")), + "closed": closure.get("closed") is True, + }, + ) + if updated.fetchone() is not None: + return True + + existing = await db.execute( + text( + """ + SELECT source_envelope -> 'closure' + FROM awooop_conversation_event + WHERE project_id = :project_id + AND channel_type = 'internal' + AND source_envelope ->> 'schema_version' = :schema_version + AND source_envelope ->> 'fingerprint' = :fingerprint + LIMIT 1 + """ + ), + { + "project_id": project_id, + "schema_version": CANDIDATE_SCHEMA_VERSION, + "fingerprint": _safe_text(candidate.get("fingerprint")), + }, + ) + row = existing.fetchone() + if row is None: + return False + stored = row[0] + if isinstance(stored, str): + stored = json.loads(stored) + return bool( + isinstance(stored, Mapping) + and stored.get("schema_version") == CLOSURE_SCHEMA_VERSION + and stored.get("candidate_fingerprint") + == closure.get("candidate_fingerprint") + and stored.get("execution_receipt_id") + == closure.get("execution_receipt_id") + and stored.get("resource_verifier_receipt_id") + == closure.get("resource_verifier_receipt_id") + and stored.get("latency_verifier_receipt_id") + == closure.get("latency_verifier_receipt_id") + and stored.get("status") == closure.get("status") + and stored.get("durable_closure_ack") is True + ) + + +async def record_cpu_p99_closure( + *, + candidate: Mapping[str, Any], + execution_receipt: Mapping[str, Any], + resource_verifier_receipt: Mapping[str, Any], + latency_verifier_receipt: Mapping[str, Any], + project_id: str = "awoooi", + persist_closure: PersistClosure | None = None, +) -> dict[str, Any]: + """Build and durably append only a fully verified closure receipt.""" + + closure = build_cpu_p99_closure( + candidate=candidate, + execution_receipt=execution_receipt, + resource_verifier_receipt=resource_verifier_receipt, + latency_verifier_receipt=latency_verifier_receipt, + ) + if closure.get("closed") is not True: + return closure + try: + closure_to_persist = {**closure, "durable_closure_ack": True} + durable = await (persist_closure or persist_cpu_p99_closure)( + candidate, closure_to_persist, project_id + ) + except Exception as exc: # noqa: BLE001 - no durable ack means not closed + logger.warning( + "cpu_p99_closure_persistence_failed", + error_type=type(exc).__name__, + trace_id=closure.get("trace_id"), + ) + durable = False + if not durable: + return { + **closure, + "status": "durable_closure_writeback_failed", + "closed": False, + "durable_closure_ack": False, + "active_blockers": list( + dict.fromkeys( + [ + *closure.get("active_blockers", []), + "durable_closure_writeback_failed", + ] + ) + ), + } + return {**closure, "durable_closure_ack": True} diff --git a/apps/api/src/services/independent_verifier_registry.py b/apps/api/src/services/independent_verifier_registry.py index a61d136b0..1ac8f703c 100644 --- a/apps/api/src/services/independent_verifier_registry.py +++ b/apps/api/src/services/independent_verifier_registry.py @@ -111,6 +111,20 @@ _EXTERNAL_READBACKS: tuple[dict[str, Any], ...] = ( "verifier": "agent99_alertmanager_pull_relay_verifier", "runtime_status": "receipt_pending", }, + { + "readback_id": "cpu_resource_postcondition", + "verifier": "cpu_resource_independent_verifier", + "required_source": "prometheus_independent_query", + "source_ref": "apps/api/src/services/cpu_p99_controlled_recovery.py", + "runtime_status": "receipt_pending", + }, + { + "readback_id": "api_p99_postcondition", + "verifier": "signoz_api_latency_independent_verifier", + "required_source": "signoz_independent_query", + "source_ref": "apps/api/src/services/cpu_p99_controlled_recovery.py", + "runtime_status": "receipt_pending", + }, ) diff --git a/apps/api/tests/test_cpu_p99_controlled_recovery.py b/apps/api/tests/test_cpu_p99_controlled_recovery.py new file mode 100644 index 000000000..fe49fca1c --- /dev/null +++ b/apps/api/tests/test_cpu_p99_controlled_recovery.py @@ -0,0 +1,552 @@ +from __future__ import annotations + +# ruff: noqa: E402, I001 + +import os +from pathlib import Path +from uuid import UUID + +import pytest +import yaml + +os.environ.setdefault("DATABASE_URL", "postgresql+asyncpg://test:test@localhost/test") + +from src.services.cpu_p99_controlled_recovery import ( # noqa: E402 + build_cpu_resource_post_verifier_receipt, + build_cpu_p99_closure, + build_cpu_p99_controlled_candidate, + build_signoz_api_p99_post_verifier_receipt, + persist_cpu_p99_closure, + record_cpu_p99_closure, +) +from src.services.independent_verifier_registry import ( # noqa: E402 + build_independent_verifier_registry, +) +from src.services.service_registry import ServiceRegistryClient # noqa: E402 + + +ROOT = Path(__file__).resolve().parents[3] +RUN_ID = "6dc24cae-46ae-51a5-8278-7a9fbdafaa34" +TRACE_ID = "cpu-p99-trace-034" +WORK_ITEM_ID = "AIA-CPU-P99-034" +OBSERVED_AT = "2026-07-22T04:00:00+00:00" + + +class _MemoryCandidateStore: + def __init__(self) -> None: + self.records: dict[str, tuple[str, dict]] = {} + + async def __call__(self, record: dict, project_id: str) -> dict: + assert project_id == "awoooi" + fingerprint = record["fingerprint"] + existing = self.records.get(fingerprint) + if existing is not None: + return { + "event_id": existing[0], + "created": False, + "record": existing[1], + } + event_id = str(UUID(int=len(self.records) + 1)) + self.records[fingerprint] = (event_id, record) + return {"event_id": event_id, "created": True, "record": record} + + +def _common(receipt_id: str) -> dict[str, object]: + return { + "receipt_id": receipt_id, + "trace_id": TRACE_ID, + "run_id": RUN_ID, + "work_item_id": WORK_ITEM_ID, + "observed_at": OBSERVED_AT, + "freshness_verified": True, + "max_age_seconds": 900, + } + + +def _resource( + *, + canonical_asset_id: str = "service:awoooi-api", + asset_domain: str = "kubernetes_workload", + alertname: str = "HostHighCpuLoad", + namespace: str = "awoooi-prod", + labels: dict[str, str] | None = None, +) -> dict[str, object]: + return { + **_common("resource-receipt-034"), + "schema_version": "cpu_resource_evidence_v1", + "source": "prometheus_readonly_query", + "durable_readback_ack": True, + "alertname": alertname, + "alert_category": "host_resource", + "canonical_asset_id": canonical_asset_id, + "asset_domain": asset_domain, + "metric_name": "cpu_usage_percent", + "observed_value": 93.2, + "threshold": 80.0, + "namespace": namespace, + "labels": labels + or { + "namespace": "awoooi-prod", + "workload_kind": "deployment", + "workload_name": "awoooi-api", + }, + } + + +def _latency() -> dict[str, object]: + return { + **_common("latency-receipt-034"), + "schema_version": "signoz_api_p99_evidence_v1", + "source": "signoz_readonly_query", + "durable_readback_ack": True, + "canonical_asset_id": "service:awoooi-api", + "service_name": "awoooi-api", + "namespace": "awoooi-prod", + "p99_seconds": 3.4, + "threshold_seconds": 2.0, + } + + +def _logs(root_asset: str = "service:awoooi-api") -> dict[str, object]: + return { + **_common("logs-receipt-034"), + "schema_version": "sanitized_log_correlation_receipt_v1", + "durable_readback_ack": True, + "sanitized": True, + "untrusted_evidence": True, + "raw_log_recorded": False, + "evidence_refs": ["incident-evidence:cpu-p99:034"], + "canonical_asset_ids": list(dict.fromkeys([root_asset, "service:awoooi-api"])), + "root_cause_asset_id": root_asset, + "rca_status": "correlated", + } + + +def _change(root_asset: str = "service:awoooi-api") -> dict[str, object]: + return { + **_common("change-receipt-034"), + "schema_version": "recent_change_correlation_receipt_v1", + "durable_readback_ack": True, + "correlation_status": "no_change_in_window", + "canonical_asset_ids": [root_asset], + "root_cause_asset_id": root_asset, + "change_refs": [], + } + + +@pytest.mark.asyncio +async def test_exact_kubernetes_cpu_p99_candidate_is_durable_and_deduplicated() -> None: + store = _MemoryCandidateStore() + first = await build_cpu_p99_controlled_candidate( + root_cause_target="awoooi-api", + resource_evidence=_resource(), + latency_evidence=_latency(), + log_evidence=_logs(), + recent_change_evidence=_change(), + persist_candidate=store, + ) + duplicate = await build_cpu_p99_controlled_candidate( + root_cause_target="awoooi-api", + resource_evidence=_resource(), + latency_evidence=_latency(), + log_evidence=_logs(), + recent_change_evidence=_change(), + persist_candidate=store, + ) + + assert first["accepted"] is True + assert first["created"] is True + assert first["status"] == "bounded_repair_candidate_recorded" + assert first["typed_domain"] == "kubernetes_workload" + assert first["executor"] == "kubernetes_controlled_executor" + assert first["verifier"] == "kubernetes_rollout_verifier" + assert first["repair_dispatch_allowed"] is True + assert first["runtime_mutation_performed"] is False + assert first["agent99_dispatch_performed"] is False + assert duplicate["created"] is False + assert duplicate["deduplicated"] is True + assert duplicate["candidate_event_id"] == first["candidate_event_id"] + assert len(store.records) == 1 + + record = next(iter(store.records.values()))[1] + assert record["repair_candidate"]["apply_authorized_by_this_service"] is False + assert record["policy"]["provider_call_allowed"] is False + assert record["policy"]["runtime_mutation_allowed"] is False + assert record["evidence"]["raw_logs_recorded"] is False + + record["repair_candidate"]["executor"] = "Agent99" + mismatched_duplicate = await build_cpu_p99_controlled_candidate( + root_cause_target="awoooi-api", + resource_evidence=_resource(), + latency_evidence=_latency(), + log_evidence=_logs(), + recent_change_evidence=_change(), + persist_candidate=store, + ) + assert mismatched_duplicate["accepted"] is False + assert mismatched_duplicate["status"] == "durable_candidate_receipt_mismatch" + assert mismatched_duplicate["runtime_mutation_performed"] is False + + +@pytest.mark.asyncio +async def test_host_pressure_uses_no_write_ansible_lane_not_kubernetes_or_agent99() -> ( + None +): + root_asset = "host:192.168.0.110" + store = _MemoryCandidateStore() + result = await build_cpu_p99_controlled_candidate( + root_cause_target="host-pressure-controller", + resource_evidence=_resource( + canonical_asset_id=root_asset, + asset_domain="host_systemd", + alertname="HostLoadAverageSustainedHigh", + namespace="", + labels={"host": "110", "host_type": "bare_metal"}, + ), + latency_evidence=_latency(), + log_evidence=_logs(root_asset), + recent_change_evidence=_change(root_asset), + persist_candidate=store, + ) + + assert result["status"] == "no_write_host_rca_candidate_recorded" + assert result["typed_domain"] == "host_systemd" + assert result["executor"] == "host_ansible_executor" + assert result["action_id"] == "ansible:110-host-pressure-readonly" + assert result["repair_dispatch_allowed"] is False + assert result["agent99_dispatch_performed"] is False + assert "kubectl" not in str(result).lower() + + +@pytest.mark.asyncio +async def test_exact_database_root_cause_beats_generic_host_route_and_stays_no_write() -> ( + None +): + root_asset = "service:signoz-clickhouse" + store = _MemoryCandidateStore() + result = await build_cpu_p99_controlled_candidate( + root_cause_target="signoz-clickhouse", + resource_evidence=_resource( + canonical_asset_id=root_asset, + asset_domain="database", + namespace="", + labels={"host": "110", "host_type": "bare_metal"}, + ), + latency_evidence=_latency(), + log_evidence=_logs(root_asset), + recent_change_evidence=_change(root_asset), + persist_candidate=store, + ) + + assert result["status"] == "no_write_database_rca_candidate_recorded" + assert result["typed_domain"] == "database" + assert result["executor"] == "db_bounded_executor" + assert result["verifier"] == "db_independent_verifier" + assert result["repair_dispatch_allowed"] is False + assert result["agent99_dispatch_performed"] is False + assert "restart" not in result["action_id"] + assert "kubectl" not in result["action_id"] + + +@pytest.mark.asyncio +async def test_unknown_asset_creates_only_durable_drift_receipt() -> None: + store = _MemoryCandidateStore() + result = await build_cpu_p99_controlled_candidate( + root_cause_target="mystery-cpu-source", + resource_evidence=_resource( + canonical_asset_id="unresolved", + asset_domain="unknown", + namespace="", + labels={"host_type": "bare_metal"}, + ), + latency_evidence=_latency(), + log_evidence=_logs("unresolved"), + recent_change_evidence=_change("unresolved"), + persist_candidate=store, + ) + + assert result["accepted"] is True + assert result["status"] == "asset_identity_unresolved" + assert result["candidate_created"] is False + assert result["repair_dispatch_allowed"] is False + assert result["executor"] is None + assert result["verifier"] is None + assert result["runtime_mutation_performed"] is False + assert len(store.records) == 1 + record = next(iter(store.records.values()))[1] + assert record["kind"] == "asset_identity_drift_work_item" + assert record["cross_domain_fallback_allowed"] is False + + +@pytest.mark.asyncio +async def test_untrusted_or_misaligned_logs_create_no_candidate() -> None: + store = _MemoryCandidateStore() + logs = _logs() + logs["raw_log_recorded"] = True + + result = await build_cpu_p99_controlled_candidate( + root_cause_target="awoooi-api", + resource_evidence=_resource(), + latency_evidence=_latency(), + log_evidence=logs, + recent_change_evidence=_change(), + persist_candidate=store, + ) + + assert result["accepted"] is False + assert result["status"] == "canonical_evidence_alignment_blocked" + assert "sanitized_log_correlation_unverified" in result["active_blockers"] + assert result["runtime_mutation_performed"] is False + assert store.records == {} + + +def _execution_receipt(candidate: dict) -> dict: + common = { + "trace_id": TRACE_ID, + "run_id": RUN_ID, + "work_item_id": WORK_ITEM_ID, + } + repair = candidate["repair_candidate"] + return { + **common, + "schema_version": "typed_bounded_execution_receipt_v1", + "receipt_id": "execution-receipt-034", + "status": "applied", + "durable_writeback_ack": True, + "applied_at": "2026-07-22T03:59:00+00:00", + "candidate_fingerprint": candidate["fingerprint"], + "canonical_asset_id": candidate["canonical_asset_id"], + "action_id": repair["action_id"], + "executor": repair["executor"], + "runtime_mutation_performed": True, + "cross_domain_fallback_performed": False, + "domain_verifier": repair["verifier"], + "domain_verifier_status": "verified", + "domain_verifier_independent": True, + "domain_verifier_receipt_id": "kubernetes-verifier-receipt-034", + "durable_domain_verifier_ack": True, + } + + +def _post_readbacks(candidate: dict, execution: dict) -> tuple[dict, dict]: + resource = { + **_common("prometheus-post-readback-034"), + "schema_version": "prometheus_cpu_resource_post_readback_v1", + "source": "prometheus_independent_query", + "trusted_server_readback": True, + "independent": True, + "durable_readback_ack": True, + "candidate_fingerprint": candidate["fingerprint"], + "execution_receipt_id": execution["receipt_id"], + "canonical_asset_id": candidate["canonical_asset_id"], + "metric_name": candidate["evidence"]["resource_metric"], + "observed_value": 55.0, + "runtime_mutation_performed": False, + "cross_domain_fallback_performed": False, + } + latency = { + **_common("signoz-post-readback-034"), + "schema_version": "signoz_api_p99_post_readback_v1", + "source": "signoz_independent_query", + "trusted_server_readback": True, + "independent": True, + "durable_readback_ack": True, + "candidate_fingerprint": candidate["fingerprint"], + "execution_receipt_id": execution["receipt_id"], + "canonical_asset_id": "service:awoooi-api", + "service_name": "awoooi-api", + "namespace": "awoooi-prod", + "p99_seconds": 1.2, + "runtime_mutation_performed": False, + "cross_domain_fallback_performed": False, + } + return resource, latency + + +def _closure_receipts(candidate: dict) -> tuple[dict, dict, dict]: + execution = _execution_receipt(candidate) + resource_readback, latency_readback = _post_readbacks(candidate, execution) + resource = build_cpu_resource_post_verifier_receipt( + candidate=candidate, + execution_receipt=execution, + metric_readback=resource_readback, + ) + latency = build_signoz_api_p99_post_verifier_receipt( + candidate=candidate, + execution_receipt=execution, + signoz_readback=latency_readback, + ) + return execution, resource, latency + + +@pytest.mark.asyncio +async def test_closure_requires_same_run_resource_and_latency_verifiers() -> None: + store = _MemoryCandidateStore() + result = await build_cpu_p99_controlled_candidate( + root_cause_target="awoooi-api", + resource_evidence=_resource(), + latency_evidence=_latency(), + log_evidence=_logs(), + recent_change_evidence=_change(), + persist_candidate=store, + ) + candidate = next(iter(store.records.values()))[1] + execution, resource, latency = _closure_receipts(candidate) + + closed = build_cpu_p99_closure( + candidate=candidate, + execution_receipt=execution, + resource_verifier_receipt=resource, + latency_verifier_receipt=latency, + ) + assert closed["closed"] is True + assert closed["status"] == "verified_closed" + assert closed["resource_verifier_closed"] is True + assert closed["latency_verifier_closed"] is True + assert closed["same_run_receipts"] is True + assert resource["verifier"] == "cpu_resource_independent_verifier" + assert latency["verifier"] == "signoz_api_latency_independent_verifier" + + _, stale_latency_readback = _post_readbacks(candidate, execution) + stale_latency_readback["freshness_verified"] = False + stale_latency = build_signoz_api_p99_post_verifier_receipt( + candidate=candidate, + execution_receipt=execution, + signoz_readback=stale_latency_readback, + ) + assert stale_latency["status"] == "blocked" + assert "post_readback_freshness_unverified" in stale_latency["active_blockers"] + + execution["runtime_mutation_performed"] = False + no_apply = build_cpu_p99_closure( + candidate=candidate, + execution_receipt=execution, + resource_verifier_receipt=resource, + latency_verifier_receipt=latency, + ) + assert no_apply["closed"] is False + assert "typed_execution_receipt_unverified" in no_apply["active_blockers"] + execution["runtime_mutation_performed"] = True + + latency["run_id"] = "91abf6c8-5001-4b17-a538-3205d6f85757" + blocked = build_cpu_p99_closure( + candidate=candidate, + execution_receipt=execution, + resource_verifier_receipt=resource, + latency_verifier_receipt=latency, + ) + assert blocked["closed"] is False + assert blocked["same_run_receipts"] is False + assert "same_run_run_id_mismatch" in blocked["active_blockers"] + assert result["runtime_mutation_performed"] is False + + +@pytest.mark.asyncio +async def test_verified_closure_requires_durable_append_ack() -> None: + store = _MemoryCandidateStore() + await build_cpu_p99_controlled_candidate( + root_cause_target="awoooi-api", + resource_evidence=_resource(), + latency_evidence=_latency(), + log_evidence=_logs(), + recent_change_evidence=_change(), + persist_candidate=store, + ) + candidate = next(iter(store.records.values()))[1] + execution, resource, latency = _closure_receipts(candidate) + persisted: list[dict] = [] + + unacknowledged = build_cpu_p99_closure( + candidate=candidate, + execution_receipt=execution, + resource_verifier_receipt=resource, + latency_verifier_receipt=latency, + ) + assert await persist_cpu_p99_closure(candidate, unacknowledged, "awoooi") is False + + async def persist_closure(candidate_record, closure, project_id): + assert candidate_record is candidate + assert project_id == "awoooi" + persisted.append(dict(closure)) + return True + + pending = await record_cpu_p99_closure( + candidate=candidate, + execution_receipt={**execution, "runtime_mutation_performed": False}, + resource_verifier_receipt=resource, + latency_verifier_receipt=latency, + persist_closure=persist_closure, + ) + assert pending["closed"] is False + assert pending["durable_closure_ack"] is False + assert persisted == [] + + closure = await record_cpu_p99_closure( + candidate=candidate, + execution_receipt=execution, + resource_verifier_receipt=resource, + latency_verifier_receipt=latency, + persist_closure=persist_closure, + ) + + assert closure["closed"] is True + assert closure["durable_closure_ack"] is True + assert persisted[0]["durable_closure_ack"] is True + + +def test_alert_and_registry_contract_require_correlation_not_single_signal_apply() -> ( + None +): + signoz = yaml.safe_load( + (ROOT / "ops/signoz/alerting/rules.yaml").read_text(encoding="utf-8") + ) + p99 = next( + rule + for group in signoz["groups"] + for rule in group.get("rules", []) + if rule.get("alert") == "APIHighLatencyP99" + ) + labels = p99["labels"] + assert labels["auto_repair"] == "false" + assert labels["canonical_asset_id"] == "service:awoooi-api" + assert labels["namespace"] == "awoooi-prod" + assert labels["workload_kind"] == "deployment" + assert labels["workload_name"] == "awoooi-api" + assert "不得因 P99 單一訊號直接 restart" in p99["annotations"]["runbook"] + + for path in ( + ROOT / "ops/monitoring/alerts-unified.yml", + ROOT / "ops/monitoring/alerts.yml", + ): + payload = yaml.safe_load(path.read_text(encoding="utf-8")) + host_cpu = next( + rule + for group in payload["groups"] + for rule in group.get("rules", []) + if rule.get("alert") == "HostHighCpuLoad" + ) + assert host_cpu["labels"]["auto_repair"] == "false" + + identity = ServiceRegistryClient( + ROOT / "ops/config/service-registry.yaml" + ).resolve_identity("awoooi-api") + assert identity["resolution_status"] == "resolved" + assert identity["canonical_id"] == "service:awoooi-api" + assert identity["kubernetes_namespace"] == "awoooi-prod" + assert identity["kubernetes_kind"] == "deployment" + assert identity["kubernetes_name"] == "awoooi-api" + + readbacks = { + row["readback_id"]: row + for row in build_independent_verifier_registry()["external_readbacks"] + } + assert readbacks["cpu_resource_postcondition"]["required_source"] == ( + "prometheus_independent_query" + ) + assert readbacks["api_p99_postcondition"]["required_source"] == ( + "signoz_independent_query" + ) + assert readbacks["cpu_resource_postcondition"]["runtime_status"] == ( + "receipt_pending" + ) + assert readbacks["api_p99_postcondition"]["runtime_status"] == ("receipt_pending") 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 156e490d6..ca14a087e 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": 31, - "in_progress": 29, + "source_implemented_runtime_pending": 32, + "in_progress": 28, "not_started_or_no_current_evidence": 2, "superseded": 2, }, diff --git a/docs/operations/sre-ai-agent-conversation-commitments.snapshot.json b/docs/operations/sre-ai-agent-conversation-commitments.snapshot.json index 10a92f129..c44f58e71 100644 --- a/docs/operations/sre-ai-agent-conversation-commitments.snapshot.json +++ b/docs/operations/sre-ai-agent-conversation-commitments.snapshot.json @@ -1,7 +1,7 @@ { "schema_version": "sre_ai_agent_conversation_commitments_v1", "program_id": "AIA-SRE-P0-20260715", - "generated_at": "2026-07-22T12:33:42+08:00", + "generated_at": "2026-07-22T15:36:56+08:00", "source_scope": "This Codex AI Agent task only; cross-task delegation envelopes are excluded.", "authoritative_for_execution_order": false, "raw_conversation_embedded": false, @@ -51,7 +51,7 @@ {"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":"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-034","category":"named_incident","status":"source_implemented_runtime_pending","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 關閉。","source_evidence":["same-run canonical CPU/P99 metrics, sanitized log and recent-change receipt correlation","exact typed Kubernetes, Host Ansible and database RCA lanes with unknown-asset fail-closed behavior","fingerprint-deduplicated durable candidate receipt without provider, Agent99 or runtime mutation","durable closure requires typed execution plus deterministic independent resource and SigNoz P99 verifier receipts on one run"]}, {"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。"}, {"id":"AIA-CONV-037","category":"named_incident","status":"source_implemented_runtime_pending","title":"AlertChainBroken_Alertmanager webhook error 大於 10% 自動旁路與修復","linked_work_items":["AIA-SRE-014","AIA-SRE-017"],"terminal_condition":"integration-scoped counters、旁路 delivery、receiver contract 與恢復 receipt 全部成立。"}, @@ -98,8 +98,8 @@ "active_or_completed_commitments": 68, "by_status": { "analysis_or_governance_complete": 6, - "source_implemented_runtime_pending": 31, - "in_progress": 29, + "source_implemented_runtime_pending": 32, + "in_progress": 28, "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 1f40fa2dc..8356692d6 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: "1a2777467c6d7b7912617f2fa6d22f590fa95a469aab0a8ff477d9725b425854" + awoooi.wooo.work/source-sha256: "b1ad3cdd9ad3e041afd57f867738027b7668f72df63d0f238d7fa4befc584f9f" data: service-registry.yaml: | # ops/config/service-registry.yaml @@ -190,6 +190,11 @@ data: 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-api containers: [] - name: awoooi-web diff --git a/ops/config/service-registry.yaml b/ops/config/service-registry.yaml index c08392609..1072efada 100644 --- a/ops/config/service-registry.yaml +++ b/ops/config/service-registry.yaml @@ -176,6 +176,11 @@ services: 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-api containers: [] - name: awoooi-web diff --git a/ops/monitoring/alerts.yml b/ops/monitoring/alerts.yml index 60b54a18d..0bbe0b77d 100644 --- a/ops/monitoring/alerts.yml +++ b/ops/monitoring/alerts.yml @@ -116,7 +116,7 @@ groups: severity: warning layer: systemd-188 team: ops - auto_repair: "true" + auto_repair: "false" # MCP Phase 2a (ADR-071, 2026-04-11 Claude Sonnet 4.6): SSH MCP 路由標籤 mcp_provider: "ssh_host" host_type: "bare_metal" diff --git a/ops/signoz/alerting/rules.yaml b/ops/signoz/alerting/rules.yaml index 75461f1df..2c9eee68b 100644 --- a/ops/signoz/alerting/rules.yaml +++ b/ops/signoz/alerting/rules.yaml @@ -51,9 +51,18 @@ groups: severity: warning source: signoz team: backend + alert_category: service_latency + auto_repair: "false" + canonical_asset_id: "service:awoooi-api" + asset_domain: kubernetes_workload + target_resource: awoooi-api + namespace: awoooi-prod + workload_kind: deployment + workload_name: awoooi-api annotations: summary: "API P99 延遲 > 2s" description: "服務 {{ $labels.service_name }} P99: {{ $value }}s" + runbook: "先建立 CPU/P99 correlation receipt,要求同一 run 的 canonical resource metrics、sanitized logs 與 recent-change evidence;不得因 P99 單一訊號直接 restart。只有 typed executor 完成 bounded repair,且 resource 與 SigNoz latency independent verifier 都通過,才可 closure。" - alert: APIHighLatencyP95 expr: |