From 10dfc6d2b07d04d459311ff1a35aae8525c1a06e Mon Sep 17 00:00:00 2001 From: Johannes du Plessis Date: Sat, 6 Jun 2026 08:11:53 -0700 Subject: [PATCH] fix: precompute usage tab caches (#1434) * fix: precompute usage tab caches Co-authored-by: open-swe[bot] * fix: dedupe usage cache refreshes Co-authored-by: open-swe[bot] --------- Co-authored-by: open-swe[bot] --- agent/dashboard/agent_usage.py | 176 ++++++++++++++++++++++++++++++++- agent/dashboard/routes.py | 8 +- agent/webapp.py | 15 ++- tests/test_agent_usage.py | 96 ++++++++++++++++++ ui/src/routes/usage.tsx | 15 ++- 5 files changed, 299 insertions(+), 11 deletions(-) diff --git a/agent/dashboard/agent_usage.py b/agent/dashboard/agent_usage.py index a53ddd7e..8dcded43 100644 --- a/agent/dashboard/agent_usage.py +++ b/agent/dashboard/agent_usage.py @@ -4,8 +4,9 @@ from __future__ import annotations import asyncio import logging +import os from collections import Counter -from collections.abc import Awaitable, Callable +from collections.abc import Awaitable, Callable, Iterable from datetime import UTC, datetime, timedelta from typing import Any, Literal @@ -21,15 +22,19 @@ USAGE_LEADERBOARD_CACHE_NAMESPACE: list[str] = ["agent_usage", "leaderboard_cach REVIEWER_STATS_CACHE_NAMESPACE: list[str] = ["agent_usage", "reviewer_stats_cache"] Period = Literal["7d", "30d", "all"] +_PERIODS: tuple[Period, ...] = ("7d", "30d", "all") _AGENT_SOURCES = frozenset({"dashboard", "github", "slack", "linear"}) _PR_REFRESH_INTERVAL_MS = 10 * 60 * 1000 _MAX_PR_REFRESH_PER_REQUEST = 25 _PR_REFRESH_CONCURRENCY = 5 _CACHE_TTL_MS = 10 * 60 * 1000 +_CACHE_WARM_INTERVAL_SECONDS = 5 * 60 _CACHE_SEARCH_LIMIT = 1000 _GITHUB_API = "https://api.github.com" logger = logging.getLogger(__name__) +_CACHE_REFRESH_IN_FLIGHT: set[tuple[tuple[str, ...], Period]] = set() +_CACHE_REFRESH_IN_FLIGHT_LOCK = asyncio.Lock() def _client(): @@ -53,6 +58,20 @@ def _normalize_period(period: str | None) -> Period: return period if period in {"7d", "30d", "all"} else "30d" +def usage_cache_warmer_enabled() -> bool: + value = os.environ.get("USAGE_CACHE_WARMER_ENABLED", "true").strip().lower() + return value not in {"0", "false", "no", "off"} + + +def usage_cache_warm_interval_seconds() -> int: + raw = os.environ.get("USAGE_CACHE_WARM_INTERVAL_SECONDS", "") + try: + value = int(raw) if raw else _CACHE_WARM_INTERVAL_SECONDS + except ValueError: + return _CACHE_WARM_INTERVAL_SECONDS + return max(60, value) + + def _record_from_item(item: Any) -> dict[str, Any] | None: if item is None: return None @@ -605,12 +624,97 @@ async def refresh_reviewer_stats_cache(period: str | None = "30d") -> dict[str, return snapshot +def _empty_usage_snapshot(period: Period) -> dict[str, Any]: + return {"period": period, "users": [], "total_members": 0} + + +def _empty_reviewer_stats_snapshot(period: Period) -> dict[str, Any]: + return { + "period": period, + "reviewed_prs": 0, + "prs_with_findings": 0, + "findings_recorded": 0, + "surfaced_findings": 0, + "addressed_findings": 0, + "resolved_after_update": 0, + "dismissed_findings": 0, + "unresolved_surfaced_findings": 0, + "resolution_rate": 0.0, + "human_replies": 0, + "severity_counts": {}, + "top_categories": [], + } + + +def _cache_refresh_key(namespace: list[str], period: Period) -> tuple[tuple[str, ...], Period]: + return tuple(namespace), period + + +async def _claim_cache_refresh(namespace: list[str], period: Period) -> bool: + key = _cache_refresh_key(namespace, period) + async with _CACHE_REFRESH_IN_FLIGHT_LOCK: + if key in _CACHE_REFRESH_IN_FLIGHT: + return False + _CACHE_REFRESH_IN_FLIGHT.add(key) + return True + + +async def _release_cache_refresh(namespace: list[str], period: Period) -> None: + key = _cache_refresh_key(namespace, period) + async with _CACHE_REFRESH_IN_FLIGHT_LOCK: + _CACHE_REFRESH_IN_FLIGHT.discard(key) + + +async def _schedule_refresh_if_idle( + namespace: list[str], + period: Period, + schedule_refresh: Callable[[Period], None], +) -> None: + if not await _claim_cache_refresh(namespace, period): + return + try: + schedule_refresh(period) + except Exception: + await _release_cache_refresh(namespace, period) + raise + + +async def _run_claimed_cache_refresh( + namespace: list[str], + period: Period, + refresh: Callable[[str | None], Awaitable[dict[str, Any]]], +) -> dict[str, Any]: + try: + return await refresh(period) + finally: + await _release_cache_refresh(namespace, period) + + +async def refresh_claimed_usage_leaderboard_cache(period: str | None = "30d") -> dict[str, Any]: + normalized_period = _normalize_period(period) + return await _run_claimed_cache_refresh( + USAGE_LEADERBOARD_CACHE_NAMESPACE, + normalized_period, + refresh_usage_leaderboard_cache, + ) + + +async def refresh_claimed_reviewer_stats_cache(period: str | None = "30d") -> dict[str, Any]: + normalized_period = _normalize_period(period) + return await _run_claimed_cache_refresh( + REVIEWER_STATS_CACHE_NAMESPACE, + normalized_period, + refresh_reviewer_stats_cache, + ) + + async def _cached_snapshot( namespace: list[str], period: Period, refresh: Callable[[str | None], Awaitable[dict[str, Any]]], *, schedule_refresh: Callable[[Period], None] | None = None, + empty_snapshot: Callable[[Period], dict[str, Any]] | None = None, ) -> tuple[dict[str, Any], int | None]: cached = await _get_value(namespace, period) if cached: @@ -620,9 +724,73 @@ async def _cached_snapshot( if not generated_at_ms or _now_ms() - generated_at_ms <= _CACHE_TTL_MS: return snapshot, generated_at_ms or None if schedule_refresh is not None: - schedule_refresh(period) + await _schedule_refresh_if_idle(namespace, period, schedule_refresh) return snapshot, generated_at_ms - return await refresh(period), _now_ms() + if schedule_refresh is not None and empty_snapshot is not None: + await _schedule_refresh_if_idle(namespace, period, schedule_refresh) + return empty_snapshot(period), None + snapshot = await refresh(period) + return snapshot, _now_ms() + + +async def _refresh_cached_snapshot_if_stale( + namespace: list[str], + period: Period, + refresh: Callable[[str | None], Awaitable[dict[str, Any]]], +) -> bool: + cached = await _get_value(namespace, period) + if cached: + snapshot = cached.get("snapshot") + generated_at_ms = _coerce_int(cached.get("generated_at_ms")) + if ( + isinstance(snapshot, dict) + and generated_at_ms + and _now_ms() - generated_at_ms <= _CACHE_TTL_MS + ): + return False + if not await _claim_cache_refresh(namespace, period): + return False + try: + await refresh(period) + return True + finally: + await _release_cache_refresh(namespace, period) + + +async def precompute_usage_caches(periods: Iterable[Period] = _PERIODS) -> dict[str, int]: + refreshed = {"usage": 0, "reviewer": 0} + for period in periods: + try: + if await _refresh_cached_snapshot_if_stale( + USAGE_LEADERBOARD_CACHE_NAMESPACE, + period, + refresh_usage_leaderboard_cache, + ): + refreshed["usage"] += 1 + except Exception: + logger.debug( + "Failed to precompute usage leaderboard cache for %s", period, exc_info=True + ) + try: + if await _refresh_cached_snapshot_if_stale( + REVIEWER_STATS_CACHE_NAMESPACE, + period, + refresh_reviewer_stats_cache, + ): + refreshed["reviewer"] += 1 + except Exception: + logger.debug("Failed to precompute reviewer stats cache for %s", period, exc_info=True) + return refreshed + + +async def run_usage_cache_warmer(interval_seconds: int | None = None) -> None: + interval = interval_seconds or usage_cache_warm_interval_seconds() + while True: + try: + await precompute_usage_caches() + except Exception: + logger.exception("Usage cache warmer failed") + await asyncio.sleep(interval) def _usage_payload_from_snapshot( @@ -696,12 +864,14 @@ async def list_agent_usage_leaderboard( normalized_period, refresh_usage_leaderboard_cache, schedule_refresh=schedule_usage_refresh, + empty_snapshot=_empty_usage_snapshot, ) reviewer_stats, reviewer_generated_at_ms = await _cached_snapshot( REVIEWER_STATS_CACHE_NAMESPACE, normalized_period, refresh_reviewer_stats_cache, schedule_refresh=schedule_reviewer_refresh, + empty_snapshot=_empty_reviewer_stats_snapshot, ) payload = _usage_payload_from_snapshot( usage_snapshot, diff --git a/agent/dashboard/routes.py b/agent/dashboard/routes.py index 3c2b46c5..112cbdac 100644 --- a/agent/dashboard/routes.py +++ b/agent/dashboard/routes.py @@ -15,8 +15,8 @@ from pydantic import BaseModel from .admin import is_admin from .agent_usage import ( list_agent_usage_leaderboard, - refresh_reviewer_stats_cache, - refresh_usage_leaderboard_cache, + refresh_claimed_reviewer_stats_cache, + refresh_claimed_usage_leaderboard_cache, ) from .analyzer_cron import remove_continual_cron from .enabled_repos import ( @@ -721,10 +721,10 @@ async def api_agent_usage_leaderboard( current_login=session["sub"], current_email=session.get("email"), schedule_usage_refresh=lambda cache_period: background_tasks.add_task( - refresh_usage_leaderboard_cache, cache_period + refresh_claimed_usage_leaderboard_cache, cache_period ), schedule_reviewer_refresh=lambda cache_period: background_tasks.add_task( - refresh_reviewer_stats_cache, cache_period + refresh_claimed_reviewer_stats_cache, cache_period ), ) diff --git a/agent/webapp.py b/agent/webapp.py index ce8d77d4..4c900625 100644 --- a/agent/webapp.py +++ b/agent/webapp.py @@ -1,5 +1,6 @@ """Custom FastAPI routes for LangGraph server.""" +import asyncio import hashlib import hmac import json @@ -7,7 +8,7 @@ import logging import os import uuid from collections.abc import AsyncIterator -from contextlib import asynccontextmanager +from contextlib import asynccontextmanager, suppress from datetime import UTC, datetime from typing import Any from urllib.parse import parse_qs, quote @@ -109,10 +110,20 @@ logger = logging.getLogger(__name__) @asynccontextmanager async def lifespan(_app: FastAPI) -> AsyncIterator[None]: + from .dashboard.agent_usage import run_usage_cache_warmer, usage_cache_warmer_enabled from .utils.sandbox import validate_sandbox_startup_config validate_sandbox_startup_config() - yield + usage_cache_task: asyncio.Task[None] | None = None + if usage_cache_warmer_enabled(): + usage_cache_task = asyncio.create_task(run_usage_cache_warmer(), name="usage-cache-warmer") + try: + yield + finally: + if usage_cache_task: + usage_cache_task.cancel() + with suppress(asyncio.CancelledError): + await usage_cache_task app = FastAPI(lifespan=lifespan) diff --git a/tests/test_agent_usage.py b/tests/test_agent_usage.py index b340558f..db206f6f 100644 --- a/tests/test_agent_usage.py +++ b/tests/test_agent_usage.py @@ -9,6 +9,7 @@ class FakeStore: def __init__(self, values: dict[tuple[tuple[str, ...], str], dict] | None = None): self.values = values or {} self.puts: list[tuple[list[str], str, dict]] = [] + self.searches: list[tuple[list[str], int]] = [] async def get_item(self, namespace: list[str], key: str) -> dict | None: value = self.values.get((tuple(namespace), key)) @@ -18,6 +19,15 @@ class FakeStore: self.puts.append((namespace, key, value)) self.values[(tuple(namespace), key)] = value + async def search_items(self, namespace: list[str], limit: int = 1000) -> dict: + self.searches.append((namespace, limit)) + items = [ + {"value": value} + for (item_namespace, _), value in self.values.items() + if item_namespace == tuple(namespace) + ] + return {"items": items[:limit]} + class FakeThreads: def __init__(self, threads: list[dict]): @@ -37,6 +47,13 @@ class FakeClient: self.threads = threads or FakeThreads([]) +@pytest.fixture(autouse=True) +def clear_cache_refresh_guard(): + agent_usage._CACHE_REFRESH_IN_FLIGHT.clear() + yield + agent_usage._CACHE_REFRESH_IN_FLIGHT.clear() + + @pytest.mark.asyncio async def test_cached_usage_payload_returns_stale_snapshot_and_schedules_refresh(monkeypatch): usage_snapshot = { @@ -122,6 +139,85 @@ async def test_cached_usage_payload_returns_stale_snapshot_and_schedules_refresh assert store.puts == [] +@pytest.mark.asyncio +async def test_cache_miss_schedules_refresh_without_blocking(monkeypatch): + store = FakeStore() + monkeypatch.setattr(agent_usage, "_client", lambda: FakeClient(store=store)) + + async def fail_refresh(period): + raise AssertionError(f"refresh should have been scheduled, not awaited: {period}") + + monkeypatch.setattr(agent_usage, "refresh_usage_leaderboard_cache", fail_refresh) + monkeypatch.setattr(agent_usage, "refresh_reviewer_stats_cache", fail_refresh) + usage_refreshes: list[str] = [] + reviewer_refreshes: list[str] = [] + + payload = await agent_usage.list_agent_usage_leaderboard( + period="7d", + limit=10, + current_login="octo", + current_email="octo@example.com", + schedule_usage_refresh=usage_refreshes.append, + schedule_reviewer_refresh=reviewer_refreshes.append, + ) + + duplicate_payload = await agent_usage.list_agent_usage_leaderboard( + period="7d", + limit=10, + current_login="octo", + current_email="octo@example.com", + schedule_usage_refresh=usage_refreshes.append, + schedule_reviewer_refresh=reviewer_refreshes.append, + ) + + assert usage_refreshes == ["7d"] + assert reviewer_refreshes == ["7d"] + assert payload["period"] == "7d" + assert payload["rows"] == [] + assert payload["generated_at_ms"] is None + assert payload["reviewer_stats"]["generated_at_ms"] is None + assert duplicate_payload["rows"] == [] + assert store.searches == [] + + +@pytest.mark.asyncio +async def test_precompute_usage_caches_refreshes_stale_snapshots(monkeypatch): + store = FakeStore() + monkeypatch.setattr(agent_usage, "_client", lambda: FakeClient(store=store)) + usage_refreshes: list[str | None] = [] + reviewer_refreshes: list[str | None] = [] + + async def refresh_usage(period): + usage_refreshes.append(period) + snapshot = agent_usage._empty_usage_snapshot(agent_usage._normalize_period(period)) + await store.put_item( + agent_usage.USAGE_LEADERBOARD_CACHE_NAMESPACE, + agent_usage._normalize_period(period), + {"generated_at_ms": agent_usage._now_ms(), "snapshot": snapshot}, + ) + return snapshot + + async def refresh_reviewer(period): + reviewer_refreshes.append(period) + snapshot = agent_usage._empty_reviewer_stats_snapshot(agent_usage._normalize_period(period)) + await store.put_item( + agent_usage.REVIEWER_STATS_CACHE_NAMESPACE, + agent_usage._normalize_period(period), + {"generated_at_ms": agent_usage._now_ms(), "snapshot": snapshot}, + ) + return snapshot + + monkeypatch.setattr(agent_usage, "refresh_usage_leaderboard_cache", refresh_usage) + monkeypatch.setattr(agent_usage, "refresh_reviewer_stats_cache", refresh_reviewer) + + refreshed = await agent_usage.precompute_usage_caches(("30d",)) + + assert refreshed == {"usage": 1, "reviewer": 1} + assert usage_refreshes == ["30d"] + assert reviewer_refreshes == ["30d"] + assert agent_usage._CACHE_REFRESH_IN_FLIGHT == set() + + @pytest.mark.asyncio async def test_reviewer_stats_snapshot_counts_surfaced_and_resolved_findings(monkeypatch): threads = [ diff --git a/ui/src/routes/usage.tsx b/ui/src/routes/usage.tsx index eb5333eb..c36d863e 100644 --- a/ui/src/routes/usage.tsx +++ b/ui/src/routes/usage.tsx @@ -48,7 +48,18 @@ function UsagePage() { queryFn: () => api.usageLeaderboard(activePeriod, 10), enabled: !!session.data, staleTime: 5 * 60 * 1000, + refetchInterval: (query) => { + const data = query.state.data + return data && + (!data.generated_at_ms || !data.reviewer_stats.generated_at_ms) + ? 2000 + : false + }, }) + const usageIsPrecomputing = + !!leaderboard.data && !leaderboard.data.generated_at_ms + const reviewerIsPrecomputing = + !!leaderboard.data && !leaderboard.data.reviewer_stats.generated_at_ms if (session.isLoading) { return ( @@ -87,7 +98,7 @@ function UsagePage() { } > - {leaderboard.isLoading ? ( + {leaderboard.isLoading || usageIsPrecomputing ? (
@@ -114,7 +125,7 @@ function UsagePage() { title="Reviewer stats" description="Issues surfaced by Open SWE Review and how often users addressed them." > - {leaderboard.isLoading ? ( + {leaderboard.isLoading || reviewerIsPrecomputing ? (