feat(telegram): enforce canonical product routing
This commit is contained in:
@@ -68,8 +68,15 @@ class FailoverAlerter:
|
||||
f"Fallback 鏈:{_escape_md(fallback_chain_str)}\n\n"
|
||||
f"自動恢復服務持續監控,3 次 HEALTHY 後自動切回"
|
||||
)
|
||||
await self._send(msg)
|
||||
logger.info("failover_alert_sent", to_provider=to_provider)
|
||||
delivered = await self._deliver_reserved(
|
||||
msg,
|
||||
dedup_key,
|
||||
ttl=DEDUP_TTL_SEC,
|
||||
)
|
||||
if delivered:
|
||||
logger.info("failover_alert_sent", to_provider=to_provider)
|
||||
else:
|
||||
logger.warning("failover_alert_no_send", to_provider=to_provider)
|
||||
|
||||
async def alert_recovery(self, event: dict[str, Any]) -> None:
|
||||
"""Ollama 自動恢復告警 — 1h dedup per host
|
||||
@@ -99,8 +106,15 @@ class FailoverAlerter:
|
||||
f"切換路徑:{_escape_md(str(from_provider))} → {_escape_md(str(to_provider))}\n\n"
|
||||
f"自動化飛輪已恢復至高效能推理模式"
|
||||
)
|
||||
await self._send(msg)
|
||||
logger.info("recovery_alert_sent", from_provider=from_provider)
|
||||
delivered = await self._deliver_reserved(
|
||||
msg,
|
||||
dedup_key,
|
||||
ttl=RECOVERY_DEDUP_TTL_SEC,
|
||||
)
|
||||
if delivered:
|
||||
logger.info("recovery_alert_sent", from_provider=from_provider)
|
||||
else:
|
||||
logger.warning("recovery_alert_no_send", from_provider=from_provider)
|
||||
|
||||
async def alert_governance(self, event_type: str, payload: dict[str, Any]) -> None:
|
||||
"""AI 治理告警(dedup 1h)
|
||||
@@ -135,8 +149,11 @@ class FailoverAlerter:
|
||||
return
|
||||
|
||||
msg = format_governance_alert_card(event_type, payload)
|
||||
await self._send(msg)
|
||||
logger.info("governance_alert_sent", event_type=event_type)
|
||||
delivered = await self._deliver_reserved(msg, dedup_key, ttl=3600)
|
||||
if delivered:
|
||||
logger.info("governance_alert_sent", event_type=event_type)
|
||||
else:
|
||||
logger.warning("governance_alert_no_send", event_type=event_type)
|
||||
|
||||
async def alert_gemini_quota_exceeded(self, event: dict[str, Any]) -> None:
|
||||
"""Gemini 每日上限觸發,降級到 188 CPU 備援 — 24h dedup(每日重置)"""
|
||||
@@ -161,8 +178,15 @@ class FailoverAlerter:
|
||||
f"進入容災模式至明日 0:00\n"
|
||||
f"建議檢查是否有異常流量,評估是否升級 Gemini 配額"
|
||||
)
|
||||
await self._send(msg)
|
||||
logger.info("quota_alert_sent", quota=quota, current_count=current_count)
|
||||
delivered = await self._deliver_reserved(
|
||||
msg,
|
||||
dedup_key,
|
||||
ttl=QUOTA_DEDUP_TTL_SEC,
|
||||
)
|
||||
if delivered:
|
||||
logger.info("quota_alert_sent", quota=quota, current_count=current_count)
|
||||
else:
|
||||
logger.warning("quota_alert_no_send", quota=quota, current_count=current_count)
|
||||
|
||||
async def alert_provider_version_changed(self, changed_providers: list[str], probed: int) -> None:
|
||||
"""AI Provider 版本變更告警 — dedup 1h/provider
|
||||
@@ -171,19 +195,20 @@ class FailoverAlerter:
|
||||
每個 provider 獨立 dedup,避免同一版本重複告警。
|
||||
"""
|
||||
now_str = datetime.now(TAIPEI_TZ).strftime("%Y-%m-%d %H:%M")
|
||||
sent: list[str] = []
|
||||
reserved: list[tuple[str, str]] = []
|
||||
|
||||
for provider in changed_providers:
|
||||
dedup_key = f"alert:provider_version_changed:{provider}"
|
||||
if not await self._check_dedup(dedup_key, ttl=3600):
|
||||
logger.debug("provider_version_alert_dedup_skipped", provider=provider)
|
||||
continue
|
||||
sent.append(provider)
|
||||
reserved.append((provider, dedup_key))
|
||||
|
||||
if not sent:
|
||||
if not reserved:
|
||||
return
|
||||
|
||||
providers_md = "\n".join(f"• {_escape_md(p)}" for p in sent)
|
||||
providers = [provider for provider, _dedup_key in reserved]
|
||||
providers_md = "\n".join(f"• {_escape_md(p)}" for p in providers)
|
||||
msg = (
|
||||
f"*AI Provider 版本變更偵測*\n\n"
|
||||
f"時間:{_escape_md(now_str)}\n"
|
||||
@@ -191,8 +216,20 @@ class FailoverAlerter:
|
||||
f"版本已變更:\n{providers_md}\n\n"
|
||||
f"系統已自動記錄版本歷史,請確認是否需要重新驗證推理品質"
|
||||
)
|
||||
await self._send(msg)
|
||||
logger.info("provider_version_alert_sent", sent=sent)
|
||||
delivered = False
|
||||
try:
|
||||
delivered = await self._send(msg)
|
||||
finally:
|
||||
for _provider, dedup_key in reserved:
|
||||
await self._finalize_dedup(
|
||||
dedup_key,
|
||||
ttl=3600,
|
||||
delivered=delivered,
|
||||
)
|
||||
if delivered:
|
||||
logger.info("provider_version_alert_sent", sent=providers)
|
||||
else:
|
||||
logger.warning("provider_version_alert_no_send", providers=providers)
|
||||
|
||||
# -------------------------------------------------------------------------
|
||||
# Dedup(Redis SET NX EX)
|
||||
@@ -200,8 +237,12 @@ class FailoverAlerter:
|
||||
|
||||
async def _check_dedup(self, key: str, ttl: int) -> bool:
|
||||
"""
|
||||
Redis SET NX EX 防止重複告警。
|
||||
True = 第一次(應送出),False = 已送過(跳過)。
|
||||
Redis SET NX EX 建立 delivery reservation。
|
||||
True = 可嘗試送出,False = 已送過或另一個 worker 正在送。
|
||||
|
||||
reservation 的 value 是 ``pending``,不是 sent receipt。外層只有收到
|
||||
affirmative provider ack 後才透過 ``_finalize_dedup`` 改成 ``sent``;
|
||||
no-send / exception 會刪除 reservation,讓下一輪可以安全重試。
|
||||
|
||||
2026-04-25 P1.5 by Claude Engineer-D — Telegram dedup 鐵律 10min/24h TTL
|
||||
2026-04-27 Wave8-X2 by Claude — dedup fail-open 修復
|
||||
@@ -212,7 +253,12 @@ class FailoverAlerter:
|
||||
# 優先嘗試 Redis
|
||||
if self._redis is not None:
|
||||
try:
|
||||
ok = await self._redis.set(f"{key}:dedup", "1", ex=ttl, nx=True)
|
||||
ok = await self._redis.set(
|
||||
f"{key}:dedup",
|
||||
"pending",
|
||||
ex=ttl,
|
||||
nx=True,
|
||||
)
|
||||
return bool(ok)
|
||||
except Exception as e:
|
||||
logger.warning("dedup_redis_failed_using_memory", error=str(e))
|
||||
@@ -233,11 +279,63 @@ class FailoverAlerter:
|
||||
self._memory_dedup[key] = now
|
||||
return True
|
||||
|
||||
async def _finalize_dedup(
|
||||
self,
|
||||
key: str,
|
||||
*,
|
||||
ttl: int,
|
||||
delivered: bool,
|
||||
) -> None:
|
||||
"""Commit sent dedup only after provider ack, otherwise release it.
|
||||
|
||||
2026-07-14 Codex v1.1 (Asia/Taipei): route acceptance and ``ok`` are
|
||||
not delivery truth. Pending reservations prevent concurrent sends but
|
||||
must never survive a structured no-send receipt as a terminal success.
|
||||
"""
|
||||
redis_finalized = False
|
||||
if self._redis is not None:
|
||||
try:
|
||||
if delivered:
|
||||
await self._redis.set(f"{key}:dedup", "sent", ex=ttl)
|
||||
else:
|
||||
await self._redis.delete(f"{key}:dedup")
|
||||
redis_finalized = True
|
||||
except Exception as e:
|
||||
logger.warning(
|
||||
"dedup_finalize_redis_failed_using_memory",
|
||||
delivered=delivered,
|
||||
error=str(e),
|
||||
)
|
||||
|
||||
if delivered:
|
||||
if not redis_finalized:
|
||||
import time
|
||||
|
||||
self._memory_dedup[key] = time.time()
|
||||
else:
|
||||
self._memory_dedup.pop(key, None)
|
||||
|
||||
async def _deliver_reserved(
|
||||
self,
|
||||
message: str,
|
||||
key: str,
|
||||
*,
|
||||
ttl: int,
|
||||
) -> bool:
|
||||
"""Deliver one reserved alert and always finalize or release its state."""
|
||||
|
||||
delivered = False
|
||||
try:
|
||||
delivered = await self._send(message)
|
||||
return delivered
|
||||
finally:
|
||||
await self._finalize_dedup(key, ttl=ttl, delivered=delivered)
|
||||
|
||||
# -------------------------------------------------------------------------
|
||||
# 發送(透過 TelegramGateway singleton)
|
||||
# -------------------------------------------------------------------------
|
||||
|
||||
async def _send(self, message: str) -> None:
|
||||
async def _send(self, message: str) -> bool:
|
||||
"""發送至 Telegram SRE_GROUP_CHAT_ID
|
||||
|
||||
使用現有 TelegramGateway singleton(不另建 HTTP client),
|
||||
@@ -248,18 +346,34 @@ class FailoverAlerter:
|
||||
2026-04-25 P1.5 by Claude Engineer-D — 告警失敗不能阻斷主流程
|
||||
"""
|
||||
try:
|
||||
from src.core.config import get_settings
|
||||
from src.services.telegram_gateway import get_telegram_gateway
|
||||
|
||||
settings = get_settings()
|
||||
chat_id = getattr(settings, "SRE_GROUP_CHAT_ID", None)
|
||||
if not chat_id:
|
||||
logger.warning("telegram_chat_id_missing_failover_alert")
|
||||
return
|
||||
from src.services.telegram_gateway import (
|
||||
_telegram_send_delivery_succeeded,
|
||||
get_telegram_gateway,
|
||||
)
|
||||
|
||||
gateway = get_telegram_gateway()
|
||||
await gateway.send_alert_notification(text=message, parse_mode="MarkdownV2")
|
||||
logger.info("telegram_failover_alert_sent", message_len=len(message))
|
||||
if not gateway.canonical_destination_chat_id(
|
||||
product_id="awoooi",
|
||||
signal_family="shared_infrastructure",
|
||||
severity="P1",
|
||||
):
|
||||
logger.warning("telegram_canonical_route_unavailable_failover_alert")
|
||||
return False
|
||||
delivery = await gateway.send_alert_notification(
|
||||
text=message,
|
||||
parse_mode="MarkdownV2",
|
||||
product_id="awoooi",
|
||||
signal_family="shared_infrastructure",
|
||||
severity="P1",
|
||||
)
|
||||
if _telegram_send_delivery_succeeded(delivery):
|
||||
return True
|
||||
else:
|
||||
logger.warning(
|
||||
"telegram_failover_alert_blocked_no_send",
|
||||
message_len=len(message),
|
||||
)
|
||||
return False
|
||||
except Exception as e:
|
||||
# 不 raise — 告警失敗不該阻斷主流程(鐵律)
|
||||
# 2026-05-06 Codex: Telegram/httpx exception 字串可能包含 bot token URL,
|
||||
@@ -269,6 +383,7 @@ class FailoverAlerter:
|
||||
error=_sanitize_telegram_error(str(e)),
|
||||
error_type=type(e).__name__,
|
||||
)
|
||||
return False
|
||||
|
||||
|
||||
# -------------------------------------------------------------------------
|
||||
|
||||
Reference in New Issue
Block a user