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
32 changes: 31 additions & 1 deletion src/berth/daemon/admin_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,9 +8,10 @@
from fastapi.responses import StreamingResponse

from berth.backends.base import Backend
from berth.daemon.admin import get_backends, router
from berth.daemon.admin import get_backends, get_conn, router
from berth.store import deployments as dep_store
from berth.store import nodes as nodes_store
from berth.store import usage_events as _usage_events


@router.get("/requests")
Expand Down Expand Up @@ -72,6 +73,35 @@ def predictor_stats(request: Request):
return task.stats_snapshot()


@router.get("/usage/series")
def usage_series(
window_s: int = 86400,
bucket_s: int = 3600,
group_by: str | None = None,
conn: sqlite3.Connection = Depends(get_conn),
):
"""Read-only aggregate request volume + tokens over time, bucketed.

Powers the Overview dashboard's volume/throughput tiles and the
traffic-over-time / top-models history. Bounded by the daemon's usage
retention window. group_by may be 'model' or 'key' (or omitted/'none').
"""
if window_s <= 0 or bucket_s <= 0:
raise HTTPException(400, "window_s and bucket_s must be positive")
if window_s // bucket_s > 1024:
raise HTTPException(400, "too many buckets requested (cap is 1024)")
if group_by not in (None, "none", "model", "key"):
raise HTTPException(400, "group_by must be one of: model, key")
gb = None if group_by in (None, "none") else group_by
data = _usage_events.series_in_window(
conn, window_s=window_s, bucket_s=bucket_s, group_by=gb,
)
base = {"window_s": window_s, "bucket_s": bucket_s, "group_by": gb}
if gb:
return {**base, "groups": data}
return {**base, "buckets": data}


@router.get("/deployments/current/logs")
def stream_current_logs(request: Request):
conn: sqlite3.Connection = request.app.state.conn
Expand Down
137 changes: 137 additions & 0 deletions src/berth/store/usage_events.py
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,143 @@ def list_recent(
return [_row(r) for r in rows]


def series_in_window(
conn: sqlite3.Connection,
*,
window_s: int,
bucket_s: int,
group_by: str | None = None,
) -> list[dict]:
"""Time-bucketed request volume + tokens over the past `window_s` seconds.

Mirrors `key_usage.bucketed_usage`: buckets are `bucket_s` wide, zero-filled,
returned oldest-first so the UI can plot left-to-right.

- ``group_by=None`` -> a flat list of bucket dicts:
``[{bucket_idx, ts_offset_s, count, tokens_in, tokens_out}, ...]``.
- ``group_by='model'|'key'`` -> one entry per group, sorted by total desc:
``[{key, label, total, tokens_in, tokens_out, buckets:[...]}, ...]`` where
each group carries its own zero-filled bucket list.

Bounded by retention: rows older than the daemon's retention window have
already been rolled up into usage_aggregates and dropped, so the effective
history can be shorter than `window_s`.
"""
bucket_s = max(1, int(bucket_s))
num_buckets = max(1, int(window_s) // bucket_s)
col = {None: None, "model": "model_name", "key": "api_key_id"}[group_by]

def _empty_buckets() -> list[dict]:
# bucket_idx 0 = most recent; emit oldest -> newest.
return [
{
"bucket_idx": i,
"ts_offset_s": i * bucket_s,
"count": 0,
"tokens_in": 0,
"tokens_out": 0,
}
for i in range(num_buckets - 1, -1, -1)
]

def _fill(buckets: list[dict], by_idx: dict) -> None:
for b in buckets:
r = by_idx.get(b["bucket_idx"])
if r is not None:
b["count"] = int(r["count"])
b["tokens_in"] = int(r["tokens_in"])
b["tokens_out"] = int(r["tokens_out"])

# Fully literal SQL per group_by case — no string interpolation of column
# names into the query (avoids SQL-injection surface entirely; the group
# column is an internal whitelist, never user input). Only the bucket width
# and window are bound parameters.
params = (bucket_s, f"-{window_s} seconds")
if col is None:
rows = conn.execute(
"""
SELECT
CAST((CAST(strftime('%s', 'now') AS INTEGER)
- CAST(strftime('%s', ts) AS INTEGER)) / ? AS INTEGER) AS bucket_idx,
COUNT(*) AS count,
COALESCE(SUM(tokens_in), 0) AS tokens_in,
COALESCE(SUM(tokens_out), 0) AS tokens_out
FROM usage_events
WHERE ts > datetime('now', ?)
GROUP BY bucket_idx
""",
params,
).fetchall()
elif col == "model_name":
rows = conn.execute(
"""
SELECT
model_name AS grp,
CAST((CAST(strftime('%s', 'now') AS INTEGER)
- CAST(strftime('%s', ts) AS INTEGER)) / ? AS INTEGER) AS bucket_idx,
COUNT(*) AS count,
COALESCE(SUM(tokens_in), 0) AS tokens_in,
COALESCE(SUM(tokens_out), 0) AS tokens_out
FROM usage_events
WHERE ts > datetime('now', ?)
GROUP BY grp, bucket_idx
""",
params,
).fetchall()
else: # col == "api_key_id"
rows = conn.execute(
"""
SELECT
api_key_id AS grp,
CAST((CAST(strftime('%s', 'now') AS INTEGER)
- CAST(strftime('%s', ts) AS INTEGER)) / ? AS INTEGER) AS bucket_idx,
COUNT(*) AS count,
COALESCE(SUM(tokens_in), 0) AS tokens_in,
COALESCE(SUM(tokens_out), 0) AS tokens_out
FROM usage_events
WHERE ts > datetime('now', ?)
GROUP BY grp, bucket_idx
""",
params,
).fetchall()

if col is None:
out = _empty_buckets()
_fill(out, {int(r["bucket_idx"]): r for r in rows})
return out

groups: dict[str, dict] = {}
for r in rows:
raw = r["grp"]
gkey = "" if raw is None else str(raw)
g = groups.setdefault(gkey, {
"key": gkey,
"label": "—" if raw is None else str(raw),
"total": 0,
"tokens_in": 0,
"tokens_out": 0,
"_by_idx": {},
})
g["total"] += int(r["count"])
g["tokens_in"] += int(r["tokens_in"])
g["tokens_out"] += int(r["tokens_out"])
g["_by_idx"][int(r["bucket_idx"])] = r

result: list[dict] = []
for g in sorted(groups.values(), key=lambda x: x["total"], reverse=True):
buckets = _empty_buckets()
_fill(buckets, g["_by_idx"])
result.append({
"key": g["key"],
"label": g["label"],
"total": g["total"],
"tokens_in": g["tokens_in"],
"tokens_out": g["tokens_out"],
"buckets": buckets,
})
return result


def purge_older_than(
conn: sqlite3.Connection, *, before_iso: str,
) -> int:
Expand Down
2 changes: 2 additions & 0 deletions src/berth/ui/assets/index-B1_2_G6r.css

Large diffs are not rendered by default.

2 changes: 0 additions & 2 deletions src/berth/ui/assets/index-BsGljYOY.css

This file was deleted.

Large diffs are not rendered by default.

4 changes: 2 additions & 2 deletions src/berth/ui/index.html
Original file line number Diff line number Diff line change
Expand Up @@ -15,8 +15,8 @@
href="https://fonts.googleapis.com/css2?family=JetBrains+Mono:wght@300;400;500;600;700&display=swap"
rel="stylesheet"
/>
<script type="module" crossorigin src="/assets/index-Drstn1oh.js"></script>
<link rel="stylesheet" crossorigin href="/assets/index-BsGljYOY.css">
<script type="module" crossorigin src="/assets/index-CNEC70C6.js"></script>
<link rel="stylesheet" crossorigin href="/assets/index-B1_2_G6r.css">
</head>
<body>
<div id="root"></div>
Expand Down
42 changes: 42 additions & 0 deletions tests/unit/test_admin_endpoints.py
Original file line number Diff line number Diff line change
Expand Up @@ -1008,3 +1008,45 @@ async def test_deploy_with_backend_adopted_returns_400(app):
)
assert r.status_code == 400, r.text
assert "reserved" in r.json()["detail"].lower()


@pytest.mark.asyncio
async def test_usage_series_ungrouped(app):
"""/admin/usage/series returns zero-filled buckets and counts seeded events."""
usage_store.record(app.state.conn, model_name="qwen", base_name="qwen", tokens_out=12)
usage_store.record(app.state.conn, model_name="qwen", base_name="qwen", tokens_out=8)
transport = httpx.ASGITransport(app=app)
async with httpx.AsyncClient(transport=transport, base_url="http://test") as c:
r = await c.get("/admin/usage/series?window_s=3600&bucket_s=3600")
assert r.status_code == 200, r.text
body = r.json()
assert body["group_by"] is None
assert body["buckets"][-1]["count"] == 2
assert body["buckets"][-1]["tokens_out"] == 20


@pytest.mark.asyncio
async def test_usage_series_group_by_model(app):
usage_store.record(app.state.conn, model_name="a", base_name="a")
usage_store.record(app.state.conn, model_name="a", base_name="a")
usage_store.record(app.state.conn, model_name="b", base_name="b")
transport = httpx.ASGITransport(app=app)
async with httpx.AsyncClient(transport=transport, base_url="http://test") as c:
r = await c.get("/admin/usage/series?window_s=3600&bucket_s=3600&group_by=model")
assert r.status_code == 200, r.text
body = r.json()
assert body["group_by"] == "model"
assert [g["label"] for g in body["groups"]] == ["a", "b"]
assert body["groups"][0]["total"] == 2


@pytest.mark.asyncio
async def test_usage_series_rejects_bad_params(app):
transport = httpx.ASGITransport(app=app)
async with httpx.AsyncClient(transport=transport, base_url="http://test") as c:
bad_window = await c.get("/admin/usage/series?window_s=0&bucket_s=60")
too_many = await c.get("/admin/usage/series?window_s=10000000&bucket_s=1")
bad_group = await c.get("/admin/usage/series?window_s=3600&bucket_s=60&group_by=bogus")
assert bad_window.status_code == 400, bad_window.text
assert too_many.status_code == 400, too_many.text
assert bad_group.status_code == 400, bad_group.text
69 changes: 69 additions & 0 deletions tests/unit/test_usage_series.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,69 @@
"""Store tests for the bucketed usage-series query that powers the Overview
dashboard's volume/throughput tiles. The /admin/usage/series route handler is
covered HTTP-level in test_admin_endpoints.py (reusing its app fixture, which
imports the daemon modules in the right order — admin_runtime imported on its
own hits a pre-existing admin<->admin_runtime circular import).
"""
from berth.store import db
from berth.store import usage_events as ue


def _fresh(tmp_path):
conn = db.connect(tmp_path / "t.db")
db.init_schema(conn)
return conn


# --- store: series_in_window -------------------------------------------------

def test_series_ungrouped_counts_and_tokens(tmp_path):
conn = _fresh(tmp_path)
for _ in range(3):
ue.record(conn, model_name="m", base_name="m", tokens_in=2, tokens_out=10)
rows = ue.series_in_window(conn, window_s=3600, bucket_s=3600, group_by=None)
assert len(rows) == 1 # window_s // bucket_s == 1 bucket
assert rows[-1]["count"] == 3
assert rows[-1]["tokens_in"] == 6
assert rows[-1]["tokens_out"] == 30


def test_series_zero_fills_empty_buckets_oldest_first(tmp_path):
conn = _fresh(tmp_path)
ue.record(conn, model_name="m", base_name="m", tokens_out=1)
rows = ue.series_in_window(conn, window_s=3600, bucket_s=600, group_by=None)
assert len(rows) == 6 # 3600 / 600
# oldest first; all but the most-recent bucket are empty
assert rows[0]["count"] == 0
assert rows[-1]["count"] == 1
# ts_offset decreases toward the present
assert rows[0]["ts_offset_s"] > rows[-1]["ts_offset_s"]


def test_series_group_by_model_sorted_by_total(tmp_path):
conn = _fresh(tmp_path)
for _ in range(3):
ue.record(conn, model_name="busy", base_name="busy", tokens_out=5)
ue.record(conn, model_name="quiet", base_name="quiet", tokens_out=7)
groups = ue.series_in_window(conn, window_s=3600, bucket_s=3600, group_by="model")
assert [g["label"] for g in groups] == ["busy", "quiet"] # busy first (higher total)
by_label = {g["label"]: g for g in groups}
assert by_label["busy"]["total"] == 3
assert by_label["quiet"]["tokens_out"] == 7
# each group still carries a zero-filled bucket list
assert by_label["busy"]["buckets"][-1]["count"] == 3


def test_series_group_by_key(tmp_path):
conn = _fresh(tmp_path)
ue.record(conn, model_name="m", base_name="m", api_key_id=None, tokens_out=1)
groups = ue.series_in_window(conn, window_s=3600, bucket_s=3600, group_by="key")
# NULL api_key id collapses to the "—" label
assert groups[0]["label"] == "—"
assert groups[0]["total"] == 1


def test_series_empty_window_returns_zeroed_buckets(tmp_path):
conn = _fresh(tmp_path)
rows = ue.series_in_window(conn, window_s=3600, bucket_s=1800, group_by=None)
assert len(rows) == 2
assert all(r["count"] == 0 for r in rows)
Loading