feat(#15): Approval Polling → SSE 即時更新
Phase 15: 將 Approval 輪詢機制改為 Server-Sent Events 後端變更: - 新增 /api/v1/approvals/stream SSE 端點 - 建立/簽核/拒絕時發布 SSE 事件 - 使用現有 EventPublisher 基礎設施 前端變更: - 新增 useApprovalSSE hook (自動連線/斷線管理) - approval.store 新增 connectSSE/disconnectSSE actions - 更新三個組件使用 SSE 取代 setInterval polling: - LiveApprovalPanel - AICommandPanel - HITLSection 效益: - 即時推送 (延遲 ~0ms vs polling 5s) - 減少 API 請求 (僅變更時推送) - 自動重連 + Fallback to polling Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
@@ -25,13 +25,23 @@ import re
|
||||
from typing import TYPE_CHECKING
|
||||
from uuid import UUID
|
||||
|
||||
from fastapi import APIRouter, BackgroundTasks, Depends, Header, HTTPException, status
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from src.services.notifications import ExecutionStatus
|
||||
|
||||
from fastapi import (
|
||||
APIRouter,
|
||||
BackgroundTasks,
|
||||
Depends,
|
||||
Header,
|
||||
HTTPException,
|
||||
Request,
|
||||
status,
|
||||
)
|
||||
from fastapi.responses import StreamingResponse
|
||||
|
||||
from src.core.config import settings
|
||||
from src.core.logging import get_logger
|
||||
from src.core.sse import EventType, SSEEvent, get_publisher
|
||||
from src.models.approval import (
|
||||
ApprovalRequest,
|
||||
ApprovalRequestCreate,
|
||||
@@ -521,6 +531,9 @@ async def create_approval(
|
||||
required_signatures=approval.required_signatures,
|
||||
)
|
||||
|
||||
# SSE: 發布建立事件
|
||||
asyncio.create_task(_publish_approval_event("created", approval))
|
||||
|
||||
return ApprovalRequestResponse.from_approval(approval)
|
||||
|
||||
|
||||
@@ -667,6 +680,10 @@ async def sign_approval(
|
||||
note="Approval has no incident_id in metadata, cannot update Incident status",
|
||||
)
|
||||
|
||||
# SSE: 發布簽核/核准事件
|
||||
event_action = "approved" if execution_triggered else "signed"
|
||||
asyncio.create_task(_publish_approval_event(event_action, approval))
|
||||
|
||||
return SignResponse(
|
||||
success=True,
|
||||
message=message,
|
||||
@@ -752,4 +769,129 @@ async def reject_approval(
|
||||
reason=request.reason,
|
||||
)
|
||||
|
||||
# SSE: 發布拒絕事件
|
||||
asyncio.create_task(_publish_approval_event("rejected", approval))
|
||||
|
||||
return ApprovalRequestResponse.from_approval(approval)
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# SSE Event Publishing (Phase 15: Polling → SSE)
|
||||
# =============================================================================
|
||||
|
||||
async def _publish_approval_event(
|
||||
action: str,
|
||||
approval: ApprovalRequest,
|
||||
) -> None:
|
||||
"""
|
||||
發布 Approval SSE 事件
|
||||
|
||||
Args:
|
||||
action: 事件動作 (created, signed, approved, rejected, executed)
|
||||
approval: 授權請求物件
|
||||
"""
|
||||
try:
|
||||
pub = await get_publisher()
|
||||
|
||||
event = SSEEvent(
|
||||
type=EventType.APPROVAL,
|
||||
data={
|
||||
"action": action,
|
||||
"approval_id": str(approval.id),
|
||||
"status": approval.status.value,
|
||||
"risk_level": approval.risk_level.value,
|
||||
"current_signatures": approval.current_signatures,
|
||||
"required_signatures": approval.required_signatures,
|
||||
"action_text": approval.action[:100],
|
||||
},
|
||||
)
|
||||
|
||||
sent_count = await pub.publish(event, topic="approvals")
|
||||
|
||||
if sent_count > 0:
|
||||
logger.debug(
|
||||
"approval_sse_published",
|
||||
action=action,
|
||||
approval_id=str(approval.id),
|
||||
sent_count=sent_count,
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(
|
||||
"approval_sse_publish_failed",
|
||||
action=action,
|
||||
approval_id=str(approval.id),
|
||||
error=str(e),
|
||||
)
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# GET /api/v1/approvals/stream (SSE)
|
||||
# =============================================================================
|
||||
|
||||
@router.get(
|
||||
"/stream",
|
||||
summary="SSE 即時授權更新",
|
||||
description="Server-Sent Events 即時推送授權狀態變更,取代輪詢機制",
|
||||
)
|
||||
async def stream_approvals(request: Request) -> StreamingResponse:
|
||||
"""
|
||||
SSE 即時授權更新 (Phase 15: Polling → SSE)
|
||||
|
||||
事件類型:
|
||||
- connected: 連線建立
|
||||
- heartbeat: 心跳 (每 15 秒)
|
||||
- approval: 授權狀態變更 (created, signed, approved, rejected, executed)
|
||||
|
||||
Client Usage (JavaScript):
|
||||
```javascript
|
||||
const es = new EventSource('/api/v1/approvals/stream');
|
||||
es.addEventListener('approval', (e) => {
|
||||
const data = JSON.parse(e.data);
|
||||
console.log('Approval update:', data.action, data.approval_id);
|
||||
// Refresh approval list or update specific item
|
||||
});
|
||||
es.addEventListener('heartbeat', () => {
|
||||
console.log('Connection alive');
|
||||
});
|
||||
```
|
||||
"""
|
||||
logger.info(
|
||||
"approval_stream_connect",
|
||||
client_ip=request.client.host if request.client else "unknown",
|
||||
)
|
||||
|
||||
pub = await get_publisher()
|
||||
|
||||
# 訂閱 approvals topic
|
||||
client = await pub.subscribe(
|
||||
topics=["approvals"],
|
||||
metadata={"ip": request.client.host if request.client else "unknown"},
|
||||
)
|
||||
|
||||
async def event_generator():
|
||||
"""SSE 事件生成器,含斷線偵測"""
|
||||
try:
|
||||
async for data in pub.stream(client):
|
||||
if await request.is_disconnected():
|
||||
logger.info("approval_stream_disconnected", client_id=client.id)
|
||||
break
|
||||
yield data
|
||||
|
||||
except asyncio.CancelledError:
|
||||
logger.info("approval_stream_cancelled", client_id=client.id)
|
||||
raise
|
||||
|
||||
finally:
|
||||
logger.info("approval_stream_cleanup", client_id=client.id)
|
||||
|
||||
return StreamingResponse(
|
||||
event_generator(),
|
||||
media_type="text/event-stream",
|
||||
headers={
|
||||
"Cache-Control": "no-cache, no-store, must-revalidate",
|
||||
"Connection": "keep-alive",
|
||||
"X-Accel-Buffering": "no",
|
||||
"Access-Control-Allow-Origin": "*",
|
||||
},
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user