diff --git a/control_plane_db.py b/control_plane_db.py index 4616cb3..5da9098 100644 --- a/control_plane_db.py +++ b/control_plane_db.py @@ -31,7 +31,7 @@ from typing import Any, Iterator, Sequence import dependency_graph -SCHEMA_VERSION = 4 +SCHEMA_VERSION = 5 # Assignable work kinds only — raw monitoring incidents are never work items. WORK_KINDS = frozenset({"issue", "pr"}) @@ -186,6 +186,34 @@ CREATE INDEX IF NOT EXISTS idx_dependency_edges_target ON dependency_edges(remote, org, repo, target_kind, target_number); CREATE INDEX IF NOT EXISTS idx_assignments_session ON assignments(session_id, status); CREATE INDEX IF NOT EXISTS idx_incident_gitea ON incident_links(gitea_org, gitea_repo, gitea_issue_number); + +-- Model usage, token cost, latency, and performance events (#651) +CREATE TABLE IF NOT EXISTS usage_events ( + usage_id INTEGER PRIMARY KEY AUTOINCREMENT, + session_id TEXT, + remote TEXT NOT NULL DEFAULT 'dadeschools', + org TEXT NOT NULL DEFAULT '', + repo TEXT NOT NULL DEFAULT '', + project_id TEXT, + role TEXT NOT NULL DEFAULT 'unknown', + model TEXT NOT NULL DEFAULT 'unknown', + issue_number INTEGER, + pr_number INTEGER, + stage TEXT NOT NULL DEFAULT 'unknown', + input_tokens INTEGER, + output_tokens INTEGER, + total_tokens INTEGER, + estimated_cost_usd REAL, + latency_ms INTEGER, + duration_ms INTEGER, + status TEXT NOT NULL DEFAULT 'success', + metadata TEXT, + created_at TEXT NOT NULL +); + +CREATE INDEX IF NOT EXISTS idx_usage_events_scope ON usage_events(remote, org, repo); +CREATE INDEX IF NOT EXISTS idx_usage_events_role_model ON usage_events(role, model); +CREATE INDEX IF NOT EXISTS idx_usage_events_stage ON usage_events(stage); """ @@ -339,6 +367,7 @@ class ControlPlaneDB: self._migrate_incident_links_null_scope(conn) self._migrate_lease_lifecycle_columns(conn) self._migrate_session_ownership_columns(conn) + self._migrate_usage_events_table(conn) conn.execute( "INSERT OR REPLACE INTO schema_meta(key, value) VALUES (?, ?)", ("schema_version", str(SCHEMA_VERSION)), @@ -518,6 +547,174 @@ class ControlPlaneDB: f"UPDATE incident_links SET {col} = '' WHERE {col} IS NULL" ) + def _migrate_usage_events_table(self, conn: sqlite3.Connection) -> None: + """Create usage_events table and indexes if they do not exist (#651).""" + conn.execute(""" + CREATE TABLE IF NOT EXISTS usage_events ( + usage_id INTEGER PRIMARY KEY AUTOINCREMENT, + session_id TEXT, + remote TEXT NOT NULL DEFAULT 'dadeschools', + org TEXT NOT NULL DEFAULT '', + repo TEXT NOT NULL DEFAULT '', + project_id TEXT, + role TEXT NOT NULL DEFAULT 'unknown', + model TEXT NOT NULL DEFAULT 'unknown', + issue_number INTEGER, + pr_number INTEGER, + stage TEXT NOT NULL DEFAULT 'unknown', + input_tokens INTEGER, + output_tokens INTEGER, + total_tokens INTEGER, + estimated_cost_usd REAL, + latency_ms INTEGER, + duration_ms INTEGER, + status TEXT NOT NULL DEFAULT 'success', + metadata TEXT, + created_at TEXT NOT NULL + ); + """) + conn.execute("CREATE INDEX IF NOT EXISTS idx_usage_events_scope ON usage_events(remote, org, repo);") + conn.execute("CREATE INDEX IF NOT EXISTS idx_usage_events_role_model ON usage_events(role, model);") + conn.execute("CREATE INDEX IF NOT EXISTS idx_usage_events_stage ON usage_events(stage);") + + def record_usage_event( + self, + *, + 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: + """Record a model usage, token cost, latency, or stage performance event (#651).""" + ts = created_at or _ts() + meta_str: str | None = None + if metadata is not None: + from webui import console_redaction + redacted_meta = console_redaction.redact_payload(metadata) + if isinstance(redacted_meta, str): + meta_str = redacted_meta + else: + try: + meta_str = json.dumps(redacted_meta, default=str) + except Exception: + meta_str = str(redacted_meta) + + if total_tokens is None and (input_tokens is not None or output_tokens is not None): + total_tokens = (input_tokens or 0) + (output_tokens or 0) + + with self._tx(immediate=True) as conn: + cursor = conn.execute( + """ + INSERT INTO usage_events ( + session_id, remote, org, repo, project_id, role, model, + issue_number, pr_number, stage, input_tokens, output_tokens, + total_tokens, estimated_cost_usd, latency_ms, duration_ms, + status, metadata, created_at + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + """, + ( + session_id, + remote, + org, + repo, + project_id, + role, + model, + issue_number, + pr_number, + stage, + input_tokens, + output_tokens, + total_tokens, + estimated_cost_usd, + latency_ms, + duration_ms, + status, + meta_str, + ts, + ), + ) + return cursor.lastrowid + + def query_usage_events( + self, + *, + 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, + issue_number: int | None = None, + pr_number: int | None = None, + stage: str | None = None, + session_id: str | None = None, + limit: int = 500, + offset: int = 0, + ) -> list[dict[str, Any]]: + """Query stored usage events matching filters (#651).""" + conditions = [] + params = [] + if remote: + conditions.append("remote = ?") + params.append(remote) + if org: + conditions.append("org = ?") + params.append(org) + if repo: + conditions.append("repo = ?") + params.append(repo) + if project_id: + conditions.append("project_id = ?") + params.append(project_id) + if role: + conditions.append("role = ?") + params.append(role) + if model: + conditions.append("model = ?") + params.append(model) + if issue_number is not None: + conditions.append("issue_number = ?") + params.append(issue_number) + if pr_number is not None: + conditions.append("pr_number = ?") + params.append(pr_number) + if stage: + conditions.append("stage = ?") + params.append(stage) + if session_id: + conditions.append("session_id = ?") + params.append(session_id) + + where_clause = f"WHERE {' AND '.join(conditions)}" if conditions else "" + sql = f""" + SELECT * FROM usage_events + {where_clause} + ORDER BY usage_id ASC + LIMIT ? OFFSET ? + """ + params.extend([limit, offset]) + + with self._tx(immediate=False) as conn: + cursor = conn.execute(sql, params) + rows = cursor.fetchall() + return [dict(row) for row in rows] + # ── sessions ────────────────────────────────────────────────────────── def upsert_session( diff --git a/docs/observability/analytics-instrumentation.md b/docs/observability/analytics-instrumentation.md new file mode 100644 index 0000000..3994045 --- /dev/null +++ b/docs/observability/analytics-instrumentation.md @@ -0,0 +1,125 @@ +# Model Usage, Token Cost, Latency, and Workflow Analytics (Phase 4) + +- **Tracking Issue:** [#651](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/651) +- **Parent Epic:** [#631](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/651) +- **Console Surface:** `/analytics`, `/api/v1/analytics`, `/api/v1/analytics/usage` + +## 1. Overview + +The Web Console Analytics module provides durable, aggregate visibility into **model usage, token cost, latency percentiles, and workflow-stage performance** across projects, worker roles, AI models, issues, and PRs. + +### Non-Goals +- No mandatory client-side telemetry that leaks prompts or secret keys. +- No third-party payment provider or billing integration. +- No automatic model routing changes without controller policy (#647). + +--- + +## 2. Event Schema (`usage_events`) + +Usage metrics are stored in the control-plane database under table `usage_events`. + +| Column | Type | Description | +|---|---|---| +| `usage_id` | `INTEGER` | Primary key (autoincrement) | +| `session_id` | `TEXT` | Optional active session identifier | +| `remote` | `TEXT` | Known Gitea instance (`dadeschools` or `prgs`) | +| `org` | `TEXT` | Repository owner / organization | +| `repo` | `TEXT` | Repository name | +| `project_id` | `TEXT` | Optional project identifier | +| `role` | `TEXT` | Active worker role (`author`, `reviewer`, `merger`, `reconciler`, `controller`) | +| `model` | `TEXT` | LLM model identifier (e.g. `gemini-3.6-flash`, `claude-3-5-sonnet`) | +| `issue_number` | `INTEGER` | Correlated Gitea issue number (optional) | +| `pr_number` | `INTEGER` | Correlated Gitea PR number (optional) | +| `stage` | `TEXT` | Workflow stage (`preflight`, `implementation`, `review`, `merge`, `reconciliation`) | +| `input_tokens` | `INTEGER` | Input token count (optional / nullable) | +| `output_tokens` | `INTEGER` | Output token count (optional / nullable) | +| `total_tokens` | `INTEGER` | Total token count (optional / nullable) | +| `estimated_cost_usd` | `REAL` | Estimated USD cost (optional / nullable) | +| `latency_ms` | `INTEGER` | Request latency in milliseconds (optional / nullable) | +| `duration_ms` | `INTEGER` | Stage execution duration in milliseconds (optional / nullable) | +| `status` | `TEXT` | Outcome status (`success`, `failure`, `timeout`) | +| `metadata` | `TEXT` | Redacted metadata or summary string | +| `created_at` | `TEXT` | ISO 8601 UTC timestamp | + +--- + +## 3. Handling of Missing Data ("Unknown" vs. Zero Fabrication) + +To ensure operational metrics accurately reflect evidence: +- **Untracked or missing metrics are displayed as `Unknown`**, never zero-fabricated. +- If an event omits `estimated_cost_usd`, `latency_ms`, or token counts, the aggregator marks those fields as missing (`None`) rather than defaulting to `0` or `$0.00`. +- Summary tables and KPI cards explicitly indicate when data is unmeasured or partially reported. + +--- + +## 4. Redaction & Security Rules + +Per `#633` security policy: +- Free-text fields (`metadata`, `prompt_summary`, `session_id`) are run through `console_redaction.redact_text` before persistence and output serialization. +- Secret tokens, keychain commands, authorization headers, passwords, and JWTs are stripped automatically. + +--- + +## 5. Opt-in Instrumentation Guide + +Applications, MCP servers, and background sessions can report usage metrics through either Python API or HTTP ingestion. + +### Python Ingestion + +```python +from webui.analytics_loader import record_usage + +record_usage( + remote="dadeschools", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + role="author", + model="gemini-3.6-flash", + issue_number=651, + stage="implementation", + input_tokens=1420, + output_tokens=380, + total_tokens=1800, + estimated_cost_usd=0.00045, + latency_ms=320, + duration_ms=4500, + status="success", + metadata={"note": "Implementation of analytics module"}, +) +``` + +### HTTP Ingestion API + +```http +POST /api/v1/analytics/usage HTTP/1.1 +Content-Type: application/json + +{ + "remote": "dadeschools", + "org": "Scaled-Tech-Consulting", + "repo": "Gitea-Tools", + "role": "author", + "model": "gemini-3.6-flash", + "issue_number": 651, + "stage": "implementation", + "input_tokens": 1420, + "output_tokens": 380, + "total_tokens": 1800, + "estimated_cost_usd": 0.00045, + "latency_ms": 320, + "duration_ms": 4500, + "status": "success", + "metadata": "Analytics schema landed" +} +``` + +--- + +## 6. Querying Analytics API + +```http +GET /api/v1/analytics?role=author&stage=implementation HTTP/1.1 +``` + +Returns `AnalyticsSnapshot` JSON containing aggregations (`by_model`, `by_stage`, `by_role`, `by_work_item`, `by_project`) and latency percentiles (`p50`, `p90`, `p95`, `p99`). diff --git a/tests/test_control_plane_db.py b/tests/test_control_plane_db.py index 8f251eb..1752299 100644 --- a/tests/test_control_plane_db.py +++ b/tests/test_control_plane_db.py @@ -36,7 +36,7 @@ class ControlPlaneDBTest(unittest.TestCase): rows = dict(conn.execute("SELECT key, value FROM schema_meta").fetchall()) finally: conn.close() - self.assertEqual(rows["schema_version"], "4") + self.assertEqual(rows["schema_version"], "5") self.assertIn("DB coordinates", rows["architecture"]) self.assertIn("bridge", rows["architecture"].lower()) diff --git a/tests/test_webui_analytics.py b/tests/test_webui_analytics.py new file mode 100644 index 0000000..c71322b --- /dev/null +++ b/tests/test_webui_analytics.py @@ -0,0 +1,196 @@ +"""Unit and integration tests for Model Usage & Performance Analytics (#651).""" + +from __future__ import annotations + +import os +import tempfile +import unittest +from starlette.testclient import TestClient + +import control_plane_db +from webui.analytics_loader import ( + ANALYTICS_SCHEMA_VERSION, + compute_percentile, + load_analytics, + record_usage, +) +from webui.app import create_app +from webui import console_redaction + + +class AnalyticsLoaderTest(unittest.TestCase): + + def setUp(self) -> None: + self.temp_dir = tempfile.TemporaryDirectory() + self.db_path = os.path.join(self.temp_dir.name, "test_control_plane.sqlite3") + os.environ["GITEA_CONTROL_PLANE_DB"] = self.db_path + self.db = control_plane_db.ControlPlaneDB(db_path=self.db_path) + + def tearDown(self) -> None: + self.temp_dir.cleanup() + + def test_compute_percentile(self) -> None: + self.assertIsNone(compute_percentile([], 50.0)) + self.assertEqual(compute_percentile([100], 50.0), 100.0) + + # 2 elements: [100, 200] + self.assertEqual(compute_percentile([100, 200], 50.0), 150.0) + + # 100 elements: 1..100 + vals = list(range(1, 101)) + self.assertEqual(compute_percentile(vals, 50.0), 50.5) + self.assertAlmostEqual(compute_percentile(vals, 90.0), 90.1) + + def test_record_and_aggregate_usage(self) -> None: + # Record event 1 (complete data) + u1 = record_usage( + db_path=self.db_path, + remote="dadeschools", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + role="author", + model="gemini-3.6-flash", + issue_number=651, + stage="implementation", + input_tokens=1000, + output_tokens=500, + estimated_cost_usd=0.0015, + latency_ms=200, + duration_ms=3000, + metadata={"secret_key": "secret123", "note": "token=secret123"}, + ) + self.assertGreater(u1, 0) + + # Record event 2 (missing tokens and cost -> unknown) + u2 = record_usage( + db_path=self.db_path, + remote="dadeschools", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + role="reviewer", + model="claude-3-5-sonnet", + pr_number=846, + stage="review", + latency_ms=500, + duration_ms=6000, + ) + self.assertGreater(u2, u1) + + snapshot = load_analytics( + db_path=self.db_path, + remote="dadeschools", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + ) + + self.assertTrue(snapshot.ok) + self.assertEqual(snapshot.schema_version, ANALYTICS_SCHEMA_VERSION) + self.assertEqual(snapshot.total_events, 2) + + # Verify overall summary + summary = snapshot.overall_summary + self.assertEqual(summary.total_events, 2) + self.assertEqual(summary.events_with_tokens, 1) + self.assertEqual(summary.total_tokens, 1500) + self.assertEqual(summary.events_with_cost, 1) + self.assertEqual(summary.estimated_cost_usd, 0.0015) + self.assertEqual(summary.events_with_latency, 2) + self.assertEqual(summary.latency_p50_ms, 350.0) + + # Verify missing data handling (AC 3: not zero-fabricated) + reviewer_model = snapshot.by_model.get("claude-3-5-sonnet") + self.assertIsNotNone(reviewer_model) + self.assertEqual(reviewer_model.total_events, 1) + self.assertEqual(reviewer_model.events_with_tokens, 0) + self.assertIsNone(reviewer_model.total_tokens) + self.assertEqual(reviewer_model.display_tokens, "Unknown") + self.assertEqual(reviewer_model.events_with_cost, 0) + self.assertIsNone(reviewer_model.estimated_cost_usd) + self.assertEqual(reviewer_model.display_cost, "Unknown") + + # Verify redaction (AC 4) + e1 = [e for e in snapshot.events if e.usage_id == u1][0] + self.assertIsNotNone(e1.metadata) + self.assertNotIn("secret123", e1.metadata) + self.assertIn("[REDACTED]", e1.metadata) + + def test_missing_db_fail_soft(self) -> None: + invalid_path = "/nonexistent_path_dir/db.sqlite3" + snapshot = load_analytics(db_path=invalid_path) + self.assertFalse(snapshot.ok) + self.assertIn("control_plane_db_unavailable", snapshot.reason) + self.assertEqual(snapshot.overall_summary.display_tokens, "Unknown") + + +class AnalyticsWebUITest(unittest.TestCase): + + def setUp(self) -> None: + self.temp_dir = tempfile.TemporaryDirectory() + self.db_path = os.path.join(self.temp_dir.name, "test_webui.sqlite3") + os.environ["GITEA_CONTROL_PLANE_DB"] = self.db_path + self.app = create_app() + self.client = TestClient(self.app) + + record_usage( + db_path=self.db_path, + remote="dadeschools", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + role="author", + model="gemini-3.6-flash", + issue_number=651, + stage="implementation", + input_tokens=2000, + output_tokens=1000, + estimated_cost_usd=0.003, + latency_ms=150, + duration_ms=2500, + ) + + def tearDown(self) -> None: + self.temp_dir.cleanup() + + def test_analytics_html_route(self) -> None: + response = self.client.get("/analytics") + self.assertEqual(response.status_code, 200) + self.assertIn("Model Usage & Performance Analytics", response.text) + self.assertIn("gemini-3.6-flash", response.text) + self.assertIn("3,000", response.text) + + def test_analytics_api_route(self) -> None: + response = self.client.get("/api/v1/analytics") + self.assertEqual(response.status_code, 200) + data = response.json() + self.assertTrue(data["ok"]) + self.assertEqual(data["total_events"], 1) + self.assertIn("gemini-3.6-flash", data["by_model"]) + + def test_analytics_ingest_endpoint(self) -> None: + payload = { + "remote": "dadeschools", + "org": "Scaled-Tech-Consulting", + "repo": "Gitea-Tools", + "role": "reviewer", + "model": "claude-3-5-sonnet", + "pr_number": 846, + "stage": "review", + "input_tokens": 500, + "output_tokens": 100, + "latency_ms": 400, + "metadata": "Review note token=secret456", + } + response = self.client.post("/api/v1/analytics/usage", json=payload) + self.assertEqual(response.status_code, 201) + res_json = response.json() + self.assertTrue(res_json["ok"]) + self.assertGreater(res_json["usage_id"], 0) + + # Check that it appears in GET /api/v1/analytics + res2 = self.client.get("/api/v1/analytics") + self.assertEqual(res2.status_code, 200) + data2 = res2.json() + self.assertEqual(data2["total_events"], 2) + + +if __name__ == "__main__": + unittest.main() diff --git a/webui/analytics_loader.py b/webui/analytics_loader.py new file mode 100644 index 0000000..d6ef3d3 --- /dev/null +++ b/webui/analytics_loader.py @@ -0,0 +1,428 @@ +"""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 + + 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=remote, + org=org, + repo=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", remote), + org=r.get("org", org), + repo=r.get("repo", repo), + project_id=r.get("project_id"), + role=r.get("role", "unknown"), + model=r.get("model", "unknown"), + issue_number=r.get("issue_number"), + pr_number=r.get("pr_number"), + stage=r.get("stage", "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", "success"), + metadata=meta, + created_at=r.get("created_at", ""), + ) + ) + + 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=remote, + org=org, + repo=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() diff --git a/webui/analytics_views.py b/webui/analytics_views.py new file mode 100644 index 0000000..b2583ca --- /dev/null +++ b/webui/analytics_views.py @@ -0,0 +1,201 @@ +"""HTML views for the Model Usage & Performance Analytics console (#651).""" + +from __future__ import annotations + +from webui.analytics_loader import AnalyticsSnapshot, GroupMetrics, UsageEvent +from webui.layout import render_page + + +def _render_badge(text: str, badge_type: str = "muted") -> str: + return f'{text}' + + +def _render_group_table(title: str, groups: dict[str, GroupMetrics], key_header: str = "Group") -> str: + if not groups: + return ( + f"
No telemetry events recorded for this dimension.
| {key_header} | +Events | +Total Tokens | +Est. Cost | +Latency (p50) | +Latency (p90) | +Avg Stage Duration | +
|---|
No individual telemetry events recorded yet. Opt-in instrumentation via session logging or POST /api/v1/analytics/usage.
| ID | +Timestamp | +Role | +Model | +Stage | +Work Item | +Tokens | +Cost | +Latency | +Duration | +Status | +
|---|
+ Durable console analytics for model usage, token cost, latency percentiles, and workflow-stage performance correlated to issues, PRs, and worker roles. +
+ + {status_notice} + +