Address review blockers and medium/low findings: - F1: HTML-escape all dynamic analytics fields (html.escape quote=True) - F2: Gate POST /api/v1/analytics/usage through console_authz record_analytics_usage (fail closed for unauthenticated / Phase 1) - F3: Enforce usage_events retention (max rows + max age) - F4: Coerce None remote/org/repo to empty strings in load_analytics Add tests for XSS escaping, unauthorized ingest deny, retention, and None scope coercion. Document authz + retention in analytics guide. Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
434 lines
14 KiB
Python
434 lines
14 KiB
Python
"""Model usage, token cost, latency, and workflow-performance analytics (#651, Phase 4).
|
|
|
|
Ingests session instrumentation metrics, aggregates usage/cost/latency percentiles
|
|
by project, role, model, issue/PR, and stage, enforcing secret redaction and
|
|
explicitly rendering missing metrics as "Unknown" without zero-fabrication.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import math
|
|
from dataclasses import asdict, dataclass
|
|
from typing import Any, Sequence
|
|
|
|
import control_plane_db
|
|
from webui import console_redaction
|
|
|
|
ANALYTICS_SCHEMA_VERSION = 1
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class UsageEvent:
|
|
usage_id: int
|
|
session_id: str | None
|
|
remote: str
|
|
org: str
|
|
repo: str
|
|
project_id: str | None
|
|
role: str
|
|
model: str
|
|
issue_number: int | None
|
|
pr_number: int | None
|
|
stage: str
|
|
input_tokens: int | None
|
|
output_tokens: int | None
|
|
total_tokens: int | None
|
|
estimated_cost_usd: float | None
|
|
latency_ms: int | None
|
|
duration_ms: int | None
|
|
status: str
|
|
metadata: str | None
|
|
created_at: str
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
d = asdict(self)
|
|
if d["metadata"]:
|
|
d["metadata"] = console_redaction.redact_text(d["metadata"])
|
|
return d
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class GroupMetrics:
|
|
name: str
|
|
total_events: int
|
|
events_with_tokens: int
|
|
input_tokens: int | None
|
|
output_tokens: int | None
|
|
total_tokens: int | None
|
|
events_with_cost: int
|
|
estimated_cost_usd: float | None
|
|
events_with_latency: int
|
|
latency_p50_ms: float | None
|
|
latency_p90_ms: float | None
|
|
latency_p95_ms: float | None
|
|
latency_p99_ms: float | None
|
|
latency_avg_ms: float | None
|
|
events_with_duration: int
|
|
duration_avg_ms: float | None
|
|
display_tokens: str
|
|
display_cost: str
|
|
display_latency_p50: str
|
|
display_latency_p90: str
|
|
display_duration_avg: str
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return asdict(self)
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class AnalyticsSnapshot:
|
|
ok: bool
|
|
reason: str
|
|
schema_version: int
|
|
remote: str
|
|
org: str
|
|
repo: str
|
|
total_events: int
|
|
overall_summary: GroupMetrics
|
|
by_project: dict[str, GroupMetrics]
|
|
by_role: dict[str, GroupMetrics]
|
|
by_model: dict[str, GroupMetrics]
|
|
by_work_item: dict[str, GroupMetrics]
|
|
by_stage: dict[str, GroupMetrics]
|
|
events: tuple[UsageEvent, ...]
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return {
|
|
"ok": self.ok,
|
|
"reason": self.reason,
|
|
"schema_version": self.schema_version,
|
|
"remote": self.remote,
|
|
"org": self.org,
|
|
"repo": self.repo,
|
|
"total_events": self.total_events,
|
|
"overall_summary": self.overall_summary.to_dict(),
|
|
"by_project": {k: v.to_dict() for k, v in self.by_project.items()},
|
|
"by_role": {k: v.to_dict() for k, v in self.by_role.items()},
|
|
"by_model": {k: v.to_dict() for k, v in self.by_model.items()},
|
|
"by_work_item": {k: v.to_dict() for k, v in self.by_work_item.items()},
|
|
"by_stage": {k: v.to_dict() for k, v in self.by_stage.items()},
|
|
"events": [e.to_dict() for e in self.events],
|
|
}
|
|
|
|
|
|
def compute_percentile(values: Sequence[float | int], percentile: float) -> float | None:
|
|
if not values:
|
|
return None
|
|
sorted_vals = sorted(values)
|
|
n = len(sorted_vals)
|
|
if n == 1:
|
|
return float(sorted_vals[0])
|
|
k = (n - 1) * (percentile / 100.0)
|
|
f = math.floor(k)
|
|
c = math.ceil(k)
|
|
if f == c:
|
|
return float(sorted_vals[int(f)])
|
|
d0 = sorted_vals[int(f)] * (c - k)
|
|
d1 = sorted_vals[int(c)] * (k - f)
|
|
return float(d0 + d1)
|
|
|
|
|
|
def aggregate_events(group_name: str, events: Sequence[UsageEvent]) -> GroupMetrics:
|
|
total_events = len(events)
|
|
if total_events == 0:
|
|
return GroupMetrics(
|
|
name=group_name,
|
|
total_events=0,
|
|
events_with_tokens=0,
|
|
input_tokens=None,
|
|
output_tokens=None,
|
|
total_tokens=None,
|
|
events_with_cost=0,
|
|
estimated_cost_usd=None,
|
|
events_with_latency=0,
|
|
latency_p50_ms=None,
|
|
latency_p90_ms=None,
|
|
latency_p95_ms=None,
|
|
latency_p99_ms=None,
|
|
latency_avg_ms=None,
|
|
events_with_duration=0,
|
|
duration_avg_ms=None,
|
|
display_tokens="Unknown",
|
|
display_cost="Unknown",
|
|
display_latency_p50="Unknown",
|
|
display_latency_p90="Unknown",
|
|
display_duration_avg="Unknown",
|
|
)
|
|
|
|
token_events = [
|
|
e for e in events
|
|
if e.total_tokens is not None or e.input_tokens is not None or e.output_tokens is not None
|
|
]
|
|
events_with_tokens = len(token_events)
|
|
if events_with_tokens > 0:
|
|
input_tokens = sum(e.input_tokens or 0 for e in token_events)
|
|
output_tokens = sum(e.output_tokens or 0 for e in token_events)
|
|
total_tokens = sum(
|
|
e.total_tokens if e.total_tokens is not None else ((e.input_tokens or 0) + (e.output_tokens or 0))
|
|
for e in token_events
|
|
)
|
|
display_tokens = f"{total_tokens:,}"
|
|
else:
|
|
input_tokens = None
|
|
output_tokens = None
|
|
total_tokens = None
|
|
display_tokens = "Unknown"
|
|
|
|
cost_events = [e for e in events if e.estimated_cost_usd is not None]
|
|
events_with_cost = len(cost_events)
|
|
if events_with_cost > 0:
|
|
estimated_cost_usd = round(sum(e.estimated_cost_usd for e in cost_events), 6)
|
|
display_cost = f"${estimated_cost_usd:.4f}"
|
|
else:
|
|
estimated_cost_usd = None
|
|
display_cost = "Unknown"
|
|
|
|
latency_vals = [e.latency_ms for e in events if e.latency_ms is not None]
|
|
events_with_latency = len(latency_vals)
|
|
if events_with_latency > 0:
|
|
latency_p50_ms = compute_percentile(latency_vals, 50.0)
|
|
latency_p90_ms = compute_percentile(latency_vals, 90.0)
|
|
latency_p95_ms = compute_percentile(latency_vals, 95.0)
|
|
latency_p99_ms = compute_percentile(latency_vals, 99.0)
|
|
latency_avg_ms = round(sum(latency_vals) / events_with_latency, 2)
|
|
display_latency_p50 = f"{round(latency_p50_ms, 1)} ms" if latency_p50_ms is not None else "Unknown"
|
|
display_latency_p90 = f"{round(latency_p90_ms, 1)} ms" if latency_p90_ms is not None else "Unknown"
|
|
else:
|
|
latency_p50_ms = None
|
|
latency_p90_ms = None
|
|
latency_p95_ms = None
|
|
latency_p99_ms = None
|
|
latency_avg_ms = None
|
|
display_latency_p50 = "Unknown"
|
|
display_latency_p90 = "Unknown"
|
|
|
|
duration_vals = [e.duration_ms for e in events if e.duration_ms is not None]
|
|
events_with_duration = len(duration_vals)
|
|
if events_with_duration > 0:
|
|
duration_avg_ms = round(sum(duration_vals) / events_with_duration, 2)
|
|
display_duration_avg = f"{round(duration_avg_ms / 1000.0, 2)} s" if duration_avg_ms >= 1000 else f"{round(duration_avg_ms, 1)} ms"
|
|
else:
|
|
duration_avg_ms = None
|
|
display_duration_avg = "Unknown"
|
|
|
|
return GroupMetrics(
|
|
name=group_name,
|
|
total_events=total_events,
|
|
events_with_tokens=events_with_tokens,
|
|
input_tokens=input_tokens,
|
|
output_tokens=output_tokens,
|
|
total_tokens=total_tokens,
|
|
events_with_cost=events_with_cost,
|
|
estimated_cost_usd=estimated_cost_usd,
|
|
events_with_latency=events_with_latency,
|
|
latency_p50_ms=latency_p50_ms,
|
|
latency_p90_ms=latency_p90_ms,
|
|
latency_p95_ms=latency_p95_ms,
|
|
latency_p99_ms=latency_p99_ms,
|
|
latency_avg_ms=latency_avg_ms,
|
|
events_with_duration=events_with_duration,
|
|
duration_avg_ms=duration_avg_ms,
|
|
display_tokens=display_tokens,
|
|
display_cost=display_cost,
|
|
display_latency_p50=display_latency_p50,
|
|
display_latency_p90=display_latency_p90,
|
|
display_duration_avg=display_duration_avg,
|
|
)
|
|
|
|
|
|
def record_usage(
|
|
*,
|
|
db_path: str | None = None,
|
|
session_id: str | None = None,
|
|
remote: str = "dadeschools",
|
|
org: str = "",
|
|
repo: str = "",
|
|
project_id: str | None = None,
|
|
role: str = "unknown",
|
|
model: str = "unknown",
|
|
issue_number: int | None = None,
|
|
pr_number: int | None = None,
|
|
stage: str = "unknown",
|
|
input_tokens: int | None = None,
|
|
output_tokens: int | None = None,
|
|
total_tokens: int | None = None,
|
|
estimated_cost_usd: float | None = None,
|
|
latency_ms: int | None = None,
|
|
duration_ms: int | None = None,
|
|
status: str = "success",
|
|
metadata: str | dict[str, Any] | None = None,
|
|
created_at: str | None = None,
|
|
) -> int:
|
|
"""Ingest/record a single usage event with optional metrics."""
|
|
db = control_plane_db.ControlPlaneDB(db_path=db_path)
|
|
return db.record_usage_event(
|
|
session_id=session_id,
|
|
remote=remote,
|
|
org=org,
|
|
repo=repo,
|
|
project_id=project_id,
|
|
role=role,
|
|
model=model,
|
|
issue_number=issue_number,
|
|
pr_number=pr_number,
|
|
stage=stage,
|
|
input_tokens=input_tokens,
|
|
output_tokens=output_tokens,
|
|
total_tokens=total_tokens,
|
|
estimated_cost_usd=estimated_cost_usd,
|
|
latency_ms=latency_ms,
|
|
duration_ms=duration_ms,
|
|
status=status,
|
|
metadata=metadata,
|
|
created_at=created_at,
|
|
)
|
|
|
|
|
|
def load_analytics(
|
|
*,
|
|
db_path: str | None = None,
|
|
remote: str | None = None,
|
|
org: str | None = None,
|
|
repo: str | None = None,
|
|
project_id: str | None = None,
|
|
role: str | None = None,
|
|
model: str | None = None,
|
|
stage: str | None = None,
|
|
issue_number: int | None = None,
|
|
pr_number: int | None = None,
|
|
limit: int = 500,
|
|
) -> AnalyticsSnapshot:
|
|
"""Load analytics snapshot aggregated by project, role, model, issue/PR, and stage."""
|
|
remote_filter = (remote or "").strip() or None
|
|
org_filter = (org or "").strip() or None
|
|
repo_filter = (repo or "").strip() or None
|
|
role_filter = (role or "").strip() or None
|
|
model_filter = (model or "").strip() or None
|
|
stage_filter = (stage or "").strip() or None
|
|
|
|
# F4: coerce optional scope filters to str so AnalyticsSnapshot never holds None.
|
|
scope_remote = (remote or "").strip()
|
|
scope_org = (org or "").strip()
|
|
scope_repo = (repo or "").strip()
|
|
|
|
try:
|
|
db = control_plane_db.ControlPlaneDB(db_path=db_path)
|
|
rows = db.query_usage_events(
|
|
remote=remote_filter,
|
|
org=org_filter,
|
|
repo=repo_filter,
|
|
project_id=project_id,
|
|
role=role_filter,
|
|
model=model_filter,
|
|
stage=stage_filter,
|
|
issue_number=issue_number,
|
|
pr_number=pr_number,
|
|
limit=limit,
|
|
)
|
|
except Exception as exc:
|
|
empty_summary = aggregate_events("Overall", [])
|
|
return AnalyticsSnapshot(
|
|
ok=False,
|
|
reason=f"control_plane_db_unavailable: {exc}",
|
|
schema_version=ANALYTICS_SCHEMA_VERSION,
|
|
remote=scope_remote,
|
|
org=scope_org,
|
|
repo=scope_repo,
|
|
total_events=0,
|
|
overall_summary=empty_summary,
|
|
by_project={},
|
|
by_role={},
|
|
by_model={},
|
|
by_work_item={},
|
|
by_stage={},
|
|
events=(),
|
|
)
|
|
|
|
parsed_events: list[UsageEvent] = []
|
|
for r in rows:
|
|
meta = console_redaction.redact_text(r.get("metadata")) if r.get("metadata") else None
|
|
parsed_events.append(
|
|
UsageEvent(
|
|
usage_id=r["usage_id"],
|
|
session_id=r.get("session_id"),
|
|
remote=r.get("remote") or scope_remote,
|
|
org=r.get("org") or scope_org,
|
|
repo=r.get("repo") or scope_repo,
|
|
project_id=r.get("project_id"),
|
|
role=r.get("role") or "unknown",
|
|
model=r.get("model") or "unknown",
|
|
issue_number=r.get("issue_number"),
|
|
pr_number=r.get("pr_number"),
|
|
stage=r.get("stage") or "unknown",
|
|
input_tokens=r.get("input_tokens"),
|
|
output_tokens=r.get("output_tokens"),
|
|
total_tokens=r.get("total_tokens"),
|
|
estimated_cost_usd=r.get("estimated_cost_usd"),
|
|
latency_ms=r.get("latency_ms"),
|
|
duration_ms=r.get("duration_ms"),
|
|
status=r.get("status") or "success",
|
|
metadata=meta,
|
|
created_at=r.get("created_at") or "",
|
|
)
|
|
)
|
|
|
|
overall_summary = aggregate_events("Overall", parsed_events)
|
|
|
|
# Group by project
|
|
groups_by_project: dict[str, list[UsageEvent]] = {}
|
|
for e in parsed_events:
|
|
key = e.project_id or (f"{e.org}/{e.repo}" if e.org and e.repo else "default")
|
|
groups_by_project.setdefault(key, []).append(e)
|
|
by_project = {k: aggregate_events(k, v) for k, v in groups_by_project.items()}
|
|
|
|
# Group by role
|
|
groups_by_role: dict[str, list[UsageEvent]] = {}
|
|
for e in parsed_events:
|
|
groups_by_role.setdefault(e.role, []).append(e)
|
|
by_role = {k: aggregate_events(k, v) for k, v in groups_by_role.items()}
|
|
|
|
# Group by model
|
|
groups_by_model: dict[str, list[UsageEvent]] = {}
|
|
for e in parsed_events:
|
|
groups_by_model.setdefault(e.model, []).append(e)
|
|
by_model = {k: aggregate_events(k, v) for k, v in groups_by_model.items()}
|
|
|
|
# Group by work item
|
|
groups_by_work_item: dict[str, list[UsageEvent]] = {}
|
|
for e in parsed_events:
|
|
if e.issue_number:
|
|
key = f"issue #{e.issue_number}"
|
|
elif e.pr_number:
|
|
key = f"pr #{e.pr_number}"
|
|
else:
|
|
key = "unlinked"
|
|
groups_by_work_item.setdefault(key, []).append(e)
|
|
by_work_item = {k: aggregate_events(k, v) for k, v in groups_by_work_item.items()}
|
|
|
|
# Group by stage
|
|
groups_by_stage: dict[str, list[UsageEvent]] = {}
|
|
for e in parsed_events:
|
|
groups_by_stage.setdefault(e.stage, []).append(e)
|
|
by_stage = {k: aggregate_events(k, v) for k, v in groups_by_stage.items()}
|
|
|
|
return AnalyticsSnapshot(
|
|
ok=True,
|
|
reason="ok",
|
|
schema_version=ANALYTICS_SCHEMA_VERSION,
|
|
remote=scope_remote,
|
|
org=scope_org,
|
|
repo=scope_repo,
|
|
total_events=len(parsed_events),
|
|
overall_summary=overall_summary,
|
|
by_project=by_project,
|
|
by_role=by_role,
|
|
by_model=by_model,
|
|
by_work_item=by_work_item,
|
|
by_stage=by_stage,
|
|
events=tuple(parsed_events),
|
|
)
|
|
|
|
|
|
def snapshot_to_dict(snapshot: AnalyticsSnapshot) -> dict[str, Any]:
|
|
return snapshot.to_dict()
|