"""Live feed: `/live` (latest active events, never cached) and `/live/stream` (SSE).""" from __future__ import annotations from typing import Any from fastapi import APIRouter, Query, Request, Response from companyatlas.api import queries as q from companyatlas.api import serializers as ser from companyatlas.api.common import NO_STORE from companyatlas.api.sse import MAX_STREAM_S, live_event_stream, sse_response from companyatlas.db import connection ORDER = 10 router = APIRouter(prefix="/api/v1", tags=["live"]) @router.get("/live", summary="Latest active events") async def live(response: Response, limit: int = Query(50, ge=1, le=200), since: str | None = None, event_type: str | None = None, min_importance: float | None = Query(None, ge=0, le=1), country: str | None = None, industry: str | None = None) -> dict[str, Any]: response.headers["cache-control"] = NO_STORE since_dt = q.parse_iso(since) where, params = q.event_filters(event_type=event_type, min_importance=min_importance, country=country, industry=industry) if since_dt is not None: where.append("e.detected_at > :live_since") params["live_since"] = since_dt async with connection() as conn: rows = await q.fetch_events(conn, where, params, sort="recent", limit=limit) items = [ser.event(r) for r in rows] return {"items": items, "count": len(items), "cursor": items[0]["detected_at"] if items else (since_dt or q.now_utc()), "server_time": q.now_utc()} @router.get("/live/stream", summary="Server-sent events stream of new events") async def live_stream(request: Request, since: str | None = None, event_type: str | None = None, min_importance: float | None = Query(None, ge=0, le=1), max_s: float = Query(MAX_STREAM_S, ge=1, le=MAX_STREAM_S)) -> Any: since_dt = q.parse_iso(since) return sse_response(live_event_stream(request, since=since_dt, event_type=event_type, min_importance=min_importance, max_s=max_s))