feat(alertmanager): track recovered lifecycle receipts

This commit is contained in:
Your Name
2026-07-22 13:51:15 +08:00
parent fb4f8e3d0c
commit 3a6f8badb7
5 changed files with 1483 additions and 6 deletions

View File

@@ -63,6 +63,9 @@ from src.services.alertmanager_llm_guard import (
ALERTMANAGER_LLM_INFLIGHT_LOCK_TTL_SECONDS,
try_acquire_alertmanager_llm_lock,
)
from src.services.alertmanager_recovery_lifecycle import (
handle_alertmanager_resolved,
)
from src.services.approval_db import get_approval_service
from src.services.auto_approve import get_auto_approve_policy
from src.services.auto_repair_service import AutoRepairService
@@ -2274,6 +2277,75 @@ def _trusted_internal_alertmanager_source(request: Request) -> tuple[bool, str]:
return True, forwarded_chain[0]
def _normalized_alertmanager_resolved_identity(
alert: AlertmanagerAlert,
) -> dict[str, str]:
"""Rebuild the same fingerprint identity used by the firing path."""
labels = dict(alert.labels or {})
annotations = dict(alert.annotations or {})
alertname = str(labels.get("alertname") or "UnknownAlert")
alert_type = get_incident_type(alertname)
if alert_type not in {
"k8s_node_failure",
"k8s_pod_crash",
"db_connection_timeout",
"service_404",
"high_cpu",
"high_memory",
"disk_full",
"ssl_expiry",
"custom",
}:
alert_type = "custom"
severity = {
"critical": "critical",
"warning": "warning",
"info": "info",
}.get(str(labels.get("severity") or "warning").casefold(), "warning")
instance = str(labels.get("instance") or "")
instance_clean = instance.split(":")[0] if ":" in instance else instance
target_resource = str(
labels.get("component")
or labels.get("pod")
or labels.get("name")
or labels.get("container")
or labels.get("job")
or (
instance_clean
if instance_clean and not is_internal_ip(instance_clean)
else ""
)
or alertname
)
namespace = str(labels.get("namespace") or "default")
message = str(
annotations.get("summary")
or annotations.get("description")
or alertname
)
normalized_alert = AlertPayload(
alert_type=alert_type,
severity=severity,
source="alertmanager",
target_resource=target_resource,
namespace=namespace,
message=message,
metrics={},
labels=labels,
)
fingerprint = AlertAnalyzer.generate_fingerprint(normalized_alert)
return {
"alertname": alertname,
"alert_type": alert_type,
"severity": severity,
"target_resource": target_resource,
"namespace": namespace,
"message": message,
"fingerprint": fingerprint,
}
async def _process_new_alert_background(
alert_context: dict,
alert_id: str,
@@ -3350,9 +3422,38 @@ async def alertmanager_webhook(
alert_count=len(payload.alerts),
)
# 只處理第一個 firing 告警 (避免告警風暴)
# 只處理第一個 firing 告警;若本批只有 resolved走同 fingerprint
# canonical lifecycle重複恢復不得重發 Telegram。
firing_alerts = [a for a in payload.alerts if a.status == "firing"]
if not firing_alerts:
resolved_alerts = [a for a in payload.alerts if a.status == "resolved"]
if resolved_alerts:
resolved_alert = resolved_alerts[0]
identity = _normalized_alertmanager_resolved_identity(resolved_alert)
source_alert_id = str(
resolved_alert.fingerprint or identity["fingerprint"]
)
recovery = await handle_alertmanager_resolved(
fingerprint=identity["fingerprint"],
alertname=identity["alertname"],
severity=identity["severity"],
namespace=identity["namespace"],
target_resource=identity["target_resource"],
source_alert_id=source_alert_id,
source_ends_at=str(resolved_alert.endsAt or ""),
)
_record_webhook_metric("success")
return AlertResponse(
success=True,
message=(
"Alertmanager resolved lifecycle: "
f"{recovery.status}; provider_ack="
f"{str(recovery.provider_acknowledged).lower()}"
),
alert_id=source_alert_id,
approval_created=False,
converged=recovery.duplicate,
)
_record_webhook_metric("success")
return AlertResponse(
success=True,

View File

@@ -0,0 +1,767 @@
"""Canonical Alertmanager resolved/recovered lifecycle with send-once Telegram ack.
The lifecycle is keyed by the normalized AWOOOI alert fingerprint. Repeated
resolved webhooks update one durable row; retries are bounded and a known
provider send is never repeated without acknowledgement reconciliation. A
recovery is acknowledged only when Telegram proves the canonical destination
and provider message identity.
"""
from __future__ import annotations
import hashlib
import json
from collections.abc import Awaitable, Callable, Mapping
from dataclasses import dataclass
from datetime import UTC, datetime
from html import escape
from typing import Any
from uuid import NAMESPACE_URL, UUID, uuid5
import structlog
from sqlalchemy import text
from src.services.audit_sink import sanitize
logger = structlog.get_logger(__name__)
_RECOVERY_RESERVATION_TTL_SECONDS = 300
_MAX_RECOVERY_DELIVERY_ATTEMPTS = 3
CorrelationLookup = Callable[[str], Awaitable[dict[str, Any] | None]]
RecoveryReservation = Callable[..., Awaitable[dict[str, Any]]]
IncidentResolver = Callable[[str], Awaitable[bool]]
RecoverySender = Callable[..., Awaitable[dict[str, Any]]]
OutcomeRecorder = Callable[..., Awaitable[None]]
@dataclass(frozen=True)
class AlertmanagerRecoveryResult:
"""Public-safe result for one resolved webhook."""
status: str
fingerprint: str
incident_id: str
lifecycle_event_id: str
duplicate: bool
provider_acknowledged: bool
provider_send_performed: bool
def _lifecycle_provider_event_id(fingerprint: str) -> str:
safe_fingerprint = str(fingerprint or "").strip()[:64] or "missing-fingerprint"
return f"alertmanager:lifecycle:{safe_fingerprint}"
def _parse_timestamp(value: object) -> datetime | None:
try:
parsed = datetime.fromisoformat(str(value or ""))
except ValueError:
return None
if parsed.tzinfo is None:
return parsed.replace(tzinfo=UTC)
return parsed.astimezone(UTC)
def _reservation_decision(
lifecycle: Mapping[str, Any] | None,
*,
incident_id: str,
now: datetime,
) -> dict[str, Any]:
"""Choose one send/retry decision without touching runtime state."""
current = dict(lifecycle or {})
status = str(current.get("status") or "")
attempts = max(0, int(current.get("delivery_attempts") or 0))
current_incident_id = str(current.get("incident_id") or "")
if incident_id and current_incident_id and incident_id != current_incident_id:
return {
"admitted": True,
"status": "recovery_pending_delivery",
"attempts": 1,
"new_incident_cycle": True,
}
if status == "recovered_acknowledged":
return {
"admitted": False,
"status": "duplicate_recovered_suppressed",
"attempts": attempts,
}
if current.get("provider_send_performed") is True:
return {
"admitted": False,
"status": "recovery_delivery_ack_unresolved",
"attempts": attempts,
}
if not incident_id:
return {
"admitted": False,
"status": "resolved_identity_unmatched",
"attempts": attempts,
}
if status == "recovery_pending_delivery":
reserved_at = _parse_timestamp(current.get("reserved_at"))
if (
reserved_at is not None
and (now - reserved_at).total_seconds() < _RECOVERY_RESERVATION_TTL_SECONDS
):
return {
"admitted": False,
"status": "duplicate_recovery_inflight_suppressed",
"attempts": attempts,
}
if attempts >= _MAX_RECOVERY_DELIVERY_ATTEMPTS:
return {
"admitted": False,
"status": "recovery_delivery_retry_exhausted",
"attempts": attempts,
}
return {
"admitted": True,
"status": "recovery_pending_delivery",
"attempts": attempts + 1,
}
def _merge_reservation_observation(
lifecycle: Mapping[str, Any] | None,
*,
decision: Mapping[str, Any],
observed: Mapping[str, Any],
now: datetime,
is_new: bool,
) -> dict[str, Any]:
"""Persist observations without letting duplicates rewrite lifecycle truth."""
updated = dict(lifecycle or {})
if decision.get("new_incident_cycle") is True:
completed_cycles = list(updated.get("completed_cycles") or [])
if updated.get("incident_id"):
completed_cycles.append(
{
"incident_id": updated.get("incident_id"),
"approval_id": updated.get("approval_id"),
"status": updated.get("status"),
"delivery_attempts": updated.get("delivery_attempts"),
"provider_send_performed": updated.get("provider_send_performed"),
"provider_acknowledged": updated.get("provider_acknowledged"),
"delivery_ack": updated.get("delivery_ack") or {},
"completed_at": updated.get("completed_at"),
}
)
updated["completed_cycles"] = completed_cycles[-8:]
updated.update(observed)
updated.update(
{
"last_observed_at": now.isoformat(),
"executor_invocation_allowed": False,
"agent99_dispatch_allowed": False,
"paid_provider_call_allowed": False,
"runtime_infrastructure_mutation_allowed": False,
}
)
if decision["admitted"]:
for stale_key in (
"completed_at",
"incident_resolution_succeeded",
"last_duplicate_at",
"last_duplicate_status",
):
updated.pop(stale_key, None)
updated.update(
{
"status": decision["status"],
"delivery_attempts": decision["attempts"],
"reserved_at": now.isoformat(),
"provider_send_performed": False,
"provider_acknowledged": False,
"delivery_ack": {},
}
)
elif is_new:
updated.update(
{
"status": decision["status"],
"delivery_attempts": decision["attempts"],
"reserved_at": None,
"provider_send_performed": False,
"provider_acknowledged": False,
}
)
else:
updated.update(
{
"last_duplicate_status": decision["status"],
"last_duplicate_at": now.isoformat(),
}
)
return updated
async def lookup_alertmanager_incident_correlation(
fingerprint: str,
) -> dict[str, Any] | None:
"""Read the latest durable approval/incident identity for one fingerprint."""
from src.db.base import get_db_context
async with get_db_context("awoooi") as db:
result = await db.execute(
text(
"""
SELECT approval.id, approval.incident_id,
approval.telegram_message_id
FROM approval_records approval
JOIN incidents incident
ON incident.incident_id = approval.incident_id
WHERE approval.fingerprint = :fingerprint
AND approval.incident_id IS NOT NULL
AND incident.resolved_at IS NULL
AND upper(coalesce(incident.status::text, ''))
NOT IN ('RESOLVED', 'CLOSED')
ORDER BY approval.last_seen_at DESC, approval.created_at DESC
LIMIT 1
"""
),
{"fingerprint": fingerprint},
)
row = result.fetchone()
if row is None:
return None
return {
"approval_id": str(row[0]),
"incident_id": str(row[1]),
"telegram_message_id": int(row[2]) if row[2] is not None else None,
}
async def reserve_alertmanager_recovery(
*,
fingerprint: str,
alertname: str,
severity: str,
namespace: str,
target_resource: str,
source_alert_id: str,
source_ends_at: str,
correlation: Mapping[str, Any] | None,
) -> dict[str, Any]:
"""Atomically reserve one recovery delivery on the canonical lifecycle row."""
from src.db.base import get_db_context
now = datetime.now(UTC)
provider_event_id = _lifecycle_provider_event_id(fingerprint)
run_id = uuid5(NAMESPACE_URL, f"awooop:alertmanager-lifecycle:{fingerprint}")
correlated_incident_id = str((correlation or {}).get("incident_id") or "")
correlated_approval_id = str((correlation or {}).get("approval_id") or "")
correlated_message_id = (correlation or {}).get("telegram_message_id")
lock_key = f"awoooi:{provider_event_id}"
async with get_db_context("awoooi") as db:
await db.execute(
text("SELECT pg_advisory_xact_lock(hashtextextended(:lock_key, 0))"),
{"lock_key": lock_key},
)
selected = await db.execute(
text(
"""
SELECT event_id, source_envelope
FROM awooop_conversation_event
WHERE project_id = 'awoooi'
AND channel_type = 'internal'
AND provider_event_id = :provider_event_id
FOR UPDATE
"""
),
{"provider_event_id": provider_event_id},
)
existing = selected.fetchone()
envelope: dict[str, Any]
if existing is None:
event_id = uuid5(
NAMESPACE_URL,
f"awooop:alertmanager-lifecycle-event:{fingerprint}",
)
envelope = {
"schema_version": "alertmanager_canonical_lifecycle_v1",
"alertmanager_lifecycle": {},
}
else:
event_id = UUID(str(existing[0]))
raw_envelope = existing[1]
if isinstance(raw_envelope, str):
raw_envelope = json.loads(raw_envelope)
envelope = (
dict(raw_envelope)
if isinstance(raw_envelope, Mapping)
else {
"schema_version": "alertmanager_canonical_lifecycle_v1",
"alertmanager_lifecycle": {},
}
)
current_lifecycle = envelope.get("alertmanager_lifecycle")
current = current_lifecycle if isinstance(current_lifecycle, Mapping) else None
incident_id = correlated_incident_id or str(
(current or {}).get("incident_id") or ""
)
approval_id = correlated_approval_id or str(
(current or {}).get("approval_id") or ""
)
telegram_message_id = (
int(correlated_message_id)
if correlated_message_id is not None
else int((current or {}).get("telegram_message_id"))
if (current or {}).get("telegram_message_id") is not None
else None
)
decision = _reservation_decision(
current,
incident_id=incident_id,
now=now,
)
lifecycle = _merge_reservation_observation(
current,
decision=decision,
observed={
"fingerprint": fingerprint,
"alertname": alertname,
"severity": severity,
"namespace": namespace,
"target_resource": target_resource,
"source_alert_id": source_alert_id,
"source_ends_at": source_ends_at,
"approval_id": approval_id or (current or {}).get("approval_id"),
"incident_id": incident_id or (current or {}).get("incident_id"),
"telegram_message_id": telegram_message_id,
},
now=now,
is_new=existing is None,
)
envelope["alertmanager_lifecycle"] = lifecycle
safe_envelope = sanitize(envelope)
serialized = json.dumps(safe_envelope, ensure_ascii=False, default=str)
raw_preview = (
f"Alertmanager lifecycle {decision['status']}: "
f"{alertname} {incident_id or 'unmatched'}"
)
preview = str(sanitize({"preview": raw_preview})["preview"])[:256]
content_hash = hashlib.sha256(fingerprint.encode("utf-8")).hexdigest()
if existing is None:
await db.execute(
text(
"""
INSERT INTO awooop_conversation_event (
event_id, project_id, channel_type, provider_event_id,
platform_subject_id, channel_user_id, channel_chat_id,
run_id, content_type, content_hash, content_preview,
content_redacted, redaction_version, source_envelope,
is_duplicate, provider_ts, received_at
) VALUES (
:event_id, 'awoooi', 'internal', :provider_event_id,
'alertmanager', 'alertmanager', :channel_chat_id,
:run_id, 'text', :content_hash, :preview,
:preview, 'audit_sink_v1', CAST(:source_envelope AS jsonb),
FALSE, NOW(), NOW()
)
"""
),
{
"event_id": event_id,
"provider_event_id": provider_event_id,
"channel_chat_id": "alertmanager:internal",
"run_id": run_id,
"content_hash": content_hash,
"preview": preview,
"source_envelope": serialized,
},
)
else:
await db.execute(
text(
"""
UPDATE awooop_conversation_event
SET source_envelope = CAST(:source_envelope AS jsonb),
content_preview = :preview,
content_redacted = :preview,
is_duplicate = :is_duplicate,
provider_ts = NOW()
WHERE event_id = :event_id
"""
),
{
"event_id": event_id,
"source_envelope": serialized,
"preview": preview,
"is_duplicate": not decision["admitted"],
},
)
return {
"event_id": str(event_id),
"admitted": bool(decision["admitted"]),
"status": str(decision["status"]),
"duplicate": not bool(decision["admitted"]),
"delivery_attempt": int(decision["attempts"]),
"incident_id": incident_id,
"telegram_message_id": telegram_message_id,
"provider_acknowledged": bool(lifecycle.get("provider_acknowledged")),
"provider_send_performed": bool(lifecycle.get("provider_send_performed")),
}
async def resolve_alertmanager_incident(incident_id: str) -> bool:
"""Resolve the correlated incident without emitting a second postmortem card."""
from src.services.incident_service import get_incident_service
incident = await get_incident_service().resolve_incident(
incident_id,
resolution_type="alertmanager_recovered",
emit_postmortem=False,
)
return incident is not None
def _delivery_ack(result: Mapping[str, Any]) -> dict[str, Any]:
route = result.get("_awoooi_canonical_route_receipt")
context = result.get("_awoooi_delivery_context")
provider_result = result.get("result")
return {
"provider_message_id": (
str(provider_result.get("message_id"))
if isinstance(provider_result, Mapping)
and provider_result.get("message_id") is not None
else None
),
"sender_bot_alias": (
route.get("sender_bot_alias") if isinstance(route, Mapping) else None
),
"destination_alias": (
route.get("destination_alias") if isinstance(route, Mapping) else None
),
"destination_binding": (
route.get("destination_binding") if isinstance(route, Mapping) else None
),
"provider_destination_binding": (
context.get("provider_destination_binding")
if isinstance(context, Mapping)
else None
),
"provider_destination_verification_method": (
context.get("provider_destination_verification_method")
if isinstance(context, Mapping)
else None
),
}
def _outcome_matches_current_cycle(
lifecycle: Mapping[str, Any],
*,
incident_id: str,
delivery_attempt: int,
acknowledged: bool,
) -> bool:
"""Reject stale outcomes while allowing a late verified ack to win."""
if str(lifecycle.get("incident_id") or "") != incident_id:
return False
current_attempt = max(0, int(lifecycle.get("delivery_attempts") or 0))
if current_attempt != delivery_attempt and not acknowledged:
return False
if lifecycle.get("status") == "recovered_acknowledged" and not acknowledged:
return False
return True
async def send_alertmanager_recovery(
*,
incident_id: str,
alertname: str,
severity: str,
target_resource: str,
fingerprint: str,
reply_to_message_id: int | None,
) -> dict[str, Any]:
"""Send one canonical recovery card and return only verified ack metadata."""
from src.services.telegram_gateway import (
_telegram_send_delivery_succeeded,
get_telegram_gateway,
)
route_severity = {
"critical": "P0",
"warning": "P1",
"info": "P2",
}.get(severity.casefold(), "P1")
body = (
"✅ <b>ALERT RECOVERED</b>\n"
"━━━━━━━━━━━━━━━━━━━\n"
f"📋 <code>{escape(incident_id)}</code>\n"
f"🔔 告警:<code>{escape(alertname)}</code>\n"
f"🎯 資產:<code>{escape(target_resource)}</code>\n"
f"🔗 fingerprint<code>{escape(fingerprint[:32])}</code>\n\n"
"Alertmanager 已回報 resolvedincident lifecycle 已更新。"
)
result = await get_telegram_gateway().send_canonical_message(
product_id="awoooi",
signal_family="incident_lifecycle",
severity=route_severity,
text=body,
reply_to_message_id=reply_to_message_id,
)
acknowledged = _telegram_send_delivery_succeeded(result)
return {
"acknowledged": acknowledged,
"provider_send_performed": bool(
isinstance(result, Mapping)
and result.get("_awooop_provider_send_performed") is True
),
"ack": _delivery_ack(result) if isinstance(result, Mapping) else {},
}
async def record_alertmanager_recovery_outcome(
*,
lifecycle_event_id: str,
fingerprint: str,
incident_id: str,
delivery_attempt: int,
incident_resolution_succeeded: bool,
delivery: Mapping[str, Any],
) -> None:
"""Finalize the same lifecycle row with incident and destination ack truth."""
from src.db.base import get_db_context
acknowledged = bool(delivery.get("acknowledged") is True)
provider_send_performed = bool(delivery.get("provider_send_performed") is True)
status = (
"incident_resolution_failed"
if not incident_resolution_succeeded
else "recovered_acknowledged"
if acknowledged
else "recovery_delivery_failed"
)
async with get_db_context("awoooi") as db:
selected = await db.execute(
text(
"""
SELECT source_envelope
FROM awooop_conversation_event
WHERE event_id = CAST(:event_id AS uuid)
AND project_id = 'awoooi'
AND provider_event_id = :provider_event_id
FOR UPDATE
"""
),
{
"event_id": lifecycle_event_id,
"provider_event_id": _lifecycle_provider_event_id(fingerprint),
},
)
row = selected.fetchone()
if row is None:
raise RuntimeError("alertmanager_lifecycle_receipt_missing")
raw_envelope = row[0]
if isinstance(raw_envelope, str):
raw_envelope = json.loads(raw_envelope)
envelope = dict(raw_envelope) if isinstance(raw_envelope, Mapping) else {}
lifecycle = envelope.get("alertmanager_lifecycle")
lifecycle = dict(lifecycle) if isinstance(lifecycle, Mapping) else {}
if not _outcome_matches_current_cycle(
lifecycle,
incident_id=incident_id,
delivery_attempt=delivery_attempt,
acknowledged=acknowledged,
):
logger.info(
"alertmanager_recovery_stale_outcome_suppressed",
fingerprint=fingerprint,
incident_id=incident_id,
delivery_attempt=delivery_attempt,
)
return
lifecycle.update(
{
"status": status,
"incident_resolution_succeeded": incident_resolution_succeeded,
"provider_send_performed": provider_send_performed,
"provider_acknowledged": acknowledged,
"delivery_ack": dict(delivery.get("ack") or {}),
"completed_at": datetime.now(UTC).isoformat(),
}
)
envelope["alertmanager_lifecycle"] = lifecycle
serialized = json.dumps(sanitize(envelope), ensure_ascii=False, default=str)
await db.execute(
text(
"""
UPDATE awooop_conversation_event
SET source_envelope = CAST(:source_envelope AS jsonb),
content_preview = :preview,
content_redacted = :preview,
is_duplicate = FALSE,
provider_ts = NOW()
WHERE event_id = CAST(:event_id AS uuid)
"""
),
{
"event_id": lifecycle_event_id,
"source_envelope": serialized,
"preview": f"Alertmanager lifecycle {status}"[:256],
},
)
async def handle_alertmanager_resolved(
*,
fingerprint: str,
alertname: str,
severity: str,
namespace: str,
target_resource: str,
source_alert_id: str,
source_ends_at: str,
correlation_lookup: CorrelationLookup = lookup_alertmanager_incident_correlation,
recovery_reservation: RecoveryReservation = reserve_alertmanager_recovery,
incident_resolver: IncidentResolver = resolve_alertmanager_incident,
recovery_sender: RecoverySender = send_alertmanager_recovery,
outcome_recorder: OutcomeRecorder = record_alertmanager_recovery_outcome,
) -> AlertmanagerRecoveryResult:
"""Resolve, notify, and acknowledge one canonical recovery transition."""
try:
correlation = await correlation_lookup(fingerprint)
reservation = await recovery_reservation(
fingerprint=fingerprint,
alertname=alertname,
severity=severity,
namespace=namespace,
target_resource=target_resource,
source_alert_id=source_alert_id,
source_ends_at=source_ends_at,
correlation=correlation,
)
except Exception as exc: # noqa: BLE001 - no durable receipt means no action
logger.warning(
"alertmanager_recovery_reservation_failed",
fingerprint=fingerprint,
error=type(exc).__name__,
)
return AlertmanagerRecoveryResult(
status="durable_lifecycle_unavailable",
fingerprint=fingerprint,
incident_id="",
lifecycle_event_id="",
duplicate=False,
provider_acknowledged=False,
provider_send_performed=False,
)
incident_id = str(
reservation.get("incident_id") or (correlation or {}).get("incident_id") or ""
)
lifecycle_event_id = str(reservation.get("event_id") or "")
delivery_attempt = int(reservation.get("delivery_attempt") or 0)
if reservation.get("admitted") is not True:
return AlertmanagerRecoveryResult(
status=str(reservation.get("status") or "recovery_suppressed"),
fingerprint=fingerprint,
incident_id=incident_id,
lifecycle_event_id=lifecycle_event_id,
duplicate=bool(reservation.get("duplicate") is True),
provider_acknowledged=bool(
reservation.get("provider_acknowledged") is True
),
provider_send_performed=bool(
reservation.get("provider_send_performed") is True
),
)
try:
incident_resolved = await incident_resolver(incident_id)
except Exception as exc: # noqa: BLE001 - persist a retryable terminal
logger.warning(
"alertmanager_incident_resolution_failed",
fingerprint=fingerprint,
incident_id=incident_id,
error=type(exc).__name__,
)
incident_resolved = False
if not incident_resolved:
await outcome_recorder(
lifecycle_event_id=lifecycle_event_id,
fingerprint=fingerprint,
incident_id=incident_id,
delivery_attempt=delivery_attempt,
incident_resolution_succeeded=False,
delivery={
"acknowledged": False,
"provider_send_performed": False,
"ack": {},
},
)
return AlertmanagerRecoveryResult(
status="incident_resolution_failed",
fingerprint=fingerprint,
incident_id=incident_id,
lifecycle_event_id=lifecycle_event_id,
duplicate=False,
provider_acknowledged=False,
provider_send_performed=False,
)
try:
delivery = await recovery_sender(
incident_id=incident_id,
alertname=alertname,
severity=severity,
target_resource=target_resource,
fingerprint=fingerprint,
reply_to_message_id=(
int(reservation["telegram_message_id"])
if reservation.get("telegram_message_id") is not None
else int(correlation["telegram_message_id"])
if correlation and correlation.get("telegram_message_id") is not None
else None
),
)
except Exception as exc: # noqa: BLE001 - durable failure remains retryable
logger.warning(
"alertmanager_recovery_delivery_failed",
fingerprint=fingerprint,
incident_id=incident_id,
error=type(exc).__name__,
)
delivery = {
"acknowledged": False,
"provider_send_performed": False,
"ack": {},
}
await outcome_recorder(
lifecycle_event_id=lifecycle_event_id,
fingerprint=fingerprint,
incident_id=incident_id,
delivery_attempt=delivery_attempt,
incident_resolution_succeeded=True,
delivery=delivery,
)
acknowledged = bool(delivery.get("acknowledged") is True)
return AlertmanagerRecoveryResult(
status=(
"recovered_acknowledged" if acknowledged else "recovery_delivery_failed"
),
fingerprint=fingerprint,
incident_id=incident_id,
lifecycle_event_id=lifecycle_event_id,
duplicate=False,
provider_acknowledged=acknowledged,
provider_send_performed=bool(delivery.get("provider_send_performed") is True),
)

View File

@@ -0,0 +1,603 @@
from __future__ import annotations
from datetime import UTC, datetime, timedelta
from unittest.mock import AsyncMock
import pytest
from fastapi import BackgroundTasks
from starlette.requests import Request
from src.api.v1 import webhooks as webhooks_module
from src.api.v1.webhooks import (
AlertmanagerAlert,
AlertmanagerPayload,
_normalized_alertmanager_resolved_identity,
alertmanager_webhook,
)
from src.services.alertmanager_recovery_lifecycle import (
AlertmanagerRecoveryResult,
_merge_reservation_observation,
_outcome_matches_current_cycle,
_reservation_decision,
handle_alertmanager_resolved,
send_alertmanager_recovery,
)
from src.services.telegram_gateway import _telegram_destination_binding
def _request() -> Request:
return Request(
{
"type": "http",
"method": "POST",
"path": "/api/v1/webhooks/alertmanager",
"scheme": "http",
"server": ("testserver", 80),
"client": ("127.0.0.1", 50000),
"query_string": b"",
"headers": [],
}
)
def test_reservation_suppresses_acknowledged_and_fresh_inflight_duplicates() -> None:
now = datetime.now(UTC)
acknowledged = _reservation_decision(
{
"status": "recovered_acknowledged",
"delivery_attempts": 1,
},
incident_id="INC-20260722-ABC123",
now=now,
)
inflight = _reservation_decision(
{
"status": "recovery_pending_delivery",
"delivery_attempts": 1,
"reserved_at": (now - timedelta(seconds=30)).isoformat(),
},
incident_id="INC-20260722-ABC123",
now=now,
)
assert acknowledged == {
"admitted": False,
"status": "duplicate_recovered_suppressed",
"attempts": 1,
}
assert inflight == {
"admitted": False,
"status": "duplicate_recovery_inflight_suppressed",
"attempts": 1,
}
def test_reservation_allows_bounded_retry_but_fails_closed_without_identity() -> None:
now = datetime.now(UTC)
retry = _reservation_decision(
{
"status": "recovery_delivery_failed",
"delivery_attempts": 1,
},
incident_id="INC-20260722-ABC123",
now=now,
)
exhausted = _reservation_decision(
{
"status": "recovery_delivery_failed",
"delivery_attempts": 3,
},
incident_id="INC-20260722-ABC123",
now=now,
)
unmatched = _reservation_decision(
None,
incident_id="",
now=now,
)
assert retry == {
"admitted": True,
"status": "recovery_pending_delivery",
"attempts": 2,
}
assert exhausted["status"] == "recovery_delivery_retry_exhausted"
assert exhausted["admitted"] is False
assert unmatched["status"] == "resolved_identity_unmatched"
assert unmatched["admitted"] is False
def test_duplicate_observation_preserves_terminal_ack_and_inflight_reservation() -> (
None
):
now = datetime.now(UTC)
acknowledged = {
"status": "recovered_acknowledged",
"delivery_attempts": 1,
"reserved_at": (now - timedelta(seconds=60)).isoformat(),
"provider_send_performed": True,
"provider_acknowledged": True,
"delivery_ack": {"provider_message_id": "991"},
}
decision = _reservation_decision(
acknowledged,
incident_id="INC-20260722-ABC123",
now=now,
)
updated = _merge_reservation_observation(
acknowledged,
decision=decision,
observed={"source_ends_at": "2026-07-22T02:00:00Z"},
now=now,
is_new=False,
)
assert updated["status"] == "recovered_acknowledged"
assert updated["reserved_at"] == acknowledged["reserved_at"]
assert updated["provider_acknowledged"] is True
assert updated["delivery_ack"] == {"provider_message_id": "991"}
assert updated["last_duplicate_status"] == "duplicate_recovered_suppressed"
inflight = {
"status": "recovery_pending_delivery",
"delivery_attempts": 1,
"reserved_at": (now - timedelta(seconds=30)).isoformat(),
"provider_send_performed": False,
"provider_acknowledged": False,
}
inflight_decision = _reservation_decision(
inflight,
incident_id="INC-20260722-ABC123",
now=now,
)
inflight_updated = _merge_reservation_observation(
inflight,
decision=inflight_decision,
observed={"source_ends_at": "2026-07-22T02:00:00Z"},
now=now,
is_new=False,
)
assert inflight_updated["status"] == "recovery_pending_delivery"
assert inflight_updated["reserved_at"] == inflight["reserved_at"]
assert inflight_updated["last_duplicate_status"] == (
"duplicate_recovery_inflight_suppressed"
)
def test_provider_send_without_destination_ack_is_never_resent_automatically() -> None:
decision = _reservation_decision(
{
"status": "recovery_delivery_failed",
"delivery_attempts": 1,
"provider_send_performed": True,
"provider_acknowledged": False,
},
incident_id="INC-20260722-ABC123",
now=datetime.now(UTC),
)
assert decision == {
"admitted": False,
"status": "recovery_delivery_ack_unresolved",
"attempts": 1,
}
def test_new_incident_reopens_same_fingerprint_without_losing_prior_ack() -> None:
now = datetime.now(UTC)
previous = {
"incident_id": "INC-20260721-OLD001",
"approval_id": "approval-old",
"status": "recovered_acknowledged",
"delivery_attempts": 1,
"provider_send_performed": True,
"provider_acknowledged": True,
"delivery_ack": {"provider_message_id": "880"},
"completed_at": (now - timedelta(days=1)).isoformat(),
}
decision = _reservation_decision(
previous,
incident_id="INC-20260722-NEW001",
now=now,
)
updated = _merge_reservation_observation(
previous,
decision=decision,
observed={
"incident_id": "INC-20260722-NEW001",
"approval_id": "approval-new",
},
now=now,
is_new=False,
)
assert decision["new_incident_cycle"] is True
assert updated["incident_id"] == "INC-20260722-NEW001"
assert updated["status"] == "recovery_pending_delivery"
assert updated["delivery_attempts"] == 1
assert updated["provider_acknowledged"] is False
assert "completed_at" not in updated
assert updated["completed_cycles"][-1]["incident_id"] == ("INC-20260721-OLD001")
assert updated["completed_cycles"][-1]["provider_acknowledged"] is True
def test_outcome_state_is_monotonic_and_cycle_bound() -> None:
lifecycle = {
"incident_id": "INC-20260722-ABC123",
"status": "recovery_pending_delivery",
"delivery_attempts": 2,
}
assert not _outcome_matches_current_cycle(
lifecycle,
incident_id="INC-20260721-OLD001",
delivery_attempt=1,
acknowledged=True,
)
assert not _outcome_matches_current_cycle(
lifecycle,
incident_id="INC-20260722-ABC123",
delivery_attempt=1,
acknowledged=False,
)
assert _outcome_matches_current_cycle(
lifecycle,
incident_id="INC-20260722-ABC123",
delivery_attempt=1,
acknowledged=True,
)
assert not _outcome_matches_current_cycle(
{**lifecycle, "status": "recovered_acknowledged"},
incident_id="INC-20260722-ABC123",
delivery_attempt=2,
acknowledged=False,
)
@pytest.mark.asyncio
async def test_resolved_handler_updates_incident_and_records_destination_ack() -> None:
correlation_lookup = AsyncMock(
return_value={
"approval_id": "approval-1",
"incident_id": "INC-20260722-ABC123",
"telegram_message_id": 881,
}
)
reservation = AsyncMock(
return_value={
"event_id": "0b7d5caf-4aa1-4e5a-8f51-987d66db7c20",
"admitted": True,
"status": "recovery_pending_delivery",
"duplicate": False,
"delivery_attempt": 1,
}
)
incident_resolver = AsyncMock(return_value=True)
sender = AsyncMock(
return_value={
"acknowledged": True,
"provider_send_performed": True,
"ack": {
"provider_message_id": "991",
"destination_binding": "binding-1",
"provider_destination_binding": "binding-1",
},
}
)
outcome_recorder = AsyncMock(return_value=None)
result = await handle_alertmanager_resolved(
fingerprint="a" * 32,
alertname="DockerContainerUnhealthy",
severity="critical",
namespace="default",
target_resource="alertmanager",
source_alert_id="source-fp-1",
source_ends_at="2026-07-22T02:00:00Z",
correlation_lookup=correlation_lookup,
recovery_reservation=reservation,
incident_resolver=incident_resolver,
recovery_sender=sender,
outcome_recorder=outcome_recorder,
)
assert result.status == "recovered_acknowledged"
assert result.provider_acknowledged is True
incident_resolver.assert_awaited_once_with("INC-20260722-ABC123")
sender.assert_awaited_once()
assert sender.await_args.kwargs["reply_to_message_id"] == 881
outcome_recorder.assert_awaited_once()
assert outcome_recorder.await_args.kwargs["delivery"]["acknowledged"] is True
@pytest.mark.asyncio
async def test_delivery_retry_uses_identity_from_durable_lifecycle_receipt() -> None:
correlation_lookup = AsyncMock(return_value=None)
reservation = AsyncMock(
return_value={
"event_id": "0b7d5caf-4aa1-4e5a-8f51-987d66db7c20",
"admitted": True,
"status": "recovery_pending_delivery",
"duplicate": False,
"delivery_attempt": 2,
"incident_id": "INC-20260722-ABC123",
"telegram_message_id": 881,
}
)
incident_resolver = AsyncMock(return_value=True)
sender = AsyncMock(
return_value={
"acknowledged": True,
"provider_send_performed": True,
"ack": {"provider_message_id": "992"},
}
)
outcome_recorder = AsyncMock(return_value=None)
result = await handle_alertmanager_resolved(
fingerprint="a" * 32,
alertname="DockerContainerUnhealthy",
severity="critical",
namespace="default",
target_resource="alertmanager",
source_alert_id="source-fp-1",
source_ends_at="2026-07-22T02:00:00Z",
correlation_lookup=correlation_lookup,
recovery_reservation=reservation,
incident_resolver=incident_resolver,
recovery_sender=sender,
outcome_recorder=outcome_recorder,
)
assert result.status == "recovered_acknowledged"
assert result.incident_id == "INC-20260722-ABC123"
incident_resolver.assert_awaited_once_with("INC-20260722-ABC123")
assert sender.await_args.kwargs["reply_to_message_id"] == 881
assert outcome_recorder.await_args.kwargs["delivery_attempt"] == 2
@pytest.mark.asyncio
async def test_duplicate_resolved_event_never_resolves_or_sends_again() -> None:
correlation_lookup = AsyncMock(
return_value={
"approval_id": "approval-1",
"incident_id": "INC-20260722-ABC123",
"telegram_message_id": 881,
}
)
reservation = AsyncMock(
return_value={
"event_id": "0b7d5caf-4aa1-4e5a-8f51-987d66db7c20",
"admitted": False,
"status": "duplicate_recovered_suppressed",
"duplicate": True,
"delivery_attempt": 1,
"provider_acknowledged": True,
"provider_send_performed": True,
}
)
incident_resolver = AsyncMock(
side_effect=AssertionError("duplicate must not resolve incident again")
)
sender = AsyncMock(
side_effect=AssertionError("duplicate must not send Telegram again")
)
outcome_recorder = AsyncMock(
side_effect=AssertionError("duplicate must not rewrite acknowledged outcome")
)
result = await handle_alertmanager_resolved(
fingerprint="a" * 32,
alertname="DockerContainerUnhealthy",
severity="critical",
namespace="default",
target_resource="alertmanager",
source_alert_id="source-fp-1",
source_ends_at="2026-07-22T02:00:00Z",
correlation_lookup=correlation_lookup,
recovery_reservation=reservation,
incident_resolver=incident_resolver,
recovery_sender=sender,
outcome_recorder=outcome_recorder,
)
assert result.status == "duplicate_recovered_suppressed"
assert result.duplicate is True
assert result.provider_acknowledged is True
assert result.provider_send_performed is True
incident_resolver.assert_not_awaited()
sender.assert_not_awaited()
outcome_recorder.assert_not_awaited()
@pytest.mark.asyncio
async def test_unmatched_resolved_identity_records_no_false_recovery() -> None:
correlation_lookup = AsyncMock(return_value=None)
reservation = AsyncMock(
return_value={
"event_id": "0b7d5caf-4aa1-4e5a-8f51-987d66db7c20",
"admitted": False,
"status": "resolved_identity_unmatched",
"duplicate": True,
"delivery_attempt": 0,
}
)
incident_resolver = AsyncMock(
side_effect=AssertionError("unmatched identity must remain fail closed")
)
sender = AsyncMock(
side_effect=AssertionError("unmatched identity must not send recovery")
)
result = await handle_alertmanager_resolved(
fingerprint="b" * 32,
alertname="UnknownAlert",
severity="warning",
namespace="default",
target_resource="unknown-target",
source_alert_id="source-fp-2",
source_ends_at="2026-07-22T02:00:00Z",
correlation_lookup=correlation_lookup,
recovery_reservation=reservation,
incident_resolver=incident_resolver,
recovery_sender=sender,
outcome_recorder=AsyncMock(),
)
assert result.status == "resolved_identity_unmatched"
assert result.provider_acknowledged is False
incident_resolver.assert_not_awaited()
sender.assert_not_awaited()
@pytest.mark.asyncio
async def test_resolved_webhook_enters_canonical_lifecycle(
monkeypatch: pytest.MonkeyPatch,
) -> None:
lifecycle = AsyncMock(
return_value=AlertmanagerRecoveryResult(
status="recovered_acknowledged",
fingerprint="a" * 32,
incident_id="INC-20260722-ABC123",
lifecycle_event_id="0b7d5caf-4aa1-4e5a-8f51-987d66db7c20",
duplicate=False,
provider_acknowledged=True,
provider_send_performed=True,
)
)
monkeypatch.setattr(
webhooks_module,
"handle_alertmanager_resolved",
lifecycle,
)
alert = AlertmanagerAlert(
status="resolved",
labels={
"alertname": "DockerContainerUnhealthy",
"severity": "critical",
"namespace": "default",
"name": "alertmanager",
},
annotations={"summary": "container unhealthy"},
startsAt="2026-07-22T01:00:00Z",
endsAt="2026-07-22T02:00:00Z",
fingerprint="provider-fingerprint",
)
identity = _normalized_alertmanager_resolved_identity(alert)
response = await alertmanager_webhook(
_request(),
AlertmanagerPayload(status="resolved", alerts=[alert]),
BackgroundTasks(),
)
assert response.success is True
assert response.alert_id == "provider-fingerprint"
assert "provider_ack=true" in response.message
lifecycle.assert_awaited_once_with(
fingerprint=identity["fingerprint"],
alertname="DockerContainerUnhealthy",
severity="critical",
namespace="default",
target_resource="alertmanager",
source_alert_id="provider-fingerprint",
source_ends_at="2026-07-22T02:00:00Z",
)
@pytest.mark.asyncio
async def test_recovery_sender_requires_destination_bound_provider_receipt(
monkeypatch: pytest.MonkeyPatch,
) -> None:
from src.services import telegram_gateway as gateway_module
destination_binding = _telegram_destination_binding(-100123)
provider_result = {
"ok": True,
"result": {
"message_id": 991,
"chat": {"id": -100123},
},
"_awooop_delivery_status": "sent",
"_awooop_provider_send_performed": True,
"_awoooi_canonical_route_receipt": {
"schema_version": "telegram_canonical_egress_receipt_v1",
"decision": "allowed",
"provider_send_performed": True,
"sender_bot_alias": "tsenyang_bot",
"destination_alias": "awoooi_sre_war_room",
"destination_binding": destination_binding,
},
"_awoooi_delivery_context": {
"schema_version": "telegram_delivery_context_v1",
"sender_bot_alias": "tsenyang_bot",
"destination_alias": "awoooi_sre_war_room",
"destination_binding": destination_binding,
"payload_destination_binding": destination_binding,
"provider_destination_binding": destination_binding,
"provider_destination_verification_method": (
"requested_chat_id_matches_provider_chat_id"
),
"destination_binding_verified": True,
},
}
class _Gateway:
send_canonical_message = AsyncMock(return_value=provider_result)
gateway = _Gateway()
monkeypatch.setattr(gateway_module, "get_telegram_gateway", lambda: gateway)
delivery = await send_alertmanager_recovery(
incident_id="INC-20260722-ABC123",
alertname="DockerContainerUnhealthy",
severity="critical",
target_resource="alertmanager",
fingerprint="a" * 32,
reply_to_message_id=881,
)
assert delivery["acknowledged"] is True
assert delivery["provider_send_performed"] is True
assert delivery["ack"]["provider_message_id"] == "991"
assert delivery["ack"]["destination_binding"] == destination_binding
assert delivery["ack"]["provider_destination_binding"] == destination_binding
assert gateway.send_canonical_message.await_args.kwargs == {
"product_id": "awoooi",
"signal_family": "incident_lifecycle",
"severity": "P0",
"text": gateway.send_canonical_message.await_args.kwargs["text"],
"reply_to_message_id": 881,
}
def test_resolved_normalization_is_stable_for_same_alert_identity() -> None:
common = {
"labels": {
"alertname": "DockerContainerUnhealthy",
"severity": "critical",
"namespace": "default",
"name": "alertmanager",
},
"annotations": {"summary": "container unhealthy"},
"startsAt": "2026-07-22T01:00:00Z",
"fingerprint": "provider-fingerprint",
}
firing = _normalized_alertmanager_resolved_identity(
AlertmanagerAlert(status="firing", **common)
)
resolved = _normalized_alertmanager_resolved_identity(
AlertmanagerAlert(
status="resolved",
endsAt="2026-07-22T02:00:00Z",
**common,
)
)
assert resolved["fingerprint"] == firing["fingerprint"]
assert resolved["target_resource"] == "alertmanager"
assert resolved["namespace"] == "default"

View File

@@ -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": 28,
"in_progress": 32,
"source_implemented_runtime_pending": 29,
"in_progress": 31,
"not_started_or_no_current_evidence": 2,
"superseded": 2,
},
@@ -176,6 +176,12 @@ def test_loader_returns_fixed_architecture_provider_order_and_agent99_bridge() -
assert "fingerprint-deduplicated" in " ".join(
commitments["AIA-CONV-030"]["source_evidence"]
)
assert commitments["AIA-CONV-031"]["status"] == (
"source_implemented_runtime_pending"
)
assert "destination-bound provider acknowledgement" in " ".join(
commitments["AIA-CONV-031"]["source_evidence"]
)
assert "Host112" in commitments["AIA-CONV-049"]["title"]
assert commitments["AIA-CONV-049"]["status"] == (
"source_implemented_runtime_pending"

View File

@@ -47,7 +47,7 @@
{"id":"AIA-CONV-028","category":"telegram","status":"source_implemented_runtime_pending","title":"串通 callback ingress、dispatch、verifier 與 durable action receipt","linked_work_items":["AIA-SRE-009","AIA-SRE-015","AIA-SRE-017"],"terminal_condition":"recipient callback 與 same-run controlled action receipt 可由 production readback 證明。"},
{"id":"AIA-CONV-029","category":"telegram","status":"source_implemented_runtime_pending","title":"AwoooI SRE 戰情室 Bot 自動理解並正確回覆使用者訊息","linked_work_items":["AIA-SRE-012","AIA-SRE-013","AIA-SRE-017"],"terminal_condition":"群組訊息經身份、意圖、RAG 與 policy 後產生可追溯回覆,不接受未授權 runtime mutation。","source_evidence":["telegram group text durable ingress receipt","deterministic intent classification","bounded sanitized KM/RAG retrieval","traceable no-runtime-mutation provider receipt"]},
{"id":"AIA-CONV-030","category":"telegram","status":"source_implemented_runtime_pending","title":"對話可啟動受控調查/修復並建立 Codex 開發 work item","linked_work_items":["AIA-SRE-001","AIA-SRE-015","AIA-SRE-017"],"terminal_condition":"討論結果轉為去重 work itemproduction action 仍走 typed executor/verifier。","source_evidence":["verified controlled_action_request ingress gate","exact canonical asset/domain resolution or asset_identity_unresolved drift item","fingerprint-deduplicated durable internal work-item receipt","deterministic no-provider no-runtime-mutation Telegram response"]},
{"id":"AIA-CONV-031","category":"telegram","status":"in_progress","title":"同 fingerprint 去重、聚合、抑噪及 resolved/recovered 更新","linked_work_items":["AIA-SRE-017"],"terminal_condition":"同事件只維護 canonical lifecycle恢復卡有 destination-bound provider acknowledgement。"},
{"id":"AIA-CONV-031","category":"telegram","status":"source_implemented_runtime_pending","title":"同 fingerprint 去重、聚合、抑噪及 resolved/recovered 更新","linked_work_items":["AIA-SRE-017"],"terminal_condition":"同事件只維護 canonical lifecycle恢復卡有 destination-bound provider acknowledgement。","source_evidence":["existing firing fingerprint convergence and parent-child alert storm aggregation","single-row durable resolved lifecycle reservation with bounded retry","duplicate resolved webhook suppression before incident or Telegram side effects","recovery card accepted only with destination-bound provider acknowledgement"]},
{"id":"AIA-CONV-032","category":"named_incident","status":"in_progress","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。"},
{"id":"AIA-CONV-033","category":"named_incident","status":"in_progress","title":"INC-20260711-D037E5 修復 stale candidate 阻擋、自動排隊與 durable closure","linked_work_items":["AIA-SRE-007","AIA-SRE-015"],"terminal_condition":"新候選可 claimapply 不重複verifier、Telegram 與 learning receipts 同 run 完成。"},
@@ -98,8 +98,8 @@
"active_or_completed_commitments": 68,
"by_status": {
"analysis_or_governance_complete": 6,
"source_implemented_runtime_pending": 28,
"in_progress": 32,
"source_implemented_runtime_pending": 29,
"in_progress": 31,
"not_started_or_no_current_evidence": 2,
"superseded": 2
},