Introduce an admin portal (React + Nginx), WebSocket routing, and API versioning middleware with `/api/v1/` prefix deprecation. Add master API key authentication, new Prometheus metrics for AI token consumption and active WebSockets, and extend S3 config with a public endpoint URL. Update test paths and fixtures to align with the new routing structure.
72 lines
2.9 KiB
Python
72 lines
2.9 KiB
Python
"""
|
|
WebSocket client for receiving real-time pipeline events.
|
|
"""
|
|
import asyncio
|
|
import json
|
|
import logging
|
|
from typing import AsyncGenerator, Dict, Any, Optional
|
|
|
|
import websockets
|
|
from websockets.exceptions import ConnectionClosed
|
|
|
|
from .exceptions import AuthError, PipelineError
|
|
|
|
logger = logging.getLogger("imago_client.websocket")
|
|
|
|
class PipelineStream:
|
|
"""Stream of real-time pipeline events for an image using WebSockets."""
|
|
|
|
def __init__(self, ws_url: str, api_key: str, image_id: int):
|
|
self.ws_url = ws_url.replace("http://", "ws://").replace("https://", "wss://")
|
|
self.api_key = api_key
|
|
self.image_id = image_id
|
|
|
|
async def stream_events(self) -> AsyncGenerator[Dict[str, Any], None]:
|
|
"""Connects to the WebSocket and yields events."""
|
|
url = f"{self.ws_url}/pipeline/{self.image_id}?token={self.api_key}"
|
|
|
|
try:
|
|
async with websockets.connect(url) as websocket:
|
|
async for message in websocket:
|
|
try:
|
|
event = json.loads(message)
|
|
yield event
|
|
|
|
# Stop iterating if the pipeline is finished
|
|
if event.get("event") in ("pipeline.done", "pipeline.error"):
|
|
break
|
|
except json.JSONDecodeError:
|
|
logger.warning(f"Failed to decode message: {message}")
|
|
|
|
except ConnectionClosed as e:
|
|
if e.code == 4001:
|
|
raise AuthError("No authentication token provided for WebSocket") from e
|
|
elif e.code == 4003:
|
|
raise AuthError("Forbidden or unauthorized to access this pipeline stream") from e
|
|
else:
|
|
raise PipelineError(f"WebSocket connection closed unexpectedly: {e.code} - {e.reason}") from e
|
|
except Exception as e:
|
|
raise PipelineError(f"WebSocket error: {e}") from e
|
|
|
|
class AdminMonitorStream:
|
|
"""Admin stream for listening to all pipeline events."""
|
|
|
|
def __init__(self, ws_url: str, api_key: str):
|
|
self.ws_url = ws_url.replace("http://", "ws://").replace("https://", "wss://")
|
|
self.api_key = api_key
|
|
|
|
async def stream_events(self) -> AsyncGenerator[Dict[str, Any], None]:
|
|
url = f"{self.ws_url}/admin/monitor?token={self.api_key}"
|
|
try:
|
|
async with websockets.connect(url) as websocket:
|
|
async for message in websocket:
|
|
try:
|
|
event = json.loads(message)
|
|
yield event
|
|
except json.JSONDecodeError:
|
|
logger.warning(f"Failed to decode message: {message}")
|
|
except ConnectionClosed as e:
|
|
if e.code in (4001, 4003):
|
|
raise AuthError("Unauthorized to access admin stream") from e
|
|
raise PipelineError(f"Connection closed: {e.code} - {e.reason}") from e
|