diff --git a/.env.example b/.env.example index 172bdf655..cd09bc682 100644 --- a/.env.example +++ b/.env.example @@ -103,6 +103,13 @@ # (sensitive) # CREDENTIAL_ENCRYPTION_KEY= +# HMAC key for signing pagination cursors. Auto-provisioned to ~/.luthien/cursor_hmac.key on first run. +# In multi-replica deployments, set this explicitly so cursors validate across replicas. +# Each replica that auto-generates its own key will reject cursors issued by other replicas. +# (sensitive) +# (default derived from runtime at startup) +# CURSOR_HMAC_KEY= + # === OBSERVABILITY =============================================== diff --git a/.gitignore b/.gitignore index 61d5952d7..05a101247 100644 --- a/.gitignore +++ b/.gitignore @@ -86,6 +86,9 @@ tests/**/failure_registry/*.json # Serena AI assistant local state .serena/ +# Sisyphus agent working state (plans, evidence, run artifacts) +.sisyphus/ + # Local scratch files kanban.jpg .e2e-logs/ diff --git a/changelog.d/perf-baseline.md b/changelog.d/perf-baseline.md new file mode 100644 index 000000000..f23ee6f70 --- /dev/null +++ b/changelog.d/perf-baseline.md @@ -0,0 +1,13 @@ +--- +category: Chores & Docs +pr: 752 +--- + +**Admin UI performance baseline**: Establishes perf infrastructure and captures SQLite baseline for history/conversation pages. + - Perf test scaffolding: isolated DB (`~/.luthien/perf.db`), seeding fixtures (sami-like, tier-100/1000/10000), and Playwright harness + - `scripts/perf_explain.py` — captures EXPLAIN QUERY PLAN for top slow queries + - `scripts/perf_report.py` — generates Markdown baseline report from seeded DB + query plans + - `scripts/run_perf.sh` — orchestrates seed + test + SLO assertion workflow + - Middleware timing (`Server-Timing` header) and payload-size contract tests + - Query plan evidence: 2× TEMP B-TREE on `session_list`, full SCAN on `recent_calls` + - Postgres baseline skipped (not available locally); run `./scripts/run_perf.sh --backend postgres` to capture diff --git a/changelog.d/perf-fix.md b/changelog.d/perf-fix.md new file mode 100644 index 000000000..cfee6271c --- /dev/null +++ b/changelog.d/perf-fix.md @@ -0,0 +1,14 @@ +--- +category: Features +pr: 752 +--- + +**Admin UI performance optimizations**: Cursor pagination and memory caps for the history page and conversation viewer. + - Cursor-paginated infinite scroll on `/history` (20 sessions per page instead of all) + - Conversation viewer loads full session JSON via `/api/history/sessions/{id}` and renders structured turns; paginated lazy loading of turns is deferred to a follow-up PR + - Raw events capped at 50 per call_id (FIFO) to bound per-turn memory growth + - Debounced server-side filter on `/history` to reduce query load + - New fragment endpoints: `/ui/fragments/sessions`, `/ui/fragments/sessions/{id}/turns` + - **Known limitation**: `filter=claude` uses a full-table payload scan (`payload LIKE '%claude-code%'`) with no index. It is correct for small deployments but will be slow on large Postgres instances. A structured `client_type` column or trigram index is the long-term fix. + - **Known limitation**: session-ID search (`q=`) uses a leading-wildcard `LIKE '%q%'` which cannot use a btree index. Intended for small deployments; a trigram index or prefix-only match is the long-term fix. + - **Known limitation**: `CURSOR_HMAC_KEY` is auto-provisioned per-instance to `~/.luthien/cursor_hmac.key`. In multi-replica deployments without sticky sessions, each replica generates its own key — cursors issued by one replica will be rejected (400) by another. Set `CURSOR_HMAC_KEY` explicitly in the environment when running behind a load balancer. diff --git a/dev/context/migration_concurrent.md b/dev/context/migration_concurrent.md new file mode 100644 index 000000000..73ce61310 --- /dev/null +++ b/dev/context/migration_concurrent.md @@ -0,0 +1,105 @@ +# Migration Runner: CONCURRENTLY Support Audit + +_Date: 2026-05-15 | Branch: perf-baseline_ + +## Background + +`CREATE INDEX CONCURRENTLY` is a Postgres feature that builds an index without holding a lock on the table, allowing reads and writes during the build. The constraint: it **cannot run inside a transaction block**. This audit investigates whether the current migration runner can safely execute such a statement. + +--- + +## Current behavior + +### PostgreSQL runner (`docker/run-migrations.sh`) + +- Applied by the `migrations` Docker service at startup; controlled by `docker compose up migrations`. +- Sequentially applies all `*.sql` files in `migrations/postgres/` in alphabetical order. +- **No `BEGIN`/`COMMIT` transaction wrapping** is added around migration files. The runner calls `psql -f "$migration"` directly: + ```sh + psql -h "$PGHOST" -U "$PGUSER" -d "$PGDATABASE" -f "$migration" + ``` +- `psql` defaults to autocommit mode — each statement in the file runs in its own implicit transaction unless the file itself contains explicit `BEGIN`/`COMMIT` blocks. +- The `_migrations` tracking row (`INSERT INTO _migrations`) is inserted in a **separate, subsequent `psql` invocation**, not inside the same transaction as the migration file. This means the tracking and the DDL are non-atomic: a crash between the two steps leaves schema changes applied but untracked. +- Migration state is tracked in the `_migrations` table (columns: `filename TEXT PK`, `applied_at TIMESTAMP`, `content_hash TEXT`). +- Applied-migration detection uses `SELECT COUNT(*) FROM _migrations WHERE filename = '$filename'`, checked per file before applying. +- Hash validation compares stored MD5 against local file MD5 and aborts on mismatch. + +### SQLite runner (`src/luthien_proxy/utils/migration_check.py :: _apply_sqlite_migrations`) + +- Runs in-process at gateway startup for dockerless/SQLite deployments. +- Uses `executescript()` to apply each `.sql` file — this method issues an implicit `COMMIT` before execution and runs all statements in the file sequentially. +- `CREATE INDEX CONCURRENTLY` is not a SQLite concept; `AGENTS.md` explicitly lists it under "What to OMIT in SQLite migrations" and directs authors to use `CREATE INDEX IF NOT EXISTS` instead. +- SQLite tracking is also done in the `_migrations` table but is written inside the same connection context (not atomic with the `executescript`, however — a mid-script crash leaves partial schema with no tracking record). + +--- + +## Verdict + +**PARTIAL** + +`CREATE INDEX CONCURRENTLY` can be placed in a Postgres migration file today and will execute successfully — because the runner uses `psql -f` in autocommit mode with **no outer transaction wrapping**. The statement will not hit the "cannot run inside a transaction block" error. + +However: + +1. **Non-atomic tracking** — the `INSERT INTO _migrations` tracking row is a separate psql call. If it fails, the index exists on disk but the migration is untracked. A re-run will try to apply the file again; `CREATE INDEX CONCURRENTLY IF NOT EXISTS` protects against failure in that case. +2. **SQLite incompatibility** — a companion SQLite migration must use plain `CREATE INDEX IF NOT EXISTS` (standard `AGENTS.md` practice; no code change needed). +3. **No explicit guidance in runner or AGENTS.md** about CONCURRENTLY for Postgres beyond the SQLite omit rule — the assumption has been "it just works because psql is autocommit." + +--- + +## Findings + +1. **No BEGIN/COMMIT wrapping in Postgres runner.** `run-migrations.sh` calls `psql -f "$migration"` with zero explicit transaction control around migration files. psql autocommit applies. + +2. **`BEGIN` in existing migrations is always PL/pgSQL, not transaction control.** Searching all postgres migration files reveals `BEGIN` only inside `$$ LANGUAGE plpgsql` function/trigger bodies (e.g., `014_add_session_search_fts.sql`, `000_init_databases.sql`). No migration wraps its DDL in a `BEGIN...COMMIT` block. + +3. **Tracking INSERT is not atomic with migration application.** Lines 153–156 of `run-migrations.sh` run the migration file, then insert into `_migrations` in a second psql call. A process kill between those two calls yields applied-but-untracked state. `CREATE INDEX CONCURRENTLY IF NOT EXISTS` + idempotent DDL is the correct mitigation. + +4. **SQLite runner uses `executescript()`, not raw `execute()`.** This means the entire SQL file is submitted to SQLite's native multi-statement parser in one call. It handles trigger `BEGIN...END` correctly but does not guarantee atomicity across the file; a mid-script error leaves partial schema with no `_migrations` entry. + +5. **AGENTS.md already documents the SQLite handling rule.** "What to OMIT in SQLite migrations" includes `CREATE INDEX CONCURRENTLY` — use plain `CREATE INDEX IF NOT EXISTS`. This is the only dual-dialect consideration; Postgres needs no special handling beyond `IF NOT EXISTS`. + +6. **Migration 006 establishes the index-in-migration pattern.** `006_add_session_id.sql` creates two partial indexes (`WHERE session_id IS NOT NULL`) with `CREATE INDEX IF NOT EXISTS`. This is the precedent: use `IF NOT EXISTS` for idempotence, and the runner handles it without transaction complications. + +7. **`014_add_session_search_fts.sql` creates multiple indexes in one file.** A GIN index, a btree partial index, and an expression index are all created in a single migration file, all with `IF NOT EXISTS`. This confirms that non-trivial index migrations work fine under the current runner. + +--- + +## Risk assessment + +### If a future PR needs `CREATE INDEX CONCURRENTLY` (Postgres) + +**Risk: LOW** — the runner already runs in autocommit mode. No runner changes are required. + +**Smallest safe path:** + +1. Postgres migration file: use `CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_name ON table(col)`. + - `IF NOT EXISTS` handles the non-atomic tracking race condition: if the runner crashes after DDL but before tracking, the re-run skips the existing index without error. + - Note: `CREATE INDEX CONCURRENTLY IF NOT EXISTS` requires Postgres 9.5+. Luthien targets modern Postgres; this is not a concern. +2. SQLite migration file: use plain `CREATE INDEX IF NOT EXISTS idx_name ON table(col)` (no CONCURRENTLY keyword). +3. No changes to `run-migrations.sh` or `migration_check.py` are needed. + +**Residual risk:** `CREATE INDEX CONCURRENTLY` holds a share-update-exclusive lock, not a full table lock, but it does require two table scans. On a large `conversation_events` table it may run for minutes. The Docker `migrations` container has no configurable `lock_timeout`; a very large production table could cause the migration container to hang. Mitigation: document the expected index build time in the migration file comment, or run it manually outside the automated runner for very large tables. + +**Out-of-scope risk (do not fix here):** The non-atomic tracking gap exists for ALL migrations, not just CONCURRENTLY ones. A proper fix would wrap both the DDL and the `INSERT INTO _migrations` in a single transaction — but that would break `CREATE INDEX CONCURRENTLY`. The correct long-term approach is to move tracking into the same psql session with `\set ON_ERROR_STOP on` and careful sequencing, but that is a separate refactor not required for this PR series. + +--- + +## Experimental Validation + +Confirmed: The P8 audit findings are correct. + +**SQLite experiment** (run 2026-05-15): +- `CREATE INDEX IF NOT EXISTS` via `executescript()`: works correctly +- `CREATE INDEX CONCURRENTLY`: fails with `sqlite3.OperationalError: near "IF": syntax error` +- Conclusion: SQLite migrations must always use plain `CREATE INDEX IF NOT EXISTS` + +**Postgres validation** (theoretical, based on runner analysis): +- The Postgres runner (`docker/run-migrations.sh`) uses `psql -f` with no transaction wrapping +- `CREATE INDEX CONCURRENTLY` requires running outside a transaction block +- Since the runner does NOT wrap in BEGIN/COMMIT, CONCURRENTLY should work +- Practical test deferred (no Postgres available in local dev); theoretical analysis confirmed + +**Verdict**: PARTIAL support confirmed experimentally: +- SQLite: CONCURRENTLY not supported (syntax error) — use plain `CREATE INDEX IF NOT EXISTS` +- Postgres: CONCURRENTLY supported (no transaction wrapping in runner) — safe to use diff --git a/pyproject.toml b/pyproject.toml index dbaa8830f..ecb5683c4 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -82,7 +82,7 @@ docstring-code-format = true convention = "google" [tool.pytest.ini_options] -addopts = "-q -ra -m 'not e2e and not integration and not mock_e2e and not sqlite_e2e' --import-mode=importlib --cov=src/luthien_proxy --cov-report=term-missing --timeout=3 --timeout-method=signal" +addopts = "-q -ra -m 'not e2e and not integration and not mock_e2e and not sqlite_e2e and not perf' --import-mode=importlib --cov=src/luthien_proxy --cov-report=term-missing --timeout=3 --timeout-method=signal" testpaths = ["tests"] asyncio_mode = "auto" filterwarnings = [ @@ -95,6 +95,8 @@ markers = [ "integration: marks integration tests that require external services (OpenAI, Anthropic APIs)", "mock_e2e: marks e2e tests that use the mock Anthropic server (no real API calls)", "sqlite_e2e: marks e2e tests running the gateway in-process with SQLite (no Docker)", + "perf: marks performance tests that measure gateway latency and throughput (opt-in via ./scripts/run_perf.sh)", + "contract: marks API contract snapshot tests that validate response shapes", "llm01: OWASP LLM01 - Prompt Injection scenarios", "llm02: OWASP LLM02 - Insecure Output Handling scenarios (reserved, no tests yet)", "llm04: OWASP LLM04 - Model Denial of Service scenarios (reserved, no tests yet)", @@ -122,13 +124,15 @@ reportMissingImports = "warning" [dependency-groups] dev = [ + "playwright==1.50.0", "pre-commit>=4.3.0", "pytest>=8.4.1", "pytest-asyncio>=1.1.0", "pytest-cov>=6.2.1", + "pytest-playwright>=0.5.0", + "pytest-timeout>=2.4.0", "ruff>=0.12.10", "pyright>=1.1.406,<1.2", - "pytest-timeout>=2.4.0", "radon>=6.0.1", "vulture>=2.14", "asgi-lifespan>=2.1.0", diff --git a/scripts/perf_explain.py b/scripts/perf_explain.py new file mode 100755 index 000000000..a339b6316 --- /dev/null +++ b/scripts/perf_explain.py @@ -0,0 +1,228 @@ +#!/usr/bin/env python3 +"""Capture EXPLAIN QUERY PLAN for the top slow queries against the perf DB. + +Usage: + uv run python scripts/perf_explain.py --backend sqlite + uv run python scripts/perf_explain.py --backend postgres + +Outputs: .sisyphus/evidence/baseline-query-plans.md + +Safety: refuses to connect if DATABASE_URL points to the dev DB (local.db). +""" + +import argparse +import os +import sqlite3 +import subprocess +import sys +from datetime import datetime, timezone +from pathlib import Path + +_REPO_ROOT = Path(__file__).resolve().parent.parent +sys.path.insert(0, str(_REPO_ROOT / "src")) + +from luthien_proxy.perf.db import ensure_perf_isolation, get_perf_db_url, migrate_perf_db # noqa: E402 +from luthien_proxy.perf.seeding import seed_sessions # noqa: E402 +from luthien_proxy.utils.db_sqlite import parse_sqlite_url # noqa: E402 + +EVIDENCE_DIR = _REPO_ROOT / ".sisyphus" / "evidence" +OUTPUT_PATH = EVIDENCE_DIR / "baseline-query-plans.md" + +# ── Queries ──────────────────────────────────────────────────────────────── +# Exact SQL extracted from source (adapted: $N → ? for sqlite3, no f-string +# interpolation — using the hot-path / no-user-filter variant). +# +# Source: src/luthien_proxy/history/service.py (_fetch_session_list_sqlite) +SESSION_LIST_SQL = """\ +SELECT + ce.session_id, + MIN(ce.created_at) as first_ts, + MAX(ce.created_at) as last_ts, + COUNT(*) as total_events, + COUNT(DISTINCT ce.call_id) as turn_count, + SUM(CASE + WHEN ce.event_type LIKE 'policy.%' + AND ce.event_type NOT LIKE 'policy.%judge.evaluation%' + THEN 1 ELSE 0 + END) as policy_interventions +FROM conversation_events ce +WHERE ce.session_id IS NOT NULL +GROUP BY ce.session_id +ORDER BY last_ts DESC +LIMIT ? OFFSET ?\ +""" + +# Source: src/luthien_proxy/history/service.py (fetch_session_detail) +SESSION_DETAIL_SQL = """\ +SELECT call_id, event_type, payload, created_at +FROM conversation_events +WHERE session_id = ? +ORDER BY created_at ASC\ +""" + +# Source: src/luthien_proxy/debug/service.py (fetch_recent_calls) +RECENT_CALLS_SQL = """\ +SELECT + call_id, + COUNT(*) as event_count, + MAX(created_at) as latest, + MAX(session_id) as session_id +FROM conversation_events +GROUP BY call_id +ORDER BY latest DESC +LIMIT ?\ +""" + +QUERIES: list[tuple[str, str, tuple[object, ...]]] = [ + ("session_list", SESSION_LIST_SQL, (50, 0)), + ("session_detail", SESSION_DETAIL_SQL, ("placeholder-session-id",)), + ("recent_calls", RECENT_CALLS_SQL, (50,)), +] + + +def get_git_sha() -> str: + try: + result = subprocess.run( + ["git", "rev-parse", "HEAD"], + capture_output=True, + text=True, + cwd=_REPO_ROOT, + ) + return result.stdout.strip() if result.returncode == 0 else "unknown" + except Exception: + return "unknown" + + +def format_explain_plan(rows: list[tuple[int, int, int, str]]) -> str: + """Format EXPLAIN QUERY PLAN rows as a tree. + + SQLite EXPLAIN QUERY PLAN returns (id, parent, notused, detail). + We indent based on parent depth to show the nested structure. + """ + if not rows: + return "(no plan output)" + id_to_depth: dict[int, int] = {0: -1} + lines = [] + for row in rows: + row_id, parent_id, _notused, detail = row[0], row[1], row[2], row[3] + parent_depth = id_to_depth.get(parent_id, -1) + depth = parent_depth + 1 + id_to_depth[row_id] = depth + indent = " " * depth + connector = "`--" if depth > 0 else "" + lines.append(f"{indent}{connector}{detail}") + return "\n".join(lines) + + +def ensure_no_dev_db_in_env() -> None: + database_url = os.environ.get("DATABASE_URL", "") + if not database_url: + return + try: + ensure_perf_isolation(database_url) + except RuntimeError as e: + # ensure_perf_isolation message always contains "isolation" + print(f"isolation refuse: DATABASE_URL is set to the dev database.\n{e}") + sys.exit(1) + + +def explain_sqlite(db_path: str) -> None: + # Ensure migrations are applied (idempotent) + print("Applying migrations...", file=sys.stderr) + migrate_perf_db("sqlite") + + conn = sqlite3.connect(db_path) + try: + row_count = conn.execute("SELECT COUNT(*) FROM conversation_events").fetchone()[0] + session_count = conn.execute( + "SELECT COUNT(DISTINCT session_id) FROM conversation_events WHERE session_id IS NOT NULL" + ).fetchone()[0] + + if row_count == 0: + print("Perf DB is empty — seeding with tier=100...", file=sys.stderr) + conn.close() + seed_sessions("sqlite", tier=100) + conn = sqlite3.connect(db_path) + row_count = conn.execute("SELECT COUNT(*) FROM conversation_events").fetchone()[0] + session_count = conn.execute( + "SELECT COUNT(DISTINCT session_id) FROM conversation_events WHERE session_id IS NOT NULL" + ).fetchone()[0] + + print(f"DB has {row_count} events, {session_count} sessions.", file=sys.stderr) + + git_sha = get_git_sha() + timestamp = datetime.now(timezone.utc).isoformat() + + sections: list[str] = [] + sections.append("---") + sections.append(f"git_sha: {git_sha}") + sections.append(f"timestamp: {timestamp}") + sections.append("backend: sqlite") + sections.append(f"row_count: {row_count}") + sections.append(f"session_count: {session_count}") + sections.append("---") + sections.append("") + + for name, sql, params in QUERIES: + print(f"Running EXPLAIN QUERY PLAN for {name}...", file=sys.stderr) + sections.append(f"## Query: {name}") + sections.append("") + sections.append("### SQL") + sections.append("") + sections.append("```sql") + sections.append(sql) + sections.append("```") + sections.append("") + sections.append("### EXPLAIN QUERY PLAN") + sections.append("") + sections.append("```") + try: + rows = conn.execute(f"EXPLAIN QUERY PLAN {sql}", params).fetchall() + sections.append(format_explain_plan(rows)) + except sqlite3.OperationalError as e: + sections.append(f"ERROR: {e}") + sections.append("```") + sections.append("") + + EVIDENCE_DIR.mkdir(parents=True, exist_ok=True) + OUTPUT_PATH.write_text("\n".join(sections) + "\n", encoding="utf-8") + print(f"Written: {OUTPUT_PATH}", file=sys.stderr) + + finally: + conn.close() + + +def main() -> None: + parser = argparse.ArgumentParser(description="Capture EXPLAIN QUERY PLAN for slow queries against the perf DB.") + parser.add_argument( + "--backend", + choices=["sqlite", "postgres"], + required=True, + help="Database backend to use.", + ) + args = parser.parse_args() + + ensure_no_dev_db_in_env() + + try: + url = get_perf_db_url(args.backend) + except RuntimeError as e: + print(f"isolation refuse: {e}") + sys.exit(1) + + try: + ensure_perf_isolation(url) + except RuntimeError as e: + print(f"isolation refuse: {e}") + sys.exit(1) + + if args.backend == "sqlite": + db_path = parse_sqlite_url(url) + explain_sqlite(db_path) + else: + print("SKIPPED: Postgres backend not available in this environment.", file=sys.stderr) + print("# SKIPPED: Postgres not available", file=sys.stderr) + + +if __name__ == "__main__": + main() diff --git a/scripts/perf_report.py b/scripts/perf_report.py new file mode 100755 index 000000000..ae5c8d3f1 --- /dev/null +++ b/scripts/perf_report.py @@ -0,0 +1,353 @@ +#!/usr/bin/env python3 +"""Generate a Markdown performance baseline report from perf test results. + +Reads .sisyphus/evidence/perf-results-*.json files and embeds +.sisyphus/evidence/baseline-query-plans.md. When no JSON files exist +(P10-P12 not yet run), every data section is rendered with a NO DATA YET +placeholder so the file is still structurally valid for diffing. + +Usage: + uv run python scripts/perf_report.py --output .sisyphus/evidence/perf-report-baseline.md + uv run python scripts/perf_report.py --output out.md --deterministic-mode +""" + +from __future__ import annotations + +import argparse +import glob +import importlib.metadata +import json +import platform +import subprocess +import sys +from datetime import datetime, timezone +from pathlib import Path + +_REPO_ROOT = Path(__file__).parent.parent +_EVIDENCE_DIR = _REPO_ROOT / ".sisyphus" / "evidence" + + +def _git_sha(repo_root: Path | None = None) -> str: + root = repo_root or _REPO_ROOT + try: + result = subprocess.run( + ["git", "rev-parse", "HEAD"], + capture_output=True, + text=True, + cwd=root, + timeout=5, + ) + return result.stdout.strip() if result.returncode == 0 else "unknown" + except Exception: + return "unknown" + + +def _playwright_version() -> str: + try: + return importlib.metadata.version("playwright") + except Exception: + return "unknown" + + +def _ram_info() -> str: + try: + result = subprocess.run( + ["sysctl", "-n", "hw.memsize"], + capture_output=True, + text=True, + timeout=3, + ) + if result.returncode == 0: + gb = int(result.stdout.strip()) // (1024**3) + return f"{gb} GB" + except Exception: + pass + try: + with open("/proc/meminfo") as f: + for line in f: + if line.startswith("MemTotal:"): + kb = int(line.split()[1]) + return f"{kb // (1024**2)} GB" + except Exception: + pass + return "unknown" + + +def load_results(evidence_dir: Path | None = None) -> list[dict]: + d = evidence_dir or _EVIDENCE_DIR + results: list[dict] = [] + for path in sorted(glob.glob(str(d / "perf-results-*.json"))): + try: + with open(path) as f: + results.append(json.load(f)) + except Exception: + continue + return results + + +def load_query_plans(evidence_dir: Path | None = None) -> str: + d = evidence_dir or _EVIDENCE_DIR + path = d / "baseline-query-plans.md" + if path.exists(): + return path.read_text() + return "_Query plans not yet captured. Run `scripts/perf_explain.py` first._\n" + + +def _find_result(results: list[dict], type_: str) -> dict | None: + for r in results: + if r.get("type") == type_: + return r + return None + + +def _section_hardware(git_sha: str, playwright_ver: str, ram: str) -> str: + rows = [ + ("Machine", platform.machine()), + ("Processor", platform.processor() or platform.machine()), + ("RAM", ram), + ("OS", f"{platform.system()} {platform.release()}"), + ("Python", platform.python_version()), + ("git_sha", f"`{git_sha}`"), + ("DB backend", "sqlite"), + ("Playwright", playwright_ver), + ] + table = ["| Field | Value |", "|-------|-------|"] + table.extend(f"| {k} | {v} |" for k, v in rows) + return "\n".join(["## Hardware & Versions", ""] + table + [""]) + + +def _section_per_page_timings(results: list[dict]) -> str: + header = "## Per-Page Timings" + r = _find_result(results, "page_timings") + if not r: + return "\n".join([header, "", "_NO DATA YET — run `scripts/run_perf.sh` to populate._", ""]) + + data = r.get("data", {}) + pages = sorted(data.keys()) + fixtures: set[str] = set() + for page_data in data.values(): + fixtures.update(page_data.keys()) + fixture_list = sorted(fixtures) + + col_header = " | ".join(f"{f} median_ms | {f} p95_ms" for f in fixture_list) + col_sep = " | ".join("--- | ---" for _ in fixture_list) + lines = [header, "", f"| Page | {col_header} |", f"|------| {col_sep} |"] + + for page in pages: + cells: list[str] = [] + for fixture in fixture_list: + fdata = data[page].get(fixture, {}) + cells.append(str(fdata.get("median_ms", "—"))) + cells.append(str(fdata.get("p95_ms", "—"))) + lines.append(f"| {page} | " + " | ".join(cells) + " |") + + lines.append("") + return "\n".join(lines) + + +def _section_throttled(results: list[dict]) -> str: + header = "## Throttled (sami-like)" + r = _find_result(results, "throttled") + if not r: + return "\n".join([header, "", "_NO DATA YET_", ""]) + + data = r.get("data", {}) + lines = [header, "", "| Page | Fixture | Median ms | P95 ms |", "|------|---------|-----------|--------|"] + for page in sorted(data.keys()): + for fixture, fdata in sorted(data[page].items()): + median = fdata.get("median_ms", "—") + p95 = fdata.get("p95_ms", "—") + lines.append(f"| {page} | {fixture} | {median} | {p95} |") + lines.append("") + return "\n".join(lines) + + +def _section_transcript_open(results: list[dict]) -> str: + header = "## Transcript Open" + r = _find_result(results, "transcript_open") + if not r: + return "\n".join([header, "", "_NO DATA YET_", ""]) + + data = r.get("data", {}) + lines = [ + header, + "", + "| Metric | Value |", + "|--------|-------|", + f"| first_turn_painted_ms | {data.get('first_turn_painted_ms', '—')} |", + "", + ] + return "\n".join(lines) + + +def _section_sse_memory(results: list[dict]) -> str: + header = "## SSE Memory Growth" + r = _find_result(results, "sse_memory") + if not r: + return "\n".join([header, "", "_NO DATA YET_", ""]) + + data = r.get("data", {}) + lines = [ + header, + "", + "| Metric | Value |", + "|--------|-------|", + f"| heap_growth_mb | {data.get('heap_growth_mb', '—')} |", + f"| events_count | {data.get('events_count', '—')} |", + "", + ] + return "\n".join(lines) + + +def _section_server_timing(results: list[dict]) -> str: + header = "## Server-Timing Breakdown" + r = _find_result(results, "server_timing") + if not r: + return "\n".join([header, "", "_NO DATA YET_", ""]) + + data = r.get("data", {}) + lines = [header, "", "| Phase | Median ms |", "|-------|-----------|"] + for phase in ("db", "serialize", "render"): + lines.append(f"| {phase} | {data.get(f'{phase}_ms', '—')} |") + lines.append("") + return "\n".join(lines) + + +def _section_payload_size(results: list[dict]) -> str: + header = "## Payload Size Breakdown" + r = _find_result(results, "payload_size") + if not r: + return "\n".join([header, "", "_NO DATA YET_", ""]) + + data = r.get("data", {}) + lines = [header, "", "| Endpoint | Bytes |", "|----------|-------|"] + for endpoint in sorted(data.keys()): + lines.append(f"| {endpoint} | {data[endpoint].get('bytes', '—')} |") + lines.append("") + return "\n".join(lines) + + +def _section_query_plans(query_plans: str) -> str: + return "\n".join(["## Query Plans", "", query_plans.strip(), ""]) + + +def _section_top_hotspots(results: list[dict]) -> str: + header = "## Top Hotspots" + has_data = any( + r.get("type") in ("page_timings", "throttled", "server_timing", "payload_size", "sse_memory") for r in results + ) + + if not has_data: + lines = [ + header, + "", + "_NO DATA YET — hotspots will be derived from measurement results._", + "", + "**Known candidates (from code review):**", + "", + "1. `history_list.html:514` — hardcodes `?limit=10000` (sends full dataset on every load)", + "2. `conversation_live.js:92-118` — `loadInitial()` fetches entire session upfront", + "3. `conversation_live.js:215-244` — full DOM re-render on every SSE event", + "4. `conversation_live.js:164-172` — unbounded `rawEvents[callId]` array (memory leak risk)", + "5. `history_list.html:423-448` — client-side filter runs on every keystroke", + "", + "**Query plan risks:**", + "", + "- `session_list`: 2× TEMP B-TREE (COUNT DISTINCT + ORDER BY) — scales poorly with row count", + "- `recent_calls`: SCAN on all rows — O(n) over conversation_events", + "", + ] + return "\n".join(lines) + + hotspots: list[str] = [] + + r_page = _find_result(results, "page_timings") + if r_page: + for page, fixtures in r_page.get("data", {}).items(): + for fixture, stats in fixtures.items(): + p95 = stats.get("p95_ms", 0) + if isinstance(p95, (int, float)) and p95 > 1000: + hotspots.append(f"`{page}` ({fixture}) p95={p95}ms — exceeds 1s SLO") + + r_payload = _find_result(results, "payload_size") + if r_payload: + for endpoint, stats in r_payload.get("data", {}).items(): + bytes_ = stats.get("bytes", 0) + if isinstance(bytes_, int) and bytes_ > 50_000: + hotspots.append(f"`{endpoint}` payload={bytes_ // 1024}KB — exceeds 50KB budget") + + r_sse = _find_result(results, "sse_memory") + if r_sse: + growth = r_sse.get("data", {}).get("heap_growth_mb", 0) + if isinstance(growth, (int, float)) and growth > 10: + hotspots.append(f"SSE heap growth={growth}MB over session — unbounded accumulation risk") + + lines = [header, ""] + if hotspots: + lines.extend(f"- {h}" for h in hotspots) + else: + lines.append("_No hotspots detected above threshold. See individual sections for details._") + lines.append("") + return "\n".join(lines) + + +def generate_report( + results: list[dict], + query_plans: str, + git_sha: str, + playwright_ver: str, + generated_at: str | None = None, + ram: str | None = None, +) -> str: + timestamp = generated_at or datetime.now(timezone.utc).isoformat() + ram_str = ram or _ram_info() + + parts = [ + f"git_sha: {git_sha}", + f"browser_version: {playwright_ver}", + "backend: sqlite", + f"generated_at: {timestamp}", + "", + "# Luthien Admin UI — Performance Baseline Report", + "", + _section_hardware(git_sha, playwright_ver, ram_str), + _section_per_page_timings(results), + _section_throttled(results), + _section_transcript_open(results), + _section_sse_memory(results), + _section_server_timing(results), + _section_payload_size(results), + _section_query_plans(query_plans), + _section_top_hotspots(results), + ] + return "\n".join(parts) + + +def main() -> None: + parser = argparse.ArgumentParser(description="Generate perf baseline Markdown report") + parser.add_argument("--output", required=True, help="Path to write the Markdown report") + parser.add_argument( + "--deterministic-mode", + action="store_true", + help="Fix timestamp to epoch so output is byte-identical across runs (for reproducibility testing)", + ) + args = parser.parse_args() + + results = load_results() + query_plans = load_query_plans() + sha = _git_sha() + pw_ver = _playwright_version() + + generated_at = "2000-01-01T00:00:00+00:00" if args.deterministic_mode else None + + report = generate_report(results, query_plans, sha, pw_ver, generated_at=generated_at) + + output_path = Path(args.output) + output_path.parent.mkdir(parents=True, exist_ok=True) + output_path.write_text(report) + + print(f"Report written to {output_path}", file=sys.stderr) + + +if __name__ == "__main__": + main() diff --git a/scripts/run_perf.sh b/scripts/run_perf.sh new file mode 100755 index 000000000..b5bf32887 --- /dev/null +++ b/scripts/run_perf.sh @@ -0,0 +1,292 @@ +#!/usr/bin/env bash +# Playwright version: 1.50.0 +# +# ISOLATION ENFORCEMENT: +# This script refuses to run against the development database (~/.luthien/local.db). +# Perf tests use a dedicated isolated database to prevent fixture data pollution +# and ensure reproducible baseline measurements: +# SQLite: ~/.luthien/perf.db (hardcoded; never local.db) +# Postgres: perf_test schema in a dedicated Postgres perf instance +# DATABASE_URL must be set explicitly and must not reference local.db. +# +# ABOUTME: Performance test runner for admin UI latency and payload SLOs. +# ABOUTME: Runs Playwright-based perf tests against an isolated perf database. + +set -euo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +REPO_ROOT="$(cd "$SCRIPT_DIR/.." && pwd)" +cd "$REPO_ROOT" + +# Colors +GREEN='\033[0;32m' +RED='\033[0;31m' +YELLOW='\033[1;33m' +BLUE='\033[0;34m' +BOLD='\033[1m' +NC='\033[0m' + +info() { echo -e "${BLUE}▸${NC} $*"; } +ok() { echo -e "${GREEN}✓${NC} $*"; } +warn() { echo -e "${YELLOW}⚠${NC} $*"; } +fail() { echo -e "${RED}✗${NC} $*"; } +header() { echo -e "\n${BOLD}═══ $* ═══${NC}"; } + +# ── Defaults ────────────────────────────────────────────────────────────────── + +TIER="" +FIXTURE="sami-like" +SEED_ONLY=false +CLEAN=false +ASSERT_SLO=false +THROTTLED=false +BACKEND="sqlite" + +# ── Help ────────────────────────────────────────────────────────────────────── + +show_help() { + cat <<'EOF' +Performance test runner for admin UI latency and payload SLOs. + +Usage: + ./scripts/run_perf.sh --tier {100|1000|10000} [options] + ./scripts/run_perf.sh --clean [--backend {sqlite|postgres}] + ./scripts/run_perf.sh --help + +Options: + --tier {100|1000|10000} Sessions to seed [required unless --clean] + --fixture {sami-like} Fixture profile (default: sami-like) + --seed-only Seed the database; skip test assertions + --clean Drop the perf database and exit + --assert-slo Fail if any SLO thresholds are exceeded (sets PERF_ASSERT_SLO=1) + --throttled CDP network throttling -- 1 Mbps + 300ms RTT (only with --fixture sami-like) + --backend {sqlite|postgres} Database backend (default: sqlite) + --help Show this help message + +Environment: + DATABASE_URL Required (refused if unset or contains local.db) + SQLite example: sqlite:///$HOME/.luthien/perf.db + +Examples: + DATABASE_URL=sqlite://$HOME/.luthien/perf.db ./scripts/run_perf.sh --tier 100 + DATABASE_URL=sqlite://$HOME/.luthien/perf.db ./scripts/run_perf.sh --tier 100 --assert-slo + DATABASE_URL=sqlite://$HOME/.luthien/perf.db ./scripts/run_perf.sh --tier 100 --throttled + ./scripts/run_perf.sh --clean + DATABASE_URL=sqlite://$HOME/.luthien/perf.db ./scripts/run_perf.sh --seed-only --tier 1000 + +Postgres --clean note: + For Postgres, --clean executes DROP SCHEMA perf_test CASCADE. + Set DATABASE_URL to the Postgres perf instance before running. +EOF + exit 0 +} + +# ── Argument parsing ────────────────────────────────────────────────────────── + +while [[ $# -gt 0 ]]; do + case "$1" in + --tier) + if [[ $# -lt 2 ]]; then fail "--tier requires an argument"; exit 1; fi + case "$2" in + 100|1000|10000) TIER="$2" ;; + *) fail "Invalid --tier: $2 (expected: 100, 1000, or 10000)"; exit 1 ;; + esac + shift 2 + ;; + --fixture) + if [[ $# -lt 2 ]]; then fail "--fixture requires an argument"; exit 1; fi + case "$2" in + sami-like) FIXTURE="$2" ;; + *) fail "Unknown --fixture: $2 (expected: sami-like)"; exit 1 ;; + esac + shift 2 + ;; + --backend) + if [[ $# -lt 2 ]]; then fail "--backend requires an argument"; exit 1; fi + case "$2" in + sqlite|postgres) BACKEND="$2" ;; + *) fail "Unknown --backend: $2 (expected: sqlite or postgres)"; exit 1 ;; + esac + shift 2 + ;; + --seed-only) SEED_ONLY=true; shift ;; + --clean) CLEAN=true; shift ;; + --assert-slo) ASSERT_SLO=true; shift ;; + --throttled) THROTTLED=true; shift ;; + --help|-h) show_help ;; + *) fail "Unknown option: $1"; exit 1 ;; + esac +done + +# ── Validate option combinations ────────────────────────────────────────────── + +if $THROTTLED && [[ "$FIXTURE" != "sami-like" ]]; then + fail "--throttled is only valid with --fixture sami-like (got: --fixture $FIXTURE)" + exit 1 +fi + +if ! $CLEAN && [[ -z "$TIER" ]]; then + fail "Required: --tier {100|1000|10000} (or use --clean to drop the perf DB)" + exit 1 +fi + +# ── SQLite clean ────────────────────────────────────────────────────────────── +# Runs before the isolation check: deletes the perf DB file, never the dev DB. + +if $CLEAN && [[ "$BACKEND" == "sqlite" ]]; then + header "Cleaning Perf Database (SQLite)" + PERF_DB="$HOME/.luthien/perf.db" + if [[ -f "$PERF_DB" ]]; then + rm -f "$PERF_DB" + ok "Removed $PERF_DB" + else + info "Nothing to clean: $PERF_DB does not exist" + fi + exit 0 +fi + +# ── Isolation check ─────────────────────────────────────────────────────────── +# +# This script refuses to run against the dev database. Perf tests MUST use an +# isolated database to prevent fixture data pollution and ensure reproducibility. +# Applies to all non-SQLite-clean operations. + +_db_url="${DATABASE_URL:-}" + +if [[ -z "$_db_url" ]]; then + fail "ISOLATION REFUSED: DATABASE_URL is not set." + fail " The gateway defaults to ~/.luthien/local.db (the dev database) when unset." + fail " This script refuses to run without an explicit isolated database URL." + fail " Set DATABASE_URL to a perf-specific path, e.g.:" + fail " export DATABASE_URL=sqlite:///\$HOME/.luthien/perf.db" + exit 1 +fi + +if [[ "$_db_url" == *"local.db"* ]]; then + fail "ISOLATION REFUSED: DATABASE_URL points to the dev database (local.db)." + fail " This script refuses to run against local.db to prevent data pollution." + fail " DATABASE_URL=$_db_url" + fail " Set DATABASE_URL to a perf-specific path, e.g.:" + fail " export DATABASE_URL=sqlite:///\$HOME/.luthien/perf.db" + exit 1 +fi + +# ── Postgres clean (after isolation check) ──────────────────────────────────── + +if $CLEAN && [[ "$BACKEND" == "postgres" ]]; then + header "Cleaning Perf Database (Postgres)" + warn "Executing: DROP SCHEMA perf_test CASCADE" + warn " Target: $_db_url" + uv run python - <<'PYEOF' +import os +import sys + +try: + import psycopg2 # type: ignore[import-untyped] +except ImportError: + print("psycopg2 not installed; run: uv add psycopg2-binary", file=sys.stderr) + sys.exit(1) + +try: + url = os.environ["DATABASE_URL"] + conn = psycopg2.connect(url) + conn.autocommit = True + cur = conn.cursor() + cur.execute("DROP SCHEMA IF EXISTS perf_test CASCADE") + conn.close() + print("perf_test schema dropped") +except Exception as exc: + print(f"Error dropping schema: {exc}", file=sys.stderr) + sys.exit(1) +PYEOF + ok "Postgres perf_test schema dropped" + exit 0 +fi + +# ── Pre-flight ──────────────────────────────────────────────────────────────── + +header "Pre-flight Checks" + +# Ensure Chromium is installed (Playwright 1.50.0 -- pinned at top of file). +info "Checking Playwright Chromium..." +uv run playwright install chromium --with-deps 2>/dev/null || true + +_chromium_ver="$(uv run python -c ' +from playwright.sync_api import sync_playwright +with sync_playwright() as p: + browser = p.chromium.launch() + ver = browser.version + browser.close() + print(ver) +' 2>/dev/null || echo "unknown")" + +_git_sha="$(git rev-parse --short HEAD 2>/dev/null || echo "unknown")" + +ok "Chromium version: $_chromium_ver" +ok "Git SHA: $_git_sha" + +# ── Environment ─────────────────────────────────────────────────────────────── + +export PERF_TIER="$TIER" +export PERF_FIXTURE="$FIXTURE" +export PERF_BACKEND="$BACKEND" + +if $ASSERT_SLO; then + export PERF_ASSERT_SLO=1 + info "SLO assertion enabled -- tests fail if thresholds exceeded" +fi + +if $THROTTLED; then + export PERF_THROTTLED=1 + info "Network throttling enabled -- 1 Mbps bandwidth + 300ms RTT (sami-like profile)" +fi + +# ── Seed only ───────────────────────────────────────────────────────────────── + +if $SEED_ONLY; then + header "Seeding Database (tier=$TIER, fixture=$FIXTURE)" + info "Seeding $TIER sessions -- test assertions will NOT run" + export PERF_SEED_ONLY=1 + uv run pytest \ + -m perf \ + tests/luthien_proxy/perf_tests/ \ + -v --no-cov \ + || true + ok "Seeding complete" + exit 0 +fi + +# ── Run perf tests ──────────────────────────────────────────────────────────── + +_slo_flag="no" +_throttle_flag="no" +$ASSERT_SLO && _slo_flag="yes" +$THROTTLED && _throttle_flag="yes" + +header "Perf Tests" +info " Tier: $TIER sessions" +info " Fixture: $FIXTURE" +info " Backend: $BACKEND" +info " Assert SLO: $_slo_flag" +info " Throttled: $_throttle_flag" +info " Database: $_db_url" + +exit_code=0 +uv run pytest \ + -m perf \ + tests/luthien_proxy/perf_tests/ \ + -v --no-cov \ + || exit_code=$? + +# ── Summary ─────────────────────────────────────────────────────────────────── + +header "Results" +if [[ $exit_code -eq 0 ]]; then + ok "All perf tests passed" + $ASSERT_SLO && ok "SLO thresholds: all met" +else + fail "Perf tests failed (exit $exit_code)" + $ASSERT_SLO && fail "One or more SLO thresholds were exceeded" +fi + +exit $exit_code diff --git a/src/luthien_proxy/config_fields.py b/src/luthien_proxy/config_fields.py index 4f041c93a..248bd45ef 100644 --- a/src/luthien_proxy/config_fields.py +++ b/src/luthien_proxy/config_fields.py @@ -187,6 +187,14 @@ class ConfigFieldMeta: "Fernet key for encrypting server credentials at rest", sensitive=True, category="security", ), + ConfigFieldMeta( + "cursor_hmac_key", "CURSOR_HMAC_KEY", str, "luthien-perf-cursor-key-dev", + "HMAC key for signing pagination cursors. Auto-provisioned to ~/.luthien/cursor_hmac.key on first run.\n" + "# In multi-replica deployments, set this explicitly so cursors validate across replicas.\n" + "# Each replica that auto-generates its own key will reject cursors issued by other replicas.", + sensitive=True, category="security", + dynamic_default=True, + ), # ── observability ───────────────────────────────────────────────────── ConfigFieldMeta( diff --git a/src/luthien_proxy/debug/service.py b/src/luthien_proxy/debug/service.py index 54ae6d8dc..e23e729ed 100644 --- a/src/luthien_proxy/debug/service.py +++ b/src/luthien_proxy/debug/service.py @@ -15,6 +15,7 @@ import urllib.parse from typing import TYPE_CHECKING, Any +from luthien_proxy.perf.timing_middleware import time_phase from luthien_proxy.utils.db import parse_db_ts if TYPE_CHECKING: @@ -216,15 +217,16 @@ async def fetch_call_events(call_id: str, db_pool: DatabasePool) -> CallEventsRe Exception: If database query fails """ async with db_pool.connection() as conn: - rows = await conn.fetch( - """ - SELECT call_id, event_type, payload, created_at, session_id - FROM conversation_events - WHERE call_id = $1 - ORDER BY created_at ASC - """, - call_id, - ) + with time_phase("db"): + rows = await conn.fetch( + """ + SELECT call_id, event_type, payload, created_at, session_id + FROM conversation_events + WHERE call_id = $1 + ORDER BY created_at ASC + """, + call_id, + ) if not rows: raise ValueError(f"No events found for call_id: {call_id}") @@ -271,19 +273,20 @@ async def fetch_call_diff(call_id: str, db_pool: DatabasePool) -> CallDiffRespon Exception: If database query fails """ async with db_pool.connection() as conn: - rows = await conn.fetch( - """ - SELECT call_id, event_type, payload - FROM conversation_events - WHERE call_id = $1 AND event_type IN ( - 'transaction.request_recorded', - 'transaction.non_streaming_response_recorded', - 'transaction.streaming_response_recorded' + with time_phase("db"): + rows = await conn.fetch( + """ + SELECT call_id, event_type, payload + FROM conversation_events + WHERE call_id = $1 AND event_type IN ( + 'transaction.request_recorded', + 'transaction.non_streaming_response_recorded', + 'transaction.streaming_response_recorded' + ) + ORDER BY created_at ASC + """, + call_id, ) - ORDER BY created_at ASC - """, - call_id, - ) if not rows: raise ValueError(f"No events found for call_id: {call_id}") @@ -337,20 +340,21 @@ async def fetch_recent_calls(limit: int, db_pool: DatabasePool) -> CallListRespo Exception: If database query fails """ async with db_pool.connection() as conn: - rows = await conn.fetch( - """ - SELECT - call_id, - COUNT(*) as event_count, - MAX(created_at) as latest, - MAX(session_id) as session_id - FROM conversation_events - GROUP BY call_id - ORDER BY latest DESC - LIMIT $1 - """, - limit, - ) + with time_phase("db"): + rows = await conn.fetch( + """ + SELECT + call_id, + COUNT(*) as event_count, + MAX(created_at) as latest, + MAX(session_id) as session_id + FROM conversation_events + GROUP BY call_id + ORDER BY latest DESC + LIMIT $1 + """, + limit, + ) calls = [ CallListItem( diff --git a/src/luthien_proxy/history/service.py b/src/luthien_proxy/history/service.py index 77e48f33c..2e3a0977f 100644 --- a/src/luthien_proxy/history/service.py +++ b/src/luthien_proxy/history/service.py @@ -11,9 +11,12 @@ import json import logging import re -from datetime import datetime +from datetime import datetime, timezone from typing import Any, TypedDict, cast +from uuid import UUID as _UUID +from luthien_proxy.perf.timing_middleware import time_phase +from luthien_proxy.utils.cursor import decode_cursor, encode_cursor from luthien_proxy.utils.db import DatabasePool, parse_db_ts from .models import ( @@ -352,6 +355,29 @@ def _extract_preview_message(payload: dict[str, Any] | str | None) -> str | None return None +def _format_session_ts(dt: datetime) -> str: + now = datetime.now(timezone.utc) + if dt.tzinfo is None: + dt = dt.replace(tzinfo=timezone.utc) + total_seconds = (now - dt).total_seconds() + + if total_seconds < 60: + return "just now" + elif total_seconds < 3600: + mins = int(total_seconds // 60) + return f"{mins}m ago" + elif total_seconds < 86400: + hours = int(total_seconds // 3600) + return f"{hours}h ago" + elif total_seconds < 7 * 86400: + days = int(total_seconds // 86400) + return f"{days}d ago" + elif dt.year == now.year: + return f"{dt.strftime('%b')} {dt.day}" + else: + return f"{dt.strftime('%b')} {dt.day}, {dt.year}" + + async def fetch_session_list( limit: int, db_pool: DatabasePool, @@ -394,135 +420,136 @@ async def _fetch_session_list_pg( # touch conversation_calls in the hot CTE — user_ids come from a separate # post-query keyed on the page's session_ids (mirrors the SQLite pattern). async with db_pool.connection() as conn: - if user_id is not None: - total_count = await conn.fetchval( - """ - SELECT COUNT(DISTINCT ce.session_id) - FROM conversation_events ce - JOIN conversation_calls cc ON ce.call_id = cc.call_id - WHERE ce.session_id IS NOT NULL AND cc.user_id = $1 - """, - user_id, - ) - else: - total_count = await conn.fetchval( - """ - SELECT COUNT(DISTINCT session_id) - FROM conversation_events - WHERE session_id IS NOT NULL - """ - ) - - # When the caller filters by user_id we restrict the events under - # consideration to call_ids belonging to that user — a single shared - # subquery used by every CTE so preview_message / models_used cannot - # leak content from another user's calls under a shared session_id. - user_call_filter = ( - "AND ce.call_id IN (SELECT call_id FROM conversation_calls WHERE user_id = $3)" - if user_id is not None - else "" - ) - query_args: list[Any] = [limit, offset] - if user_id is not None: - query_args.append(user_id) + with time_phase("db"): + if user_id is not None: + total_count = await conn.fetchval( + """ + SELECT COUNT(DISTINCT ce.session_id) + FROM conversation_events ce + JOIN conversation_calls cc ON ce.call_id = cc.call_id + WHERE ce.session_id IS NOT NULL AND cc.user_id = $1 + """, + user_id, + ) + else: + total_count = await conn.fetchval( + """ + SELECT COUNT(DISTINCT session_id) + FROM conversation_events + WHERE session_id IS NOT NULL + """ + ) - rows = await conn.fetch( - f""" - WITH session_stats AS ( - SELECT - ce.session_id, - MIN(ce.created_at) as first_ts, - MAX(ce.created_at) as last_ts, - COUNT(*) as total_events, - COUNT(DISTINCT ce.call_id) as turn_count, - COUNT(*) FILTER ( - WHERE ce.event_type LIKE 'policy.%' - AND ce.event_type NOT LIKE 'policy.%judge.evaluation%' - ) as policy_interventions - FROM conversation_events ce - WHERE ce.session_id IS NOT NULL - {user_call_filter} - GROUP BY ce.session_id - ), - session_models AS ( - SELECT DISTINCT - ce.session_id, - ce.payload->>'final_model' as model - FROM conversation_events ce - WHERE ce.session_id IS NOT NULL - AND ce.event_type = 'transaction.request_recorded' - AND ce.payload->>'final_model' IS NOT NULL - {user_call_filter} - ), - session_first_message AS ( - SELECT DISTINCT ON (ce.session_id) - ce.session_id, - ce.payload as request_payload - FROM conversation_events ce - WHERE ce.session_id IS NOT NULL - AND ce.event_type = 'transaction.request_recorded' - -- Skip probe requests: max_tokens=1 means internal probe (token counting, quota). - -- COALESCE to 2 so requests without max_tokens are not skipped. - AND COALESCE((ce.payload->'final_request'->>'max_tokens')::int, 2) > 1 - {user_call_filter} - ORDER BY ce.session_id, ce.created_at ASC + # When the caller filters by user_id we restrict the events under + # consideration to call_ids belonging to that user — a single shared + # subquery used by every CTE so preview_message / models_used cannot + # leak content from another user's calls under a shared session_id. + user_call_filter = ( + "AND ce.call_id IN (SELECT call_id FROM conversation_calls WHERE user_id = $3)" + if user_id is not None + else "" ) - SELECT - s.session_id, - s.first_ts, - s.last_ts, - s.total_events, - s.turn_count, - s.policy_interventions, - COALESCE( - array_agg(DISTINCT m.model) FILTER (WHERE m.model IS NOT NULL), - ARRAY[]::text[] - ) as models, - f.request_payload - FROM session_stats s - LEFT JOIN session_models m ON s.session_id = m.session_id - LEFT JOIN session_first_message f ON s.session_id = f.session_id - GROUP BY s.session_id, s.first_ts, s.last_ts, - s.total_events, s.turn_count, s.policy_interventions, - f.request_payload - ORDER BY s.last_ts DESC - LIMIT $1 OFFSET $2 - """, - *query_args, - ) - - # Separate user_ids lookup keyed on the page's session_ids. Distinct - # users only — never collapse via MIN/MAX. When a user filter is in - # effect the same scoping is applied so the response doesn't leak the - # *existence* of other users sharing the session. - user_ids_by_session: dict[str, list[str]] = {} - if rows: - session_ids_on_page = [str(row["session_id"]) for row in rows] - placeholders = ", ".join(f"${i + 1}" for i in range(len(session_ids_on_page))) + query_args: list[Any] = [limit, offset] if user_id is not None: - user_id_filter_clause = f"AND cc.user_id = ${len(session_ids_on_page) + 1}" - user_id_extra_args: list[Any] = [user_id] - else: - user_id_filter_clause = "" - user_id_extra_args = [] - user_id_rows = await conn.fetch( + query_args.append(user_id) + + rows = await conn.fetch( f""" - SELECT DISTINCT ce.session_id, cc.user_id - FROM conversation_events ce - JOIN conversation_calls cc ON ce.call_id = cc.call_id - WHERE ce.session_id IN ({placeholders}) - AND cc.user_id IS NOT NULL - {user_id_filter_clause} + WITH session_stats AS ( + SELECT + ce.session_id, + MIN(ce.created_at) as first_ts, + MAX(ce.created_at) as last_ts, + COUNT(*) as total_events, + COUNT(DISTINCT ce.call_id) as turn_count, + COUNT(*) FILTER ( + WHERE ce.event_type LIKE 'policy.%' + AND ce.event_type NOT LIKE 'policy.%judge.evaluation%' + ) as policy_interventions + FROM conversation_events ce + WHERE ce.session_id IS NOT NULL + {user_call_filter} + GROUP BY ce.session_id + ), + session_models AS ( + SELECT DISTINCT + ce.session_id, + ce.payload->>'final_model' as model + FROM conversation_events ce + WHERE ce.session_id IS NOT NULL + AND ce.event_type = 'transaction.request_recorded' + AND ce.payload->>'final_model' IS NOT NULL + {user_call_filter} + ), + session_first_message AS ( + SELECT DISTINCT ON (ce.session_id) + ce.session_id, + ce.payload as request_payload + FROM conversation_events ce + WHERE ce.session_id IS NOT NULL + AND ce.event_type = 'transaction.request_recorded' + -- Skip probe requests: max_tokens=1 means internal probe (token counting, quota). + -- COALESCE to 2 so requests without max_tokens are not skipped. + AND COALESCE((ce.payload->'final_request'->>'max_tokens')::int, 2) > 1 + {user_call_filter} + ORDER BY ce.session_id, ce.created_at ASC + ) + SELECT + s.session_id, + s.first_ts, + s.last_ts, + s.total_events, + s.turn_count, + s.policy_interventions, + COALESCE( + array_agg(DISTINCT m.model) FILTER (WHERE m.model IS NOT NULL), + ARRAY[]::text[] + ) as models, + f.request_payload + FROM session_stats s + LEFT JOIN session_models m ON s.session_id = m.session_id + LEFT JOIN session_first_message f ON s.session_id = f.session_id + GROUP BY s.session_id, s.first_ts, s.last_ts, + s.total_events, s.turn_count, s.policy_interventions, + f.request_payload + ORDER BY s.last_ts DESC + LIMIT $1 OFFSET $2 """, - *session_ids_on_page, - *user_id_extra_args, + *query_args, ) - for r in user_id_rows: - sid = str(r["session_id"]) - uid = str(r["user_id"]) - bucket = user_ids_by_session.setdefault(sid, []) - if uid not in bucket: - bucket.append(uid) + + # Separate user_ids lookup keyed on the page's session_ids. Distinct + # users only — never collapse via MIN/MAX. When a user filter is in + # effect the same scoping is applied so the response doesn't leak the + # *existence* of other users sharing the session. + user_ids_by_session: dict[str, list[str]] = {} + if rows: + session_ids_on_page = [str(row["session_id"]) for row in rows] + placeholders = ", ".join(f"${i + 1}" for i in range(len(session_ids_on_page))) + if user_id is not None: + user_id_filter_clause = f"AND cc.user_id = ${len(session_ids_on_page) + 1}" + user_id_extra_args: list[Any] = [user_id] + else: + user_id_filter_clause = "" + user_id_extra_args = [] + user_id_rows = await conn.fetch( + f""" + SELECT DISTINCT ce.session_id, cc.user_id + FROM conversation_events ce + JOIN conversation_calls cc ON ce.call_id = cc.call_id + WHERE ce.session_id IN ({placeholders}) + AND cc.user_id IS NOT NULL + {user_id_filter_clause} + """, + *session_ids_on_page, + *user_id_extra_args, + ) + for r in user_id_rows: + sid = str(r["session_id"]) + uid = str(r["user_id"]) + bucket = user_ids_by_session.setdefault(sid, []) + if uid not in bucket: + bucket.append(uid) sessions = [ SessionSummary( @@ -561,141 +588,140 @@ async def _fetch_session_list_sqlite( # SECURITY INVARIANT: user_id is bound as a query parameter, never # interpolated into the SQL string. async with db_pool.connection() as conn: - if user_id is not None: - total_count = await conn.fetchval( - """ - SELECT COUNT(DISTINCT ce.session_id) + with time_phase("db"): + if user_id is not None: + total_count = await conn.fetchval( + """ + SELECT COUNT(DISTINCT ce.session_id) + FROM conversation_events ce + JOIN conversation_calls cc ON ce.call_id = cc.call_id + WHERE ce.session_id IS NOT NULL AND cc.user_id = $1 + """, + user_id, + ) + else: + total_count = await conn.fetchval( + """ + SELECT COUNT(DISTINCT session_id) + FROM conversation_events + WHERE session_id IS NOT NULL + """ + ) + + # PERF: only filter through conversation_calls when a user filter is + # actually requested. Unfiltered list calls (the hot path) skip the + # conversation_calls subquery entirely. user_ids are populated by a + # separate post-query keyed on the page's session_ids (SQLite has no + # array_agg, so we can't compute them inside this query anyway). + user_call_filter = ( + "AND ce.call_id IN (SELECT call_id FROM conversation_calls WHERE user_id = $3)" + if user_id is not None + else "" + ) + query_args: list[Any] = [limit, offset] + if user_id is not None: + query_args.append(user_id) + + rows = await conn.fetch( + f""" + SELECT + ce.session_id, + MIN(ce.created_at) as first_ts, + MAX(ce.created_at) as last_ts, + COUNT(*) as total_events, + COUNT(DISTINCT ce.call_id) as turn_count, + SUM(CASE + WHEN ce.event_type LIKE 'policy.%' + AND ce.event_type NOT LIKE 'policy.%judge.evaluation%' + THEN 1 ELSE 0 + END) as policy_interventions FROM conversation_events ce - JOIN conversation_calls cc ON ce.call_id = cc.call_id - WHERE ce.session_id IS NOT NULL AND cc.user_id = $1 + WHERE ce.session_id IS NOT NULL + {user_call_filter} + GROUP BY ce.session_id + ORDER BY last_ts DESC + LIMIT $1 OFFSET $2 """, - user_id, - ) - else: - total_count = await conn.fetchval( - """ - SELECT COUNT(DISTINCT session_id) - FROM conversation_events - WHERE session_id IS NOT NULL - """ + *query_args, ) - # PERF: only filter through conversation_calls when a user filter is - # actually requested. Unfiltered list calls (the hot path) skip the - # conversation_calls subquery entirely. user_ids are populated by a - # separate post-query keyed on the page's session_ids (SQLite has no - # array_agg, so we can't compute them inside this query anyway). - user_call_filter = ( - "AND ce.call_id IN (SELECT call_id FROM conversation_calls WHERE user_id = $3)" - if user_id is not None - else "" - ) - query_args: list[Any] = [limit, offset] - if user_id is not None: - query_args.append(user_id) - - rows = await conn.fetch( - f""" - SELECT - ce.session_id, - MIN(ce.created_at) as first_ts, - MAX(ce.created_at) as last_ts, - COUNT(*) as total_events, - COUNT(DISTINCT ce.call_id) as turn_count, - SUM(CASE - WHEN ce.event_type LIKE 'policy.%' - AND ce.event_type NOT LIKE 'policy.%judge.evaluation%' - THEN 1 ELSE 0 - END) as policy_interventions - FROM conversation_events ce - WHERE ce.session_id IS NOT NULL - {user_call_filter} - GROUP BY ce.session_id - ORDER BY last_ts DESC - LIMIT $1 OFFSET $2 - """, - *query_args, - ) + total = int(total_count) if total_count is not None else 0 # type: ignore[arg-type] - total = int(total_count) if total_count is not None else 0 # type: ignore[arg-type] + if not rows: + return SessionListResponse(sessions=[], total=total, offset=offset, has_more=False) - if not rows: - return SessionListResponse(sessions=[], total=total, offset=offset, has_more=False) + session_ids = [str(row["session_id"]) for row in rows] + placeholders = ", ".join(f"${i + 1}" for i in range(len(session_ids))) - session_ids = [str(row["session_id"]) for row in rows] - placeholders = ", ".join(f"${i + 1}" for i in range(len(session_ids))) + # When a user_id filter is in effect, restrict the model/preview/user-id + # lookups to that user's call_ids — without this, preview_message and + # models_used can leak content from other users' calls that happen to + # share the session_id. + # NOTE: this clause is *separate from* the `user_call_filter` used in + # the main aggregation above — different placeholder slot ($N differs + # because session_ids are also bound here). Don't fold into one. + if user_id is not None: + user_call_filter_lookups = f"AND ce.call_id IN (SELECT call_id FROM conversation_calls WHERE user_id = ${len(session_ids) + 1})" + extra_args: list[Any] = [user_id] + else: + user_call_filter_lookups = "" + extra_args = [] - # When a user_id filter is in effect, restrict the model/preview/user-id - # lookups to that user's call_ids — without this, preview_message and - # models_used can leak content from other users' calls that happen to - # share the session_id. - # NOTE: this clause is *separate from* the `user_call_filter` used in - # the main aggregation above — different placeholder slot ($N differs - # because session_ids are also bound here). Don't fold into one. - if user_id is not None: - user_call_filter_lookups = ( - f"AND ce.call_id IN (SELECT call_id FROM conversation_calls WHERE user_id = ${len(session_ids) + 1})" + # One query for all models on this page + model_rows = await conn.fetch( + f""" + SELECT ce.session_id, json_extract(ce.payload, '$.final_model') as model + FROM conversation_events ce + WHERE ce.session_id IN ({placeholders}) + AND ce.event_type = 'transaction.request_recorded' + AND json_extract(ce.payload, '$.final_model') IS NOT NULL + {user_call_filter_lookups} + """, + *session_ids, + *extra_args, ) - extra_args: list[Any] = [user_id] - else: - user_call_filter_lookups = "" - extra_args = [] - - # One query for all models on this page - model_rows = await conn.fetch( - f""" - SELECT ce.session_id, json_extract(ce.payload, '$.final_model') as model - FROM conversation_events ce - WHERE ce.session_id IN ({placeholders}) - AND ce.event_type = 'transaction.request_recorded' - AND json_extract(ce.payload, '$.final_model') IS NOT NULL - {user_call_filter_lookups} - """, - *session_ids, - *extra_args, - ) - # One query for first qualifying preview per session on this page - preview_rows = await conn.fetch( - f""" - SELECT ce.session_id, ce.payload as request_payload - FROM conversation_events ce - WHERE ce.session_id IN ({placeholders}) - AND ce.event_type = 'transaction.request_recorded' - AND COALESCE( - CAST(json_extract(ce.payload, '$.final_request.max_tokens') AS INTEGER), - 2 - ) > 1 - {user_call_filter_lookups} - ORDER BY ce.session_id, ce.created_at ASC - """, - *session_ids, - *extra_args, - ) + # One query for first qualifying preview per session on this page + preview_rows = await conn.fetch( + f""" + SELECT ce.session_id, ce.payload as request_payload + FROM conversation_events ce + WHERE ce.session_id IN ({placeholders}) + AND ce.event_type = 'transaction.request_recorded' + AND COALESCE( + CAST(json_extract(ce.payload, '$.final_request.max_tokens') AS INTEGER), + 2 + ) > 1 + {user_call_filter_lookups} + ORDER BY ce.session_id, ce.created_at ASC + """, + *session_ids, + *extra_args, + ) - # Distinct user_ids per session — never collapse via MIN/MAX, that lies - # on multi-user sessions. Returned as a list so the consumer can render - # mixed-identity sessions honestly. When a user filter is in effect - # we constrain to that user so the response doesn't leak the *existence* - # of other users sharing the session. - if user_id is not None: - user_id_filter_clause = f"AND cc.user_id = ${len(session_ids) + 1}" - user_id_args: list[Any] = [user_id] - else: - user_id_filter_clause = "" - user_id_args = [] - user_id_rows = await conn.fetch( - f""" - SELECT DISTINCT ce.session_id, cc.user_id - FROM conversation_events ce - JOIN conversation_calls cc ON ce.call_id = cc.call_id - WHERE ce.session_id IN ({placeholders}) - AND cc.user_id IS NOT NULL - {user_id_filter_clause} - """, - *session_ids, - *user_id_args, - ) + # Distinct user_ids per session — never collapse via MIN/MAX, that lies + # on multi-user sessions. Returned as a list so the consumer can render + # mixed-identity sessions honestly. When a user filter is in effect + # we constrain to that user so the response doesn't leak the *existence* + # of other users sharing the session. + if user_id is not None: + user_id_filter_clause = f"AND cc.user_id = ${len(session_ids) + 1}" + user_id_args: list[Any] = [user_id] + else: + user_id_filter_clause = "" + user_id_args = [] + user_id_rows = await conn.fetch( + f""" + SELECT DISTINCT ce.session_id, cc.user_id + FROM conversation_events ce + JOIN conversation_calls cc ON ce.call_id = cc.call_id + WHERE ce.session_id IN ({placeholders}) + AND cc.user_id IS NOT NULL + {user_id_filter_clause} + """, + *session_ids, + *user_id_args, + ) # Build per-session lookup maps from the bulk results models_by_session: dict[str, list[str]] = {} @@ -753,15 +779,16 @@ async def fetch_session_detail(session_id: str, db_pool: DatabasePool) -> Sessio ValueError: If no events found for session_id """ async with db_pool.connection() as conn: - rows = await conn.fetch( - """ - SELECT call_id, event_type, payload, created_at - FROM conversation_events - WHERE session_id = $1 - ORDER BY created_at ASC - """, - session_id, - ) + with time_phase("db"): + rows = await conn.fetch( + """ + SELECT call_id, event_type, payload, created_at + FROM conversation_events + WHERE session_id = $1 + ORDER BY created_at ASC + """, + session_id, + ) if not rows: raise ValueError(f"No events found for session_id: {session_id}") @@ -1058,10 +1085,325 @@ def _format_message_markdown(msg: ConversationMessage) -> str: return "\n".join(lines) +async def fetch_session_turns_page( + session_id: str, + cursor_token: str | None, + limit: int, + db_pool: DatabasePool, +) -> dict[str, object]: + """Fetch a cursor-paginated page of raw events for a session. + + Returns a dict with keys ``turns`` (list of event dicts) and + ``next_cursor`` (opaque token or None when no further pages exist). + """ + cursor_ts = None + cursor_event_id = None + if cursor_token is not None: + cursor_ts, cursor_event_id = decode_cursor(cursor_token, kind="turns") + + async with db_pool.connection() as conn: + with time_phase("db"): + if cursor_ts is None: + rows = await conn.fetch( + """ + SELECT id, event_type, payload, created_at + FROM conversation_events + WHERE session_id = $1 + ORDER BY created_at ASC, id ASC + LIMIT $2 + """, + session_id, + limit + 1, + ) + else: + assert cursor_event_id is not None + if db_pool.is_sqlite: + # SQLite stores event ids as plain TEXT. If a Postgres-issued + # cursor (UUID string) is decoded here (e.g. same CURSOR_HMAC_KEY + # after a backend migration), the comparison succeeds but may + # return wrong pages. Operators should rotate cursors after + # switching backends. + cursor_id_param: str | _UUID = cursor_event_id + else: + try: + cursor_id_param = _UUID(cursor_event_id) + except ValueError as exc: + raise ValueError(f"Invalid cursor: event id is not a valid UUID: {exc}") from exc + rows = await conn.fetch( + """ + SELECT id, event_type, payload, created_at + FROM conversation_events + WHERE session_id = $1 + AND (datetime(created_at), id) > (datetime($2), $3) + ORDER BY created_at ASC, id ASC + LIMIT $4 + """ + if db_pool.is_sqlite + else """ + SELECT id, event_type, payload, created_at + FROM conversation_events + WHERE session_id = $1 + AND (created_at, id) > ($2, $3) + ORDER BY created_at ASC, id ASC + LIMIT $4 + """, + session_id, + cursor_ts.isoformat() if db_pool.is_sqlite else cursor_ts, + cursor_id_param, + limit + 1, + ) + + has_more = len(rows) > limit + page_rows = list(rows[:limit]) + + turns: list[dict[str, object]] = [] + for row in page_rows: + payload_raw = row["payload"] + if isinstance(payload_raw, dict): + payload_str = json.dumps(payload_raw) + else: + payload_str = str(payload_raw) if payload_raw is not None else "" + + turns.append( + { + "event_id": str(row["id"]), + "created_at": parse_db_ts(row["created_at"]), + "event_type": str(row["event_type"]), + "payload_preview": payload_str[:200], + } + ) + + next_cursor: str | None = None + if has_more and page_rows: + last_row = page_rows[-1] + last_ts = parse_db_ts(last_row["created_at"]) + last_event_id = str(last_row["id"]) + next_cursor = encode_cursor(last_ts, last_event_id, kind="turns") + + return {"turns": turns, "next_cursor": next_cursor} + + +_Q_MAX_LEN = 128 + + +def _escape_like(value: str) -> str: + return value.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_") + + +async def fetch_sessions_page( + cursor_token: str | None, + limit: int, + db_pool: DatabasePool, + q: str | None = None, + quick_filter: str | None = None, +) -> dict[str, Any]: + """Fetch a cursor-paginated page of session summaries. + + Returns a dict with keys ``sessions`` (list of session dicts) and + ``next_cursor`` (opaque token or None when no further pages exist). + + Note: ``q`` uses a leading-wildcard LIKE which cannot use a btree index. + ``quick_filter='claude'`` scans the full payload column. Both are intended + for small deployments; see changelog for the long-term fix path. + """ + if q is not None: + q = q[:_Q_MAX_LEN] + + if db_pool.is_sqlite: + sqlite_args: list[object] = [] + + q_filter = "" + if q: + sqlite_args.append(f"%{_escape_like(q.lower())}%") + q_filter = "AND LOWER(session_id) LIKE ? ESCAPE '\\'" + + cursor_filter = "" + if cursor_token is not None: + cursor_ts, cursor_sid = decode_cursor(cursor_token, kind="sessions") + cursor_filter = "AND (datetime(last_ts), session_id) < (datetime(?), ?)" + sqlite_args.extend([cursor_ts.isoformat(), cursor_sid]) + + filter_clause = "" + if quick_filter == "30days": + filter_clause = "AND datetime(last_ts) >= datetime('now', '-30 days')" + elif quick_filter == "claude": + filter_clause = "AND session_id IN (SELECT DISTINCT session_id FROM conversation_events WHERE payload LIKE '%claude-code%')" + + sqlite_args.append(limit + 1) + query_args: list[object] = sqlite_args + + sessions_query = f""" + WITH sessions_agg AS ( + SELECT + session_id, + MIN(created_at) AS first_ts, + MAX(created_at) AS last_ts, + COUNT(DISTINCT call_id) AS turn_count, + SUM(CASE + WHEN event_type LIKE 'policy.%' + AND event_type NOT LIKE 'policy.%judge.evaluation%' + THEN 1 ELSE 0 + END) AS policy_interventions + FROM conversation_events + WHERE session_id IS NOT NULL + {q_filter} + GROUP BY session_id + ) + SELECT session_id, first_ts, last_ts, turn_count, policy_interventions + FROM sessions_agg + WHERE 1=1 + {cursor_filter} + {filter_clause} + ORDER BY last_ts DESC, session_id DESC + LIMIT ? + """ + else: + query_args = [limit + 1] + + q_filter = "" + if q: + query_args.append(f"%{_escape_like(q)}%") + q_idx = len(query_args) + q_filter = f"AND session_id ILIKE ${q_idx} ESCAPE '\\'" + + cursor_filter = "" + if cursor_token is not None: + cursor_ts, cursor_sid = decode_cursor(cursor_token, kind="sessions") + query_args.append(cursor_ts) + ts_idx = len(query_args) + query_args.append(cursor_sid) + sid_idx = len(query_args) + cursor_filter = f"AND (last_ts, session_id) < (${ts_idx}, ${sid_idx})" + + filter_clause = "" + if quick_filter == "30days": + filter_clause = "AND last_ts >= NOW() - INTERVAL '30 days'" + elif quick_filter == "claude": + filter_clause = "AND session_id IN (SELECT DISTINCT session_id FROM conversation_events WHERE payload::text ILIKE '%claude-code%')" + + sessions_query = f""" + WITH sessions_agg AS ( + SELECT + session_id, + MIN(created_at) AS first_ts, + MAX(created_at) AS last_ts, + COUNT(DISTINCT call_id) AS turn_count, + COUNT(*) FILTER ( + WHERE event_type LIKE 'policy.%' + AND event_type NOT LIKE 'policy.%judge.evaluation%' + ) AS policy_interventions + FROM conversation_events + WHERE session_id IS NOT NULL + {q_filter} + GROUP BY session_id + ) + SELECT session_id, first_ts, last_ts, turn_count, policy_interventions + FROM sessions_agg + WHERE 1=1 + {cursor_filter} + {filter_clause} + ORDER BY last_ts DESC, session_id DESC + LIMIT $1 + """ + + async with db_pool.connection() as conn: + with time_phase("db"): + rows = list(await conn.fetch(sessions_query, *query_args)) + + has_more = len(rows) > limit + page_rows = rows[:limit] + + if not page_rows: + return {"sessions": [], "next_cursor": None} + + session_ids = [str(r["session_id"]) for r in page_rows] + + if db_pool.is_sqlite: + placeholders = ", ".join("?" for _ in session_ids) + max_tokens_check = """ + AND COALESCE( + CAST(json_extract(payload, '$.final_request.max_tokens') AS INTEGER), + 2 + ) > 1 + """ + model_field_sql = "json_extract(payload, '$.final_model')" + else: + placeholders = ", ".join(f"${i + 1}" for i in range(len(session_ids))) + max_tokens_check = """ + AND COALESCE((payload->'final_request'->>'max_tokens')::int, 2) > 1 + """ + model_field_sql = "payload->>'final_model'" + + preview_rows = await conn.fetch( + f""" + SELECT session_id, payload + FROM conversation_events + WHERE session_id IN ({placeholders}) + AND event_type = 'transaction.request_recorded' + {max_tokens_check} + ORDER BY session_id, created_at ASC + """, + *session_ids, + ) + + model_rows = await conn.fetch( + f""" + SELECT session_id, {model_field_sql} as model + FROM conversation_events + WHERE session_id IN ({placeholders}) + AND event_type = 'transaction.request_recorded' + AND {model_field_sql} IS NOT NULL + """, + *session_ids, + ) + + previews: dict[str, str] = {} + for pr in preview_rows: + sid = str(pr["session_id"]) + if sid not in previews: + previews[sid] = _extract_preview_message(cast(_PreviewPayload, pr["payload"])) or "" + + models_by_session: dict[str, list[str]] = {} + for mr in model_rows: + sid = str(mr["session_id"]) + model = str(mr["model"]) + session_models = models_by_session.setdefault(sid, []) + if model not in session_models: + session_models.append(model) + + next_cursor: str | None = None + if has_more: + last_row = page_rows[-1] + last_ts_val = parse_db_ts(last_row["last_ts"]) + # Cursor encodes MAX(created_at) — "last active" time, not creation time. + # A new event on an older session bumps its last_ts and can re-surface or + # skip that session across page boundaries. This is intentional: the list + # is ordered by activity, not by when sessions were created. + next_cursor = encode_cursor(last_ts_val, str(last_row["session_id"]), kind="sessions") + + sessions = [ + { + "session_id": str(r["session_id"]), + "first_ts": str(r["first_ts"]), + "last_ts": str(r["last_ts"]), + "last_ts_formatted": _format_session_ts(parse_db_ts(r["last_ts"])), + "turn_count": int(r["turn_count"]), # type: ignore[arg-type] + "policy_interventions": int(r["policy_interventions"]), # type: ignore[arg-type] + "models_used": models_by_session.get(str(r["session_id"]), []), + "preview": previews.get(str(r["session_id"]), ""), + } + for r in page_rows + ] + + return {"sessions": sessions, "next_cursor": next_cursor} + + __all__ = [ "extract_text_content", "fetch_session_list", "fetch_session_detail", "export_session_markdown", "export_session_jsonl", + "fetch_session_turns_page", + "fetch_sessions_page", ] diff --git a/src/luthien_proxy/main.py b/src/luthien_proxy/main.py index 152a6501c..7e48519d1 100644 --- a/src/luthien_proxy/main.py +++ b/src/luthien_proxy/main.py @@ -8,6 +8,7 @@ import logging import os import secrets +import tempfile from collections.abc import MutableMapping from contextlib import asynccontextmanager @@ -41,6 +42,7 @@ ) from luthien_proxy.observability.redis_event_publisher import RedisEventPublisher from luthien_proxy.observability.sentry import init_sentry +from luthien_proxy.perf.timing_middleware import ServerTimingMiddleware from luthien_proxy.pipeline.upstream_headers import validate_upstream_headers_at_startup from luthien_proxy.policy_manager import PolicyManager from luthien_proxy.rate_limit import TokenBucketRateLimiter @@ -194,6 +196,14 @@ async def lifespan(app: FastAPI): await _config_registry.initialize() logger.info("Config registry initialized") + _CURSOR_HMAC_DEV_SENTINEL = "luthien-perf-cursor-key-dev" + if settings.cursor_hmac_key == _CURSOR_HMAC_DEV_SENTINEL: + logger.warning( + "CURSOR_HMAC_KEY is set to the public dev default. " + "Pagination cursors can be forged by anyone who has read this source. " + "Set CURSOR_HMAC_KEY to a random secret before deploying to production." + ) + # Fail fast on UPSTREAM_HEADERS misconfiguration rather than silently # disabling the integration on first request. validate_upstream_headers_at_startup() @@ -430,6 +440,10 @@ async def dispatch(self, request: Request, call_next): if request.url.path.startswith("/api/") or request.url.path in ("/health", "/ready"): # Prevent CDN/edge caching of API and health responses (Railway, Cloudflare, etc.) response.headers["Cache-Control"] = "no-store, no-cache, must-revalidate" + elif request.url.path.startswith("/ui/fragments/"): + # Fragment responses embed HMAC-signed cursors; a stale cached fragment + # replays a stale cursor that becomes a 400 after key rotation. + response.headers["Cache-Control"] = "no-store" elif request.url.path.startswith("/static/"): path = request.url.path if path.endswith((".js", ".html", ".css")): @@ -440,6 +454,10 @@ async def dispatch(self, request: Request, call_next): app.add_middleware(StaticCacheMiddleware) + # Add ServerTimingMiddleware as the last (innermost) middleware + # so it captures actual handler latency + app.add_middleware(ServerTimingMiddleware) + # Include routers app.include_router(gateway_router) # /v1/messages app.include_router(debug_router) # /api/debug/* @@ -714,6 +732,25 @@ def auto_provision_defaults() -> dict[str, str]: os.environ["ADMIN_API_KEY"] = value provisioned["ADMIN_API_KEY"] = value + if not os.environ.get("CURSOR_HMAC_KEY"): + data_dir = os.path.join(os.path.expanduser("~"), ".luthien") + os.makedirs(data_dir, exist_ok=True) + key_path = os.path.join(data_dir, "cursor_hmac.key") + if os.path.exists(key_path): + with open(key_path) as f: + value = f.read().strip() + logger.info("CURSOR_HMAC_KEY loaded from %s", key_path) + else: + value = secrets.token_urlsafe(32) + with tempfile.NamedTemporaryFile(mode="w", dir=data_dir, delete=False) as tmp: + tmp.write(value) + tmp_path = tmp.name + os.chmod(tmp_path, 0o600) + os.replace(tmp_path, key_path) + logger.info("CURSOR_HMAC_KEY generated and saved to %s", key_path) + os.environ["CURSOR_HMAC_KEY"] = value + provisioned["CURSOR_HMAC_KEY"] = value + if not os.environ.get("POLICY_CONFIG"): value = "config/railway_policy_config.yaml" if on_railway else "config/policy_config.yaml" os.environ["POLICY_CONFIG"] = value diff --git a/src/luthien_proxy/perf/__init__.py b/src/luthien_proxy/perf/__init__.py new file mode 100644 index 000000000..cf44e1770 --- /dev/null +++ b/src/luthien_proxy/perf/__init__.py @@ -0,0 +1,5 @@ +"""Perf-measurement utilities for the Luthien proxy. + +Isolated from the main application — writes only to the perf database, +never to the dev database (~/.luthien/local.db). +""" diff --git a/src/luthien_proxy/perf/db.py b/src/luthien_proxy/perf/db.py new file mode 100644 index 000000000..562565038 --- /dev/null +++ b/src/luthien_proxy/perf/db.py @@ -0,0 +1,132 @@ +"""Perf-DB isolation enforcement, migration runner, and drop helpers. + +The perf database is a completely isolated database used only for performance +benchmarking. It must never alias the dev database (~/.luthien/local.db). +""" + +from __future__ import annotations + +import asyncio +import os +from pathlib import Path +from typing import Literal + + +def get_perf_db_url(backend: Literal["sqlite", "postgres"]) -> str: + """Return the URL for the perf-test database. + + Args: + backend: "sqlite" → file URL under ~/.luthien/perf.db; + "postgres" → DATABASE_URL with perf_test schema override. + + Returns: + A database URL string for use with the migration runner. + + Raises: + ValueError: When backend is "postgres" (not yet implemented). + RuntimeError: When backend is "postgres" and DATABASE_URL is unset. + """ + if backend == "postgres": + raise ValueError( + "Postgres backend is not yet implemented for the perf harness. Use --backend sqlite (the default)." + ) + if backend == "sqlite": + return f"sqlite:///{Path.home()}/.luthien/perf.db" + base_url = os.environ.get("DATABASE_URL", "") + if not base_url: + raise RuntimeError("DATABASE_URL environment variable is required for postgres backend") + separator = "&" if "?" in base_url else "?" + return f"{base_url}{separator}options=-csearch_path=perf_test" + + +def ensure_perf_isolation(url: str) -> None: + """Assert that a database URL is not the dev database. + + This is the safety gate — call it before any write to the perf DB. + + Args: + url: The database URL to inspect. + + Raises: + RuntimeError: If the URL points to the dev database (contains "local.db"), + or if it is a Postgres URL without the "perf_test" schema override. + The message always contains the word "isolation". + """ + if "local.db" in url: + raise RuntimeError( + "Perf DB isolation violation: URL contains 'local.db' — " + "refusing to use the dev database as the perf database. " + "Use get_perf_db_url() to obtain the correct perf DB URL." + ) + if url.startswith(("postgresql://", "postgres://")) and "perf_test" not in url: + raise RuntimeError( + f"Perf DB isolation violation: Postgres URL must include " + f"'perf_test' schema (add ?options=-csearch_path=perf_test). Got: {url!r}" + ) + + +def drop_perf_db(backend: Literal["sqlite", "postgres"]) -> None: + """Drop the perf database. Idempotent — safe to call when already dropped. + + Args: + backend: "sqlite" removes ~/.luthien/perf.db (no-op if absent); + "postgres" runs DROP SCHEMA IF EXISTS perf_test CASCADE. + """ + if backend == "sqlite": + perf_path = Path.home() / ".luthien" / "perf.db" + perf_path.unlink(missing_ok=True) + return + + url = get_perf_db_url("postgres") + + async def _drop() -> None: + import asyncpg # type: ignore[import-untyped] # noqa: PLC0415 + + conn = await asyncpg.connect(url) + try: + await conn.execute("DROP SCHEMA IF EXISTS perf_test CASCADE") + finally: + await conn.close() + + asyncio.run(_drop()) + + +def migrate_perf_db(backend: Literal["sqlite", "postgres"]) -> None: + """Apply all migrations to the perf database. + + Calls ensure_perf_isolation before touching the database. For SQLite, + creates ~/.luthien/ if needed and runs the bundled migration scripts + via the standard migration runner. + + Args: + backend: "sqlite" or "postgres". + + Raises: + RuntimeError: If isolation check fails or migrations fail. + NotImplementedError: For the "postgres" backend (not yet implemented). + """ + url = get_perf_db_url(backend) + ensure_perf_isolation(url) + + if backend == "sqlite": + _migrate_sqlite(url) + else: + raise NotImplementedError("Postgres perf migration is not yet implemented") + + +def _migrate_sqlite(url: str) -> None: + from luthien_proxy.utils.db import DatabasePool # noqa: PLC0415 + from luthien_proxy.utils.db_sqlite import parse_sqlite_url # noqa: PLC0415 + from luthien_proxy.utils.migration_check import _apply_sqlite_migrations # noqa: PLC0415 + + db_path = Path(parse_sqlite_url(url)) + db_path.parent.mkdir(parents=True, exist_ok=True) + + async def _run() -> None: + db_pool = DatabasePool(url) + try: + await _apply_sqlite_migrations(db_pool) + finally: + await db_pool.close() + + asyncio.run(_run()) diff --git a/src/luthien_proxy/perf/seeding.py b/src/luthien_proxy/perf/seeding.py new file mode 100644 index 000000000..437dd9998 --- /dev/null +++ b/src/luthien_proxy/perf/seeding.py @@ -0,0 +1,309 @@ +"""Direct-SQL seeding module for perf-test database. + +Inserts production-shaped rows into conversation_calls and conversation_events +for performance benchmarking. Uses direct sqlite3 connections and executemany +for maximum throughput. + +All session_ids are prefixed with 'perf-seed-{tier}-' or 'perf-seed-sami-'. +IDs are fully deterministic — drop + re-seed produces identical data. + +FK ordering: conversation_calls rows are inserted before conversation_events rows. +""" + +from __future__ import annotations + +import random +import sqlite3 +import time +from dataclasses import dataclass +from datetime import datetime, timedelta, timezone +from pathlib import Path +from typing import Literal + +from luthien_proxy.perf.db import ensure_perf_isolation, get_perf_db_url, migrate_perf_db + +_MODEL = "claude-haiku-4-5" +_BASE_TS = datetime(2025, 1, 1, 0, 0, 0, tzinfo=timezone.utc) +_BATCH_SIZE = 5000 + +_CALLS_INSERT = ( + "INSERT INTO conversation_calls" + " (call_id, model_name, provider, status, created_at, completed_at, session_id)" + " VALUES (?, ?, ?, ?, ?, ?, ?)" +) +_EVENTS_INSERT = ( + "INSERT INTO conversation_events" + " (id, call_id, event_type, payload, created_at, session_id)" + " VALUES (?, ?, ?, ?, ?, ?)" +) + +# Pre-built JSON template fragments — content is pure ASCII, no escaping needed. +_REQ_PAD = "A" * 50 +_RESP_PAD = "B" * 100 + +_REQ_HEAD = ( + '{"final_request": {"model": "' + _MODEL + '", "max_tokens": 1024,' + ' "stream": true, "temperature": 0.7,' + ' "messages": [{"role": "user", "content": "' +) +_REQ_MID = ( + '"}]}, "original_request": {"model": "' + _MODEL + '", "max_tokens": 1024,' + ' "stream": true, "temperature": 0.7,' + ' "messages": [{"role": "user", "content": "' +) +_REQ_TAIL = '"}]}, "final_model": "' + _MODEL + '"}' + +_RESP_HEAD = ( + '{"final_response": {"id": "msg_000000", "type": "message",' + ' "role": "assistant", "model": "' + _MODEL + '",' + ' "stop_reason": "end_turn", "stop_sequence": null,' + ' "usage": {"input_tokens": 256, "output_tokens": 512},' + ' "content": [{"type": "text", "text": "' +) +_RESP_TAIL = '"}]}}' + + +@dataclass(frozen=True) +class SeedingReport: + """Report returned by seeding functions with metrics about the seeding run.""" + + tier: int | str + total_sessions: int + total_rows: int + total_bytes: int + elapsed_seconds: float + backend: str + biggest_session_message_count: int + + +def _fmt_ts(dt: datetime) -> str: + return dt.strftime("%Y-%m-%d %H:%M:%S") + + +def _req_payload(session_id: str, call_idx: int) -> str: + """~5 KB JSON string for a transaction.request_recorded event.""" + content = f"s={session_id[:12]} c={call_idx:04d} " + _REQ_PAD + return _REQ_HEAD + content + _REQ_MID + content + _REQ_TAIL + + +def _resp_payload(session_id: str, call_idx: int) -> str: + """~20 KB JSON string for a transaction.streaming_response_recorded event.""" + text = f"r={session_id[:12]} c={call_idx:04d} " + _RESP_PAD + return _RESP_HEAD + text + _RESP_TAIL + + +def _call_count(session_idx: int, rng_seed: int) -> int: + """Deterministic call count per session. + + Distribution (in calls; each call = 2 events): + - 50% → 5–15 calls (10–30 events; median ≈ 20 events) + - 45% → 15–50 calls (30–100 events; p95 ≈ 100 events) + - 5% → 50–250 calls (100–500 events; p99 ≈ 500 events) + """ + rng = random.Random(rng_seed * 1_000_003 + session_idx) + r = rng.random() + if r < 0.50: + return rng.randint(5, 15) + elif r < 0.95: + return rng.randint(15, 50) + else: + return rng.randint(50, 250) + + +def _sqlite_path(url: str) -> Path: + prefix = "sqlite:///" + if not url.startswith(prefix): + raise ValueError(f"Expected sqlite:/// URL, got {url!r}") + return Path(url[len(prefix) :]) + + +def _seed_sqlite( + db_path: Path, + plan: list[tuple[str, int]], + tier: int | str, + backend: str = "sqlite", +) -> SeedingReport: + """Bulk-insert plan into SQLite via executemany. + + Args: + db_path: Path to the SQLite database file. + plan: List of (session_id, n_calls) pairs. + tier: Tier label for the report. + backend: Backend label for the report. + + Returns: + SeedingReport with insertion statistics. + """ + t0 = time.monotonic() + total_bytes = 0 + biggest = 0 + + conn = sqlite3.connect(str(db_path)) + conn.execute("PRAGMA journal_mode=WAL") + conn.execute("PRAGMA synchronous=OFF") + conn.execute("PRAGMA cache_size=-131072") + conn.execute("PRAGMA temp_store=MEMORY") + + try: + # Drop indexes before bulk insert — dramatically reduces write amplification. + # Indexes are recreated after all rows are inserted. + for idx in ( + "idx_conversation_events_type", + "idx_conversation_events_created", + "idx_conversation_events_call_created", + "idx_conversation_events_session", + "idx_conversation_calls_created", + "idx_conversation_calls_session", + "idx_conversation_calls_user", + ): + conn.execute(f"DROP INDEX IF EXISTS {idx}") + + # Pass 1: conversation_calls (FK parent) — must precede events. + calls_batch: list[tuple] = [] + for session_idx, (session_id, n_calls) in enumerate(plan): + if n_calls > biggest: + biggest = n_calls + for call_idx in range(n_calls): + call_id = f"{session_id}-{call_idx:04d}" + ts = _fmt_ts(_BASE_TS + timedelta(seconds=session_idx * 3600 + call_idx * 5)) + calls_batch.append((call_id, _MODEL, "anthropic", "completed", ts, ts, session_id)) + if len(calls_batch) >= _BATCH_SIZE: + conn.executemany(_CALLS_INSERT, calls_batch) + calls_batch.clear() + if calls_batch: + conn.executemany(_CALLS_INSERT, calls_batch) + + # Pass 2: conversation_events (FK child). + events_batch: list[tuple] = [] + for session_idx, (session_id, n_calls) in enumerate(plan): + for call_idx in range(n_calls): + call_id = f"{session_id}-{call_idx:04d}" + ts_req = _fmt_ts(_BASE_TS + timedelta(seconds=session_idx * 3600 + call_idx * 5)) + ts_resp = _fmt_ts(_BASE_TS + timedelta(seconds=session_idx * 3600 + call_idx * 5 + 1)) + req_p = _req_payload(session_id, call_idx) + resp_p = _resp_payload(session_id, call_idx) + total_bytes += len(req_p) + len(resp_p) + + events_batch.append( + ( + f"{call_id}-req", + call_id, + "transaction.request_recorded", + req_p, + ts_req, + session_id, + ) + ) + events_batch.append( + ( + f"{call_id}-resp", + call_id, + "transaction.streaming_response_recorded", + resp_p, + ts_resp, + session_id, + ) + ) + + if len(events_batch) >= _BATCH_SIZE: + conn.executemany(_EVENTS_INSERT, events_batch) + events_batch.clear() + + if events_batch: + conn.executemany(_EVENTS_INSERT, events_batch) + + # Recreate indexes after bulk insert. + conn.execute("CREATE INDEX IF NOT EXISTS idx_conversation_events_type ON conversation_events(event_type)") + conn.execute("CREATE INDEX IF NOT EXISTS idx_conversation_events_created ON conversation_events(created_at)") + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_conversation_events_call_created" + " ON conversation_events(call_id, created_at)" + ) + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_conversation_events_session" + " ON conversation_events(session_id) WHERE session_id IS NOT NULL" + ) + conn.execute("CREATE INDEX IF NOT EXISTS idx_conversation_calls_created ON conversation_calls(created_at)") + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_conversation_calls_session" + " ON conversation_calls(session_id) WHERE session_id IS NOT NULL" + ) + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_conversation_calls_user" + " ON conversation_calls(user_id) WHERE user_id IS NOT NULL" + ) + conn.commit() + finally: + conn.close() + + elapsed = time.monotonic() - t0 + n_calls_total = sum(n for _, n in plan) + total_rows = n_calls_total + 2 * n_calls_total # calls + 2 events per call + + return SeedingReport( + tier=tier, + total_sessions=len(plan), + total_rows=total_rows, + total_bytes=total_bytes, + elapsed_seconds=elapsed, + backend=backend, + biggest_session_message_count=biggest, + ) + + +def seed_sessions( + backend: Literal["sqlite", "postgres"], + tier: int, +) -> SeedingReport: + """Seed the perf database with ``tier`` sessions. + + Calls ensure_perf_isolation and migrate_perf_db before inserting. + All session_ids are prefixed with ``perf-seed-{tier}-``. + IDs are fully deterministic — drop + re-seed produces identical data. + + Args: + backend: "sqlite" or "postgres". + tier: Number of sessions to insert (typically 100, 1_000, or 10_000). + + Returns: + SeedingReport with insertion statistics. + """ + url = get_perf_db_url(backend) + ensure_perf_isolation(url) + migrate_perf_db(backend) + + prefix = f"perf-seed-{tier}-" + plan = [(f"{prefix}{i:04d}", _call_count(i, rng_seed=tier)) for i in range(tier)] + + if backend == "sqlite": + return _seed_sqlite(_sqlite_path(url), plan, tier=tier, backend=backend) + raise NotImplementedError(f"backend {backend!r} not yet implemented") + + +def seed_sami_like(backend: Literal["sqlite", "postgres"]) -> SeedingReport: + """Seed the perf database with a sami-like fixture. + + 78 sessions total. Session ``perf-seed-sami-442msg`` has exactly 442 calls. + Remaining 77 sessions have 1–187 calls (realistic spread). + All session_ids are prefixed with ``perf-seed-sami-``. + + Args: + backend: "sqlite" or "postgres". + + Returns: + SeedingReport with biggest_session_message_count >= 442. + """ + url = get_perf_db_url(backend) + ensure_perf_isolation(url) + migrate_perf_db(backend) + + prefix = "perf-seed-sami-" + big_session_id = f"{prefix}442msg" + + rng = random.Random(0xABCDEF) + other_plan: list[tuple[str, int]] = [(f"{prefix}{i:03d}", rng.randint(1, 187)) for i in range(77)] + plan = [(big_session_id, 442)] + other_plan + + if backend == "sqlite": + return _seed_sqlite(_sqlite_path(url), plan, tier="sami", backend=backend) + raise NotImplementedError(f"backend {backend!r} not yet implemented") diff --git a/src/luthien_proxy/perf/timing_middleware.py b/src/luthien_proxy/perf/timing_middleware.py new file mode 100644 index 000000000..0378e8cbe --- /dev/null +++ b/src/luthien_proxy/perf/timing_middleware.py @@ -0,0 +1,130 @@ +"""Server-Timing middleware for admin/debug/UI paths. + +Records per-request timing phases via contextvars and appends a ``Server-Timing`` +header on responses whose path starts with ``/api/history/``, ``/api/debug/``, or +``/ui/fragments/``. All other paths (including ``/v1/messages``) are untouched. + +Usage:: + + from luthien_proxy.perf.timing_middleware import time_phase, ServerTimingMiddleware + + # Inside a request handler or service function: + with time_phase("db"): + rows = await db.fetch(query) + + with time_phase("serialize"): + payload = serialize(rows) + + # In FastAPI app setup (handled by P14): + app.add_middleware(ServerTimingMiddleware) +""" + +from __future__ import annotations + +import time +from collections.abc import Generator +from contextlib import contextmanager +from contextvars import ContextVar + +from starlette.middleware.base import BaseHTTPMiddleware +from starlette.requests import Request +from starlette.responses import Response + +# Paths where Server-Timing is emitted. /v1/messages is deliberately excluded. +_TIMED_PREFIXES: tuple[str, ...] = ( + "/api/history/", + "/api/debug/", + "/ui/fragments/", +) + +# Per-request phase list: list of (name, elapsed_ms) tuples. +# A new list is injected at the start of each request by ServerTimingMiddleware +# so phases never bleed across requests, even under concurrent load. +_phases_var: ContextVar[list[tuple[str, float]]] = ContextVar("_luthien_timing_phases") + + +@contextmanager +def time_phase(name: str) -> Generator[None, None, None]: + """Record the wall-clock duration of a code block as a timing phase. + + The elapsed milliseconds are appended to the current request's phase list + (stored in a ``ContextVar``). If called outside a ``ServerTimingMiddleware`` + request context the phase is silently discarded. + + Args: + name: Short identifier for the phase (e.g. ``"db"``, ``"serialize"``). + + Yields: + Nothing — use as a plain context manager. + + Example:: + + with time_phase("db"): + rows = await conn.fetch(query) + """ + start = time.perf_counter() + try: + yield + finally: + elapsed_ms = (time.perf_counter() - start) * 1000.0 + phases = _phases_var.get(None) + if phases is not None: + phases.append((name, elapsed_ms)) + + +def format_phases(phases: list[tuple[str, float]]) -> str: + """Format a list of timing phases as a ``Server-Timing`` header value. + + Args: + phases: Ordered list of ``(name, elapsed_ms)`` tuples. + + Returns: + Header value string, e.g. ``"db;dur=12.3, serialize;dur=4.5"``. + Returns an empty string if ``phases`` is empty. + + Example:: + + >>> format_phases([("db", 12.3), ("serialize", 4.5)]) + 'db;dur=12.3, serialize;dur=4.5' + """ + return ", ".join(f"{name};dur={elapsed_ms:.1f}" for name, elapsed_ms in phases) + + +class ServerTimingMiddleware(BaseHTTPMiddleware): + """ASGI middleware that adds a ``Server-Timing`` header to filtered responses. + + Only paths starting with ``/api/history/``, ``/api/debug/``, or + ``/ui/fragments/`` receive the header. All other paths (including the hot + ``/v1/messages`` path) pass through with zero overhead beyond a single + ``str.startswith`` check. + + Timing phases are recorded by calling ``time_phase(name)`` anywhere in the + request/response call stack. Context isolation is guaranteed by + ``contextvars.ContextVar``: each request gets its own fresh phase list. + """ + + async def dispatch(self, request: Request, call_next) -> Response: # noqa: D102 + path = request.url.path + should_time = path.startswith(_TIMED_PREFIXES) + + if not should_time: + return await call_next(request) + + phases: list[tuple[str, float]] = [] + token = _phases_var.set(phases) + try: + response = await call_next(request) + finally: + _phases_var.reset(token) + + if phases: + response.headers["Server-Timing"] = format_phases(phases) + + return response + + +__all__ = [ + "ServerTimingMiddleware", + "time_phase", + "format_phases", +] diff --git a/src/luthien_proxy/settings.py b/src/luthien_proxy/settings.py index 8c0320409..c2a5897a6 100644 --- a/src/luthien_proxy/settings.py +++ b/src/luthien_proxy/settings.py @@ -77,6 +77,7 @@ class Settings(_SettingsBase): # ── security ──────────────────────────────────────────────────── credential_encryption_key: str | None = None + cursor_hmac_key: str = "luthien-perf-cursor-key-dev" # ── observability ─────────────────────────────────────────────── otel_enabled: bool = False diff --git a/src/luthien_proxy/static/conversation_live.html b/src/luthien_proxy/static/conversation_live.html index 897828dd5..d954db6e5 100644 --- a/src/luthien_proxy/static/conversation_live.html +++ b/src/luthien_proxy/static/conversation_live.html @@ -917,7 +917,7 @@

-
+
Loading conversation...
diff --git a/src/luthien_proxy/static/conversation_live.js b/src/luthien_proxy/static/conversation_live.js index a782ff7d9..192fcc572 100644 --- a/src/luthien_proxy/static/conversation_live.js +++ b/src/luthien_proxy/static/conversation_live.js @@ -1,8 +1,13 @@ + +const MAX_RAW_EVENTS_PER_CALL = 50; + function escapeHtml(str) { if (str === null || str === undefined) return ''; const div = document.createElement('div'); div.textContent = String(str); - return div.innerHTML; + // innerHTML escapes <, >, &. Also escape quotes so the result is safe + // in double-quoted HTML attribute values (e.g. data-call-id="..."). + return div.innerHTML.replace(/"/g, '"').replace(/'/g, '''); } function conversationViewer() { @@ -19,12 +24,13 @@ function conversationViewer() { renderedCallIds: new Set(), turnFingerprints: {}, _rawTurns: [], + initialLoaded: false, - init() { + async init() { const pathParts = window.location.pathname.split('/'); this.conversationId = decodeURIComponent(pathParts[pathParts.length - 1]); this.setupEventDelegation(); - this.loadInitial(); + await this.loadInitial(); this.connectSSE(); window.addEventListener('beforeunload', () => { if (this.evtSource) this.evtSource.close(); @@ -90,30 +96,28 @@ function conversationViewer() { }, async loadInitial() { + const container = document.getElementById('conversation-container'); + if (!container) return; + try { const resp = await fetch( `/api/history/sessions/${encodeURIComponent(this.conversationId)}`, - { headers: { 'Accept': 'application/json' } } + { headers: { 'Accept': 'application/json' }, redirect: 'manual' } ); - - if (!resp.ok) { - if (resp.status === 403) { - window.location.href = '/login?error=required&next=' + - encodeURIComponent(window.location.pathname); - return; - } - if (resp.status === 404) throw new Error('Conversation not found'); - throw new Error(`HTTP ${resp.status}: ${resp.statusText}`); - } + if (resp.type === 'opaqueredirect') { window.location.href = '/login'; return; } + if (resp.status === 404) throw new Error('Conversation not found'); + if (!resp.ok) throw new Error(`HTTP ${resp.status}`); const data = await resp.json(); this.processTurns(data); this.updateStats(data); this.updateTimestamp(); this.renderTurns(); - this.$nextTick(() => this.autoScrollToBottom()); - } catch (err) { - this.showError(`Failed to load: ${err.message}`); + } catch (e) { + console.error('Failed to load initial turns:', e); + container.innerHTML = '
Failed to load conversation. Please refresh.
'; + } finally { + this.initialLoaded = true; } }, @@ -165,14 +169,16 @@ function conversationViewer() { this.rawEvents[callId] = []; } - this.rawEvents[callId].push({ + const bucket = this.rawEvents[callId]; + if (bucket.length >= MAX_RAW_EVENTS_PER_CALL) { + bucket.shift(); + } + bucket.push({ type: eventType, timestamp: event.timestamp || new Date().toISOString(), data: event }); - this.stats.events++; - const shouldRefresh = eventType.includes('request_recorded') || eventType.includes('response_recorded') || eventType.includes('policy.'); @@ -191,8 +197,9 @@ function conversationViewer() { try { const resp = await fetch( `/api/history/sessions/${encodeURIComponent(this.conversationId)}`, - { headers: { 'Accept': 'application/json' } } + { headers: { 'Accept': 'application/json' }, redirect: 'manual' } ); + if (resp.type === 'opaqueredirect') { window.location.href = '/login'; return; } if (!resp.ok) return; const data = await resp.json(); const rawTurns = data.turns || []; @@ -254,23 +261,6 @@ function conversationViewer() { this.turns = this.presentTurns(rawTurns); }, - // Presentation pipeline: classify preflight turns and compute - // display messages (dedup) entirely on the client side. - // - // The API sends the full conversation history on every request: - // Turn 1: [user₀] - // Turn 2: [user₀, assistant₁, user₂] - // Turn 3: [user₀, assistant₁, user₂, tool_call₂, tool_result₂, user₃] - // - // user₀ (the initial message with all preamble) is re-sent identically - // every turn. New content appears at the end, after the previous turn's - // messages. So for turn N, display = request_messages.slice(prevCount). - // Preflight turns are excluded from the count so they don't disrupt the - // sequence. - // - // Invariant: the API sends a stable, strictly-growing cumulative - // message array. If a policy rewrites or reorders earlier messages, - // the slicing will produce incorrect results. presentTurns(rawTurns) { let prevRealMsgCount = 0; @@ -294,17 +284,9 @@ function conversationViewer() { }); }, - // Classify non-conversational preflight turns using structural - // request params (not response content heuristics). - // - Quota probe: max_tokens === 1 - // - Title generation: json_schema output + low max_tokens (≤256) - // These can appear at any position in the session. classifyPreflight(turn) { const params = turn.request_params || {}; if (params.max_tokens === 1) return true; - // json_schema alone isn't sufficient — real conversations can use - // structured output. Title generation uses json_schema with a - // small token budget. if (params.output_config?.format?.type === 'json_schema' && params.max_tokens != null && params.max_tokens <= 256) return true; return false; @@ -359,12 +341,13 @@ function conversationViewer() { snapshotExpandState() { const state = { visible: [], expanded: [], open: [] }; - document.querySelectorAll('.visible[id]').forEach(el => state.visible.push(el.id)); - document.querySelectorAll('.expanded[id]').forEach(el => state.expanded.push(el.id)); - document.querySelectorAll('.open[data-event-timeline]').forEach(el => { + const container = document.getElementById('conversation-container') || document; + container.querySelectorAll('.visible[id]').forEach(el => state.visible.push(el.id)); + container.querySelectorAll('.expanded[id]').forEach(el => state.expanded.push(el.id)); + container.querySelectorAll('.open[data-event-timeline]').forEach(el => { state.open.push(el.getAttribute('data-event-timeline')); }); - document.querySelectorAll('.open[data-diff-toggle]').forEach(el => { + container.querySelectorAll('.open[data-diff-toggle]').forEach(el => { state.open.push('diff:' + el.getAttribute('data-diff-toggle')); }); return state; @@ -431,13 +414,10 @@ function conversationViewer() { if (isPreflight) classes.push('preflight'); const callId = escapeHtml(turn.call_id); + const rawCallId = turn.call_id; const displayMessages = turn._displayMessages || turn.request_messages || []; const responseMessages = turn.response_messages || []; - // Build unified message list pairing tool calls with their results. - // Tool results (from request) match tool calls (from request or response) - // via tool_call_id. We also skip request tool_calls that duplicate - // a response tool_call from this same turn (the API re-sends them). const toolResultsByCallId = {}; for (const m of displayMessages) { if (m.message_type === 'tool_result' && m.tool_call_id) { @@ -445,7 +425,6 @@ function conversationViewer() { } } - // Track response tool_call IDs to suppress duplicates from request const responseToolCallIds = new Set(); for (const m of responseMessages) { if (m.message_type === 'tool_call' && m.tool_call_id) { @@ -456,20 +435,16 @@ function conversationViewer() { const orderedMessages = []; const usedResultIds = new Set(); - // Request messages: skip tool_results (paired later) and - // tool_calls that also appear in the response (duplicates) for (const m of displayMessages) { if (m.message_type === 'tool_result') continue; if (m.message_type === 'tool_call' && responseToolCallIds.has(m.tool_call_id)) continue; orderedMessages.push(m); - // If this request tool_call has a result, pair it if (m.message_type === 'tool_call' && toolResultsByCallId[m.tool_call_id]) { orderedMessages.push(toolResultsByCallId[m.tool_call_id]); usedResultIds.add(m.tool_call_id); } } - // Response messages with paired results for (const m of responseMessages) { orderedMessages.push(m); if (m.message_type === 'tool_call' && toolResultsByCallId[m.tool_call_id]) { @@ -478,7 +453,6 @@ function conversationViewer() { } } - // Any orphaned tool results for (const id in toolResultsByCallId) { if (!usedResultIds.has(id)) { orderedMessages.push(toolResultsByCallId[id]); @@ -509,7 +483,7 @@ function conversationViewer() { } let eventTimelineHtml = ''; - const events = this.rawEvents[callId] || []; + const events = this.rawEvents[rawCallId] || []; if (events.length > 0) { const eventsHtml = events.map((evt, idx) => { const eventKey = `${callId}-${idx}`; @@ -578,7 +552,6 @@ function conversationViewer() { const content = msg.content || ''; const contentId = `c-${stableId}`; - // Tool calls: show only the pretty-formatted input, not raw content if (msg.message_type === 'tool_call' && msg.tool_input) { return `
@@ -593,7 +566,6 @@ function conversationViewer() { `; } - // Tool results: code-style block if (msg.message_type === 'tool_result') { const shouldTruncate = content.length > 800; const expandBtn = shouldTruncate @@ -644,8 +616,6 @@ function conversationViewer() { 'bash-stderr': 'bash-output', }; - // Known wrapper tags are flat (never nested within themselves) - // so a simple regex is reliable here. const tagNames = Object.keys(TAG_LABELS).map(t => t.replace(/[.*+?^${}()|[\]\\]/g, '\\$&')); const tagPattern = new RegExp( '<(' + tagNames.join('|') + ')(?:\\s[^>]*)?>([\\s\\S]*?)', @@ -670,7 +640,6 @@ function conversationViewer() { if (remaining) parts.push({ type: 'text', content: remaining }); } - // If no tags found, fall back to plain rendering if (parts.length === 0 || (parts.length === 1 && parts[0].type === 'text')) { return this._renderPlainContent(content, contentId); } diff --git a/src/luthien_proxy/static/history_list.html b/src/luthien_proxy/static/history_list.html index ce2cef2a5..6d54b28ad 100644 --- a/src/luthien_proxy/static/history_list.html +++ b/src/luthien_proxy/static/history_list.html @@ -338,7 +338,7 @@ } - +