feat(sre): persist alert-chain telegram queue receipts
This commit is contained in:
@@ -12,8 +12,9 @@ is transport and coordination only; it never executes the Linux repair.
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import html
|
||||
import json
|
||||
from collections.abc import Awaitable, Callable
|
||||
from collections.abc import Awaitable, Callable, Mapping
|
||||
from datetime import UTC, datetime
|
||||
from typing import Any, Literal
|
||||
from uuid import NAMESPACE_URL, UUID, uuid5
|
||||
@@ -252,6 +253,88 @@ def _incident_payload_for_candidate(
|
||||
return payload
|
||||
|
||||
|
||||
def _alert_chain_candidate_queued_message(
|
||||
*,
|
||||
owner: Mapping[str, Any],
|
||||
) -> str:
|
||||
"""Render a compact, truthful queue receipt for the SRE war room."""
|
||||
|
||||
incident_id = html.escape(str(owner["incident_id"]))
|
||||
run_id = html.escape(str(owner["run_id"]))
|
||||
return (
|
||||
"🚨 <b>[ALERT-CHAIN | CRITICAL]</b>\n"
|
||||
"<code>ai_automation_alert_card_v1</code>\n\n"
|
||||
"<b>Alertmanager 告警鏈路中斷</b>\n"
|
||||
f"├ 資產:<code>{ALERT_CHAIN_CANONICAL_ASSET_ID}</code>\n"
|
||||
"├ 現況:修復候選已持久化並排入受控佇列\n"
|
||||
f"├ Incident:<code>{incident_id}</code>\n"
|
||||
f"└ Run:<code>{run_id}</code>\n\n"
|
||||
"<b>AI Agent 動作</b>\n"
|
||||
"├ Agent99:已完成只讀旁路探測與驗證轉送\n"
|
||||
"├ Policy:已用 deterministic typed-domain policy 選定執行器\n"
|
||||
"├ Host Ansible:已排入 check → bounded apply → verify\n"
|
||||
"└ LLM:未介入此緊急路由(不偽造 Ollama/Claude/Gemini 執行)\n\n"
|
||||
"⏭ <b>自動後續</b>:獨立驗證 → Incident closure "
|
||||
"→ KM/RAG/MCP/PlayBook\n"
|
||||
"<code>runtime_write_gate=controlled</code>"
|
||||
)
|
||||
|
||||
|
||||
async def _send_alert_chain_candidate_queued_telegram_receipt(
|
||||
*,
|
||||
owner: Mapping[str, Any],
|
||||
project_id: str,
|
||||
) -> dict[str, Any]:
|
||||
"""Send the queue receipt and reduce provider evidence to a safe receipt."""
|
||||
|
||||
from src.services.awooop_deeplinks import incident_truth_chain_reply_markup
|
||||
from src.services.telegram_gateway import (
|
||||
_telegram_send_delivery_succeeded,
|
||||
get_telegram_gateway,
|
||||
)
|
||||
|
||||
delivery = await get_telegram_gateway().send_alert_notification(
|
||||
_alert_chain_candidate_queued_message(owner=owner),
|
||||
product_id=project_id,
|
||||
signal_family="incident_lifecycle",
|
||||
severity="P1",
|
||||
reply_markup=incident_truth_chain_reply_markup(
|
||||
str(owner["incident_id"]),
|
||||
project_id=project_id,
|
||||
),
|
||||
)
|
||||
provider_result = (
|
||||
delivery.get("result")
|
||||
if isinstance(delivery, Mapping)
|
||||
and isinstance(delivery.get("result"), Mapping)
|
||||
else {}
|
||||
)
|
||||
route_receipt = (
|
||||
delivery.get("_awoooi_canonical_route_receipt")
|
||||
if isinstance(delivery, Mapping)
|
||||
and isinstance(
|
||||
delivery.get("_awoooi_canonical_route_receipt"), Mapping
|
||||
)
|
||||
else {}
|
||||
)
|
||||
return {
|
||||
"schema_version": "alert_chain_telegram_queue_receipt_v1",
|
||||
"attempted": True,
|
||||
"durable_acknowledged": _telegram_send_delivery_succeeded(delivery),
|
||||
"delivery_status": str(
|
||||
delivery.get("_awooop_delivery_status") or "unknown"
|
||||
)
|
||||
if isinstance(delivery, Mapping)
|
||||
else "invalid_response",
|
||||
"provider_message_id": str(provider_result.get("message_id") or ""),
|
||||
"destination_alias": str(route_receipt.get("destination_alias") or ""),
|
||||
"classification": "alert_chain_health",
|
||||
"agent99_activity": "read_only_relay_completed",
|
||||
"executor_activity": "host_ansible_candidate_queued",
|
||||
"model_activity": "not_invoked_deterministic_policy",
|
||||
}
|
||||
|
||||
|
||||
async def process_alert_chain_emergency_relay(
|
||||
payload: AlertChainEmergencyRelayRequest,
|
||||
*,
|
||||
@@ -263,6 +346,7 @@ async def process_alert_chain_emergency_relay(
|
||||
),
|
||||
candidate_enqueuer: Callable[..., Awaitable[dict[str, Any]]] | None = None,
|
||||
event_recorder: Callable[..., Awaitable[Any]] = record_alertmanager_event,
|
||||
telegram_notifier: Callable[..., Awaitable[dict[str, Any]]] | None = None,
|
||||
) -> dict[str, Any]:
|
||||
"""Persist the deterministic owner and queue the exact Ansible candidate."""
|
||||
|
||||
@@ -272,6 +356,8 @@ async def process_alert_chain_emergency_relay(
|
||||
)
|
||||
|
||||
candidate_enqueuer = enqueue_ai_decision_ansible_candidate
|
||||
if telegram_notifier is None:
|
||||
telegram_notifier = _send_alert_chain_candidate_queued_telegram_receipt
|
||||
|
||||
owner = build_alert_chain_emergency_owner(payload, project_id=project_id)
|
||||
approvals = approval_service or get_approval_service()
|
||||
@@ -454,10 +540,77 @@ async def process_alert_chain_emergency_relay(
|
||||
incident_id=owner["incident_id"],
|
||||
)
|
||||
|
||||
telegram_receipt: dict[str, Any] = {
|
||||
"schema_version": "alert_chain_telegram_queue_receipt_v1",
|
||||
"attempted": False,
|
||||
"durable_acknowledged": False,
|
||||
"delivery_status": "candidate_not_queued",
|
||||
"provider_message_id": "",
|
||||
"destination_alias": "",
|
||||
"classification": "alert_chain_health",
|
||||
"agent99_activity": "read_only_relay_completed",
|
||||
"executor_activity": (
|
||||
"host_ansible_candidate_queued"
|
||||
if queued
|
||||
else "host_ansible_candidate_not_queued"
|
||||
),
|
||||
"model_activity": "not_invoked_deterministic_policy",
|
||||
}
|
||||
if queued:
|
||||
try:
|
||||
notification_result = await telegram_notifier(
|
||||
owner=owner,
|
||||
project_id=project_id,
|
||||
)
|
||||
if isinstance(notification_result, Mapping):
|
||||
delivery_fields = (
|
||||
"attempted",
|
||||
"durable_acknowledged",
|
||||
"delivery_status",
|
||||
"provider_message_id",
|
||||
"destination_alias",
|
||||
)
|
||||
telegram_receipt.update(
|
||||
{
|
||||
key: notification_result.get(key)
|
||||
for key in delivery_fields
|
||||
if key in notification_result
|
||||
}
|
||||
)
|
||||
except Exception as exc:
|
||||
telegram_receipt.update(
|
||||
{
|
||||
"attempted": True,
|
||||
"durable_acknowledged": False,
|
||||
"delivery_status": f"failed:{type(exc).__name__}",
|
||||
}
|
||||
)
|
||||
logger.warning(
|
||||
"alert_chain_candidate_queued_telegram_receipt_failed",
|
||||
error_type=type(exc).__name__,
|
||||
incident_id=owner["incident_id"],
|
||||
)
|
||||
|
||||
telegram_acknowledged = bool(
|
||||
queued
|
||||
and telegram_receipt.get("attempted") is True
|
||||
and telegram_receipt.get("durable_acknowledged") is True
|
||||
)
|
||||
queued_blockers = [
|
||||
"independent_post_verifier_pending",
|
||||
"incident_closure_and_learning_writeback_pending",
|
||||
]
|
||||
if queued and not telegram_acknowledged:
|
||||
queued_blockers.insert(0, "telegram_durable_queue_receipt_missing")
|
||||
|
||||
return {
|
||||
"schema_version": "alert_chain_emergency_ingress_result_v1",
|
||||
"status": (
|
||||
"ansible_candidate_queued_verifier_pending"
|
||||
(
|
||||
"ansible_candidate_queued_verifier_pending"
|
||||
if telegram_acknowledged
|
||||
else "ansible_candidate_queued_telegram_receipt_pending"
|
||||
)
|
||||
if queued
|
||||
else str(
|
||||
enqueue_receipt.get("status")
|
||||
@@ -494,13 +647,12 @@ async def process_alert_chain_emergency_relay(
|
||||
for value in enqueue_receipt.get("active_blockers") or []
|
||||
][:8],
|
||||
},
|
||||
"telegram_receipt": telegram_receipt,
|
||||
"alert_to_telegram_receipt_verified": telegram_acknowledged,
|
||||
"runtime_write_performed": False,
|
||||
"runtime_closure_verified": False,
|
||||
"active_blockers": (
|
||||
[
|
||||
"independent_post_verifier_pending",
|
||||
"incident_closure_and_learning_writeback_pending",
|
||||
]
|
||||
queued_blockers
|
||||
if queued
|
||||
else [
|
||||
str(value)[:160]
|
||||
@@ -510,7 +662,13 @@ async def process_alert_chain_emergency_relay(
|
||||
][:8]
|
||||
),
|
||||
"safe_next_action": (
|
||||
"ansible_worker_check_apply_verify_writeback_same_run"
|
||||
(
|
||||
"ansible_worker_check_apply_verify_writeback_same_run"
|
||||
if telegram_acknowledged
|
||||
else (
|
||||
"ansible_worker_continue_same_run_and_retry_telegram_receipt"
|
||||
)
|
||||
)
|
||||
if queued
|
||||
else "retry_same_relay_receipt_after_candidate_queue_blocker"
|
||||
),
|
||||
|
||||
@@ -24,6 +24,7 @@ from src.services.alert_chain_emergency_ingress import (
|
||||
ALERT_CHAIN_CATALOG_ID,
|
||||
ALERT_CHAIN_EXECUTOR,
|
||||
AlertChainEmergencyRelayRequest,
|
||||
_alert_chain_candidate_queued_message,
|
||||
build_alert_chain_emergency_owner,
|
||||
process_alert_chain_emergency_relay,
|
||||
)
|
||||
@@ -254,6 +255,24 @@ def test_owner_identity_is_stable_and_binds_same_run_trace_work_item() -> None:
|
||||
assert first["canonical_asset_id"] == ALERT_CHAIN_CANONICAL_ASSET_ID
|
||||
|
||||
|
||||
def test_queue_receipt_message_is_compact_truthful_and_actionable() -> None:
|
||||
owner = build_alert_chain_emergency_owner(
|
||||
AlertChainEmergencyRelayRequest.model_validate(payload())
|
||||
)
|
||||
|
||||
message = _alert_chain_candidate_queued_message(owner=owner)
|
||||
|
||||
assert "[ALERT-CHAIN | CRITICAL]" in message
|
||||
assert "ai_automation_alert_card_v1" in message
|
||||
assert owner["incident_id"] in message
|
||||
assert owner["run_id"] in message
|
||||
assert "Agent99:已完成只讀旁路探測與驗證轉送" in message
|
||||
assert "Host Ansible:已排入 check → bounded apply → verify" in message
|
||||
assert "未介入此緊急路由" in message
|
||||
assert "KM/RAG/MCP/PlayBook" in message
|
||||
assert "建議下一步" not in message
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_secondary_ingress_persists_and_queues_exact_ansible_candidate() -> None:
|
||||
request = AlertChainEmergencyRelayRequest.model_validate(payload())
|
||||
@@ -262,6 +281,7 @@ async def test_secondary_ingress_persists_and_queues_exact_ansible_candidate() -
|
||||
incidents = FakeIncidentService()
|
||||
enqueue_calls = []
|
||||
events = []
|
||||
telegram_calls = []
|
||||
|
||||
async def create_incident(**kwargs):
|
||||
assert kwargs["canonical_incident_id"] == owner["incident_id"]
|
||||
@@ -290,6 +310,18 @@ async def test_secondary_ingress_persists_and_queues_exact_ansible_candidate() -
|
||||
async def record_event(**kwargs):
|
||||
events.append(kwargs)
|
||||
|
||||
async def notify_telegram(**kwargs):
|
||||
telegram_calls.append(kwargs)
|
||||
assert kwargs["owner"] == owner
|
||||
assert kwargs["project_id"] == "awoooi"
|
||||
return {
|
||||
"attempted": True,
|
||||
"durable_acknowledged": True,
|
||||
"delivery_status": "sent",
|
||||
"provider_message_id": "4431",
|
||||
"destination_alias": "awoooi_sre_war_room",
|
||||
}
|
||||
|
||||
result = await process_alert_chain_emergency_relay(
|
||||
request,
|
||||
approval_service=approvals,
|
||||
@@ -297,6 +329,7 @@ async def test_secondary_ingress_persists_and_queues_exact_ansible_candidate() -
|
||||
incident_creator=create_incident,
|
||||
candidate_enqueuer=enqueue,
|
||||
event_recorder=record_event,
|
||||
telegram_notifier=notify_telegram,
|
||||
)
|
||||
|
||||
assert result["status"] == "ansible_candidate_queued_verifier_pending"
|
||||
@@ -308,7 +341,21 @@ async def test_secondary_ingress_persists_and_queues_exact_ansible_candidate() -
|
||||
assert result["executor"] == "host_ansible_executor"
|
||||
assert result["runtime_write_performed"] is False
|
||||
assert result["runtime_closure_verified"] is False
|
||||
assert result["alert_to_telegram_receipt_verified"] is True
|
||||
assert result["telegram_receipt"] == {
|
||||
"schema_version": "alert_chain_telegram_queue_receipt_v1",
|
||||
"attempted": True,
|
||||
"durable_acknowledged": True,
|
||||
"delivery_status": "sent",
|
||||
"provider_message_id": "4431",
|
||||
"destination_alias": "awoooi_sre_war_room",
|
||||
"classification": "alert_chain_health",
|
||||
"agent99_activity": "read_only_relay_completed",
|
||||
"executor_activity": "host_ansible_candidate_queued",
|
||||
"model_activity": "not_invoked_deterministic_policy",
|
||||
}
|
||||
assert len(enqueue_calls) == 1
|
||||
assert len(telegram_calls) == 1
|
||||
assert events[0]["stage"] == "secondary_ingress_ansible_candidate_queued"
|
||||
|
||||
repeated = await process_alert_chain_emergency_relay(
|
||||
@@ -318,9 +365,11 @@ async def test_secondary_ingress_persists_and_queues_exact_ansible_candidate() -
|
||||
incident_creator=create_incident,
|
||||
candidate_enqueuer=enqueue,
|
||||
event_recorder=record_event,
|
||||
telegram_notifier=notify_telegram,
|
||||
)
|
||||
assert repeated["identity"] == result["identity"]
|
||||
assert approvals.incremented == 1
|
||||
assert len(telegram_calls) == 2
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@@ -348,6 +397,9 @@ async def test_secondary_ingress_keeps_persisted_candidate_partial_when_queue_bl
|
||||
async def no_event(**_kwargs):
|
||||
return None
|
||||
|
||||
async def must_not_notify(**_kwargs):
|
||||
raise AssertionError("telegram must not send before exact queue ack")
|
||||
|
||||
result = await process_alert_chain_emergency_relay(
|
||||
request,
|
||||
approval_service=approvals,
|
||||
@@ -355,14 +407,70 @@ async def test_secondary_ingress_keeps_persisted_candidate_partial_when_queue_bl
|
||||
incident_creator=create_incident,
|
||||
candidate_enqueuer=blocked_enqueue,
|
||||
event_recorder=no_event,
|
||||
telegram_notifier=must_not_notify,
|
||||
)
|
||||
|
||||
assert result["candidate_persisted"] is True
|
||||
assert result["candidate_queued"] is False
|
||||
assert result["telegram_receipt"]["attempted"] is False
|
||||
assert result["alert_to_telegram_receipt_verified"] is False
|
||||
assert result["runtime_closure_verified"] is False
|
||||
assert result["safe_next_action"].startswith("retry_same_relay_receipt")
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_queued_candidate_continues_when_telegram_receipt_is_missing() -> None:
|
||||
request = AlertChainEmergencyRelayRequest.model_validate(payload())
|
||||
owner = build_alert_chain_emergency_owner(request)
|
||||
approvals = FakeApprovalService()
|
||||
incidents = FakeIncidentService()
|
||||
|
||||
async def create_incident(**_kwargs):
|
||||
incidents.incident = FakeIncident(owner["incident_id"])
|
||||
return owner["incident_id"]
|
||||
|
||||
async def enqueue(**_kwargs):
|
||||
return {
|
||||
"status": "controlled_check_mode_queued",
|
||||
"automation_run_id": owner["run_id"],
|
||||
"trace_id": owner["trace_id"],
|
||||
"work_item_id": owner["work_item_id"],
|
||||
"queued": True,
|
||||
"side_effect_performed": False,
|
||||
"active_blockers": [],
|
||||
}
|
||||
|
||||
async def record_event(**_kwargs):
|
||||
return None
|
||||
|
||||
async def notify_without_durable_ack(**_kwargs):
|
||||
return {
|
||||
"attempted": True,
|
||||
"durable_acknowledged": False,
|
||||
"delivery_status": "durable_mirror_failed",
|
||||
}
|
||||
|
||||
result = await process_alert_chain_emergency_relay(
|
||||
request,
|
||||
approval_service=approvals,
|
||||
incident_service=incidents,
|
||||
incident_creator=create_incident,
|
||||
candidate_enqueuer=enqueue,
|
||||
event_recorder=record_event,
|
||||
telegram_notifier=notify_without_durable_ack,
|
||||
)
|
||||
|
||||
assert result["status"] == (
|
||||
"ansible_candidate_queued_telegram_receipt_pending"
|
||||
)
|
||||
assert result["candidate_queued"] is True
|
||||
assert result["alert_to_telegram_receipt_verified"] is False
|
||||
assert "telegram_durable_queue_receipt_missing" in result["active_blockers"]
|
||||
assert result["safe_next_action"] == (
|
||||
"ansible_worker_continue_same_run_and_retry_telegram_receipt"
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_secondary_ingress_rejects_deterministic_owner_collision() -> None:
|
||||
request = AlertChainEmergencyRelayRequest.model_validate(payload())
|
||||
|
||||
Reference in New Issue
Block a user