mirror of
https://github.com/rennf93/roboco.git
synced 2026-08-03 07:23:24 +02:00
* feat(metrics): capture per-session turns + tool_calls (phase 1)
Persist LLM iterations (turns) and tool invocations per agent spawn session,
the raw signal the granular per-member performance metrics build on (real
effort/iterations vs wall-clock).
- sum_transcript_usage returns a 5-tuple adding turns = unique assistant
message-id count; _usage_from_transcript + _resolve_active_tokens updated to
the 5-tuple (active-tokens keeps its 4-tuple contract by slicing).
- SDK: _SessionState.turns, set by /usage/sync; /usage/status (TokenUsageStatus)
now carries turns + tool_calls (= total_calls).
- orchestrator: new _resolve_final_turns_tools (SDK primary, transcript fallback
for turns only; Grok -> 0/0) wired into _finalize_spawn_session, which writes
turns + tool_calls to agent_spawn_sessions.
- migration 055 adds turns + tool_calls (BigInteger DEFAULT 0 -> historical/Grok
rows read 0, surfaced as n/a). Verified real alembic upgrade/downgrade.
Part of metrics-granularity (v0.15.0); recon-adjusted plan on disk.
* feat(metrics): pure compute_stage_effort helper (phase 2, part 1)
Foundation-layer overlap math (no DB): split each task status window into
active (merged wall-clock overlap of spawn stints — concurrent stints counted
once, so active <= window) vs wait (queue/review idle). Distinct from summed
effort. The per-task metrics service will feed it audit-log windows + spawn
stints. 9 unit tests (disjoint/nested/partial/merged/clamped/zero/multi-window).
* feat(metrics): per-task live metrics + GET /metrics/task/{id} (phase 2)
TaskMetrics dataclass + MetricsService.get_task_metrics: summed spawn effort
(vs wall-clock), turns/tool_calls/tokens/cost, per-stage active-vs-wait
(compute_stage_effort over audit windows x spawn stints), and who-caused-rework
(revision_count + named qa/pr fail events). Open stints and the open final
stage window close at completed_at for a terminal task (else now), so stages
don't grow past completion. Exposed at GET /dashboard/metrics/task/{task_id}
(404 if absent). Real-PG tests (compose/none/in-flight) + route tests (200/404).
* feat(metrics): CEO-as-member scorecard + ceo_reject audit regression (phase 3)
The human CEO is a measured member, read purely from audit_log (agent_role='ceo'
serializes from the CEO StrEnum): approval dwell (awaiting_ceo_approval -> a CEO
decision, incl. the coordination-root reject that lands in pending), unblock
dwell (blocked -> a CEO revive), and god-mode action count (every CEO-attributed
transition). CeoScorecard + MetricsService.get_ceo_scorecard (p50/p90 via
PERCENTILE_CONT, expanding IN for the decision sets) + GET
/dashboard/metrics/member/ceo (declared before any future member/{id} route).
The ceo_reject coordination-root audit gap the plan meant to close was already
closed by the gap-sweep (routes through admin_set_status -> agent_role='ceo'
audit); locked with a regression assertion in the existing coordination-reject
test. Real-PG tests: approval/unblock/godmode, non-ceo exclusion, empty->zeros.
* feat(metrics): audit instrumentation for escalations/blocked-others/idle (phase 4a)
The three extra per-member metrics that had no data source get durable,
in-session audit events (additive; never gate the underlying action):
- apply_escalation -> task.escalated (details.escalator_slug) on both the
normal block path and the pool-divert path -> escalations count.
- _unblock_dependents -> task.unblocked_dependents (details.count) on the
completed BLOCKER task, captured before the dependency edges are pruned ->
blocked-others count (sweeper attributes to the blocker's owner).
- mark_agent_idle -> agent.idle (details.agent_slug) -> idle/utilization (the
sweeper pairs an idle mark to the member's next spawn for idle duration).
(QA pass-rate needs no new event — reuses task.awaiting_documentation[qa] +
task.qa_fail.) Real-PG tests for each; 111 transition tests still green.
* feat(metrics): member_performance_daily rollup table + migration 056 (phase 4b)
The per-member scorecard rollup: one row per (date, member_kind, agent_slug),
CEO as a first-class member_kind='ceo' row (agent_slug='' NOT NULL so the
NULL-distinct UNIQUE keeps it unique). Full column set + the four CEO-approved
extras (qa_reviews_total/passed, escalations, blocked_others, idle_seconds) plus
blocked_seconds. Overwrite-upsert on (date, member_kind, agent_slug) for an
idempotent sweep. Migration 056 verified real up/down (24 cols, 4 indexes).
* feat(metrics): _sweep_member_performance rollup sweeper (phase 4c)
The daily per-member rollup sweep (mirrors _sweep_daily_rollup): a trailing
7-day, idempotent overwrite-upsert wired into _run_sweep. One focused query per
metric merges into a (date, agent_slug) accumulator — spawn effort/turns/tokens/
cost, completed/first-pass/revisions-received, revisions-caused (qa/pr fails),
QA pass-rate (passed + total), escalations (by escalator_slug), blocked-others
(unblocked_dependents by blocker owner), idle_seconds (idle mark -> next spawn),
blocked_seconds (blocked dwell) — plus one CEO row/day (approval/unblock dwell +
god-mode). Real-PG test asserts every facet + idempotency (a 2nd sweep
overwrites, never doubles); spawn-day != completion-day split is by-design.
* feat(metrics): member/org rollup scorecards + endpoints + live overlay (phase 5)
MemberScorecard + OrgScorecard with derived rates (FPY, effort-throughput,
turns/tool-calls per task, QA pass-rate, utilization) — all division-guarded to
None. get_member_scorecard reads member_performance_daily by slug and overlays
the member's live in-flight (non-terminal) tasks' effort via get_task_metrics
(disjoint by status: completion counts stay rollup-only, overlay only enriches
effort/turns/cost; includes_live_inflight flags it). get_org_scorecard
aggregates the cell (?team=) or whole org. Routes: GET /metrics/member/{agent_id}
(404 if absent, after the ceo literal route) + GET /metrics/org?team=. Real-PG
tests (derived rates, overlay no double-count, guards, org) + route tests.
* feat(metrics): granular CEO completion notification (phase 6)
There was no CEO completion notification at all (EventType.TASK_COMPLETED was
defined but never emitted). Add notify_ceo_of_completion in
NotificationDeliveryService — a granular body (real effort vs wall-clock +
stints/turns/tool-calls/revisions[QA/PR]/cost from get_task_metrics; degrades to
wall-clock-only, turns 'n/a', when there are no spawn sessions). Reuses the
existing ALERT type (no enum migration; the notificationtype PG enum is fixed at
001). ceo_approve now emits TASK_COMPLETED + fires the notification (best-effort
via _notify_completion — never blocks completion); complete() emits
TASK_COMPLETED too (closes the dead-code gap; the WS bridge can forward it).
Pure formatter tests + real-PG notification test.
* [metrics-granularity] Phase 7: panel Scorecards tab + dashboard overview
Add the CEO-facing metrics surfaces for the granularity feature:
- New "Scorecards" tab on the Metrics page: org rollup headline, the
CEO-as-member card (approval/unblock dwell + god-mode count), and a
per-member table (completed, first-pass yield, active effort, turns/task,
QA pass-rate, escalations, blocked-others, utilization). Each member row
self-fetches its rollup scorecard; live in-flight rows carry a "live" badge.
- New dashboard overview card (ScorecardOverviewPanel): org-wide 30-day
headline (completed, FPY, throughput/hr, active effort, cost) deep-linking
into the Scorecards tab.
- Plumbing: TaskMetrics/MemberScorecard/OrgScorecard/CeoScorecard types,
observability API client methods + empty fallbacks, and the four
useCeoScorecard/useMemberScorecard/useOrgScorecard/useTaskMetrics hooks.
Panel gate green: tsc, eslint, prettier, vitest (175 tests, +6 new).
* [metrics-granularity] test: make completion-notification robust to shared-DB CEO
test_notify_ceo_of_completion_creates_alert errored in the full suite (passed
in isolation): the session-scoped test DB is shared across the run, and the
sibling real-DB board-gate test commits a role=CEO agent (slug="ceo") without
cleanup — so my env fixture's hardcoded slug="ceo" insert hit a unique-constraint
violation, and a second role=CEO row would also make _get_ceo_agent()'s
scalar_one_or_none() raise. Reuse an existing CEO when present (the singleton the
production system actually has), else create one with a unique slug. Order-
independent. Also reflow test_metrics_instrumentation.py to ruff format.
* chore(release): 0.15.0
Metrics granularity: per-member/per-task/org + CEO-as-member scorecards,
turn/tool-call capture (migration 055), member_performance_daily rollup
(migration 056) with QA pass-rate / escalations / blocked-others / utilization,
per-task active-vs-wait metrics, granular completion notification, panel
Scorecards tab + dashboard Performance card, and the ceo_reject audit fix.
Version bump across the canonical set + CHANGELOG.
* [metrics-granularity] fix pre-tag audit findings (overlay double-count + panel error states)
Adversarial review before the v0.15.0 tag surfaced two real logical gaps:
- MAJOR (backend): the live in-flight overlay re-summed ALL sessions of every
non-terminal task via get_task_metrics, but _msweep_spawn already rolls up
every CLOSED session regardless of task status — so a closed session on a
still-open task was counted twice (rollup + overlay), permanently inflating a
member's effort/turns/tokens/cost on the common reap/respawn path. The overlay
now sums only OPEN sessions (ended_at IS NULL), which the closed-only rollup
can never contain — disjoint by construction. A just-closed session lands in
the rollup on the next ~60s sweep (no gap of note). Aggregated in SQL to mirror
_msweep_spawn. Regression test reproduces the double-count (turns 10→5).
- MAJOR (panel): the four new scorecard surfaces used `isLoading || !data` with
no isError branch, so a failed query span forever on a skeleton. They now
surface a load error. Tests added.
Also: OrgSummary active-effort formatting no longer round-trips hours→seconds→
hours; dashboard grid uses xl:grid-cols-4 (was 2xl) so 4 panels show at 1280px;
corrected the inaccurate "NULL distinct" CEO-row uniqueness comment (agent_slug
is NOT NULL; the '' tuple is simply distinct from agent rows).
make quality GREEN (cov 95.31%); panel GREEN (vitest 178).
* [metrics-granularity] fix: decode bytes stream message-id before XCLAIM
StreamEventBus._recover_stream passed the pending message id to XCLAIM via
str() on the raw bytes the client returns (redis client has no
decode_responses), producing "b'1782066556728-0'". Redis rejects that with
"Unrecognized XCLAIM option", so pending-message recovery threw on every
reclaim tick and unacked messages from crashed/slow consumers were never
reclaimed (leaking in the PEL on every stream, spamming the error log). Decode
via the existing _to_str helper — the fix the sibling claim path already uses.
Pre-existing in v0.14.0 (unrelated to metrics granularity); folded into this
release per CEO. TDD regression test + CHANGELOG entry. make quality GREEN.
---------
Co-authored-by: Renn F <rennf93@users.noreply.github.com>
248 lines
8.7 KiB
Python
248 lines
8.7 KiB
Python
"""Token-usage capture — /usage/sync parses the transcript and sets totals.
|
|
|
|
The agent SDK exposes /usage/report (additive) and /usage/status (read),
|
|
but nothing ever fed token counts in, so every session reported zero and the
|
|
cost dashboard rendered all-zeros. The fix: the usage-report hook hands the
|
|
SDK the Claude Code transcript path; /usage/sync parses the per-message
|
|
``usage`` blocks and *sets* the cumulative totals absolutely. These tests pin
|
|
that contract — correct summation, idempotency (no double-count on re-sync),
|
|
graceful handling of a missing/partial transcript, and growth on re-sync.
|
|
|
|
Expected totals are derived from the input rows (no magic literals), so the
|
|
assertions track whatever the fixtures declare.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
from typing import TYPE_CHECKING
|
|
|
|
import pytest
|
|
import roboco.agent_sdk.server as srv
|
|
from fastapi.testclient import TestClient
|
|
|
|
if TYPE_CHECKING:
|
|
from collections.abc import Iterator, Sequence
|
|
from pathlib import Path
|
|
|
|
_OK = 200
|
|
|
|
# Each row is (input, output, cache_read, cache_write).
|
|
_UsageRow = tuple[int, int, int, int]
|
|
|
|
|
|
@pytest.fixture(autouse=True)
|
|
def _reset_state() -> Iterator[None]:
|
|
srv._state.reset()
|
|
yield
|
|
srv._state.reset()
|
|
|
|
|
|
@pytest.fixture(autouse=True)
|
|
def _transcript_base(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None:
|
|
# /usage/sync contains transcript_path under ROBOCO_TRANSCRIPT_DIR; point
|
|
# it at the test's tmp_path so the under-base transcripts here pass the guard.
|
|
monkeypatch.setenv("ROBOCO_TRANSCRIPT_DIR", str(tmp_path))
|
|
|
|
|
|
@pytest.fixture
|
|
def client() -> TestClient:
|
|
return TestClient(srv.app)
|
|
|
|
|
|
def _assistant_line(row: _UsageRow) -> str:
|
|
inp, out, cread, cwrite = row
|
|
return json.dumps(
|
|
{
|
|
"type": "assistant",
|
|
"message": {
|
|
"role": "assistant",
|
|
"usage": {
|
|
"input_tokens": inp,
|
|
"output_tokens": out,
|
|
"cache_read_input_tokens": cread,
|
|
"cache_creation_input_tokens": cwrite,
|
|
},
|
|
},
|
|
}
|
|
)
|
|
|
|
|
|
def _write(path: Path, *lines: str) -> None:
|
|
path.write_text("\n".join(lines) + "\n", encoding="utf-8")
|
|
|
|
|
|
def _expected(rows: Sequence[_UsageRow]) -> dict[str, int]:
|
|
return {
|
|
"tokens_input": sum(r[0] for r in rows),
|
|
"tokens_output": sum(r[1] for r in rows),
|
|
"tokens_cache_read": sum(r[2] for r in rows),
|
|
"tokens_cache_write": sum(r[3] for r in rows),
|
|
}
|
|
|
|
|
|
def test_sums_usage_across_assistant_messages(
|
|
client: TestClient, tmp_path: Path
|
|
) -> None:
|
|
rows: list[_UsageRow] = [(100, 20, 5, 3), (50, 10, 2, 1)]
|
|
transcript = tmp_path / "session.jsonl"
|
|
_write(transcript, *(_assistant_line(r) for r in rows))
|
|
resp = client.post("/usage/sync", json={"transcript_path": str(transcript)})
|
|
assert resp.status_code == _OK
|
|
body = resp.json()
|
|
for key, value in _expected(rows).items():
|
|
assert body[key] == value
|
|
|
|
|
|
def test_status_reflects_synced_totals(client: TestClient, tmp_path: Path) -> None:
|
|
rows: list[_UsageRow] = [(200, 40, 0, 0)]
|
|
transcript = tmp_path / "session.jsonl"
|
|
_write(transcript, *(_assistant_line(r) for r in rows))
|
|
client.post("/usage/sync", json={"transcript_path": str(transcript)})
|
|
status = client.get("/usage/status").json()
|
|
for key, value in _expected(rows).items():
|
|
assert status[key] == value
|
|
|
|
|
|
def test_resync_is_idempotent_not_additive(client: TestClient, tmp_path: Path) -> None:
|
|
"""The set is absolute — syncing the same transcript twice must not double."""
|
|
rows: list[_UsageRow] = [(100, 20, 0, 0)]
|
|
transcript = tmp_path / "session.jsonl"
|
|
_write(transcript, *(_assistant_line(r) for r in rows))
|
|
client.post("/usage/sync", json={"transcript_path": str(transcript)})
|
|
client.post("/usage/sync", json={"transcript_path": str(transcript)})
|
|
status = client.get("/usage/status").json()
|
|
for key, value in _expected(rows).items():
|
|
assert status[key] == value
|
|
|
|
|
|
def test_resync_after_growth_overwrites_with_new_total(
|
|
client: TestClient, tmp_path: Path
|
|
) -> None:
|
|
first: list[_UsageRow] = [(100, 20, 0, 0)]
|
|
grown: list[_UsageRow] = [(100, 20, 0, 0), (80, 15, 0, 0)]
|
|
transcript = tmp_path / "session.jsonl"
|
|
_write(transcript, *(_assistant_line(r) for r in first))
|
|
client.post("/usage/sync", json={"transcript_path": str(transcript)})
|
|
# The transcript grows as the turn continues.
|
|
_write(transcript, *(_assistant_line(r) for r in grown))
|
|
client.post("/usage/sync", json={"transcript_path": str(transcript)})
|
|
status = client.get("/usage/status").json()
|
|
for key, value in _expected(grown).items():
|
|
assert status[key] == value
|
|
|
|
|
|
def test_missing_transcript_returns_zero_without_error(
|
|
client: TestClient, tmp_path: Path
|
|
) -> None:
|
|
resp = client.post(
|
|
"/usage/sync", json={"transcript_path": str(tmp_path / "nope.jsonl")}
|
|
)
|
|
assert resp.status_code == _OK
|
|
body = resp.json()
|
|
for key, value in _expected([]).items():
|
|
assert body[key] == value
|
|
assert body["turns"] == 0
|
|
assert body["tool_calls"] == 0
|
|
|
|
|
|
def test_malformed_lines_are_skipped(client: TestClient, tmp_path: Path) -> None:
|
|
rows: list[_UsageRow] = [(100, 20, 0, 0), (50, 10, 0, 0)]
|
|
transcript = tmp_path / "session.jsonl"
|
|
_write(
|
|
transcript,
|
|
"not json at all",
|
|
_assistant_line(rows[0]),
|
|
json.dumps({"type": "user", "message": {"role": "user"}}), # no usage
|
|
"{ broken",
|
|
_assistant_line(rows[1]),
|
|
)
|
|
body = client.post("/usage/sync", json={"transcript_path": str(transcript)}).json()
|
|
exp = _expected(rows)
|
|
assert body["tokens_input"] == exp["tokens_input"]
|
|
assert body["tokens_output"] == exp["tokens_output"]
|
|
|
|
|
|
def test_parser_handles_entries_without_message(tmp_path: Path) -> None:
|
|
rows: list[_UsageRow] = [(10, 5, 0, 0)]
|
|
transcript = tmp_path / "session.jsonl"
|
|
_write(
|
|
transcript,
|
|
json.dumps({"type": "system", "subtype": "init"}),
|
|
_assistant_line(rows[0]),
|
|
)
|
|
tin, tout, cread, cwrite, turns = srv._sum_transcript_usage(transcript)
|
|
exp = _expected(rows)
|
|
assert (tin, tout, cread, cwrite) == (
|
|
exp["tokens_input"],
|
|
exp["tokens_output"],
|
|
exp["tokens_cache_read"],
|
|
exp["tokens_cache_write"],
|
|
)
|
|
assert turns == 0 # _assistant_line carries no message id
|
|
|
|
|
|
def _assistant_line_with_id(row: _UsageRow, message_id: str) -> str:
|
|
"""An assistant line carrying a message id (for de-duplication tests)."""
|
|
inp, out, cread, cwrite = row
|
|
return json.dumps(
|
|
{
|
|
"type": "assistant",
|
|
"message": {
|
|
"id": message_id,
|
|
"role": "assistant",
|
|
"usage": {
|
|
"input_tokens": inp,
|
|
"output_tokens": out,
|
|
"cache_read_input_tokens": cread,
|
|
"cache_creation_input_tokens": cwrite,
|
|
},
|
|
},
|
|
}
|
|
)
|
|
|
|
|
|
def test_parser_dedupes_repeated_message_id(tmp_path: Path) -> None:
|
|
"""One message logged across several lines (same id) is counted once.
|
|
|
|
Claude Code emits one transcript line per content block (thinking, text,
|
|
tool_use), each repeating the message's ``usage``. Summing every line would
|
|
roughly double the totals, so the parser must de-duplicate by message id.
|
|
"""
|
|
msg = (100, 20, 5, 3)
|
|
other = (7, 2, 1, 0)
|
|
transcript = tmp_path / "session.jsonl"
|
|
_write(
|
|
transcript,
|
|
_assistant_line_with_id(msg, "msg_aaa"), # thinking block
|
|
_assistant_line_with_id(msg, "msg_aaa"), # text block (same id)
|
|
_assistant_line_with_id(msg, "msg_aaa"), # tool_use block (same id)
|
|
_assistant_line_with_id(other, "msg_bbb"),
|
|
)
|
|
tin, tout, cread, cwrite, turns = srv._sum_transcript_usage(transcript)
|
|
# Counted once per id: msg + other, NOT msg * 3 + other.
|
|
exp = _expected([msg, other])
|
|
assert (tin, tout, cread, cwrite) == (
|
|
exp["tokens_input"],
|
|
exp["tokens_output"],
|
|
exp["tokens_cache_read"],
|
|
exp["tokens_cache_write"],
|
|
)
|
|
expected_turns = 2 # two unique message ids
|
|
assert turns == expected_turns
|
|
|
|
|
|
def test_sync_response_surfaces_turns(client: TestClient, tmp_path: Path) -> None:
|
|
"""/usage/sync (and thus /usage/status) reports the LLM turn count."""
|
|
transcript = tmp_path / "session.jsonl"
|
|
_write(
|
|
transcript,
|
|
_assistant_line_with_id((10, 5, 0, 0), "msg_a"),
|
|
_assistant_line_with_id((10, 5, 0, 0), "msg_a"), # same id
|
|
_assistant_line_with_id((2, 1, 0, 0), "msg_b"),
|
|
)
|
|
body = client.post("/usage/sync", json={"transcript_path": str(transcript)}).json()
|
|
expected_turns = 2
|
|
assert body["turns"] == expected_turns
|
|
assert "tool_calls" in body
|