refactor(api): Phase 17 P0 Router 層違規全部修復
消除 Router 層直接存取 Redis/DB 的違規: incidents.py (6 處): - 改用 IncidentService.get_active_incidents() - 改用 IncidentService.get_from_working_memory() - 改用 IncidentService.update_outcome() - 改用 IncidentService.resolve_incident() - 改用 IncidentService.find_by_proposal_id() stats.py (8 處): - 新增 StatsService 封裝快取邏輯 - 移除直接 Redis 存取 audit_logs.py (7 處): - 新增 AuditLogRepository 封裝 DB 操作 - Router 改用 Repository 層 webhooks.py (2 處): - 新增 SignalProducerService 封裝 Redis Stream - 改用 IncidentService.save_to_working_memory() 符合 leWOOOgo 積木化規範: Router → Service → Repository → DB/Redis Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
@@ -520,6 +520,193 @@ class IncidentService:
|
||||
return {}
|
||||
return {}
|
||||
|
||||
# =========================================================================
|
||||
# Phase 17 P0: Router 層違規修復 - 新增方法
|
||||
# =========================================================================
|
||||
|
||||
async def update_outcome(
|
||||
self,
|
||||
incident_id: str,
|
||||
effectiveness_score: int | None = None,
|
||||
human_feedback: str | None = None,
|
||||
learning_notes: str | None = None,
|
||||
should_remember: bool = True,
|
||||
) -> Incident | None:
|
||||
"""
|
||||
更新 Incident 的 outcome (人類回饋)
|
||||
|
||||
Phase 17: 從 Router 層遷移至 Service 層
|
||||
|
||||
Args:
|
||||
incident_id: 事件 ID
|
||||
effectiveness_score: 有效性評分 (1-5)
|
||||
human_feedback: 文字回饋
|
||||
learning_notes: 學習筆記
|
||||
should_remember: 是否納入長期記憶
|
||||
|
||||
Returns:
|
||||
Incident | None: 更新後的事件,失敗返回 None
|
||||
"""
|
||||
from src.models.incident import IncidentOutcome
|
||||
from src.repositories.incident_repository import get_incident_repository
|
||||
from src.utils.timezone import now_taipei
|
||||
|
||||
# 1. 從 Working Memory 讀取
|
||||
incident = await self.get_from_working_memory(incident_id)
|
||||
if incident is None:
|
||||
logger.warning("incident_not_found_for_outcome", incident_id=incident_id)
|
||||
return None
|
||||
|
||||
# 2. 更新 outcome
|
||||
if incident.outcome is None:
|
||||
incident.outcome = IncidentOutcome()
|
||||
|
||||
if effectiveness_score is not None:
|
||||
incident.outcome.effectiveness_score = effectiveness_score
|
||||
if human_feedback is not None:
|
||||
incident.outcome.human_feedback = human_feedback
|
||||
if learning_notes is not None:
|
||||
incident.outcome.learning_notes = learning_notes
|
||||
incident.outcome.should_remember = should_remember
|
||||
incident.updated_at = now_taipei()
|
||||
|
||||
# 3. 寫入 Working Memory
|
||||
redis_success = await self.save_to_working_memory(incident)
|
||||
if not redis_success:
|
||||
logger.error("outcome_redis_write_failed", incident_id=incident_id)
|
||||
return None
|
||||
|
||||
# 4. 同步到 Episodic Memory (PostgreSQL)
|
||||
try:
|
||||
repo = get_incident_repository()
|
||||
await repo.update_outcome(
|
||||
incident_id=incident_id,
|
||||
outcome=incident.outcome.model_dump(mode="json"),
|
||||
updated_at=now_taipei(),
|
||||
)
|
||||
logger.info("outcome_db_updated", incident_id=incident_id)
|
||||
except Exception as e:
|
||||
logger.warning(
|
||||
"outcome_db_update_failed",
|
||||
incident_id=incident_id,
|
||||
error=str(e),
|
||||
)
|
||||
# DB 失敗不影響主流程
|
||||
|
||||
return incident
|
||||
|
||||
async def resolve_incident(self, incident_id: str) -> Incident | None:
|
||||
"""
|
||||
將 Incident 狀態更新為 RESOLVED
|
||||
|
||||
Phase 17: 從 Router 層遷移至 Service 層
|
||||
|
||||
Args:
|
||||
incident_id: 事件 ID
|
||||
|
||||
Returns:
|
||||
Incident | None: 更新後的事件,失敗返回 None
|
||||
"""
|
||||
from src.repositories.incident_repository import get_incident_repository
|
||||
from src.utils.timezone import now_taipei
|
||||
|
||||
# 1. 從 Working Memory 讀取
|
||||
incident = await self.get_from_working_memory(incident_id)
|
||||
if incident is None:
|
||||
logger.warning("incident_not_found_for_resolve", incident_id=incident_id)
|
||||
return None
|
||||
|
||||
# 2. 更新狀態
|
||||
incident.status = IncidentStatus.RESOLVED
|
||||
incident.resolved_at = now_taipei()
|
||||
incident.updated_at = now_taipei()
|
||||
|
||||
# 3. 寫入 Working Memory
|
||||
redis_success = await self.save_to_working_memory(incident)
|
||||
if not redis_success:
|
||||
logger.error("resolve_redis_write_failed", incident_id=incident_id)
|
||||
return None
|
||||
|
||||
# 4. 同步到 Episodic Memory
|
||||
try:
|
||||
repo = get_incident_repository()
|
||||
await repo.update_status(
|
||||
incident_id=incident_id,
|
||||
status="resolved",
|
||||
updated_at=now_taipei(),
|
||||
)
|
||||
logger.info("resolve_db_updated", incident_id=incident_id)
|
||||
except Exception as e:
|
||||
logger.warning(
|
||||
"resolve_db_update_failed",
|
||||
incident_id=incident_id,
|
||||
error=str(e),
|
||||
)
|
||||
|
||||
return incident
|
||||
|
||||
async def find_by_proposal_id(self, proposal_id: str) -> Incident | None:
|
||||
"""
|
||||
根據 proposal_id 查找關聯的 Incident
|
||||
|
||||
Phase 17: 從 Router 層遷移至 Service 層
|
||||
|
||||
Args:
|
||||
proposal_id: 提案 ID (UUID 字串)
|
||||
|
||||
Returns:
|
||||
Incident | None: 找到的事件,未找到返回 None
|
||||
"""
|
||||
from uuid import UUID
|
||||
|
||||
redis_client = get_redis()
|
||||
|
||||
try:
|
||||
target_uuid = UUID(proposal_id)
|
||||
|
||||
async for key in redis_client.scan_iter(
|
||||
match=f"{INCIDENT_KEY_PREFIX}INC-*",
|
||||
count=100,
|
||||
):
|
||||
data = await redis_client.get(key)
|
||||
if data is None:
|
||||
continue
|
||||
|
||||
try:
|
||||
# 方案 C: 正規化舊格式 Enum 值
|
||||
incident_dict = json.loads(data)
|
||||
if "status" in incident_dict:
|
||||
incident_dict["status"] = normalize_status(incident_dict["status"])
|
||||
if "severity" in incident_dict:
|
||||
incident_dict["severity"] = normalize_severity(incident_dict["severity"])
|
||||
|
||||
for signal in incident_dict.get("signals", []):
|
||||
if "severity" in signal:
|
||||
signal["severity"] = normalize_severity(signal["severity"])
|
||||
|
||||
incident = Incident.model_validate(incident_dict)
|
||||
|
||||
if target_uuid in incident.proposal_ids:
|
||||
return incident
|
||||
|
||||
except Exception as e:
|
||||
logger.warning(
|
||||
"incident_parse_error_in_find",
|
||||
key=key,
|
||||
error=str(e),
|
||||
)
|
||||
continue
|
||||
|
||||
return None
|
||||
|
||||
except Exception as e:
|
||||
logger.exception(
|
||||
"find_by_proposal_id_error",
|
||||
proposal_id=proposal_id,
|
||||
error=str(e),
|
||||
)
|
||||
return None
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# Singleton
|
||||
|
||||
120
apps/api/src/services/signal_producer.py
Normal file
120
apps/api/src/services/signal_producer.py
Normal file
@@ -0,0 +1,120 @@
|
||||
"""
|
||||
Signal Producer Service - Phase 17 P0 Router 層違規修復
|
||||
======================================================
|
||||
|
||||
封裝 Redis Stream 的 Signal 生產邏輯,消除 Router 層直接存取 Redis。
|
||||
|
||||
功能:
|
||||
- XADD 寫入 Redis Stream
|
||||
- Trace Context 傳遞
|
||||
|
||||
符合 leWOOOgo 積木化規範:
|
||||
- Router -> Service -> Redis
|
||||
"""
|
||||
|
||||
from dataclasses import dataclass
|
||||
from typing import Any
|
||||
|
||||
import structlog
|
||||
|
||||
from src.core.redis_client import get_redis
|
||||
from src.core.telemetry import get_trace_context
|
||||
from src.utils.timezone import now_taipei
|
||||
|
||||
logger = structlog.get_logger(__name__)
|
||||
|
||||
# Redis Stream 配置
|
||||
SIGNAL_STREAM_KEY = "awoooi:signals"
|
||||
SIGNAL_STREAM_MAXLEN = 10000
|
||||
|
||||
|
||||
@dataclass
|
||||
class SignalData:
|
||||
"""Signal 資料結構"""
|
||||
source: str
|
||||
alert_name: str
|
||||
severity: str
|
||||
namespace: str
|
||||
target: str
|
||||
message: str
|
||||
labels: dict[str, Any] | None = None
|
||||
annotations: dict[str, Any] | None = None
|
||||
|
||||
|
||||
class SignalProducerService:
|
||||
"""
|
||||
Signal 生產者服務
|
||||
|
||||
封裝 Redis Stream 的 XADD 操作
|
||||
"""
|
||||
|
||||
async def produce(self, signal: SignalData) -> str:
|
||||
"""
|
||||
將 Signal 寫入 Redis Stream
|
||||
|
||||
Phase 17: 從 Router 層遷移至 Service 層
|
||||
|
||||
使用 Redis Streams (XADD) 實現非同步事件總線:
|
||||
- MAXLEN ~10000: 限制 Stream 長度,自動裁剪舊訊息
|
||||
- *: 自動生成 Message ID
|
||||
|
||||
Args:
|
||||
signal: Signal 資料
|
||||
|
||||
Returns:
|
||||
str: Redis Stream Message ID
|
||||
"""
|
||||
redis_client = get_redis()
|
||||
|
||||
# 組裝 Signal 字典 (所有值必須是字串)
|
||||
signal_dict = {
|
||||
"source": signal.source,
|
||||
"alert_name": signal.alert_name,
|
||||
"severity": signal.severity,
|
||||
"namespace": signal.namespace,
|
||||
"target": signal.target,
|
||||
"message": signal.message,
|
||||
"labels": str(signal.labels or {}),
|
||||
"annotations": str(signal.annotations or {}),
|
||||
"received_at": now_taipei().isoformat(),
|
||||
}
|
||||
|
||||
# Phase 15.2: 注入 Trace Context
|
||||
trace_ctx = get_trace_context()
|
||||
if trace_ctx:
|
||||
signal_dict["_trace_id"] = trace_ctx.get("trace_id", "")
|
||||
signal_dict["_span_id"] = trace_ctx.get("span_id", "")
|
||||
|
||||
# XADD 寫入 Stream
|
||||
message_id = await redis_client.xadd(
|
||||
SIGNAL_STREAM_KEY,
|
||||
signal_dict,
|
||||
maxlen=SIGNAL_STREAM_MAXLEN,
|
||||
approximate=True, # ~MAXLEN 近似裁剪,效能更好
|
||||
)
|
||||
|
||||
logger.info(
|
||||
"signal_produced",
|
||||
message_id=message_id,
|
||||
source=signal.source,
|
||||
alert_name=signal.alert_name,
|
||||
severity=signal.severity,
|
||||
trace_id=trace_ctx.get("trace_id") if trace_ctx else None,
|
||||
)
|
||||
|
||||
return message_id
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# Singleton
|
||||
# =============================================================================
|
||||
|
||||
_signal_producer: SignalProducerService | None = None
|
||||
|
||||
|
||||
def get_signal_producer() -> SignalProducerService:
|
||||
"""取得 SignalProducerService 實例 (Singleton)"""
|
||||
global _signal_producer
|
||||
if _signal_producer is None:
|
||||
_signal_producer = SignalProducerService()
|
||||
return _signal_producer
|
||||
109
apps/api/src/services/stats_service.py
Normal file
109
apps/api/src/services/stats_service.py
Normal file
@@ -0,0 +1,109 @@
|
||||
"""
|
||||
Stats Service - Phase 17 P0 Router 層違規修復
|
||||
=============================================
|
||||
|
||||
封裝統計 API 的快取邏輯,消除 Router 層直接存取 Redis。
|
||||
|
||||
功能:
|
||||
- 快取包裝器 (Redis)
|
||||
- 統計計算 (透過 Repository)
|
||||
|
||||
符合 leWOOOgo 積木化規範:
|
||||
- Router -> Service -> Redis/Repository
|
||||
"""
|
||||
|
||||
import json
|
||||
from typing import Any, Callable, Coroutine
|
||||
|
||||
import structlog
|
||||
|
||||
from src.core.redis_client import get_redis
|
||||
|
||||
logger = structlog.get_logger(__name__)
|
||||
|
||||
# 快取 TTL (秒)
|
||||
STATS_CACHE_TTL = 300 # 5 分鐘
|
||||
|
||||
|
||||
class StatsService:
|
||||
"""
|
||||
統計服務
|
||||
|
||||
封裝統計 API 的快取邏輯
|
||||
"""
|
||||
|
||||
async def get_cached_or_compute(
|
||||
self,
|
||||
cache_key: str,
|
||||
compute_fn: Callable[[], Coroutine[Any, Any, dict[str, Any]]],
|
||||
ttl: int = STATS_CACHE_TTL,
|
||||
) -> dict[str, Any]:
|
||||
"""
|
||||
快取包裝器: 先查 Redis,沒有則計算並快取
|
||||
|
||||
Phase 17: 從 Router 層遷移至 Service 層
|
||||
|
||||
Args:
|
||||
cache_key: Redis key
|
||||
compute_fn: 計算函數 (async callable)
|
||||
ttl: 快取時間 (秒)
|
||||
|
||||
Returns:
|
||||
快取或計算結果
|
||||
"""
|
||||
redis_client = get_redis()
|
||||
|
||||
# 嘗試從快取取得
|
||||
try:
|
||||
cached = await redis_client.get(cache_key)
|
||||
if cached:
|
||||
logger.debug("stats_cache_hit", key=cache_key)
|
||||
return json.loads(cached)
|
||||
except Exception as e:
|
||||
logger.warning("stats_cache_read_error", key=cache_key, error=str(e))
|
||||
|
||||
# 計算結果
|
||||
result = await compute_fn()
|
||||
|
||||
# 寫入快取
|
||||
try:
|
||||
await redis_client.set(cache_key, json.dumps(result), ex=ttl)
|
||||
logger.debug("stats_cache_set", key=cache_key, ttl=ttl)
|
||||
except Exception as e:
|
||||
logger.warning("stats_cache_write_error", key=cache_key, error=str(e))
|
||||
|
||||
return result
|
||||
|
||||
async def invalidate_cache(self, cache_key: str) -> bool:
|
||||
"""
|
||||
清除指定快取
|
||||
|
||||
Args:
|
||||
cache_key: Redis key
|
||||
|
||||
Returns:
|
||||
是否成功清除
|
||||
"""
|
||||
redis_client = get_redis()
|
||||
try:
|
||||
await redis_client.delete(cache_key)
|
||||
logger.info("stats_cache_invalidated", key=cache_key)
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.warning("stats_cache_invalidate_error", key=cache_key, error=str(e))
|
||||
return False
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# Singleton
|
||||
# =============================================================================
|
||||
|
||||
_stats_service: StatsService | None = None
|
||||
|
||||
|
||||
def get_stats_service() -> StatsService:
|
||||
"""取得 StatsService 實例 (Singleton)"""
|
||||
global _stats_service
|
||||
if _stats_service is None:
|
||||
_stats_service = StatsService()
|
||||
return _stats_service
|
||||
Reference in New Issue
Block a user