mirror of
https://github.com/rennf93/roboco.git
synced 2026-08-03 07:23:24 +02:00
v0.15.0: Metrics granularity — per-member / per-task / org + CEO scorecards (#289)
* 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>
This commit is contained in:
@@ -0,0 +1,210 @@
|
||||
"""CEO completion notification — granular effort breakdown (phase 6).
|
||||
|
||||
The pure body formatter (real effort vs wall-clock; degrades to wall-clock-only)
|
||||
+ notify_ceo_of_completion end to end against real PG.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from types import SimpleNamespace
|
||||
from typing import TYPE_CHECKING, Any, cast
|
||||
from uuid import UUID, uuid4
|
||||
|
||||
import pytest
|
||||
import pytest_asyncio
|
||||
from roboco.db.tables import (
|
||||
AgentSpawnSessionTable,
|
||||
AgentTable,
|
||||
NotificationTable,
|
||||
ProjectTable,
|
||||
TaskTable,
|
||||
)
|
||||
from roboco.models.base import (
|
||||
AgentRole,
|
||||
AgentStatus,
|
||||
Complexity,
|
||||
NotificationType,
|
||||
TaskNature,
|
||||
TaskStatus,
|
||||
TaskType,
|
||||
Team,
|
||||
)
|
||||
from roboco.models.metrics import TaskMetrics
|
||||
from roboco.services.notification_delivery import (
|
||||
_format_completion_body,
|
||||
get_notification_delivery_service,
|
||||
)
|
||||
from sqlalchemy import select
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from collections.abc import AsyncIterator
|
||||
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
|
||||
def _metrics(**over: Any) -> TaskMetrics:
|
||||
base: dict[str, Any] = {
|
||||
"task_id": str(uuid4()),
|
||||
"active_runtime_seconds": 3600,
|
||||
"wall_clock_seconds": 7200,
|
||||
"turns": 42,
|
||||
"tool_calls": 99,
|
||||
"tokens": 1000,
|
||||
"cost_usd": 4.2,
|
||||
"revision_count": 2,
|
||||
"qa_fails": 1,
|
||||
"pr_fails": 1,
|
||||
"stints": 3,
|
||||
"stages": [],
|
||||
}
|
||||
base.update(over)
|
||||
return TaskMetrics(**base)
|
||||
|
||||
|
||||
def test_format_body_with_metrics() -> None:
|
||||
body = _format_completion_body(
|
||||
cast("TaskTable", SimpleNamespace(title="Auth flow")), _metrics()
|
||||
)
|
||||
assert "Auth flow" in body
|
||||
assert "Active effort: 1.0h across 3 stint(s)" in body
|
||||
assert "42 turns" in body
|
||||
assert "Wall-clock: 2.0h" in body
|
||||
assert "2 (1 QA / 1 PR)" in body
|
||||
assert "$4.2" in body
|
||||
|
||||
|
||||
def test_format_body_turns_na_when_zero() -> None:
|
||||
body = _format_completion_body(
|
||||
cast("TaskTable", SimpleNamespace(title="T")), _metrics(turns=0)
|
||||
)
|
||||
assert "n/a turns" in body # pre-turns-migration / Grok
|
||||
|
||||
|
||||
def test_format_body_degrades_without_metrics() -> None:
|
||||
body = _format_completion_body(cast("TaskTable", SimpleNamespace(title="T")), None)
|
||||
assert body == "Task 'T' completed."
|
||||
|
||||
|
||||
@pytest_asyncio.fixture
|
||||
async def env(db_session: AsyncSession) -> AsyncIterator[dict]:
|
||||
# The CEO is a singleton in the real system, and `_get_ceo_agent()` resolves
|
||||
# it by `role == CEO` with `scalar_one_or_none()`. The session-scoped test DB
|
||||
# is shared across the run, and a sibling real-DB test commits a role=CEO
|
||||
# agent (slug="ceo") without cleanup, so it can already be present here.
|
||||
# Reuse an existing CEO rather than inserting a second one — creating another
|
||||
# would both collide on the unique slug and make `_get_ceo_agent()` raise
|
||||
# MultipleResultsFound. Order-independent: in isolation we create one.
|
||||
existing_ceo = (
|
||||
(
|
||||
await db_session.execute(
|
||||
select(AgentTable).where(AgentTable.role == AgentRole.CEO)
|
||||
)
|
||||
)
|
||||
.scalars()
|
||||
.first()
|
||||
)
|
||||
ceo = existing_ceo or AgentTable(
|
||||
id=uuid4(),
|
||||
name="CEO",
|
||||
slug=f"ceo-{uuid4().hex[:6]}",
|
||||
role=AgentRole.CEO,
|
||||
team=None,
|
||||
status=AgentStatus.ACTIVE,
|
||||
model_config={},
|
||||
system_prompt="x",
|
||||
capabilities=[],
|
||||
permissions={},
|
||||
metrics={},
|
||||
)
|
||||
if existing_ceo is None:
|
||||
db_session.add(ceo)
|
||||
dev = AgentTable(
|
||||
id=uuid4(),
|
||||
name="dev",
|
||||
slug=f"be-dev-{uuid4().hex[:6]}",
|
||||
role=AgentRole.DEVELOPER,
|
||||
team=Team.BACKEND,
|
||||
status=AgentStatus.ACTIVE,
|
||||
model_config={},
|
||||
system_prompt="x",
|
||||
capabilities=[],
|
||||
permissions={},
|
||||
metrics={},
|
||||
)
|
||||
db_session.add(dev)
|
||||
await db_session.flush()
|
||||
project = ProjectTable(
|
||||
id=uuid4(),
|
||||
name="P",
|
||||
slug=f"p-{uuid4().hex[:6]}",
|
||||
git_url="https://example.com/r.git",
|
||||
assigned_cell=Team.BACKEND,
|
||||
created_by=dev.id,
|
||||
)
|
||||
db_session.add(project)
|
||||
await db_session.flush()
|
||||
yield {"db": db_session, "ceo": ceo, "dev": dev, "project_id": project.id}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_notify_ceo_of_completion_creates_alert(env: dict) -> None:
|
||||
db = env["db"]
|
||||
base = datetime.now(UTC) - timedelta(hours=2)
|
||||
task = TaskTable(
|
||||
id=uuid4(),
|
||||
title="Ship it",
|
||||
description="d",
|
||||
acceptance_criteria=["ac"],
|
||||
task_type=TaskType.CODE,
|
||||
nature=TaskNature.TECHNICAL,
|
||||
status=TaskStatus.COMPLETED,
|
||||
team=Team.BACKEND,
|
||||
project_id=env["project_id"],
|
||||
created_by=env["dev"].id,
|
||||
assigned_to=env["dev"].id,
|
||||
estimated_complexity=Complexity.MEDIUM,
|
||||
started_at=base,
|
||||
completed_at=base + timedelta(seconds=600),
|
||||
)
|
||||
db.add(task)
|
||||
await db.flush()
|
||||
db.add(
|
||||
AgentSpawnSessionTable(
|
||||
id=uuid4(),
|
||||
agent_slug=env["dev"].slug,
|
||||
team="backend",
|
||||
role="developer",
|
||||
model="claude",
|
||||
task_id=str(task.id),
|
||||
started_at=base,
|
||||
ended_at=base + timedelta(seconds=300),
|
||||
turns=7,
|
||||
tool_calls=12,
|
||||
tokens_input=10,
|
||||
tokens_output=5,
|
||||
estimated_cost_usd=0.5,
|
||||
)
|
||||
)
|
||||
await db.flush()
|
||||
|
||||
delivery = get_notification_delivery_service(db)
|
||||
await delivery.notify_ceo_of_completion(task=task, task_id=cast("UUID", task.id))
|
||||
|
||||
rows = (
|
||||
(
|
||||
await db.execute(
|
||||
select(NotificationTable).where(
|
||||
NotificationTable.related_task_id == task.id
|
||||
)
|
||||
)
|
||||
)
|
||||
.scalars()
|
||||
.all()
|
||||
)
|
||||
assert len(rows) == 1
|
||||
note = rows[0]
|
||||
assert note.type == NotificationType.ALERT
|
||||
assert env["ceo"].id in note.to_agents
|
||||
assert "Active effort" in note.body
|
||||
assert "Ship it" in note.subject
|
||||
@@ -14,10 +14,18 @@ from httpx import ASGITransport, AsyncClient
|
||||
from roboco.api.deps import get_agent_context, get_db
|
||||
from roboco.api.routes.dashboard import get_main_pm_kanban
|
||||
from roboco.api.routes.dashboard import router as dashboard_router
|
||||
from roboco.db.tables import AgentTable
|
||||
from roboco.db.tables import AgentTable, ProjectTable, TaskTable
|
||||
from roboco.models import AgentRole, AgentStatus
|
||||
from roboco.models.base import (
|
||||
Complexity,
|
||||
TaskNature,
|
||||
TaskStatus,
|
||||
TaskType,
|
||||
Team,
|
||||
)
|
||||
from roboco.models.permissions import AgentContext
|
||||
from roboco.services.dashboard import reset_storage
|
||||
from sqlalchemy import select
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from collections.abc import AsyncGenerator, AsyncIterator
|
||||
@@ -408,3 +416,110 @@ async def test_get_main_pm_kanban_function_directly(
|
||||
"""
|
||||
result = await get_main_pm_kanban(db_session)
|
||||
assert isinstance(result, dict)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_ceo_scorecard_endpoint(dashboard_client: AsyncClient) -> None:
|
||||
resp = await dashboard_client.get(
|
||||
"/api/dashboard/metrics/member/ceo?days=30", headers=_HDR
|
||||
)
|
||||
assert resp.status_code == HTTPStatus.OK
|
||||
body = resp.json()
|
||||
assert body["member_kind"] == "ceo"
|
||||
assert set(body) >= {
|
||||
"approval_p50_seconds",
|
||||
"approval_count",
|
||||
"unblock_p50_seconds",
|
||||
"unblock_count",
|
||||
"godmode_actions",
|
||||
}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_member_scorecard_404_when_absent(dashboard_client: AsyncClient) -> None:
|
||||
resp = await dashboard_client.get(
|
||||
f"/api/dashboard/metrics/member/{uuid4()}", headers=_HDR
|
||||
)
|
||||
assert resp.status_code == HTTPStatus.NOT_FOUND
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_ceo_route_wins_over_member_uuid_route(
|
||||
dashboard_client: AsyncClient,
|
||||
) -> None:
|
||||
# The literal "ceo" must resolve to the CEO route, not the {agent_id} route.
|
||||
resp = await dashboard_client.get("/api/dashboard/metrics/member/ceo", headers=_HDR)
|
||||
assert resp.status_code == HTTPStatus.OK
|
||||
assert resp.json()["member_kind"] == "ceo"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_org_scorecard_endpoint(dashboard_client: AsyncClient) -> None:
|
||||
resp = await dashboard_client.get("/api/dashboard/metrics/org", headers=_HDR)
|
||||
assert resp.status_code == HTTPStatus.OK
|
||||
body = resp.json()
|
||||
assert body["scope"] == "org"
|
||||
assert set(body) >= {"member_count", "tasks_completed", "first_pass_yield"}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_task_metrics_404_for_missing_task(
|
||||
dashboard_client: AsyncClient,
|
||||
) -> None:
|
||||
resp = await dashboard_client.get(
|
||||
f"/api/dashboard/metrics/task/{uuid4()}", headers=_HDR
|
||||
)
|
||||
assert resp.status_code == HTTPStatus.NOT_FOUND
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_task_metrics_returns_shape_for_existing_task(
|
||||
db_session: AsyncSession, dashboard_client: AsyncClient
|
||||
) -> None:
|
||||
creator = (await db_session.execute(select(AgentTable).limit(1))).scalar_one()
|
||||
project = ProjectTable(
|
||||
id=uuid4(),
|
||||
name="P",
|
||||
slug=f"p-{uuid4().hex[:6]}",
|
||||
git_url="https://example.com/r.git",
|
||||
assigned_cell=Team.BACKEND,
|
||||
created_by=creator.id,
|
||||
)
|
||||
db_session.add(project)
|
||||
await db_session.flush()
|
||||
task = TaskTable(
|
||||
id=uuid4(),
|
||||
title="t",
|
||||
description="d",
|
||||
acceptance_criteria=["ac"],
|
||||
task_type=TaskType.CODE,
|
||||
nature=TaskNature.TECHNICAL,
|
||||
status=TaskStatus.IN_PROGRESS,
|
||||
team=Team.BACKEND,
|
||||
project_id=project.id,
|
||||
created_by=creator.id,
|
||||
estimated_complexity=Complexity.MEDIUM,
|
||||
)
|
||||
db_session.add(task)
|
||||
await db_session.flush()
|
||||
|
||||
resp = await dashboard_client.get(
|
||||
f"/api/dashboard/metrics/task/{task.id}", headers=_HDR
|
||||
)
|
||||
assert resp.status_code == HTTPStatus.OK
|
||||
body = resp.json()
|
||||
assert body["task_id"] == str(task.id)
|
||||
assert set(body) >= {
|
||||
"active_runtime_seconds",
|
||||
"wall_clock_seconds",
|
||||
"turns",
|
||||
"tool_calls",
|
||||
"tokens",
|
||||
"cost_usd",
|
||||
"revision_count",
|
||||
"qa_fails",
|
||||
"pr_fails",
|
||||
"stints",
|
||||
"stages",
|
||||
}
|
||||
assert isinstance(body["stages"], list)
|
||||
|
||||
@@ -0,0 +1,299 @@
|
||||
"""_sweep_member_performance — the granular per-member rollup, against real PG.
|
||||
|
||||
Seeds a day of spawn sessions + completed tasks + the new audit events
|
||||
(escalated / unblocked_dependents / agent.idle / qa pass+fail / CEO decisions),
|
||||
runs the sweep, and asserts the agent + CEO rows — plus idempotency (a second
|
||||
sweep overwrites, never doubles).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from contextlib import asynccontextmanager
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from typing import TYPE_CHECKING, Any
|
||||
from unittest.mock import patch
|
||||
from uuid import uuid4
|
||||
|
||||
import pytest
|
||||
import pytest_asyncio
|
||||
from roboco.db.tables import (
|
||||
AgentSpawnSessionTable,
|
||||
AgentTable,
|
||||
AuditLogTable,
|
||||
MemberPerformanceDailyTable,
|
||||
ProjectTable,
|
||||
TaskTable,
|
||||
)
|
||||
from roboco.models.base import (
|
||||
AgentRole,
|
||||
AgentStatus,
|
||||
Complexity,
|
||||
TaskNature,
|
||||
TaskStatus,
|
||||
TaskType,
|
||||
Team,
|
||||
)
|
||||
from roboco.runtime.orchestrator import AgentOrchestrator
|
||||
from sqlalchemy import select
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from collections.abc import AsyncIterator
|
||||
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
# Yesterday at noon UTC: safely inside the 7-day window and far from midnight,
|
||||
# so started_at / completed_at (+1h) and all audit events share ONE date
|
||||
# (the sweep legitimately splits cross-midnight work across days — not under test
|
||||
# here). This keeps the seeded member's rollup on a single row.
|
||||
_BASE = (datetime.now(UTC) - timedelta(days=1)).replace(
|
||||
hour=12, minute=0, second=0, microsecond=0
|
||||
)
|
||||
|
||||
_ACTIVE_SECONDS = 600
|
||||
_APPROVAL_DWELL_SECONDS = 300
|
||||
_COMPLETED_TASKS = 2
|
||||
_BLOCKED_OTHERS = 2
|
||||
_QA_TOTAL = 2
|
||||
|
||||
|
||||
class _NoCommitSession:
|
||||
"""Delegates to the real test session but turns commit() into flush() so the
|
||||
per-test rollback isolation holds while the sweep still 'commits'."""
|
||||
|
||||
def __init__(self, real: Any) -> None:
|
||||
self._real = real
|
||||
|
||||
def __getattr__(self, name: str) -> Any:
|
||||
return getattr(self._real, name)
|
||||
|
||||
async def commit(self) -> None:
|
||||
await self._real.flush()
|
||||
|
||||
|
||||
def _agent(role: AgentRole, slug: str) -> AgentTable:
|
||||
return AgentTable(
|
||||
id=uuid4(),
|
||||
name=slug,
|
||||
slug=slug,
|
||||
role=role,
|
||||
team=Team.BACKEND,
|
||||
status=AgentStatus.ACTIVE,
|
||||
model_config={},
|
||||
system_prompt="x",
|
||||
capabilities=[],
|
||||
permissions={},
|
||||
metrics={},
|
||||
)
|
||||
|
||||
|
||||
def _audit(
|
||||
target_id: Any,
|
||||
event_type: str,
|
||||
ts: datetime,
|
||||
*,
|
||||
agent_id: Any = None,
|
||||
details: dict[str, Any] | None = None,
|
||||
) -> AuditLogTable:
|
||||
return AuditLogTable(
|
||||
id=uuid4(),
|
||||
event_type=event_type,
|
||||
agent_id=agent_id,
|
||||
target_type="task",
|
||||
target_id=target_id,
|
||||
severity="info",
|
||||
details=details or {},
|
||||
timestamp=ts,
|
||||
)
|
||||
|
||||
|
||||
@pytest_asyncio.fixture
|
||||
async def seeded(db_session: AsyncSession) -> AsyncIterator[dict]:
|
||||
dev = _agent(AgentRole.DEVELOPER, f"be-dev-{uuid4().hex[:6]}")
|
||||
qa = _agent(AgentRole.QA, f"be-qa-{uuid4().hex[:6]}")
|
||||
db_session.add_all([dev, qa])
|
||||
await db_session.flush()
|
||||
project = ProjectTable(
|
||||
id=uuid4(),
|
||||
name="P",
|
||||
slug=f"p-{uuid4().hex[:6]}",
|
||||
git_url="https://example.com/r.git",
|
||||
assigned_cell=Team.BACKEND,
|
||||
created_by=dev.id,
|
||||
)
|
||||
db_session.add(project)
|
||||
await db_session.flush()
|
||||
|
||||
task = TaskTable(
|
||||
id=uuid4(),
|
||||
title="t",
|
||||
description="d",
|
||||
acceptance_criteria=["ac"],
|
||||
task_type=TaskType.CODE,
|
||||
nature=TaskNature.TECHNICAL,
|
||||
status=TaskStatus.COMPLETED,
|
||||
team=Team.BACKEND,
|
||||
project_id=project.id,
|
||||
created_by=dev.id,
|
||||
assigned_to=dev.id,
|
||||
revision_count=1,
|
||||
estimated_complexity=Complexity.MEDIUM,
|
||||
started_at=_BASE,
|
||||
completed_at=_BASE + timedelta(hours=1),
|
||||
)
|
||||
blocker = TaskTable(
|
||||
id=uuid4(),
|
||||
title="b",
|
||||
description="d",
|
||||
acceptance_criteria=["ac"],
|
||||
task_type=TaskType.CODE,
|
||||
nature=TaskNature.TECHNICAL,
|
||||
status=TaskStatus.COMPLETED,
|
||||
team=Team.BACKEND,
|
||||
project_id=project.id,
|
||||
created_by=dev.id,
|
||||
assigned_to=dev.id,
|
||||
estimated_complexity=Complexity.MEDIUM,
|
||||
started_at=_BASE,
|
||||
completed_at=_BASE + timedelta(hours=1),
|
||||
)
|
||||
db_session.add_all([task, blocker])
|
||||
await db_session.flush()
|
||||
|
||||
db_session.add(
|
||||
AgentSpawnSessionTable(
|
||||
id=uuid4(),
|
||||
agent_slug=dev.slug,
|
||||
team="backend",
|
||||
role="developer",
|
||||
model="claude",
|
||||
task_id=str(task.id),
|
||||
started_at=_BASE,
|
||||
ended_at=_BASE + timedelta(seconds=600),
|
||||
turns=5,
|
||||
tool_calls=10,
|
||||
tokens_input=100,
|
||||
tokens_output=80,
|
||||
estimated_cost_usd=1.5,
|
||||
)
|
||||
)
|
||||
db_session.add_all(
|
||||
[
|
||||
_audit(
|
||||
task.id,
|
||||
"task.qa_fail",
|
||||
_BASE,
|
||||
agent_id=qa.id,
|
||||
details={"agent_role": "qa"},
|
||||
),
|
||||
_audit(
|
||||
task.id,
|
||||
"task.awaiting_documentation",
|
||||
_BASE + timedelta(minutes=1),
|
||||
agent_id=qa.id,
|
||||
details={"agent_role": "qa"},
|
||||
),
|
||||
_audit(
|
||||
task.id,
|
||||
"task.escalated",
|
||||
_BASE,
|
||||
details={"escalator_slug": dev.slug, "target_slug": qa.slug},
|
||||
),
|
||||
_audit(
|
||||
blocker.id, "task.unblocked_dependents", _BASE, details={"count": 2}
|
||||
),
|
||||
_audit(
|
||||
dev.id,
|
||||
"agent.idle",
|
||||
_BASE + timedelta(minutes=5),
|
||||
agent_id=dev.id,
|
||||
details={"agent_slug": dev.slug},
|
||||
),
|
||||
# CEO decision: approval + god-mode (to_status is what the pairing reads).
|
||||
_audit(
|
||||
task.id,
|
||||
"task.awaiting_ceo_approval",
|
||||
_BASE,
|
||||
details={"to_status": "awaiting_ceo_approval", "agent_role": "main_pm"},
|
||||
),
|
||||
_audit(
|
||||
task.id,
|
||||
"task.completed",
|
||||
_BASE + timedelta(seconds=300),
|
||||
details={"to_status": "completed", "agent_role": "ceo"},
|
||||
),
|
||||
]
|
||||
)
|
||||
await db_session.flush()
|
||||
yield {"db": db_session, "dev": dev.slug, "qa": qa.slug}
|
||||
|
||||
|
||||
async def _rows(db: AsyncSession, slug: str) -> list[Any]:
|
||||
return list(
|
||||
(
|
||||
await db.execute(
|
||||
select(MemberPerformanceDailyTable).where(
|
||||
MemberPerformanceDailyTable.agent_slug == slug
|
||||
)
|
||||
)
|
||||
)
|
||||
.scalars()
|
||||
.all()
|
||||
)
|
||||
|
||||
|
||||
async def _run_sweep(db: AsyncSession) -> None:
|
||||
orch = AgentOrchestrator.__new__(AgentOrchestrator)
|
||||
|
||||
@asynccontextmanager
|
||||
async def _cm() -> Any:
|
||||
yield _NoCommitSession(db)
|
||||
|
||||
with patch("roboco.db.base.get_session_factory", return_value=_cm):
|
||||
await orch._sweep_member_performance()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_sweep_populates_agent_and_ceo_rows(seeded: dict) -> None:
|
||||
db = seeded["db"]
|
||||
await _run_sweep(db)
|
||||
|
||||
dev_rows = await _rows(db, seeded["dev"])
|
||||
assert len(dev_rows) == 1
|
||||
dev = dev_rows[0]
|
||||
assert dev.member_kind == "agent"
|
||||
assert dev.active_runtime_seconds == _ACTIVE_SECONDS
|
||||
assert (dev.turns, dev.tool_calls, dev.tokens) == (5, 10, 180)
|
||||
assert dev.cost_usd == pytest.approx(1.5)
|
||||
assert dev.tasks_completed == _COMPLETED_TASKS # task + blocker
|
||||
assert dev.tasks_first_pass == 1 # blocker had 0 revisions
|
||||
assert dev.revisions_received == 1
|
||||
assert dev.escalations == 1
|
||||
assert dev.blocked_others == _BLOCKED_OTHERS
|
||||
assert dev.idle_seconds > 0
|
||||
|
||||
qa_rows = await _rows(db, seeded["qa"])
|
||||
assert len(qa_rows) == 1
|
||||
qa = qa_rows[0]
|
||||
assert qa.revisions_caused == 1 # the qa_fail
|
||||
assert qa.qa_reviews_passed == 1
|
||||
assert qa.qa_reviews_total == _QA_TOTAL # 1 pass + 1 fail
|
||||
|
||||
ceo_rows = await _rows(db, "")
|
||||
assert len(ceo_rows) == 1
|
||||
ceo = ceo_rows[0]
|
||||
assert ceo.member_kind == "ceo"
|
||||
assert ceo.godmode_actions == 1
|
||||
assert ceo.ceo_approval_dwell_seconds == _APPROVAL_DWELL_SECONDS
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_sweep_is_idempotent(seeded: dict) -> None:
|
||||
db = seeded["db"]
|
||||
await _run_sweep(db)
|
||||
await _run_sweep(db) # second pass must overwrite, not double
|
||||
|
||||
dev = (await _rows(db, seeded["dev"]))[0]
|
||||
assert dev.active_runtime_seconds == _ACTIVE_SECONDS # not 1200
|
||||
assert dev.tasks_completed == _COMPLETED_TASKS # not 4
|
||||
assert dev.escalations == 1
|
||||
assert len(await _rows(db, seeded["dev"])) == 1
|
||||
@@ -0,0 +1,118 @@
|
||||
"""get_ceo_scorecard — the human CEO as a measured member (audit-log only).
|
||||
|
||||
Seeds CEO-attributed audit transitions and asserts approval dwell (incl. the
|
||||
coordination-root reject that lands in `pending`), unblock dwell, and the
|
||||
god-mode action count. The CEO never runs an LLM, so this reads only audit_log
|
||||
(agent_role='ceo').
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from typing import TYPE_CHECKING, Any
|
||||
from uuid import uuid4
|
||||
|
||||
import pytest
|
||||
import pytest_asyncio
|
||||
from roboco.db.tables import AuditLogTable
|
||||
from roboco.services.metrics import MetricsService
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from collections.abc import AsyncIterator
|
||||
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
_NOW = datetime.now(UTC)
|
||||
|
||||
|
||||
def _audit(
|
||||
task_id: Any,
|
||||
status: str,
|
||||
ts: datetime,
|
||||
*,
|
||||
agent_role: str | None = None,
|
||||
) -> AuditLogTable:
|
||||
return AuditLogTable(
|
||||
id=uuid4(),
|
||||
event_type=f"task.{status}",
|
||||
agent_id=None,
|
||||
target_type="task",
|
||||
target_id=task_id,
|
||||
severity="info",
|
||||
details={"to_status": status, "from_status": "prev", "agent_role": agent_role},
|
||||
timestamp=ts,
|
||||
)
|
||||
|
||||
|
||||
@pytest_asyncio.fixture
|
||||
async def svc(db_session: AsyncSession) -> AsyncIterator[MetricsService]:
|
||||
yield MetricsService(db_session)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_empty_window_returns_zeros(svc: MetricsService) -> None:
|
||||
card = await svc.get_ceo_scorecard(days=30)
|
||||
assert card.approval_count == 0
|
||||
assert card.unblock_count == 0
|
||||
assert card.godmode_actions == 0
|
||||
assert card.approval_p50_seconds == 0.0
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_approval_unblock_dwell_and_godmode(
|
||||
svc: MetricsService, db_session: AsyncSession
|
||||
) -> None:
|
||||
base = _NOW - timedelta(hours=2)
|
||||
t_approve = uuid4()
|
||||
t_reject = uuid4()
|
||||
t_block = uuid4()
|
||||
db_session.add_all(
|
||||
[
|
||||
# Approval: awaiting -> completed(ceo) after 300s.
|
||||
_audit(t_approve, "awaiting_ceo_approval", base),
|
||||
_audit(
|
||||
t_approve, "completed", base + timedelta(seconds=300), agent_role="ceo"
|
||||
),
|
||||
# Coordination-root reject: awaiting -> pending(ceo) after 120s.
|
||||
_audit(t_reject, "awaiting_ceo_approval", base),
|
||||
_audit(
|
||||
t_reject, "pending", base + timedelta(seconds=120), agent_role="ceo"
|
||||
),
|
||||
# Unblock: blocked -> in_progress(ceo) after 600s.
|
||||
_audit(t_block, "blocked", base),
|
||||
_audit(
|
||||
t_block, "in_progress", base + timedelta(seconds=600), agent_role="ceo"
|
||||
),
|
||||
]
|
||||
)
|
||||
await db_session.flush()
|
||||
|
||||
card = await svc.get_ceo_scorecard(days=30)
|
||||
# Two approval decisions (completed + coordination pending), median of 300/120.
|
||||
expected_approvals = 2
|
||||
assert card.approval_count == expected_approvals
|
||||
assert card.approval_p50_seconds == pytest.approx(210.0)
|
||||
# One unblock, 600s.
|
||||
assert card.unblock_count == 1
|
||||
assert card.unblock_p50_seconds == pytest.approx(600.0)
|
||||
# God-mode = every ceo-attributed transition: completed + pending + in_progress.
|
||||
expected_godmode = 3
|
||||
assert card.godmode_actions == expected_godmode
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_non_ceo_transitions_are_not_counted(
|
||||
svc: MetricsService, db_session: AsyncSession
|
||||
) -> None:
|
||||
tid = uuid4()
|
||||
db_session.add_all(
|
||||
[
|
||||
_audit(tid, "awaiting_ceo_approval", _NOW - timedelta(minutes=10)),
|
||||
# A QA fail (not the CEO) must not count as an approval or god-mode.
|
||||
_audit(tid, "needs_revision", _NOW - timedelta(minutes=5), agent_role="qa"),
|
||||
]
|
||||
)
|
||||
await db_session.flush()
|
||||
card = await svc.get_ceo_scorecard(days=30)
|
||||
assert card.approval_count == 0
|
||||
assert card.godmode_actions == 0
|
||||
@@ -0,0 +1,164 @@
|
||||
"""New audit instrumentation feeding the extra per-member metrics (phase 4).
|
||||
|
||||
- apply_escalation -> task.escalated (escalations metric)
|
||||
- _unblock_dependents -> task.unblocked_dependents (blocked-others metric)
|
||||
- mark_agent_idle -> agent.idle (idle/utilization metric)
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import TYPE_CHECKING, Any
|
||||
from uuid import uuid4
|
||||
|
||||
import pytest
|
||||
import pytest_asyncio
|
||||
from roboco.db.tables import AgentTable, AuditLogTable, ProjectTable, TaskTable
|
||||
from roboco.models.base import (
|
||||
AgentRole,
|
||||
AgentStatus,
|
||||
Complexity,
|
||||
TaskNature,
|
||||
TaskStatus,
|
||||
TaskType,
|
||||
Team,
|
||||
)
|
||||
from roboco.services.task import TaskService
|
||||
from sqlalchemy import select
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from collections.abc import AsyncIterator
|
||||
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
|
||||
def _agent(role: AgentRole, slug: str) -> AgentTable:
|
||||
return AgentTable(
|
||||
id=uuid4(),
|
||||
name=slug,
|
||||
slug=slug,
|
||||
role=role,
|
||||
team=Team.BACKEND,
|
||||
status=AgentStatus.ACTIVE,
|
||||
model_config={},
|
||||
system_prompt="x",
|
||||
capabilities=[],
|
||||
permissions={},
|
||||
metrics={},
|
||||
)
|
||||
|
||||
|
||||
def _task(project_id: Any, created_by: Any, **over: Any) -> TaskTable:
|
||||
base: dict[str, Any] = {
|
||||
"id": uuid4(),
|
||||
"title": "t",
|
||||
"description": "d",
|
||||
"acceptance_criteria": ["ac"],
|
||||
"task_type": TaskType.CODE,
|
||||
"nature": TaskNature.TECHNICAL,
|
||||
"status": TaskStatus.IN_PROGRESS,
|
||||
"team": Team.BACKEND,
|
||||
"project_id": project_id,
|
||||
"created_by": created_by,
|
||||
"estimated_complexity": Complexity.MEDIUM,
|
||||
}
|
||||
base.update(over)
|
||||
return TaskTable(**base)
|
||||
|
||||
|
||||
async def _audit_of(db: AsyncSession, event_type: str, target_id: Any) -> list[Any]:
|
||||
rows = (
|
||||
(
|
||||
await db.execute(
|
||||
select(AuditLogTable).where(
|
||||
AuditLogTable.event_type == event_type,
|
||||
AuditLogTable.target_id == target_id,
|
||||
)
|
||||
)
|
||||
)
|
||||
.scalars()
|
||||
.all()
|
||||
)
|
||||
return list(rows)
|
||||
|
||||
|
||||
@pytest_asyncio.fixture
|
||||
async def env(db_session: AsyncSession) -> AsyncIterator[dict]:
|
||||
dev = _agent(AgentRole.DEVELOPER, f"be-dev-{uuid4().hex[:6]}")
|
||||
pm = _agent(AgentRole.CELL_PM, f"be-pm-{uuid4().hex[:6]}")
|
||||
db_session.add_all([dev, pm])
|
||||
await db_session.flush()
|
||||
project = ProjectTable(
|
||||
id=uuid4(),
|
||||
name="P",
|
||||
slug=f"p-{uuid4().hex[:6]}",
|
||||
git_url="https://example.com/r.git",
|
||||
assigned_cell=Team.BACKEND,
|
||||
created_by=dev.id,
|
||||
)
|
||||
db_session.add(project)
|
||||
await db_session.flush()
|
||||
yield {
|
||||
"svc": TaskService(db_session),
|
||||
"db": db_session,
|
||||
"project_id": project.id,
|
||||
"dev": dev,
|
||||
"pm": pm,
|
||||
}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_apply_escalation_emits_task_escalated(env: dict) -> None:
|
||||
db = env["db"]
|
||||
task = _task(
|
||||
env["project_id"],
|
||||
env["dev"].id,
|
||||
assigned_to=env["dev"].id,
|
||||
claimed_by=env["dev"].id,
|
||||
)
|
||||
db.add(task)
|
||||
await db.flush()
|
||||
ok = await env["svc"].apply_escalation(
|
||||
task=task,
|
||||
target_agent_id=env["pm"].id,
|
||||
escalator_slug=env["dev"].slug,
|
||||
target_slug=env["pm"].slug,
|
||||
reason="need help with the seam",
|
||||
)
|
||||
assert ok is True
|
||||
rows = await _audit_of(db, "task.escalated", task.id)
|
||||
assert len(rows) == 1
|
||||
assert rows[0].details["escalator_slug"] == env["dev"].slug
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_unblock_dependents_emits_count(env: dict) -> None:
|
||||
db = env["db"]
|
||||
blocker = _task(env["project_id"], env["dev"].id, status=TaskStatus.COMPLETED)
|
||||
db.add(blocker)
|
||||
await db.flush()
|
||||
dependent = _task(
|
||||
env["project_id"],
|
||||
env["dev"].id,
|
||||
status=TaskStatus.BLOCKED,
|
||||
dependency_ids=[blocker.id],
|
||||
assigned_to=env["dev"].id,
|
||||
claimed_by=env["dev"].id,
|
||||
)
|
||||
db.add(dependent)
|
||||
await db.flush()
|
||||
await env["svc"]._unblock_dependents(blocker.id)
|
||||
rows = await _audit_of(db, "task.unblocked_dependents", blocker.id)
|
||||
assert len(rows) == 1
|
||||
assert rows[0].details["count"] == 1
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_mark_agent_idle_emits_agent_idle(env: dict) -> None:
|
||||
db = env["db"]
|
||||
await env["svc"].mark_agent_idle(env["dev"].id)
|
||||
rows = await _audit_of(db, "agent.idle", env["dev"].id)
|
||||
assert len(rows) == 1
|
||||
assert rows[0].details["agent_slug"] == env["dev"].slug
|
||||
refreshed = await db.get(AgentTable, env["dev"].id)
|
||||
assert refreshed is not None
|
||||
assert refreshed.status == AgentStatus.IDLE
|
||||
@@ -0,0 +1,316 @@
|
||||
"""Member / org rollup scorecards + the live in-flight overlay (real PG).
|
||||
|
||||
Seeds member_performance_daily rows (the rollup source) and asserts the derived
|
||||
rates (FPY, effort-throughput, turns/task, qa pass-rate, utilization), the live
|
||||
in-flight overlay (enriches effort but not completion counts — disjoint by
|
||||
status), and the division guards.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from typing import TYPE_CHECKING, Any, cast
|
||||
from uuid import uuid4
|
||||
|
||||
import pytest
|
||||
import pytest_asyncio
|
||||
from roboco.db.tables import (
|
||||
AgentSpawnSessionTable,
|
||||
AgentTable,
|
||||
MemberPerformanceDailyTable,
|
||||
ProjectTable,
|
||||
TaskTable,
|
||||
)
|
||||
from roboco.models.base import (
|
||||
AgentRole,
|
||||
AgentStatus,
|
||||
Complexity,
|
||||
TaskNature,
|
||||
TaskStatus,
|
||||
TaskType,
|
||||
Team,
|
||||
)
|
||||
from roboco.services.metrics import MetricsService
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from collections.abc import AsyncIterator
|
||||
from uuid import UUID
|
||||
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
_TODAY = datetime.now(UTC).date()
|
||||
_TOTAL_COMPLETED = 3
|
||||
_ROLLUP_COMPLETED = 2
|
||||
_OVERLAY_TURNS = 3
|
||||
_ROLLUP_ONLY_TURNS = 5 # closed session already in the rollup, not re-added
|
||||
_ORG_MEMBERS = 2
|
||||
_ORG_COMPLETED = 3
|
||||
|
||||
|
||||
def _daily(slug: str, **over: Any) -> MemberPerformanceDailyTable:
|
||||
base: dict[str, Any] = {
|
||||
"id": uuid4(),
|
||||
"date": _TODAY,
|
||||
"member_kind": "agent",
|
||||
"agent_slug": slug,
|
||||
"team": Team.BACKEND.value,
|
||||
"role": "developer",
|
||||
}
|
||||
base.update(over)
|
||||
return MemberPerformanceDailyTable(**base)
|
||||
|
||||
|
||||
def _agent(role: AgentRole, slug: str) -> AgentTable:
|
||||
return AgentTable(
|
||||
id=uuid4(),
|
||||
name=slug,
|
||||
slug=slug,
|
||||
role=role,
|
||||
team=Team.BACKEND,
|
||||
status=AgentStatus.ACTIVE,
|
||||
model_config={},
|
||||
system_prompt="x",
|
||||
capabilities=[],
|
||||
permissions={},
|
||||
metrics={},
|
||||
)
|
||||
|
||||
|
||||
@pytest_asyncio.fixture
|
||||
async def svc(db_session: AsyncSession) -> AsyncIterator[MetricsService]:
|
||||
yield MetricsService(db_session)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_member_scorecard_rollup_and_derived(
|
||||
svc: MetricsService, db_session: AsyncSession
|
||||
) -> None:
|
||||
dev = _agent(AgentRole.DEVELOPER, f"be-dev-{uuid4().hex[:6]}")
|
||||
db_session.add(dev)
|
||||
await db_session.flush()
|
||||
db_session.add_all(
|
||||
[
|
||||
_daily(
|
||||
dev.slug,
|
||||
tasks_completed=2,
|
||||
tasks_first_pass=1,
|
||||
active_runtime_seconds=1800,
|
||||
turns=6,
|
||||
tool_calls=12,
|
||||
tokens=100,
|
||||
cost_usd=1.0,
|
||||
qa_reviews_total=3,
|
||||
qa_reviews_passed=2,
|
||||
escalations=1,
|
||||
blocked_others=1,
|
||||
idle_seconds=600,
|
||||
revisions_caused=1,
|
||||
revisions_received=1,
|
||||
),
|
||||
_daily(
|
||||
dev.slug,
|
||||
date=_TODAY - timedelta(days=1),
|
||||
tasks_completed=1,
|
||||
tasks_first_pass=1,
|
||||
active_runtime_seconds=1800,
|
||||
turns=4,
|
||||
tool_calls=8,
|
||||
tokens=50,
|
||||
cost_usd=0.5,
|
||||
qa_reviews_total=2,
|
||||
qa_reviews_passed=2,
|
||||
idle_seconds=1200,
|
||||
),
|
||||
]
|
||||
)
|
||||
await db_session.flush()
|
||||
|
||||
card = await svc.get_member_scorecard(cast("UUID", dev.id), days=30)
|
||||
assert card is not None
|
||||
assert card.tasks_completed == _TOTAL_COMPLETED
|
||||
assert card.first_pass_yield == pytest.approx(2 / 3, abs=1e-4) # 2 of 3
|
||||
# 3 tasks over 3600s = 1h -> 3.0/hr.
|
||||
assert card.effort_throughput_per_hour == pytest.approx(3.0)
|
||||
assert (card.turns, card.tool_calls) == (10, 20)
|
||||
assert card.turns_per_task == pytest.approx(10 / 3, abs=1e-4)
|
||||
assert card.qa_pass_rate == pytest.approx(4 / 5, abs=1e-4) # 4 of 5
|
||||
assert card.escalations == 1
|
||||
assert card.blocked_others == 1
|
||||
# util = 3600 active / (3600 + 1800 idle) = 0.6667.
|
||||
assert card.utilization == pytest.approx(3600 / 5400, abs=1e-4)
|
||||
assert card.includes_live_inflight is False
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_live_overlay_enriches_effort_not_completion(
|
||||
svc: MetricsService, db_session: AsyncSession
|
||||
) -> None:
|
||||
dev = _agent(AgentRole.DEVELOPER, f"be-dev-{uuid4().hex[:6]}")
|
||||
db_session.add(dev)
|
||||
await db_session.flush()
|
||||
project = ProjectTable(
|
||||
id=uuid4(),
|
||||
name="P",
|
||||
slug=f"p-{uuid4().hex[:6]}",
|
||||
git_url="https://example.com/r.git",
|
||||
assigned_cell=Team.BACKEND,
|
||||
created_by=dev.id,
|
||||
)
|
||||
db_session.add(project)
|
||||
await db_session.flush()
|
||||
db_session.add(_daily(dev.slug, tasks_completed=2, active_runtime_seconds=100))
|
||||
# A non-terminal (in-flight) task with a spawn stint -> overlay effort.
|
||||
inflight = TaskTable(
|
||||
id=uuid4(),
|
||||
title="t",
|
||||
description="d",
|
||||
acceptance_criteria=["ac"],
|
||||
task_type=TaskType.CODE,
|
||||
nature=TaskNature.TECHNICAL,
|
||||
status=TaskStatus.IN_PROGRESS,
|
||||
team=Team.BACKEND,
|
||||
project_id=project.id,
|
||||
created_by=dev.id,
|
||||
assigned_to=dev.id,
|
||||
estimated_complexity=Complexity.MEDIUM,
|
||||
started_at=datetime.now(UTC) - timedelta(hours=1),
|
||||
)
|
||||
db_session.add(inflight)
|
||||
await db_session.flush()
|
||||
now = datetime.now(UTC)
|
||||
# OPEN (still-running) session — the live delta the rollup cannot hold yet.
|
||||
db_session.add(
|
||||
AgentSpawnSessionTable(
|
||||
id=uuid4(),
|
||||
agent_slug=dev.slug,
|
||||
team="backend",
|
||||
role="developer",
|
||||
model="claude",
|
||||
task_id=str(inflight.id),
|
||||
started_at=now - timedelta(seconds=200),
|
||||
ended_at=None,
|
||||
turns=3,
|
||||
tool_calls=4,
|
||||
tokens_input=10,
|
||||
tokens_output=5,
|
||||
estimated_cost_usd=0.2,
|
||||
)
|
||||
)
|
||||
await db_session.flush()
|
||||
|
||||
card = await svc.get_member_scorecard(cast("UUID", dev.id), days=30)
|
||||
assert card is not None
|
||||
assert card.tasks_completed == _ROLLUP_COMPLETED # in-flight NOT completed
|
||||
assert card.includes_live_inflight is True
|
||||
assert card.active_runtime_hours > 100 / 3600 # rollup + overlay effort
|
||||
assert card.turns == _OVERLAY_TURNS # from the overlay (rollup row had 0)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_live_overlay_excludes_closed_session_no_double_count(
|
||||
svc: MetricsService, db_session: AsyncSession
|
||||
) -> None:
|
||||
"""A CLOSED session on a non-terminal task is already in the daily rollup
|
||||
(via _msweep_spawn, which counts ended_at IS NOT NULL). The overlay must NOT
|
||||
re-add it, or the member's effort/turns double-count on the common
|
||||
reap/respawn path."""
|
||||
dev = _agent(AgentRole.DEVELOPER, f"be-dev-{uuid4().hex[:6]}")
|
||||
db_session.add(dev)
|
||||
await db_session.flush()
|
||||
project = ProjectTable(
|
||||
id=uuid4(),
|
||||
name="P",
|
||||
slug=f"p-{uuid4().hex[:6]}",
|
||||
git_url="https://example.com/r.git",
|
||||
assigned_cell=Team.BACKEND,
|
||||
created_by=dev.id,
|
||||
)
|
||||
db_session.add(project)
|
||||
await db_session.flush()
|
||||
# Rollup row already reflects the closed session (turns=5, active=300s).
|
||||
db_session.add(
|
||||
_daily(dev.slug, tasks_completed=0, active_runtime_seconds=300, turns=5)
|
||||
)
|
||||
inflight = TaskTable(
|
||||
id=uuid4(),
|
||||
title="t",
|
||||
description="d",
|
||||
acceptance_criteria=["ac"],
|
||||
task_type=TaskType.CODE,
|
||||
nature=TaskNature.TECHNICAL,
|
||||
status=TaskStatus.IN_PROGRESS,
|
||||
team=Team.BACKEND,
|
||||
project_id=project.id,
|
||||
created_by=dev.id,
|
||||
assigned_to=dev.id,
|
||||
estimated_complexity=Complexity.MEDIUM,
|
||||
started_at=datetime.now(UTC) - timedelta(hours=1),
|
||||
)
|
||||
db_session.add(inflight)
|
||||
await db_session.flush()
|
||||
now = datetime.now(UTC)
|
||||
db_session.add(
|
||||
AgentSpawnSessionTable(
|
||||
id=uuid4(),
|
||||
agent_slug=dev.slug,
|
||||
team="backend",
|
||||
role="developer",
|
||||
model="claude",
|
||||
task_id=str(inflight.id),
|
||||
started_at=now - timedelta(seconds=300),
|
||||
ended_at=now, # CLOSED — already counted by the rollup
|
||||
turns=5,
|
||||
tool_calls=4,
|
||||
tokens_input=10,
|
||||
tokens_output=5,
|
||||
estimated_cost_usd=0.2,
|
||||
)
|
||||
)
|
||||
await db_session.flush()
|
||||
|
||||
card = await svc.get_member_scorecard(cast("UUID", dev.id), days=30)
|
||||
assert card is not None
|
||||
assert card.turns == _ROLLUP_ONLY_TURNS # closed session is NOT re-added
|
||||
assert card.active_runtime_hours == pytest.approx(300 / 3600, abs=1e-4)
|
||||
assert card.includes_live_inflight is False # no OPEN session
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_member_scorecard_404_and_guards(
|
||||
svc: MetricsService, db_session: AsyncSession
|
||||
) -> None:
|
||||
assert await svc.get_member_scorecard(uuid4()) is None
|
||||
# An agent with no rollup rows: division guards -> None, no crash.
|
||||
dev = _agent(AgentRole.DEVELOPER, f"be-dev-{uuid4().hex[:6]}")
|
||||
db_session.add(dev)
|
||||
await db_session.flush()
|
||||
card = await svc.get_member_scorecard(cast("UUID", dev.id), days=30)
|
||||
assert card is not None
|
||||
assert card.tasks_completed == 0
|
||||
assert card.first_pass_yield is None
|
||||
assert card.effort_throughput_per_hour is None
|
||||
assert card.qa_pass_rate is None
|
||||
assert card.utilization is None
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_org_scorecard_aggregates_members(
|
||||
svc: MetricsService, db_session: AsyncSession
|
||||
) -> None:
|
||||
s1, s2 = f"be-dev-{uuid4().hex[:6]}", f"be-dev-{uuid4().hex[:6]}"
|
||||
db_session.add_all(
|
||||
[
|
||||
_daily(
|
||||
s1, tasks_completed=2, tasks_first_pass=2, active_runtime_seconds=3600
|
||||
),
|
||||
_daily(
|
||||
s2, tasks_completed=1, tasks_first_pass=0, active_runtime_seconds=3600
|
||||
),
|
||||
]
|
||||
)
|
||||
await db_session.flush()
|
||||
org = await svc.get_org_scorecard(team=Team.BACKEND, days=30)
|
||||
assert org.scope == "team"
|
||||
assert org.member_count == _ORG_MEMBERS
|
||||
assert org.tasks_completed == _ORG_COMPLETED
|
||||
assert org.first_pass_yield == pytest.approx(2 / 3, abs=1e-4)
|
||||
@@ -0,0 +1,244 @@
|
||||
"""get_task_metrics — granular per-task effort against a real Postgres.
|
||||
|
||||
Seeds a task's audit-log journey + agent spawn stints (with turns/tool_calls/
|
||||
tokens/cost) + named qa/pr fail events, then asserts the composed metrics:
|
||||
summed effort vs wall-clock, turns/tool_calls/tokens/cost, per-stage
|
||||
active-vs-wait, and who-caused-rework.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from typing import TYPE_CHECKING, Any, NamedTuple
|
||||
from uuid import uuid4
|
||||
|
||||
import pytest
|
||||
import pytest_asyncio
|
||||
from roboco.db.tables import (
|
||||
AgentSpawnSessionTable,
|
||||
AgentTable,
|
||||
AuditLogTable,
|
||||
ProjectTable,
|
||||
TaskTable,
|
||||
)
|
||||
from roboco.models.base import (
|
||||
AgentRole,
|
||||
AgentStatus,
|
||||
Complexity,
|
||||
TaskNature,
|
||||
TaskStatus,
|
||||
TaskType,
|
||||
Team,
|
||||
)
|
||||
from roboco.services.metrics import MetricsService
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from collections.abc import AsyncIterator
|
||||
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
_T0 = datetime(2026, 6, 20, 12, 0, 0, tzinfo=UTC)
|
||||
|
||||
|
||||
def _sec(n: int) -> datetime:
|
||||
return _T0 + timedelta(seconds=n)
|
||||
|
||||
|
||||
def _agent(role: AgentRole, slug: str) -> AgentTable:
|
||||
return AgentTable(
|
||||
id=uuid4(),
|
||||
name=slug,
|
||||
slug=slug,
|
||||
role=role,
|
||||
team=Team.BACKEND,
|
||||
status=AgentStatus.ACTIVE,
|
||||
model_config={},
|
||||
system_prompt="x",
|
||||
capabilities=[],
|
||||
permissions={},
|
||||
metrics={},
|
||||
)
|
||||
|
||||
|
||||
def _audit(
|
||||
task_id: Any,
|
||||
status: str,
|
||||
ts: datetime,
|
||||
*,
|
||||
agent_id: Any = None,
|
||||
event_type: str | None = None,
|
||||
) -> AuditLogTable:
|
||||
return AuditLogTable(
|
||||
id=uuid4(),
|
||||
event_type=event_type or f"task.{status}",
|
||||
agent_id=agent_id,
|
||||
target_type="task",
|
||||
target_id=task_id,
|
||||
severity="info",
|
||||
details={"to_status": status, "from_status": "prev", "team": "backend"},
|
||||
timestamp=ts,
|
||||
)
|
||||
|
||||
|
||||
class _Usage(NamedTuple):
|
||||
turns: int
|
||||
tool_calls: int
|
||||
tokens_in: int
|
||||
tokens_out: int
|
||||
cost: float
|
||||
|
||||
|
||||
def _spawn(
|
||||
task_id: str,
|
||||
started: datetime,
|
||||
ended: datetime | None,
|
||||
usage: _Usage,
|
||||
) -> AgentSpawnSessionTable:
|
||||
return AgentSpawnSessionTable(
|
||||
id=uuid4(),
|
||||
agent_slug="be-dev-1",
|
||||
team="backend",
|
||||
role="developer",
|
||||
model="claude",
|
||||
task_id=task_id,
|
||||
started_at=started,
|
||||
ended_at=ended,
|
||||
turns=usage.turns,
|
||||
tool_calls=usage.tool_calls,
|
||||
tokens_input=usage.tokens_in,
|
||||
tokens_output=usage.tokens_out,
|
||||
estimated_cost_usd=usage.cost,
|
||||
)
|
||||
|
||||
|
||||
@pytest_asyncio.fixture
|
||||
async def setup(db_session: AsyncSession) -> AsyncIterator[dict]:
|
||||
dev = _agent(AgentRole.DEVELOPER, f"be-dev-{uuid4().hex[:6]}")
|
||||
qa = _agent(AgentRole.QA, f"be-qa-{uuid4().hex[:6]}")
|
||||
db_session.add_all([dev, qa])
|
||||
await db_session.flush()
|
||||
project = ProjectTable(
|
||||
id=uuid4(),
|
||||
name="P",
|
||||
slug=f"p-{uuid4().hex[:6]}",
|
||||
git_url="https://example.com/r.git",
|
||||
assigned_cell=Team.BACKEND,
|
||||
created_by=dev.id,
|
||||
)
|
||||
db_session.add(project)
|
||||
await db_session.flush()
|
||||
yield {
|
||||
"svc": MetricsService(db_session),
|
||||
"db": db_session,
|
||||
"project_id": project.id,
|
||||
"dev_id": dev.id,
|
||||
"qa_id": qa.id,
|
||||
}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_returns_none_for_missing_task(setup: dict) -> None:
|
||||
assert await setup["svc"].get_task_metrics(uuid4()) is None
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_composes_effort_turns_stages_and_rework(setup: dict) -> None:
|
||||
db = setup["db"]
|
||||
tid = uuid4()
|
||||
db.add(
|
||||
TaskTable(
|
||||
id=tid,
|
||||
title="t",
|
||||
description="d",
|
||||
acceptance_criteria=["ac"],
|
||||
task_type=TaskType.CODE,
|
||||
nature=TaskNature.TECHNICAL,
|
||||
status=TaskStatus.COMPLETED,
|
||||
team=Team.BACKEND,
|
||||
project_id=setup["project_id"],
|
||||
created_by=setup["dev_id"],
|
||||
assigned_to=setup["dev_id"],
|
||||
revision_count=2,
|
||||
estimated_complexity=Complexity.MEDIUM,
|
||||
started_at=_T0,
|
||||
completed_at=_sec(7200),
|
||||
)
|
||||
)
|
||||
db.add_all(
|
||||
[
|
||||
_audit(tid, "claimed", _T0),
|
||||
_audit(tid, "in_progress", _sec(60)),
|
||||
_audit(tid, "awaiting_qa", _sec(3660)),
|
||||
_audit(tid, "completed", _sec(7200)),
|
||||
_audit(tid, "needs_revision", _sec(3660), event_type="task.qa_fail"),
|
||||
_audit(tid, "needs_revision", _sec(3000), event_type="task.pr_fail"),
|
||||
]
|
||||
)
|
||||
db.add_all(
|
||||
[
|
||||
_spawn(str(tid), _T0, _sec(600), _Usage(5, 10, 100, 50, 1.0)),
|
||||
_spawn(str(tid), _sec(3600), _sec(3660), _Usage(3, 4, 20, 10, 0.5)),
|
||||
]
|
||||
)
|
||||
await db.flush()
|
||||
|
||||
m = await setup["svc"].get_task_metrics(tid)
|
||||
assert m is not None
|
||||
# summed effort (600 + 60) vs wall-clock (2h).
|
||||
expected_active_s = 660
|
||||
expected_wall_s = 7200
|
||||
assert m.active_runtime_seconds == expected_active_s
|
||||
assert m.wall_clock_seconds == expected_wall_s
|
||||
assert (m.turns, m.tool_calls, m.tokens) == (8, 14, 180)
|
||||
assert m.cost_usd == pytest.approx(1.5)
|
||||
assert (m.revision_count, m.qa_fails, m.pr_fails, m.stints) == (2, 1, 1, 2)
|
||||
|
||||
stages = {s.status: s for s in m.stages}
|
||||
# claimed [0,60): stint1 covers it fully.
|
||||
assert (stages["claimed"].active_seconds, stages["claimed"].wait_seconds) == (60, 0)
|
||||
# in_progress [60,3660): stint1 60..600 (540) + stint2 3600..3660 (60) = 600 active.
|
||||
assert (
|
||||
stages["in_progress"].active_seconds,
|
||||
stages["in_progress"].wait_seconds,
|
||||
) == (600, 3000)
|
||||
# awaiting_qa [3660,7200): no stint running -> all wait.
|
||||
assert (
|
||||
stages["awaiting_qa"].active_seconds,
|
||||
stages["awaiting_qa"].wait_seconds,
|
||||
) == (0, 3540)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_in_flight_open_stint_and_open_window_decompose(setup: dict) -> None:
|
||||
db = setup["db"]
|
||||
tid = uuid4()
|
||||
db.add(
|
||||
TaskTable(
|
||||
id=tid,
|
||||
title="t",
|
||||
description="d",
|
||||
acceptance_criteria=["ac"],
|
||||
task_type=TaskType.CODE,
|
||||
nature=TaskNature.TECHNICAL,
|
||||
status=TaskStatus.IN_PROGRESS,
|
||||
team=Team.BACKEND,
|
||||
project_id=setup["project_id"],
|
||||
created_by=setup["dev_id"],
|
||||
assigned_to=setup["dev_id"],
|
||||
estimated_complexity=Complexity.MEDIUM,
|
||||
started_at=_T0,
|
||||
completed_at=None,
|
||||
)
|
||||
)
|
||||
db.add_all([_audit(tid, "claimed", _T0), _audit(tid, "in_progress", _sec(60))])
|
||||
# An OPEN stint (ended_at=None) -> runs to now.
|
||||
db.add(_spawn(str(tid), _T0, None, _Usage(2, 3, 1, 1, 0.1)))
|
||||
await db.flush()
|
||||
|
||||
m = await setup["svc"].get_task_metrics(tid)
|
||||
assert m is not None
|
||||
assert m.stints == 1
|
||||
assert m.active_runtime_seconds > 0 # open stint ran to now
|
||||
assert m.wall_clock_seconds > 0 # open task -> now
|
||||
# The open final window (in_progress) still decomposes.
|
||||
assert "in_progress" in {s.status for s in m.stages}
|
||||
@@ -15,6 +15,7 @@ import pytest
|
||||
import pytest_asyncio
|
||||
from roboco.db.tables import (
|
||||
AgentTable,
|
||||
AuditLogTable,
|
||||
JournalEntryTable,
|
||||
ProductTable,
|
||||
ProjectTable,
|
||||
@@ -1199,6 +1200,29 @@ async def test_ceo_reject_routes_coordination_task_to_main_pm(
|
||||
assert rejected.assigned_to == main_pm_id
|
||||
assert rejected.claimed_by is None
|
||||
|
||||
# Regression (metrics-granularity Phase 3): the coordination-root reject
|
||||
# routes through admin_set_status, so it MUST still emit a CEO-attributed
|
||||
# audit row transitioning OUT of awaiting_ceo_approval — the signal the CEO
|
||||
# scorecard pairs for approval latency + counts as a god-mode action. If a
|
||||
# future refactor drops the audit (the old gap), this fails.
|
||||
audit_rows = (
|
||||
(
|
||||
await db_session.execute(
|
||||
select(AuditLogTable).where(AuditLogTable.target_id == task.id)
|
||||
)
|
||||
)
|
||||
.scalars()
|
||||
.all()
|
||||
)
|
||||
ceo_rows = [
|
||||
a
|
||||
for a in audit_rows
|
||||
if (a.details or {}).get("agent_role") == "ceo"
|
||||
and (a.details or {}).get("from_status") == "awaiting_ceo_approval"
|
||||
]
|
||||
assert ceo_rows, "coordination ceo_reject must emit a ceo-attributed audit row"
|
||||
assert ceo_rows[0].details.get("to_status") == "pending"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_ceo_reject_routes_batch_umbrella_to_main_pm(
|
||||
|
||||
@@ -0,0 +1,95 @@
|
||||
"""sum_transcript_usage — token + turn counts from a Claude Code JSONL transcript.
|
||||
|
||||
The 5th return value is the LLM turn count: the number of UNIQUE assistant
|
||||
``message.id``s (Claude Code logs one line per content block, all sharing the
|
||||
message id, so naive line-counting would inflate both tokens and turns).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from roboco.agent_sdk.transcript_usage import sum_transcript_usage
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from pathlib import Path
|
||||
|
||||
_EXPECTED_TUPLE_LEN = 5
|
||||
|
||||
|
||||
def _line(msg_id: str | None, **usage: int) -> str:
|
||||
msg: dict[str, object] = {"usage": usage}
|
||||
if msg_id is not None:
|
||||
msg["id"] = msg_id
|
||||
return json.dumps({"message": msg})
|
||||
|
||||
|
||||
def _write(path: Path, lines: list[str]) -> None:
|
||||
path.write_text("\n".join(lines) + "\n", encoding="utf-8")
|
||||
|
||||
|
||||
def test_returns_five_tuple(tmp_path: Path) -> None:
|
||||
f = tmp_path / "t.jsonl"
|
||||
_write(f, [_line("m1", input_tokens=10, output_tokens=5)])
|
||||
result = sum_transcript_usage(f)
|
||||
assert len(result) == _EXPECTED_TUPLE_LEN
|
||||
|
||||
|
||||
def test_turns_counts_unique_message_ids(tmp_path: Path) -> None:
|
||||
f = tmp_path / "t.jsonl"
|
||||
_write(
|
||||
f,
|
||||
[
|
||||
_line("m1", input_tokens=10, output_tokens=5),
|
||||
_line("m2", input_tokens=20, output_tokens=7),
|
||||
_line("m3", input_tokens=1, output_tokens=1),
|
||||
],
|
||||
)
|
||||
_in, _out, _cr, _cw, turns = sum_transcript_usage(f)
|
||||
expected_turns = 3
|
||||
assert turns == expected_turns
|
||||
|
||||
|
||||
def test_repeated_message_id_counts_one_turn_and_one_usage(tmp_path: Path) -> None:
|
||||
# Claude Code emits one line per content block of the SAME assistant message,
|
||||
# each repeating the usage — must count once for tokens AND turns.
|
||||
f = tmp_path / "t.jsonl"
|
||||
_write(
|
||||
f,
|
||||
[
|
||||
_line("m1", input_tokens=10, output_tokens=5),
|
||||
_line("m1", input_tokens=10, output_tokens=5),
|
||||
_line("m1", input_tokens=10, output_tokens=5),
|
||||
],
|
||||
)
|
||||
tin, tout, _cr, _cw, turns = sum_transcript_usage(f)
|
||||
assert (tin, tout, turns) == (10, 5, 1)
|
||||
|
||||
|
||||
def test_malformed_lines_skipped_without_losing_turn_count(tmp_path: Path) -> None:
|
||||
f = tmp_path / "t.jsonl"
|
||||
_write(
|
||||
f,
|
||||
[
|
||||
_line("m1", input_tokens=10, output_tokens=5),
|
||||
"not json at all {{{",
|
||||
"",
|
||||
_line("m2", input_tokens=2, output_tokens=2),
|
||||
],
|
||||
)
|
||||
tin, _out, _cr, _cw, turns = sum_transcript_usage(f)
|
||||
assert (tin, turns) == (12, 2)
|
||||
|
||||
|
||||
def test_usage_line_without_id_sums_tokens_but_not_a_turn(tmp_path: Path) -> None:
|
||||
f = tmp_path / "t.jsonl"
|
||||
_write(
|
||||
f,
|
||||
[
|
||||
_line(None, input_tokens=4, output_tokens=1),
|
||||
_line("m1", input_tokens=6, output_tokens=1),
|
||||
],
|
||||
)
|
||||
tin, _out, _cr, _cw, turns = sum_transcript_usage(f)
|
||||
assert (tin, turns) == (10, 1)
|
||||
@@ -139,7 +139,11 @@ def test_missing_transcript_returns_zero_without_error(
|
||||
"/usage/sync", json={"transcript_path": str(tmp_path / "nope.jsonl")}
|
||||
)
|
||||
assert resp.status_code == _OK
|
||||
assert resp.json() == _expected([])
|
||||
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:
|
||||
@@ -167,7 +171,7 @@ def test_parser_handles_entries_without_message(tmp_path: Path) -> None:
|
||||
json.dumps({"type": "system", "subtype": "init"}),
|
||||
_assistant_line(rows[0]),
|
||||
)
|
||||
tin, tout, cread, cwrite = srv._sum_transcript_usage(transcript)
|
||||
tin, tout, cread, cwrite, turns = srv._sum_transcript_usage(transcript)
|
||||
exp = _expected(rows)
|
||||
assert (tin, tout, cread, cwrite) == (
|
||||
exp["tokens_input"],
|
||||
@@ -175,6 +179,7 @@ def test_parser_handles_entries_without_message(tmp_path: Path) -> None:
|
||||
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:
|
||||
@@ -214,7 +219,7 @@ def test_parser_dedupes_repeated_message_id(tmp_path: Path) -> None:
|
||||
_assistant_line_with_id(msg, "msg_aaa"), # tool_use block (same id)
|
||||
_assistant_line_with_id(other, "msg_bbb"),
|
||||
)
|
||||
tin, tout, cread, cwrite = srv._sum_transcript_usage(transcript)
|
||||
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) == (
|
||||
@@ -223,3 +228,20 @@ def test_parser_dedupes_repeated_message_id(tmp_path: Path) -> None:
|
||||
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
|
||||
|
||||
@@ -0,0 +1,60 @@
|
||||
"""The member_performance_daily rollup table — schema shape."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import TYPE_CHECKING, cast
|
||||
|
||||
from roboco.db.tables import MemberPerformanceDailyTable
|
||||
from sqlalchemy import UniqueConstraint
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from sqlalchemy import Table
|
||||
|
||||
_TABLE = cast("Table", MemberPerformanceDailyTable.__table__)
|
||||
|
||||
|
||||
def test_table_name() -> None:
|
||||
assert MemberPerformanceDailyTable.__tablename__ == "member_performance_daily"
|
||||
|
||||
|
||||
def test_has_all_metric_columns_including_extras() -> None:
|
||||
cols = set(_TABLE.columns.keys())
|
||||
assert {
|
||||
# core
|
||||
"date",
|
||||
"member_kind",
|
||||
"agent_slug",
|
||||
"team",
|
||||
"role",
|
||||
"tasks_completed",
|
||||
"tasks_first_pass",
|
||||
"revisions_caused",
|
||||
"revisions_received",
|
||||
"active_runtime_seconds",
|
||||
"turns",
|
||||
"tool_calls",
|
||||
"tokens",
|
||||
"cost_usd",
|
||||
"ceo_approval_dwell_seconds",
|
||||
"ceo_unblock_dwell_seconds",
|
||||
"godmode_actions",
|
||||
# the 4 CEO-approved extras + blocked_seconds
|
||||
"qa_reviews_total",
|
||||
"qa_reviews_passed",
|
||||
"escalations",
|
||||
"blocked_others",
|
||||
"idle_seconds",
|
||||
"blocked_seconds",
|
||||
} <= cols
|
||||
|
||||
|
||||
def test_natural_key_is_unique() -> None:
|
||||
uniques = [c for c in _TABLE.constraints if isinstance(c, UniqueConstraint)]
|
||||
key_sets = [{col.name for col in u.columns} for u in uniques]
|
||||
assert {"date", "member_kind", "agent_slug"} in key_sets
|
||||
|
||||
|
||||
def test_agent_slug_not_nullable() -> None:
|
||||
# NOT NULL DEFAULT '' — else the CEO row (agent_slug NULL) would duplicate
|
||||
# under Postgres' NULL-distinct UNIQUE semantics.
|
||||
assert _TABLE.columns["agent_slug"].nullable is False
|
||||
@@ -267,6 +267,43 @@ async def test_undecodable_message_is_acked_and_dead_lettered() -> None:
|
||||
assert invoked == []
|
||||
|
||||
|
||||
class _FakeRecoverRedis:
|
||||
"""Fake whose xpending_range returns the message id as BYTES (the real
|
||||
client has no decode_responses), and which captures the ids XCLAIM gets."""
|
||||
|
||||
def __init__(self, message_id: bytes) -> None:
|
||||
self._message_id = message_id
|
||||
self.claimed_ids: list[object] = []
|
||||
|
||||
async def xpending(self, *args: object, **kwargs: object) -> dict:
|
||||
del args, kwargs
|
||||
return {"pending": 1}
|
||||
|
||||
async def xpending_range(self, *args: object, **kwargs: object) -> list:
|
||||
del args, kwargs
|
||||
return [{"message_id": self._message_id, "time_since_delivered": 10_000}]
|
||||
|
||||
async def xclaim(self, *args: object, **kwargs: object) -> list:
|
||||
del args
|
||||
self.claimed_ids = cast("list[object]", kwargs.get("message_ids") or [])
|
||||
return [] # nothing claimed back → no handling
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_recover_stream_decodes_bytes_message_id_for_xclaim() -> None:
|
||||
"""xpending_range returns the message id as bytes; _recover_stream must
|
||||
decode it before XCLAIM. A raw ``str(bytes)`` yields ``"b'1782..-0'"``,
|
||||
which Redis rejects with "Unrecognized XCLAIM option", so pending-message
|
||||
recovery silently fails every reclaim tick."""
|
||||
bus = StreamEventBus()
|
||||
fake = _FakeRecoverRedis(b"1782066556728-0")
|
||||
bus._redis = cast("Redis", fake)
|
||||
|
||||
await bus._recover_stream("roboco:stream:usage", idle_time_ms=0)
|
||||
|
||||
assert fake.claimed_ids == ["1782066556728-0"] # decoded, not "b'...'"
|
||||
|
||||
|
||||
# --- periodic reclaim: a runtime handler failure is retried without a restart ---
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,100 @@
|
||||
"""compute_stage_effort — split each stage window into active vs wait seconds.
|
||||
|
||||
Pure overlap math (no DB): given a stage's [start, end) window and the agent
|
||||
spawn stints that ran during the task, ``active`` is the wall-clock time during
|
||||
which AT LEAST ONE stint was running (overlapping stints merged, so active can
|
||||
never exceed the window), and ``wait`` is the remainder. This is the wall-clock
|
||||
decomposition — distinct from summed effort (Σ stint durations), which can
|
||||
exceed wall-clock when stints run concurrently.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import UTC, datetime, timedelta
|
||||
|
||||
from roboco.foundation.policy.stage_effort import StageEffort, compute_stage_effort
|
||||
|
||||
_BASE = datetime(2026, 7, 1, 12, 0, 0, tzinfo=UTC)
|
||||
|
||||
|
||||
def _at(seconds: int) -> datetime:
|
||||
return _BASE + timedelta(seconds=seconds)
|
||||
|
||||
|
||||
def _window(status: str, start_s: int, end_s: int) -> tuple[str, datetime, datetime]:
|
||||
return (status, _at(start_s), _at(end_s))
|
||||
|
||||
|
||||
def _stint(start_s: int, end_s: int) -> tuple[datetime, datetime]:
|
||||
return (_at(start_s), _at(end_s))
|
||||
|
||||
|
||||
def _only(windows: list, stints: list) -> StageEffort:
|
||||
result = compute_stage_effort(windows, stints)
|
||||
assert len(result) == 1
|
||||
return result[0]
|
||||
|
||||
|
||||
def test_disjoint_stint_is_all_wait() -> None:
|
||||
eff = _only([_window("in_progress", 0, 100)], [_stint(200, 300)])
|
||||
assert (eff.active_seconds, eff.wait_seconds) == (0, 100)
|
||||
|
||||
|
||||
def test_fully_nested_stint() -> None:
|
||||
eff = _only([_window("in_progress", 0, 100)], [_stint(20, 50)])
|
||||
assert (eff.active_seconds, eff.wait_seconds) == (30, 70)
|
||||
|
||||
|
||||
def test_partial_overlap_clips_to_window() -> None:
|
||||
# stint runs 80..150 but window ends at 100 -> only 20s active in-window.
|
||||
eff = _only([_window("in_progress", 0, 100)], [_stint(80, 150)])
|
||||
assert (eff.active_seconds, eff.wait_seconds) == (20, 80)
|
||||
|
||||
|
||||
def test_multiple_nonoverlapping_stints_sum() -> None:
|
||||
eff = _only(
|
||||
[_window("in_progress", 0, 100)],
|
||||
[_stint(0, 10), _stint(40, 60)],
|
||||
)
|
||||
assert (eff.active_seconds, eff.wait_seconds) == (30, 70)
|
||||
|
||||
|
||||
def test_overlapping_stints_are_merged_not_double_counted() -> None:
|
||||
# [10,40) and [30,60) overlap -> merged union is [10,60) = 50s, NOT 60s.
|
||||
eff = _only(
|
||||
[_window("in_progress", 0, 100)],
|
||||
[_stint(10, 40), _stint(30, 60)],
|
||||
)
|
||||
# merged union [10,60) = 50s active, 50s wait (NOT 60s from double-count).
|
||||
assert (eff.active_seconds, eff.wait_seconds) == (50, 50)
|
||||
|
||||
|
||||
def test_active_never_exceeds_window_length() -> None:
|
||||
eff = _only(
|
||||
[_window("in_progress", 0, 100)],
|
||||
[_stint(-50, 500)], # stint dwarfs the window
|
||||
)
|
||||
assert (eff.active_seconds, eff.wait_seconds) == (100, 0)
|
||||
|
||||
|
||||
def test_zero_length_window() -> None:
|
||||
eff = _only([_window("claimed", 50, 50)], [_stint(0, 100)])
|
||||
assert (eff.active_seconds, eff.wait_seconds) == (0, 0)
|
||||
|
||||
|
||||
def test_each_window_decomposes_independently() -> None:
|
||||
windows = [_window("claimed", 0, 100), _window("in_progress", 100, 300)]
|
||||
stints = [_stint(50, 250)] # spans both windows
|
||||
result = compute_stage_effort(windows, stints)
|
||||
by_status = {e.status: e for e in result}
|
||||
claimed = by_status["claimed"]
|
||||
in_progress = by_status["in_progress"]
|
||||
assert (claimed.active_seconds, claimed.wait_seconds) == (50, 50)
|
||||
# in_progress: stint covers 100..250 of the 100..300 window.
|
||||
assert (in_progress.active_seconds, in_progress.wait_seconds) == (150, 50)
|
||||
|
||||
|
||||
def test_to_dict_shape() -> None:
|
||||
eff = _only([_window("in_progress", 0, 100)], [_stint(20, 50)])
|
||||
d = eff.to_dict()
|
||||
assert d == {"status": "in_progress", "active_seconds": 30, "wait_seconds": 70}
|
||||
@@ -269,7 +269,7 @@ async def test_finalize_spawn_session_http_error_uses_zero_tokens() -> None:
|
||||
|
||||
with (
|
||||
patch("roboco.runtime.orchestrator.httpx.AsyncClient", _client_cls),
|
||||
patch.object(orch, "_usage_from_transcript", return_value=(0, 0, 0, 0)),
|
||||
patch.object(orch, "_usage_from_transcript", return_value=(0, 0, 0, 0, 0)),
|
||||
patch("roboco.db.base.get_session_factory", return_value=db_factory),
|
||||
patch("roboco.billing.pricing.calculate_cost", return_value=0.0) as mock_cost,
|
||||
):
|
||||
@@ -303,7 +303,7 @@ async def test_finalize_spawn_session_non_200_uses_zero_tokens() -> None:
|
||||
|
||||
with (
|
||||
patch("roboco.runtime.orchestrator.httpx.AsyncClient", _client_cls),
|
||||
patch.object(orch, "_usage_from_transcript", return_value=(0, 0, 0, 0)),
|
||||
patch.object(orch, "_usage_from_transcript", return_value=(0, 0, 0, 0, 0)),
|
||||
patch("roboco.db.base.get_session_factory", return_value=db_factory),
|
||||
patch("roboco.billing.pricing.calculate_cost", return_value=0.0) as mock_cost,
|
||||
):
|
||||
@@ -389,7 +389,7 @@ async def test_sweep_token_snapshots_skips_zero_token_agents() -> None:
|
||||
|
||||
with (
|
||||
patch("roboco.runtime.orchestrator.httpx.AsyncClient", _client_cls),
|
||||
patch.object(orch, "_usage_from_transcript", return_value=(0, 0, 0, 0)),
|
||||
patch.object(orch, "_usage_from_transcript", return_value=(0, 0, 0, 0, 0)),
|
||||
patch("roboco.db.base.get_session_factory", return_value=db_factory),
|
||||
):
|
||||
await orch._sweep_token_snapshots()
|
||||
@@ -765,7 +765,9 @@ async def test_resolve_active_tokens_falls_back_to_transcript() -> None:
|
||||
)
|
||||
|
||||
client = _FakeHTTPClient(_handler)
|
||||
with patch.object(orch, "_usage_from_transcript", return_value=(6, 514, 100, 50)):
|
||||
with patch.object(
|
||||
orch, "_usage_from_transcript", return_value=(6, 514, 100, 50, 3)
|
||||
):
|
||||
tokens = await orch._resolve_active_tokens(
|
||||
cast("httpx.AsyncClient", client), _AGENT_ID
|
||||
)
|
||||
@@ -790,7 +792,7 @@ async def test_resolve_active_tokens_prefers_sdk() -> None:
|
||||
|
||||
client = _FakeHTTPClient(_handler)
|
||||
with patch.object(
|
||||
orch, "_usage_from_transcript", return_value=(999, 999, 999, 999)
|
||||
orch, "_usage_from_transcript", return_value=(999, 999, 999, 999, 0)
|
||||
) as mock_tx:
|
||||
tokens = await orch._resolve_active_tokens(
|
||||
cast("httpx.AsyncClient", client), _AGENT_ID
|
||||
@@ -800,6 +802,50 @@ async def test_resolve_active_tokens_prefers_sdk() -> None:
|
||||
mock_tx.assert_not_called()
|
||||
|
||||
|
||||
async def test_resolve_final_turns_tools_from_sdk() -> None:
|
||||
"""turns + tool_calls come from the SDK /usage/status when present."""
|
||||
orch = _make_orchestrator()
|
||||
|
||||
def _handler(_url: str) -> Any:
|
||||
return _mock_response(200, {"turns": 7, "tool_calls": 42, "tokens_input": 1})
|
||||
|
||||
with patch(
|
||||
"roboco.runtime.orchestrator.httpx.AsyncClient",
|
||||
lambda **_kw: _FakeHTTPClient(_handler),
|
||||
):
|
||||
turns, tool_calls = await orch._resolve_final_turns_tools(_AGENT_ID)
|
||||
|
||||
assert (turns, tool_calls) == (7, 42)
|
||||
|
||||
|
||||
async def test_resolve_final_turns_tools_transcript_fallback_for_turns() -> None:
|
||||
"""When the SDK reports 0 turns, fall back to the transcript turn count.
|
||||
|
||||
tool_calls has no transcript equivalent and stays 0 ("n/a").
|
||||
"""
|
||||
orch = _make_orchestrator()
|
||||
|
||||
def _handler(_url: str) -> Any:
|
||||
return _mock_response(200, {"turns": 0, "tool_calls": 0})
|
||||
|
||||
transcript_turns = 9
|
||||
with (
|
||||
patch(
|
||||
"roboco.runtime.orchestrator.httpx.AsyncClient",
|
||||
lambda **_kw: _FakeHTTPClient(_handler),
|
||||
),
|
||||
patch.object(
|
||||
orch,
|
||||
"_usage_from_transcript",
|
||||
return_value=(1, 2, 3, 4, transcript_turns),
|
||||
),
|
||||
):
|
||||
turns, tool_calls = await orch._resolve_final_turns_tools(_AGENT_ID)
|
||||
|
||||
assert turns == transcript_turns # recovered from the transcript
|
||||
assert tool_calls == 0
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# _usage_from_transcript — locate by session id across any project dir
|
||||
# ---------------------------------------------------------------------------
|
||||
@@ -836,7 +882,7 @@ def test_usage_from_transcript_finds_by_session_id_in_shared_app_dir(
|
||||
monkeypatch.setattr(Path, "home", lambda: tmp_path)
|
||||
|
||||
result = AgentOrchestrator._usage_from_transcript("main-pm", sid)
|
||||
assert result == (exp_in, exp_out, exp_cr, exp_cw)
|
||||
assert result == (exp_in, exp_out, exp_cr, exp_cw, 1) # one message => 1 turn
|
||||
|
||||
|
||||
def test_usage_from_transcript_without_session_id_uses_slug_glob(
|
||||
@@ -859,4 +905,4 @@ def test_usage_from_transcript_without_session_id_uses_slug_glob(
|
||||
monkeypatch.setattr(Path, "home", lambda: tmp_path)
|
||||
|
||||
result = AgentOrchestrator._usage_from_transcript("be-dev-1")
|
||||
assert result == (exp_in, exp_out, 0, 0)
|
||||
assert result == (exp_in, exp_out, 0, 0, 1) # one message => 1 turn
|
||||
|
||||
Reference in New Issue
Block a user