SPB Git forge
28commits 1branches 0releases
7.7 MBsize
maindefault branch
10 days agolast push
Python 66.3% TypeScript 22.7% JavaScript 8.6% HTML 1.4% CSS 0.7%
3.3 KB · 67 lines python
Raw Blame History
1"""Cross-company intelligence: `/signals`, `/trends`, `/map`."""2from __future__ import annotations34from typing import Any56from fastapi import APIRouter, Query, Request, Response78from companyatlas.api import aggregates as agg9from companyatlas.api import serializers as ser10from companyatlas.api.common import cached, cached_response, public_cache_value11from companyatlas.db import connection, fetch_all1213ORDER = 2014router = APIRouter(prefix="/api/v1", tags=["signals"])15TREND_WINDOW_DAYS = {"7d": 7, "30d": 30, "90d": 90}161718@router.get("/signals", summary="Active signals (labelled as signals, never facts)")19async def signals(response: Response, kind: str | None = Query(None, max_length=60), scope: str | None = Query(None, pattern="^(company|industry|country|global)$"),20                  scope_key: str | None = Query(None, max_length=80), company: str | None = Query(None, max_length=200),21                  min_strength: float | None = Query(None, ge=0, le=1), limit: int = Query(50, ge=1, le=500)) -> dict[str, Any]:22    response.headers["cache-control"] = public_cache_value(60)23    where = ["s.status = 'active'"]24    params: dict[str, Any] = {"limit": limit}25    if kind:26        where.append("s.kind = cast(:kind as text)")27        params["kind"] = kind28    if scope:29        where.append("s.scope = cast(:scope as text)")30        params["scope"] = scope31    if scope_key:32        where.append("s.scope_key = cast(:scope_key as text)")33        params["scope_key"] = scope_key34    if company:35        where.append("(c.slug = cast(:company as text) or c.id = cast(:company as text))")36        params["company"] = company37    if min_strength is not None:38        where.append("s.strength >= :min_strength")39        params["min_strength"] = min_strength40    async with connection() as conn:41        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, "42                                     "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 "43                                     f"where {' and '.join(where)} order by s.detected_at desc, s.strength desc limit :limit", **params)44    items = []45    for r in rows:46        s = ser.signal(r)47        s["company"] = ser.company_ref(r) if r.get("company_slug") else None48        items.append(s)49    return {"items": items}505152@router.get("/trends", summary="Trending terms with momentum (cached 300 s)")53async def trends(request: Request, window: str = Query("7d", pattern="^(7d|30d|90d)$"), limit: int = Query(30, ge=1, le=200)) -> Any:54    async def produce() -> dict[str, Any]:55        async with connection() as conn:56            return {"window": window, "items": await agg.trend_rows(conn, TREND_WINDOW_DAYS[window], limit)}57    return cached_response(request, await cached(f"trends:{window}:{limit}", 300, produce), 300)585960@router.get("/map", summary="Clustered map buckets (≤ 600, cached 300 s)")61async def map_(request: Request, metric: str = Query("events_30d", pattern="^(events_30d|companies|hiring)$")) -> Any:62    async def produce() -> dict[str, Any]:63        async with connection() as conn:64            buckets = await agg.map_buckets(conn, metric)65        return {"metric": metric, "buckets": buckets}66    return cached_response(request, await cached(f"map:{metric}", 300, produce), 300)67