"""SSE event endpoints for real-time notifications.""" import asyncio import json from typing import AsyncGenerator from fastapi import APIRouter from sse_starlette.sse import EventSourceResponse router = APIRouter() # Global event queue for broadcasting events event_queues: list[asyncio.Queue] = [] class EventBroadcaster: """Broadcasts events to all connected SSE clients.""" @staticmethod async def broadcast(event_type: str, data: dict): """Broadcast an event to all connected clients. Args: event_type: Type of event (e.g., "update_available", "backup_complete") data: Event data to send """ message = { "type": event_type, "data": data } # Send to all connected clients dead_queues = [] for queue in event_queues: try: await queue.put(message) except Exception: # Queue is dead, mark for removal dead_queues.append(queue) # Remove dead queues for dead_queue in dead_queues: event_queues.remove(dead_queue) broadcaster = EventBroadcaster() async def event_generator() -> AsyncGenerator[dict, None]: """Generate SSE events for a client connection.""" # Create a new queue for this client queue = asyncio.Queue() event_queues.append(queue) try: # Send initial connection event yield { "event": "connected", "data": json.dumps({"message": "Connected to Dangerous Pi event stream"}) } # Keep connection alive and send events while True: try: # Wait for events with timeout to send keep-alive message = await asyncio.wait_for(queue.get(), timeout=30.0) yield { "event": message["type"], "data": json.dumps(message["data"]) } except asyncio.TimeoutError: # Send keep-alive comment yield { "comment": "keep-alive" } except asyncio.CancelledError: # Client disconnected pass finally: # Remove queue when client disconnects if queue in event_queues: event_queues.remove(queue) @router.get("/events") async def stream_events(): """SSE endpoint for streaming events to clients. Events include: - update_available: New version available on GitHub - update_downloading: Update download in progress - update_complete: Update downloaded and ready to install - backup_complete: Backup operation completed - ups_low_battery: UPS battery below threshold - pm3_rebuild_required: PM3 client rebuild needed - command_complete: Long-running command completed """ return EventSourceResponse(event_generator()) # Helper functions for sending specific event types async def notify_update_available(version: str, url: str): """Notify clients that an update is available.""" await broadcaster.broadcast("update_available", { "version": version, "url": url }) async def notify_update_progress(progress: float, status: str): """Notify clients of update download progress.""" await broadcaster.broadcast("update_downloading", { "progress": progress, "status": status }) async def notify_backup_complete(backup_path: str, size_bytes: int): """Notify clients that backup is complete.""" await broadcaster.broadcast("backup_complete", { "path": backup_path, "size": size_bytes }) async def notify_ups_battery(percentage: float, voltage: float): """Notify clients of UPS battery status.""" await broadcaster.broadcast("ups_battery", { "percentage": percentage, "voltage": voltage }) async def notify_ups_warning(percentage: float, threshold: float): """Notify clients of low battery warning.""" await broadcaster.broadcast("ups_warning", { "percentage": percentage, "threshold": threshold, "message": f"Battery low: {percentage}%" }) async def notify_ups_critical(percentage: float): """Notify clients of critical battery level.""" await broadcaster.broadcast("ups_critical", { "percentage": percentage, "message": f"Battery critical: {percentage}%" }) async def notify_ups_shutdown(delay: int, percentage: float): """Notify clients that shutdown has been initiated.""" await broadcaster.broadcast("ups_shutdown", { "delay": delay, "percentage": percentage, "message": f"Shutdown initiated due to low battery ({percentage}%)" }) async def notify_pm3_status(connected: bool, message: str): """Notify clients of PM3 status changes.""" await broadcaster.broadcast("pm3_status", { "connected": connected, "message": message })