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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
76 changes: 73 additions & 3 deletions backend/app/api/fleet.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
"""Fleet summary and consolidation endpoints."""
from __future__ import annotations

import logging
from concurrent.futures import ThreadPoolExecutor, as_completed
from datetime import datetime, timedelta
from typing import Any
Expand All @@ -16,6 +17,8 @@
from ..models import Alert, FailureEvent, PoolLoadSample, SyncLog, Worker
from ..sync.taskcluster import ALL_WORKER_POOLS, HW_WORKER_POOLS

log = logging.getLogger(__name__)

# Job-source sampling is served through the SWR cache so page loads never block on
# live Taskcluster fan-out; a task's project/user is immutable, so it's cached too.
_POOL_SOURCES_TTL = 90
Expand Down Expand Up @@ -763,9 +766,11 @@ def android_pools() -> dict[str, Any]:
_ANDROID_PROVISIONER: dict[str, str] = {wt: prov for prov, wt in ANDROID_WORKER_POOLS}


@router.get("/android-pool-sources")
def android_pool_sources(pool: str) -> dict[str, Any]:
"""Sample running tasks for an Android pool via the TC workers endpoint (no DB)."""
def _compute_android_pool_sources(pool: str) -> dict[str, Any]:
"""Sample running tasks for an Android pool via the TC workers endpoint (no DB).

Runs from a background thread (SWR refresh), so it must not touch a
request-scoped session — it uses none."""
provisioner = _ANDROID_PROVISIONER.get(pool)
if not provisioner:
return {"pool": pool, "sample_size": 0, "by_project": {}, "by_user": {}}
Expand Down Expand Up @@ -828,6 +833,18 @@ def _fetch(task_id: str) -> dict[str, str]:
}


@router.get("/android-pool-sources")
def android_pool_sources(pool: str) -> dict[str, Any]:
"""Job-source breakdown for one Android pool.

Cached: the uncached version made one live Taskcluster call per pool per page
load (limit=1000 workers, then one task fetch per running task)."""
return cache.swr(
f"android-pool-sources:{pool}", _POOL_SOURCES_TTL,
lambda: _compute_android_pool_sources(pool),
)


# ── Android per-device health (creds-free) ───────────────────────────────────
# Android devices (Bitbar / Lambda) have no rows in the Worker table and we don't
# yet have Bitbar/Lambda API access. But the Taskcluster queue is fully public, and
Expand Down Expand Up @@ -1080,6 +1097,59 @@ def pool_sources(pool: str) -> dict[str, Any]:
return cache.swr(f"pool-sources:{pool}", _POOL_SOURCES_TTL, lambda: _compute_pool_sources(pool))


# One request per pool was half the dashboard's entire Cloud Armor budget: a single
# Pools page load fired 35 /pool-sources + 14 /android-pool-sources, and the limit is
# per-IP for the whole origin. Both are individually cached, so the cost was purely
# request COUNT — which is exactly what the rate limiter charges for. These batch
# variants collapse each loop into one request with identical payloads.
_MAX_BATCH_POOLS = 120


def _parse_pools(pools: str) -> list[str]:
"""Comma-separated pool names → de-duplicated, order-preserving, bounded list."""
seen: dict[str, None] = {}
for raw in pools.split(","):
name = raw.strip()
if name:
seen.setdefault(name)
return list(seen)[:_MAX_BATCH_POOLS]


def _sources_batch(pools: str, key_prefix: str, compute: Any) -> dict[str, Any]:
names = _parse_pools(pools)
if not names:
return {"sources": {}}
# Bounded parallelism: a warm batch is pure cache reads, but a cold one would
# otherwise sample Taskcluster serially. The per-pool compute functions each
# open their own DB session (or none), so they are safe off-thread.
out: dict[str, Any] = {}
with ThreadPoolExecutor(max_workers=min(8, len(names))) as ex:
futures = {
ex.submit(cache.swr, f"{key_prefix}:{n}", _POOL_SOURCES_TTL, lambda n=n: compute(n)): n
for n in names
}
for fut in as_completed(futures):
name = futures[fut]
try:
out[name] = fut.result()
except Exception:
# One bad pool must not fail the batch; the card just stays empty.
log.warning("pool-sources batch: %s failed for %s", key_prefix, name, exc_info=True)
return {"sources": out}


@router.get("/pool-sources-batch")
def pool_sources_batch(pools: str) -> dict[str, Any]:
"""Job-source breakdown for many pools in one request. `pools` is comma-separated."""
return _sources_batch(pools, "pool-sources", _compute_pool_sources)


@router.get("/android-pool-sources-batch")
def android_pool_sources_batch(pools: str) -> dict[str, Any]:
"""Android job-source breakdown for many pools in one request."""
return _sources_batch(pools, "android-pool-sources", _compute_android_pool_sources)


def warm_pool_sources() -> int:
"""Background job: refresh the job-source cache for every pool with running
tasks, so the first page load already has warm data. Returns pools warmed."""
Expand Down
155 changes: 155 additions & 0 deletions backend/tests/test_pool_sources_batch.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,155 @@
"""Batched job-source endpoints: one request per page, not one per pool.

These exist because the per-pool loops were measured, on 2026-09-17, as 49 of the
100 requests in a single minute that tripped Cloud Armor and replaced the
dashboard with a bare "429 Too Many Requests" page:

/api/fleet/pool-sources 35
/api/fleet/android-pool-sources 14

Both per-pool computations were already cached server-side, so the cost was purely
request COUNT -- which is exactly what a per-IP rate limiter charges for. The batch
variants must therefore collapse N requests into 1 while returning byte-identical
per-pool payloads, so no card changes behaviour.

Dependency-light on purpose (no DB, no Taskcluster): the compute functions are
monkeypatched, since what is under test is batching, parsing and fan-out -- not
the sampling itself.
"""
from __future__ import annotations

import pytest
from app import cache
from app.api import fleet


@pytest.fixture(autouse=True)
def _clear_cache():
"""The SWR cache is module-global; a leaked key would make these order-dependent."""
cache._cache.clear()
cache._refreshing.clear()
yield
cache._cache.clear()
cache._refreshing.clear()


# ── Parsing ─────────────────────────────────────────────────────────────────────

def test_parses_and_preserves_order():
assert fleet._parse_pools("b,a,c") == ["b", "a", "c"]


def test_tolerates_whitespace_and_empty_segments():
assert fleet._parse_pools(" a , ,b ,") == ["a", "b"]


def test_deduplicates_repeated_pools():
"""A repeated name must not become a repeated compute -- that is the whole point."""
assert fleet._parse_pools("a,b,a,b,a") == ["a", "b"]


def test_empty_input_yields_no_pools():
assert fleet._parse_pools("") == []
assert fleet._parse_pools(" , , ") == []


def test_batch_size_is_bounded():
"""An unbounded batch would just move the fan-out server-side."""
names = ",".join(f"pool-{i}" for i in range(fleet._MAX_BATCH_POOLS + 50))
assert len(fleet._parse_pools(names)) == fleet._MAX_BATCH_POOLS


# ── Batching ────────────────────────────────────────────────────────────────────

def test_batch_returns_one_entry_per_pool(monkeypatch):
calls: list[str] = []

def fake(pool: str) -> dict:
calls.append(pool)
return {"pool": pool, "sample_size": 7, "by_project": {}, "by_user": {}}

monkeypatch.setattr(fleet, "_compute_pool_sources", fake)
out = fleet.pool_sources_batch("alpha,beta,gamma")

assert sorted(out["sources"]) == ["alpha", "beta", "gamma"]
assert sorted(calls) == ["alpha", "beta", "gamma"]
assert out["sources"]["beta"]["sample_size"] == 7


def test_batch_payload_is_identical_to_the_single_pool_route(monkeypatch):
"""The frontend swapped N single calls for one batch; the per-pool value must not drift."""
payload = {"pool": "alpha", "sample_size": 3,
"by_project": {"autoland": 3}, "by_user": {"someone": 3}}
monkeypatch.setattr(fleet, "_compute_pool_sources", lambda pool: dict(payload, pool=pool))

single = fleet.pool_sources("alpha")
cache._cache.clear()
batched = fleet.pool_sources_batch("alpha")["sources"]["alpha"]

assert batched == single


def test_batch_is_served_from_the_existing_cache(monkeypatch):
"""A warm batch must do zero computation -- these keys share the single route's cache."""
calls: list[str] = []
monkeypatch.setattr(fleet, "_compute_pool_sources",
lambda pool: (calls.append(pool), {"pool": pool})[1])

fleet.pool_sources_batch("alpha,beta")
assert sorted(calls) == ["alpha", "beta"]

calls.clear()
fleet.pool_sources_batch("alpha,beta")
assert calls == []


def test_one_failing_pool_does_not_fail_the_batch(monkeypatch):
"""Previously each pool was its own request, so one failure cost one card."""
def flaky(pool: str) -> dict:
if pool == "bad":
raise RuntimeError("taskcluster said no")
return {"pool": pool, "sample_size": 1, "by_project": {}, "by_user": {}}

monkeypatch.setattr(fleet, "_compute_pool_sources", flaky)
out = fleet.pool_sources_batch("good,bad,alsogood")

assert sorted(out["sources"]) == ["alsogood", "good"]
assert "bad" not in out["sources"]


def test_empty_batch_short_circuits(monkeypatch):
called = False

def fake(pool: str) -> dict:
nonlocal called
called = True
return {}

monkeypatch.setattr(fleet, "_compute_pool_sources", fake)
assert fleet.pool_sources_batch("") == {"sources": {}}
assert called is False


def test_android_batch_uses_the_android_compute(monkeypatch):
"""The two batches must not share a cache namespace, or pools with the same
name in both would collide."""
monkeypatch.setattr(fleet, "_compute_pool_sources",
lambda pool: {"pool": pool, "which": "hardware"})
monkeypatch.setattr(fleet, "_compute_android_pool_sources",
lambda pool: {"pool": pool, "which": "android"})

assert fleet.pool_sources_batch("p")["sources"]["p"]["which"] == "hardware"
assert fleet.android_pool_sources_batch("p")["sources"]["p"]["which"] == "android"


def test_android_single_route_is_now_cached(monkeypatch):
"""It was the one uncached path: one live TC workers call (limit=1000) per pool,
per page load, 14 times on the Pools page."""
calls: list[str] = []
monkeypatch.setattr(fleet, "_compute_android_pool_sources",
lambda pool: (calls.append(pool), {"pool": pool})[1])

fleet.android_pool_sources("gecko-t-lambda-test-1")
fleet.android_pool_sources("gecko-t-lambda-test-1")

assert calls == ["gecko-t-lambda-test-1"]
52 changes: 49 additions & 3 deletions frontend/src/api.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,39 @@

const BASE = "/api";

/**
* GET requests in flight, keyed by full URL. Several components mount together and
* ask for the same endpoint on one page load (Layout + Overview both want
* /fleet/summary; Workers + CommandPalette both want /fleet/pools). Cloud Armor
* rate-limits per source IP across the whole origin, and corp NAT means a shared
* budget, so a duplicate request is never free. Identical concurrent GETs share
* one response; nothing is cached past settlement, so this changes no semantics.
*/
const inflight = new Map<string, Promise<unknown>>();

/** Retry once per step on the statuses that mean "try again", not "you're wrong". */
const RETRY_STATUSES = new Set([429, 502, 503, 504]);
const RETRY_BACKOFF_MS = [400, 1200];

const sleep = (ms: number) => new Promise(r => setTimeout(r, ms));

async function getOnce<T>(url: string): Promise<T> {
for (let attempt = 0; ; attempt++) {
const res = await fetch(url);
if (res.ok) return res.json() as Promise<T>;
if (!RETRY_STATUSES.has(res.status) || attempt >= RETRY_BACKOFF_MS.length) {
throw new Error(`${res.status} ${res.statusText} — ${new URL(url).pathname}`);
}
// Honour Retry-After when the server sends one (seconds or HTTP-date).
const hdr = res.headers.get("Retry-After");
const hinted = hdr ? (/^\d+$/.test(hdr) ? Number(hdr) * 1000 : Date.parse(hdr) - Date.now()) : NaN;
const wait = Number.isFinite(hinted) && hinted > 0
? Math.min(hinted, 5_000)
: RETRY_BACKOFF_MS[attempt] * (0.5 + Math.random()); // jitter: don't resynchronise a burst
await sleep(wait);
}
}

async function get<T>(path: string, params?: Record<string, string | number | boolean | undefined>): Promise<T> {
const url = new URL(path, window.location.origin);
url.pathname = BASE + path;
Expand All @@ -10,9 +43,13 @@ async function get<T>(path: string, params?: Record<string, string | number | bo
if (v !== undefined) url.searchParams.set(k, String(v));
});
}
const res = await fetch(url.toString());
if (!res.ok) throw new Error(`${res.status} ${res.statusText} — ${url.pathname}`);
return res.json();
const key = url.toString();
const existing = inflight.get(key);
if (existing) return existing as Promise<T>;

const p = getOnce<T>(key).finally(() => inflight.delete(key));
inflight.set(key, p);
return p;
}

async function post<T>(path: string, body?: unknown): Promise<T> {
Expand Down Expand Up @@ -173,6 +210,11 @@ export interface PoolSources {
by_user: Record<string, number>;
}

/** Many pools' job-source breakdowns in one response, keyed by pool name. */
export interface PoolSourcesBatch {
sources: Record<string, PoolSources>;
}

export interface CloudPool {
id: string;
name: string;
Expand Down Expand Up @@ -649,12 +691,16 @@ export const api = {
pools: () => get<PoolsResponse>("/fleet/pools"),
pendingCounts: () => get<PendingCountsResponse>("/fleet/pending-counts"),
poolSources: (pool: string) => get<PoolSources>("/fleet/pool-sources", { pool }),
poolSourcesBatch: (pools: string[]) =>
get<PoolSourcesBatch>("/fleet/pool-sources-batch", { pools: pools.join(",") }),
failures: (days = 7, platform?: string) => get<FailureInsights>("/fleet/failures", { days, platform }),
loadHistory: (hours = 48, includeSeries = false) =>
get<LoadHistory>("/fleet/load-history", { hours, include_series: includeSeries || undefined }),
cloudPools: () => get<CloudPoolsResponse>("/fleet/cloud-pools"),
androidPools: () => get<CloudPoolsResponse>("/fleet/android-pools"),
androidPoolSources: (pool: string) => get<PoolSources>("/fleet/android-pool-sources", { pool }),
androidPoolSourcesBatch: (pools: string[]) =>
get<PoolSourcesBatch>("/fleet/android-pool-sources-batch", { pools: pools.join(",") }),
androidDevices: () => get<AndroidDevicesResponse>("/fleet/android-devices"),
showcase: () => get<ShowcaseData>("/fleet/showcase"),
},
Expand Down
8 changes: 5 additions & 3 deletions frontend/src/pages/Overview.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -172,9 +172,11 @@ function MonitoredPools({ load }: { load: LoadHistory | null }) {
// fetch them only when the pinned set changes rather than on every poll.
const monitoredKey = monitored.join(",");
useEffect(() => {
for (const name of monitoredKey.split(",").filter(Boolean)) {
api.fleet.poolSources(name).then(s => setSources(prev => ({ ...prev, [name]: s }))).catch(() => {});
}
const names = monitoredKey.split(",").filter(Boolean);
if (!names.length) return;
api.fleet.poolSourcesBatch(names)
.then(d => setSources(prev => ({ ...prev, ...d.sources })))
.catch(() => {});
}, [monitoredKey]);

const byName = new Map((load?.pools ?? []).map(p => [p.pool, p] as const));
Expand Down
Loading
Loading