Files
roboco/tests/e2e_smoke/harness.py
T
10f039c36f feat(eval): golden-task eval harness + doctrine cohort stamp (#655)
* fix(notifications): exponential backoff + CAS claim for expired-unacked re-escalation

The sweep re-escalated every expired unacked ack-required notification
on every ~60s tick, forever — the live incident: 3 fresh blocker
escalations + Telegram DMs per minute from a static stale pile. Now
each notification carries reescalation_count / last_reescalated_at /
reescalation_delivered_count (migration 079): first fire at expiry,
then doubling intervals from 1h capped at 24h, hard stop after
ROBOCO_NOTIFICATION_MAX_REESCALATIONS (default 5) with one permanent
log carrying attempts-vs-delivered so 'seen and ignored' is
distinguishable from 'route never worked'. The due/wait/capped decision
is a pure function in foundation/policy/communications.py.

Per adversarial review, the attempt slot is claimed by compare-and-set
(UPDATE ... WHERE reescalation_count = :n) BEFORE delivery — the
previous draft leaned on the 60s dedup window, which never engages for
BLOCKER_ESCALATION (_LOOP_PRONE_TYPES excludes it), so concurrent
sweeps would have double-delivered. A lost claim skips delivery
outright. Legacy rows read as count=0 and keep today's first-fire
semantics. 61 tests incl. a two-session CAS race and a real alembic
upgrade/downgrade round trip.

* feat(budgets): per-task and per-project cost budgets (flag-gated)

tasks.budget_usd + projects.monthly_budget_usd (migration 080, chained
on 079; adds ix_agent_spawn_sessions_task_id since both enforcement
seams filter on bare task_id). Behind ROBOCO_TASK_BUDGETS_ENABLED
(default off, feature-flags card) — verifiably inert when off.

Claim-time: a project-month-spend guard applies to WORK-STARTING claims
only (i_will_work_on / i_will_plan) — per adversarial review, review/
doc/gate/inbound-PR claims are exempt so in-flight work can always
finish reviewing and merging at cap. Spend counts closed sessions'
estimated_cost_usd PLUS open sessions priced live from token snapshots
(the original closed-only sum read parallel long sessions as $0).

Sweep-side: the existing budget sweep also prices the active task's
spend vs budget_usd (TaskType defaults when null); on breach the task
is BLOCKED (HUMAN resolver, budget marker) BEFORE the graceful stop so
the unclaim no-ops and the dispatcher never respawns onto it, and the
CEO notification names both recovery steps. unblock on a budget-blocked
task re-checks live spend and refuses while still over — no silent
re-breach loop. Panel: budget inputs in both dialogs (0 rejected — a
zero budget silently blocks everything), spend logic consolidated in
TaskService.task_spend_usd. 42 new tests incl. a real-DB spend-query
suite and a two-tick non-refire sweep test.

* feat(eval): golden-task eval harness + doctrine cohort stamp

roboco/eval: 6 BenchTaskSpec fixtures run through the real lifecycle in
a disposable environment (the e2e_smoke harness's fake GitHub + local
git origin + throwaway DB catalog — real isolation, not convention),
scored deterministically (terminal status, revision_count, cycle time,
tokens/cost via the agent_spawn_sessions task_id join) plus a local-
model judge whose output is nested under a non_deterministic-marked
object so cohort diffs don't read judge noise as regression. CLI:
python -m roboco.eval run --role <slug> --cohort <name>. Source-
checkout-only by declared posture (deptry-scoped ignore + a hard
ImportError guard naming why; tests/ never ships in images or wheels).

agent_spawn_sessions.doctrine_version (migration 081, chained on 080)
is stamped at spawn-session finalize from the composed prompt layers —
with the session's model column it identifies a cohort durably.

Per adversarial review: bench runs patch the vault flags off (they were
writing real markdown into the operator's vault), and the real-spawn
OrchestratorStageSpawner is deliberately cut to NotImplementedError —
spawned containers' MCP wiring resolves to the production orchestrator
under real agent UUIDs, so real spawns wait for a dedicated follow-up;
the injectable scripted spawner is the working path. Full suite 13852
passed / 94% coverage in the source worktree; deptry/mypy/xenon clean.

---------

Co-authored-by: Renn F <rennf93@users.noreply.github.com>
2026-07-23 00:06:50 +02:00

555 lines
21 KiB
Python

"""e2e smoke harness — in-process RoboCo stack + scripted-agent driver.
Pieces (all REAL except GitHub and the LLM):
- The API: the real v1 flow/do routers + real middleware/exception handlers,
served by uvicorn in a thread, over the ephemeral test Postgres (the app's
own lazy engine is pointed at it by patching ``settings.database_*`` and
resetting ``_DbHolder``).
- Git: a local bare origin whose path CONTAINS ``github.com/<owner>/<repo>``
— ``_parse_git_url`` extracts owner/repo from it while clone/fetch/push
run tokenless over the local protocol.
- GitHub REST: a fake ``/_github`` router mounted on the same app
(``settings.github_api_base_url`` points at it). PR state lives in memory;
merges perform REAL git merges (squash included) on the bare origin, so
downstream git logic (cherry checks, freshness, branch sync) sees reality.
- Agents: ``ScriptedAgent`` reloads the REAL ``roboco.mcp.flow_server`` /
``do_server`` modules with that agent's env (id, role, role-scoped
manifest built from the real ``role_config``) and calls the REAL tool
functions, which POST to the in-process API over loopback HTTP.
"""
from __future__ import annotations
import asyncio
import importlib
import json
import os
import socket
import subprocess
import threading
import time
from contextlib import suppress
from dataclasses import dataclass, field
from typing import TYPE_CHECKING, Any
import pytest
import uvicorn
from cryptography.fernet import Fernet
from fastapi import APIRouter, FastAPI, Request
from fastapi.responses import JSONResponse
from sqlalchemy.engine.url import make_url
from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
if TYPE_CHECKING:
from collections.abc import Iterator
from pathlib import Path
from types import ModuleType
from typing import Protocol
from uuid import UUID
class TmpPathFactory(Protocol):
"""Structural stand-in for the one ``pytest.TempPathFactory`` method
``build_e2e_stack`` uses. pytest's real fixture value already
satisfies this shape, so it needs no adapter — but it lets
``roboco/eval/runner.py`` (an offline CLI, not a pytest session) drive
this same stack-building machinery with a plain temp-dir factory
instead of constructing a real ``pytest.Config``."""
def mktemp(self, basename: str, numbered: bool = True) -> Path: ...
_OWNER = "e2e-smoke"
_REPO = "proj"
def _git(cwd: Path, *args: str) -> str:
res = subprocess.run(
["git", "-C", str(cwd), *args],
capture_output=True,
text=True,
check=True,
)
return res.stdout.strip()
# ---------------------------------------------------------------------------
# Fake GitHub REST — PR state in memory, merges as REAL git ops on the origin
# ---------------------------------------------------------------------------
@dataclass
class _FakeGitHub:
origin: Path
admin_clone: Path
prs: dict[int, dict[str, Any]] = field(default_factory=dict)
comments: list[dict[str, Any]] = field(default_factory=list)
next_number: int = 1
def create_pr(self, title: str, body: str, head: str, base: str) -> dict[str, Any]:
number = self.next_number
self.next_number += 1
pr = {
"number": number,
"html_url": f"https://github.com/{_OWNER}/{_REPO}/pull/{number}",
"title": title,
"body": body,
"state": "open",
"merged": False,
"head": {
"ref": head,
"sha": self._sha_of(head),
"repo": {"full_name": f"{_OWNER}/{_REPO}"},
},
"base": {"ref": base},
"user": {"login": "e2e-bot"},
"author_association": "MEMBER",
}
self.prs[number] = pr
return pr
def _sha_of(self, branch: str) -> str:
try:
return _git(self.origin, "rev-parse", branch)
except subprocess.CalledProcessError:
return "0" * 40
def merge_pr(self, number: int, merge_method: str) -> dict[str, Any]:
pr = self.prs[number]
head, base = pr["head"]["ref"], pr["base"]["ref"]
admin = self.admin_clone
_git(admin, "fetch", "origin", "--prune")
_git(admin, "checkout", "-B", base, f"origin/{base}")
if merge_method == "squash":
_git(admin, "merge", "--squash", f"origin/{head}")
_git(admin, "commit", "-m", f"{pr['title']} (#{number})")
else:
_git(
admin,
"merge",
"--no-ff",
"-m",
f"Merge pull request #{number} from {head}",
f"origin/{head}",
)
_git(admin, "push", "origin", base)
sha = _git(admin, "rev-parse", "HEAD")
pr["merged"] = True
pr["state"] = "closed"
return {
"merged": True,
"sha": sha,
"message": "Pull Request successfully merged",
}
def open_prs(self, head: str | None, base: str | None) -> list[dict[str, Any]]:
out = []
for pr in self.prs.values():
if pr["state"] != "open":
continue
if head and pr["head"]["ref"] != head.split(":", 1)[-1]:
continue
if base and pr["base"]["ref"] != base:
continue
out.append(pr)
return out
def _fake_github_router(gh: _FakeGitHub) -> APIRouter:
r = APIRouter(prefix="/_github")
@r.get("/repos/{owner}/{repo}")
async def repo_caps(owner: str, repo: str) -> dict[str, Any]:
return {
"allow_squash_merge": True,
"allow_merge_commit": True,
"allow_rebase_merge": False,
}
@r.get("/repos/{owner}/{repo}/pulls/{number}")
async def get_pr(owner: str, repo: str, number: int) -> JSONResponse:
pr = gh.prs.get(number)
if pr is None:
return JSONResponse({"message": "Not Found"}, status_code=404)
# Real GitHub recomputes head.sha as the branch advances; a stale
# creation-time snapshot broke the unchanged-PR gate's semantics.
pr["head"]["sha"] = gh._sha_of(pr["head"]["ref"])
return JSONResponse(pr)
@r.get("/repos/{owner}/{repo}/pulls")
async def list_prs(
owner: str,
repo: str,
head: str | None = None,
base: str | None = None,
state: str = "open",
) -> list[dict[str, Any]]:
return gh.open_prs(head, base)
@r.post("/repos/{owner}/{repo}/pulls", status_code=201)
async def create_pr(owner: str, repo: str, request: Request) -> dict[str, Any]:
body = await request.json()
return gh.create_pr(
body["title"], body.get("body", ""), body["head"], body["base"]
)
@r.patch("/repos/{owner}/{repo}/pulls/{number}")
async def patch_pr(
owner: str, repo: str, number: int, request: Request
) -> JSONResponse:
pr = gh.prs.get(number)
if pr is None:
return JSONResponse({"message": "Not Found"}, status_code=404)
body = await request.json()
for key in ("title", "body", "state"):
if key in body:
pr[key] = body[key]
return JSONResponse(pr)
@r.put("/repos/{owner}/{repo}/pulls/{number}/merge")
async def merge_pr(
owner: str, repo: str, number: int, request: Request
) -> JSONResponse:
if number not in gh.prs:
return JSONResponse({"message": "Not Found"}, status_code=404)
body = await request.json()
try:
result = gh.merge_pr(number, body.get("merge_method", "merge"))
except subprocess.CalledProcessError as exc:
return JSONResponse(
{"message": f"Merge conflict: {exc.stderr}"}, status_code=409
)
return JSONResponse(result)
@r.post("/repos/{owner}/{repo}/pulls/{number}/requested_reviewers", status_code=201)
async def request_reviewers(
owner: str, repo: str, number: int, request: Request
) -> dict[str, Any]:
return gh.prs.get(number, {})
@r.post("/repos/{owner}/{repo}/issues/{number}/comments", status_code=201)
async def comment(
owner: str, repo: str, number: int, request: Request
) -> dict[str, Any]:
gh.comments.append({"number": number, "body": (await request.json())})
return {"id": len(gh.comments)}
@r.delete("/repos/{owner}/{repo}/git/refs/heads/{branch:path}", status_code=204)
async def delete_branch(owner: str, repo: str, branch: str) -> None:
with suppress(subprocess.CalledProcessError):
_git(gh.origin, "branch", "-D", branch)
# A minimal green check-runs signal for any commit — real-life NAS
# projects have CI, so the pr_pass CI-status guard should see a
# `success` state and exercise its pass-through branch, not the
# no_ci_configured 404 path.
@r.get("/repos/{owner}/{repo}/commits/{sha}/check-runs")
async def check_runs(owner: str, repo: str, sha: str) -> dict[str, Any]:
return {
"total_count": 1,
"check_runs": [
{"name": "ci", "status": "completed", "conclusion": "success"}
],
}
@r.get("/repos/{owner}/{repo}/actions/workflows")
async def workflows(owner: str, repo: str) -> dict[str, Any]:
return {"total_count": 1, "workflows": [{"id": 1, "name": "ci"}]}
return r
# ---------------------------------------------------------------------------
# Stack: settings patches + origin + app + uvicorn thread
# ---------------------------------------------------------------------------
@dataclass
class E2EStack:
base_url: str
root: Path
origin: Path
workspaces_root: Path
db_url: str
github: _FakeGitHub
def workspace_of(self, project_slug: str, team: str, agent_slug: str) -> Path:
return self.workspaces_root / project_slug / team / agent_slug
def run_db(self, coro_fn: Any) -> Any:
"""Run ``coro_fn(session)`` against a fresh engine/session and return."""
async def _run() -> Any:
engine = create_async_engine(self.db_url)
factory = async_sessionmaker(engine, expire_on_commit=False)
try:
async with factory() as session:
result = await coro_fn(session)
await session.commit()
return result
finally:
await engine.dispose()
return asyncio.run(_run())
def _free_port() -> int:
with socket.socket() as s:
s.bind(("127.0.0.1", 0))
return int(s.getsockname()[1])
def _seed_origin(root: Path) -> Path:
"""Bare origin at a path _parse_git_url can read owner/repo from."""
origin = root / "github.com" / _OWNER / f"{_REPO}.git"
origin.parent.mkdir(parents=True)
subprocess.run(
["git", "init", "--bare", "--initial-branch=master", str(origin)],
check=True,
capture_output=True,
)
seed = root / "seed-clone"
subprocess.run(
["git", "clone", str(origin), str(seed)], check=True, capture_output=True
)
_git(seed, "config", "user.name", "roboco-e2e")
_git(seed, "config", "user.email", "e2e@roboco.local")
(seed / "README.md").write_text("# e2e smoke project\n")
_git(seed, "add", "README.md")
_git(seed, "commit", "-m", "Initial commit")
_git(seed, "push", "origin", "master")
return origin
def _make_admin_clone(root: Path, origin: Path) -> Path:
admin = root / "gh-admin-clone"
subprocess.run(
["git", "clone", str(origin), str(admin)], check=True, capture_output=True
)
_git(admin, "config", "user.name", "fake-github")
_git(admin, "config", "user.email", "merge@github.local")
return admin
def _build_app(gh: _FakeGitHub) -> FastAPI:
from roboco.api.middleware import setup_middleware
from roboco.api.routes.dashboard import router as dashboard_router
from roboco.api.routes.health import router as health_router
from roboco.api.routes.notifications import router as notifications_router
from roboco.api.routes.orchestrator import router as orchestrator_router
from roboco.api.routes.settings import router as settings_router
from roboco.api.routes.tasks import router as tasks_router
from roboco.api.routes.v1 import do as do_module
from roboco.api.routes.v1 import flow_auditor as fa
from roboco.api.routes.v1 import flow_board as fb
from roboco.api.routes.v1 import flow_cell_pm as fcp
from roboco.api.routes.v1 import flow_dev as fd
from roboco.api.routes.v1 import flow_doc as fdoc
from roboco.api.routes.v1 import flow_main_pm as fmp
from roboco.api.routes.v1 import flow_pr_reviewer as fpr
from roboco.api.routes.v1 import flow_qa as fq
app = FastAPI(title="roboco-e2e-smoke")
setup_middleware(app)
app.include_router(health_router)
for module in (fd, fq, fdoc, fcp, fmp, fb, fa, fpr):
app.include_router(module.router)
app.include_router(do_module.router)
# The REST task surface — scenario 3 drives the real CEO
# approve-and-merge endpoint (the human gate) through it.
app.include_router(tasks_router, prefix="/api/tasks")
# The notifications router was already mounted here before the auditor
# revival diff; it remains so reactive audit dispatch (``_dispatch_audit_work``)
# can poll real ALERT rows end-to-end.
app.include_router(notifications_router, prefix="/api/notifications")
# Cloud-auth gate coverage smoke exercises the real _require_ceo and
# require_panel_token dep paths on these routers.
app.include_router(orchestrator_router, prefix="/api/orchestrator")
app.include_router(settings_router, prefix="/api/settings")
app.include_router(dashboard_router, prefix="/api/dashboard")
app.include_router(_fake_github_router(gh))
return app
def build_e2e_stack(
_test_database_url: str, tmp_path_factory: TmpPathFactory
) -> Iterator[E2EStack]:
"""Generator behind the ``e2e_stack`` fixture (defined in conftest).
``tmp_path_factory`` only needs ``.mktemp()`` (see ``TmpPathFactory``
above) — pytest's real fixture satisfies it structurally, and
``roboco/eval/runner.py`` drives this same function with a plain
non-pytest factory to reuse this stack outside a test session.
"""
from roboco.config import settings
from roboco.db import base as db_base
mp = pytest.MonkeyPatch()
root = tmp_path_factory.mktemp("e2e")
origin = _seed_origin(root)
admin = _make_admin_clone(root, origin)
gh = _FakeGitHub(origin=origin, admin_clone=admin)
workspaces = root / "workspaces"
workspaces.mkdir()
url = make_url(_test_database_url)
mp.setattr(settings, "database_host", url.host or "localhost")
mp.setattr(settings, "database_port", url.port or 5432)
mp.setattr(settings, "database_user", url.username or "")
mp.setattr(settings, "database_password", url.password or "")
mp.setattr(settings, "database_name", url.database or "")
mp.setattr(settings, "workspaces_root", str(workspaces))
mp.setattr(settings, "workspace_auto_clone", True)
mp.setattr(settings, "encryption_key", Fernet.generate_key().decode())
# The app's lazy engine must bind to the patched settings, not a leftover.
db_base._DbHolder.engine = None
db_base._DbHolder.session_factory = None
port = _free_port()
base_url = f"http://127.0.0.1:{port}"
mp.setattr(settings, "github_api_base_url", f"{base_url}/_github")
app = _build_app(gh)
# loop=settings.uvicorn_loop ("asyncio" by default): uvicorn auto-selects
# uvloop when installed, and this in-thread server has crashed CI with a
# uvloop/asyncpg segfault (uvloop 0.22 + asyncpg 0.31 + Python 3.13) —
# mirror the production default instead of picking up uvloop implicitly.
server = uvicorn.Server(
uvicorn.Config(
app,
host="127.0.0.1",
port=port,
log_level="warning",
loop=settings.uvicorn_loop,
)
)
thread = threading.Thread(target=server.run, daemon=True)
thread.start()
import httpx
deadline = time.time() + 30
while time.time() < deadline:
try:
# Any HTTP response at all means the server thread is up.
httpx.get(f"{base_url}/health", timeout=1)
break
except httpx.HTTPError:
time.sleep(0.1)
else:
raise RuntimeError("e2e app server did not become ready")
try:
yield E2EStack(
base_url=base_url,
root=root,
origin=origin,
workspaces_root=workspaces,
db_url=_test_database_url,
github=gh,
)
finally:
server.should_exit = True
thread.join(timeout=10)
db_base._DbHolder.engine = None
db_base._DbHolder.session_factory = None
mp.undo()
# ---------------------------------------------------------------------------
# Scripted agents — the REAL MCP tool functions, per-agent module reloads
# ---------------------------------------------------------------------------
class ScriptedAgent:
"""Drives the real flow/do MCP tool functions as one seeded agent."""
def __init__(self, stack: E2EStack, agent_id: UUID, slug: str, role: str) -> None:
self.stack = stack
self.agent_id = agent_id
self.slug = slug
self.role = role
self._manifest_path = stack.root / f"manifest-{slug}.json"
self._manifest_path.write_text(json.dumps(self._manifest()))
def _manifest(self) -> dict[str, Any]:
from roboco.services.gateway.role_config import get_role_config
cfg = get_role_config(self.role)
return {
"agent_id": str(self.agent_id),
"role": self.role,
"team": "backend",
"workspace_path": str(self.stack.workspaces_root),
"flow_tools": list(cfg.flow_tools),
"do_tools": list(cfg.do_tools),
"read_tools": ["Read", "Glob", "Grep"],
"write_tools": ["Edit", "Write"] if cfg.allows_write else [],
"bash_allowed": True,
"subagent_allowed": False,
"subagent_model": None,
"env": {},
}
def _module(self, name: str) -> ModuleType:
os.environ["ROBOCO_AGENT_ID"] = str(self.agent_id)
os.environ["ROBOCO_AGENT_ROLE"] = self.role
os.environ["ROBOCO_ORCHESTRATOR_URL"] = self.stack.base_url
os.environ["ROBOCO_TOOL_MANIFEST_PATH"] = str(self._manifest_path)
# The host agent environment may carry a real ROBOCO_AGENT_TOKEN issued
# for the test runner's identity. flow_server reads it before each call
# and forwards it in X-Agent-Token; the token won't match the ephemeral
# test agent IDs and causes 401s. Drop it so tests run in the same
# unsigned-token mode as CI.
os.environ.pop("ROBOCO_AGENT_TOKEN", None)
# Same leakage class for the per-verb circuit breaker: flow_server /
# do_server post rejections to ROBOCO_SDK_URL (default
# http://localhost:9000), which is this repo's own live agent SDK
# loopback port when the suite happens to run inside a real spawned
# agent container. That breaker then records genuine attempts for
# the ephemeral test-agent IDs and can trip circuit_open mid-test
# (surfaced by test_request_sandbox_guard_chain_over_real_api's 3
# deliberately-rejected calls). Point it at a guaranteed-refused
# loopback port so every environment sees the same fail-open
# bypass CI gets (no SDK process listening at all).
os.environ["ROBOCO_SDK_URL"] = "http://127.0.0.1:1"
module = importlib.import_module(name)
if getattr(module, "AGENT_ID", None) != str(self.agent_id):
module = importlib.reload(module)
return module
def flow(self, verb: str, /, **kwargs: Any) -> dict[str, Any]:
result: dict[str, Any] = getattr(self._module("roboco.mcp.flow_server"), verb)(
**kwargs
)
return result
def do(self, tool: str, /, **kwargs: Any) -> dict[str, Any]:
result: dict[str, Any] = getattr(self._module("roboco.mcp.do_server"), tool)(
**kwargs
)
return result
def expect_error(env: dict[str, Any], kind: str, context: str) -> dict[str, Any]:
"""Assert an envelope is the EXPECTED rejection kind."""
assert env.get("error") == kind, (
f"{context}: expected rejection {kind!r}, got error={env.get('error')!r}\n"
f" full: {json.dumps(env, default=str, indent=2)[:4000]}"
)
return env
def expect_ok(env: dict[str, Any], context: str) -> dict[str, Any]:
"""Assert an envelope is a success; on failure show the whole envelope."""
assert isinstance(env, dict), f"{context}: non-dict envelope: {env!r}"
assert not env.get("error"), (
f"{context}: rejected with error={env.get('error')!r}\n"
f" message : {env.get('message')}\n"
f" remediate: {env.get('remediate')}\n"
f" missing : {env.get('missing')}\n"
f" full : {json.dumps(env, default=str, indent=2)[:4000]}"
)
return env