Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 40 additions & 0 deletions backend/src/control_center/analytics/aggregator.py
Original file line number Diff line number Diff line change
Expand Up @@ -119,6 +119,18 @@ def apply_interaction_event(event: AnalyticsEvent) -> bool:
pipe.expire(key, AGG_TTL_SECONDS)
pipe.execute()

if event.org_id is not None:
# A small, low-cardinality index of "which services has this org
# actually seen traffic from" -- read_known_services below is
# what /analytics/services (a later PR) enumerates, since there
# is no other way to discover which analytics:agg:{org}:{service}:
# {date} keys exist without this index (Redis has no efficient
# "list keys matching a pattern" primitive worth using here).
# No TTL: a service name, once seen for an org, stays a legitimate
# thing to report zero-for on a quiet day, same reasoning
# AGG_TTL_SECONDS's own 90-day retention already applies to.
_redis.sadd(f"analytics:services:{event.org_id}", event.service)

if event.user_id is not None:
active_keys = [f"analytics:active_users:{date_str}"]
if event.org_id is not None:
Expand Down Expand Up @@ -271,6 +283,34 @@ def read_daily_active_users(date_strs: list[str], org_id: Optional[int] = None)
return [{"date": d, "count": read_active_user_count([d], org_id=org_id)} for d in date_strs]


def read_active_user_ids(date_strs: list[str], org_id: Optional[int] = None) -> set[str]:
"""Raw user_id set (union across date_strs) -- INTERNAL ONLY. Never
returned directly by any API response (task brief: "do not expose
raw user IDs through normal analytics endpoints") -- the one
legitimate caller is a later PR's team-roster intersection, which
only ever surfaces a resulting *count*, never this set itself."""
prefix = f"analytics:active_users:{org_id}:" if org_id is not None else "analytics:active_users:"
keys = [f"{prefix}{d}" for d in date_strs]
if not keys:
return set()
try:
if len(keys) == 1:
return set(_redis.smembers(keys[0]))
return set(_redis.sunion(*keys))
except Exception:
return set()


def read_known_services(org_id: int) -> list[str]:
"""The service names apply_interaction_event has ever recorded
traffic from for this org -- see that function's own comment on the
`analytics:services:{org_id}` index this reads."""
try:
return sorted(_redis.smembers(f"analytics:services:{org_id}"))
except Exception:
return []


def read_user_activity(org_id: int, date_str: str) -> dict[str, int]:
"""user_id -> query_count for one org/day -- never returned directly
by any API response (task brief: "do not expose raw user IDs"); only
Expand Down
39 changes: 38 additions & 1 deletion backend/src/control_center/analytics/cache.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@

import json
import os
from typing import Any, Callable
from typing import Any, Awaitable, Callable

import redis

Expand Down Expand Up @@ -65,6 +65,43 @@ def get_or_set(
return value


async def get_or_set_async(
key: str,
endpoint: str,
compute_coro: Callable[[], Awaitable[Any]],
ttl: int = DEFAULT_TTL_SECONDS,
) -> Any:
"""Async twin of `get_or_set` -- PR-C's service layer talks to other
upstreams (billing/TES/team-roster) via `httpx.AsyncClient` from
inside async FastAPI route handlers, so the compute side has to be
awaitable. Redis itself is still touched synchronously (redis-py's
sync client, same as the rest of this package) -- fine here since
those calls are fast, in-network, and already wrapped in their own
try/except; only `compute_coro` needs to be awaited.
"""
from control_center.analytics.metrics import CACHE_HITS, CACHE_MISSES

try:
cached = _redis.get(key)
except Exception:
cached = None

if cached is not None:
CACHE_HITS.labels(endpoint=endpoint).inc()
try:
return json.loads(cached)
except (TypeError, ValueError):
pass

CACHE_MISSES.labels(endpoint=endpoint).inc()
value = await compute_coro()
try:
_redis.setex(key, ttl, json.dumps(value, default=str))
except Exception:
pass
return value


def invalidate(key: str) -> None:
"""Best-effort explicit invalidation. Not required for correctness
(every cache entry has a TTL and naturally expires -- Section 9's own
Expand Down
189 changes: 189 additions & 0 deletions backend/src/control_center/analytics/router.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,189 @@
"""The /analytics/* HTTP surface (task brief Section 7). Every route:
1. resolves+enforces scope via `Depends(require_analytics_scope)` (401/403
handled entirely by that dependency -- no authorization decision is
made in this file);
2. resolves the date range (`service.resolve_date_range`, default last
30 days);
3. reads-through a 5-minute Redis cache (`cache.get_or_set_async`) keyed
by endpoint+scope+range, never computing twice for the same request
shape within the TTL;
4. delegates all computation to service.py -- no route function contains
aggregation logic itself.

No new authentication path: the same `require_analytics_scope`
dependency (built on `core.jwt_verify.verify_token`, the one JWT
verifier this whole service already uses) gates every route here.
"""
from __future__ import annotations

import csv
import io
from datetime import date
from typing import Any, Callable, Coroutine, Optional

from fastapi import APIRouter, Depends, Header, Query
from fastapi.responses import JSONResponse, StreamingResponse

from control_center.analytics import cache, service
from control_center.analytics.metrics import API_ERRORS, API_REQUESTS
from control_center.analytics.permissions import AnalyticsScope, require_analytics_scope

router = APIRouter(prefix="/analytics", tags=["analytics"])


def _cache_key(endpoint: str, scope: AnalyticsScope, from_date: date, to_date: date) -> str:
return f"analytics:{endpoint}:{scope.org_id}:{scope.team_id}:{from_date.isoformat()}:{to_date.isoformat()}"


async def _run(endpoint: str, cache_key: str, compute_coro: Callable[[], Coroutine[Any, Any, dict]]) -> JSONResponse:
API_REQUESTS.labels(endpoint=endpoint).inc()
try:
data = await cache.get_or_set_async(cache_key, endpoint, compute_coro, ttl=cache.DEFAULT_TTL_SECONDS)
except Exception:
API_ERRORS.labels(endpoint=endpoint, status_code="500").inc()
raise
return JSONResponse(data)


@router.get("/overview")
async def overview(
from_date: Optional[date] = Query(default=None),
to_date: Optional[date] = Query(default=None),
authorization: Optional[str] = Header(default=None),
scope: AnalyticsScope = Depends(require_analytics_scope),
) -> JSONResponse:
resolved_from, resolved_to = service.resolve_date_range(from_date, to_date)
key = _cache_key("overview", scope, resolved_from, resolved_to)
return await _run("overview", key, lambda: service.get_overview(scope, resolved_from, resolved_to, authorization))


@router.get("/queries")
async def queries(
from_date: Optional[date] = Query(default=None),
to_date: Optional[date] = Query(default=None),
authorization: Optional[str] = Header(default=None),
scope: AnalyticsScope = Depends(require_analytics_scope),
) -> JSONResponse:
resolved_from, resolved_to = service.resolve_date_range(from_date, to_date)
key = _cache_key("queries", scope, resolved_from, resolved_to)
return await _run("queries", key, lambda: service.get_queries(scope, resolved_from, resolved_to, authorization))


@router.get("/users")
async def users(
from_date: Optional[date] = Query(default=None),
to_date: Optional[date] = Query(default=None),
authorization: Optional[str] = Header(default=None),
scope: AnalyticsScope = Depends(require_analytics_scope),
) -> JSONResponse:
resolved_from, resolved_to = service.resolve_date_range(from_date, to_date)
key = _cache_key("users", scope, resolved_from, resolved_to)
return await _run("users", key, lambda: service.get_users(scope, resolved_from, resolved_to, authorization))


@router.get("/services")
async def services(
from_date: Optional[date] = Query(default=None),
to_date: Optional[date] = Query(default=None),
scope: AnalyticsScope = Depends(require_analytics_scope),
) -> JSONResponse:
resolved_from, resolved_to = service.resolve_date_range(from_date, to_date)
key = _cache_key("services", scope, resolved_from, resolved_to)
return await _run("services", key, lambda: service.get_services(scope, resolved_from, resolved_to))


@router.get("/workflows")
async def workflows(
from_date: Optional[date] = Query(default=None),
to_date: Optional[date] = Query(default=None),
authorization: Optional[str] = Header(default=None),
scope: AnalyticsScope = Depends(require_analytics_scope),
) -> JSONResponse:
resolved_from, resolved_to = service.resolve_date_range(from_date, to_date)
key = _cache_key("workflows", scope, resolved_from, resolved_to)
return await _run("workflows", key, lambda: service.get_workflows(scope, resolved_from, resolved_to, authorization))


@router.get("/performance")
async def performance(
from_date: Optional[date] = Query(default=None),
to_date: Optional[date] = Query(default=None),
# Still gated by require_analytics_scope (any admin role -- 403 for a
# regular user) even though the response itself is platform-wide and
# ignores scope.org_id/team_id entirely -- an authenticated admin of
# any organization may see platform-wide operational health, but a
# regular user still may not reach analytics at all.
_scope: AnalyticsScope = Depends(require_analytics_scope),
) -> JSONResponse:
resolved_from, resolved_to = service.resolve_date_range(from_date, to_date)
key = f"analytics:performance:platform:{resolved_from.isoformat()}:{resolved_to.isoformat()}"
return await _run("performance", key, lambda: service.get_performance(resolved_from, resolved_to))


@router.get("/usage")
async def usage(
authorization: Optional[str] = Header(default=None),
scope: AnalyticsScope = Depends(require_analytics_scope),
) -> JSONResponse:
key = f"analytics:usage:{scope.org_id}:{scope.team_id}"
return await _run("usage", key, lambda: service.get_usage(scope, authorization))


_EXPORTABLE: dict[str, Callable[[AnalyticsScope, date, date, Optional[str]], Coroutine[Any, Any, dict]]] = {
"overview": lambda scope, f, t, auth: service.get_overview(scope, f, t, auth),
"queries": lambda scope, f, t, auth: service.get_queries(scope, f, t, auth),
"users": lambda scope, f, t, auth: service.get_users(scope, f, t, auth),
"services": lambda scope, f, t, auth: service.get_services(scope, f, t),
"workflows": lambda scope, f, t, auth: service.get_workflows(scope, f, t, auth),
"performance": lambda scope, f, t, auth: service.get_performance(f, t),
}


def _rows_for_export(export_type: str, data: dict) -> tuple[list[str], list[list[Any]]]:
"""Turns one service.py response dict into (header, rows) for CSV --
the "daily trend" shape (queries/users) becomes one row per date, the
"services" shape becomes one row per service, everything else
(overview/workflows/performance) becomes one summary row."""
if export_type in ("queries", "users") and isinstance(data.get("daily"), list):
header = ["date", "count"]
rows = [[row.get("date"), row.get("count")] for row in data["daily"]]
return header, rows
if export_type == "services":
header = ["service", "total_calls", "errors", "error_rate", "avg_latency_ms"]
rows = [
[r.get("service"), r.get("total_calls"), r.get("errors"), r.get("error_rate"), r.get("avg_latency_ms")]
for r in data.get("services", [])
]
return header, rows
scalar_items = [(k, v) for k, v in data.items() if not isinstance(v, (list, dict))]
return [k for k, _ in scalar_items], [[v for _, v in scalar_items]]


@router.get("/export")
async def export(
type: str = Query(...),
from_date: Optional[date] = Query(default=None),
to_date: Optional[date] = Query(default=None),
authorization: Optional[str] = Header(default=None),
scope: AnalyticsScope = Depends(require_analytics_scope),
) -> StreamingResponse:
builder = _EXPORTABLE.get(type)
if builder is None:
return JSONResponse({"error": f"unknown export type: {type!r}"}, status_code=400)

resolved_from, resolved_to = service.resolve_date_range(from_date, to_date)
API_REQUESTS.labels(endpoint="export").inc()
data = await builder(scope, resolved_from, resolved_to, authorization)

header, rows = _rows_for_export(type, data)
buffer = io.StringIO()
writer = csv.writer(buffer)
writer.writerow(header)
writer.writerows(rows)
buffer.seek(0)

return StreamingResponse(
iter([buffer.getvalue()]),
media_type="text/csv",
headers={"Content-Disposition": f'attachment; filename="analytics_{type}.csv"'},
)
Loading
Loading