"""Cross-company intelligence: `/signals`, `/trends`, `/map`.""" from __future__ import annotations from typing import Any from fastapi import APIRouter, Query, Request, Response from companyatlas.api import aggregates as agg from companyatlas.api import serializers as ser from companyatlas.api.common import cached, cached_response, public_cache_value from companyatlas.db import connection, fetch_all ORDER = 20 router = APIRouter(prefix="/api/v1", tags=["signals"]) TREND_WINDOW_DAYS = {"7d": 7, "30d": 30, "90d": 90} @router.get("/signals", summary="Active signals (labelled as signals, never facts)") async def signals(response: Response, kind: str | None = Query(None, max_length=60), scope: str | None = Query(None, pattern="^(company|industry|country|global)$"), scope_key: str | None = Query(None, max_length=80), company: str | None = Query(None, max_length=200), min_strength: float | None = Query(None, ge=0, le=1), limit: int = Query(50, ge=1, le=500)) -> dict[str, Any]: response.headers["cache-control"] = public_cache_value(60) where = ["s.status = 'active'"] params: dict[str, Any] = {"limit": limit} if kind: where.append("s.kind = cast(:kind as text)") params["kind"] = kind if scope: where.append("s.scope = cast(:scope as text)") params["scope"] = scope if scope_key: where.append("s.scope_key = cast(:scope_key as text)") params["scope_key"] = scope_key if company: where.append("(c.slug = cast(:company as text) or c.id = cast(:company as text))") params["company"] = company if min_strength is not None: where.append("s.strength >= :min_strength") params["min_strength"] = min_strength async with connection() as conn: rows = await fetch_all(conn, "select s.*, c.slug as company_slug, c.display_name as company_display_name, c.canonical_domain as company_domain, " "c.country as company_country, c.logo_url as company_logo_url from signals s left join companies c on c.id = s.company_id " f"where {' and '.join(where)} order by s.detected_at desc, s.strength desc limit :limit", **params) items = [] for r in rows: s = ser.signal(r) s["company"] = ser.company_ref(r) if r.get("company_slug") else None items.append(s) return {"items": items} @router.get("/trends", summary="Trending terms with momentum (cached 300 s)") async def trends(request: Request, window: str = Query("7d", pattern="^(7d|30d|90d)$"), limit: int = Query(30, ge=1, le=200)) -> Any: async def produce() -> dict[str, Any]: async with connection() as conn: return {"window": window, "items": await agg.trend_rows(conn, TREND_WINDOW_DAYS[window], limit)} return cached_response(request, await cached(f"trends:{window}:{limit}", 300, produce), 300) @router.get("/map", summary="Clustered map buckets (≤ 600, cached 300 s)") async def map_(request: Request, metric: str = Query("events_30d", pattern="^(events_30d|companies|hiring)$")) -> Any: async def produce() -> dict[str, Any]: async with connection() as conn: buckets = await agg.map_buckets(conn, metric) return {"metric": metric, "buckets": buckets} return cached_response(request, await cached(f"map:{metric}", 300, produce), 300)