From 91249efb86dfc3a8bd4875f3ee48c2681fb49426 Mon Sep 17 00:00:00 2001 From: Tin Dang Date: Thu, 13 Aug 2026 13:22:57 +0700 Subject: [PATCH 1/3] =?UTF-8?q?feat(evals):=20eval-run-executor=20?= =?UTF-8?q?=E2=80=94=20governed=20set=20replay=20(R7=20L1)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Replay an eval set against a model THROUGH the existing governance path, so a run is a burst of billed upstream calls that enters the same budget/credit/tier/rate guards a live request does — and never trips a shared breaker that would degrade another tenant. Two new tenant-scoped tables (four-manifest rule: migration f5b2d8c41a37 + EXPECTED_TABLES + alembic env.py import + guardrails NOT-IN allow-list): - eval_runs one row per launched run; status DERIVED from its cases (M7) - eval_case_results per-case outcome; response_text is the ZDR-gated payload-at-rest; UNIQUE(eval_run_id, eval_case_id) makes a resumed drive idempotent The executor REUSES CompletionUseCase.complete() per case (never reimplements governance), so M1 (9-guard order) + M2 (one usage_record per dial, on the served model id) hold by construction. It injects its OWN per-tenant breaker + concurrency semaphore (TenantExecutionRegistry) via complete(upstream=...) — NEVER the global app.state.circuit_breaker — closing the recurring cross-tenant-DoS defect (R:GLOBAL_BREAKER). Every dial carries a per-call timeout; breaker-open and timeout both fail the case closed (errored), the run continues. Security: - ZDR refused OUTRIGHT at launch (403 ERR_ZDR_PAYLOAD_BLOCKED) AND re-checked ATOMICALLY (raise_if_zdr_locked) at the first payload write — a flip mid-run persists nothing further and marks the run blocked. - Auth-scoped resume: the raw API key is NEVER persisted at rest (only key_id). In-process crash resumes from the in-memory key; a cross-redeploy resume must re-supply the key via a fresh authenticated request — never a forged identity. - Uniform 404 for absent/cross-tenant runs (ERR_EVAL_RUN_NOT_FOUND) — no oracle. API (one OpenAI-wire envelope, reused from the eval-set-store surface): POST /v1/evals/sets/{set_id}/runs -> 201 { id:"er_..", status, case_count, ... } GET /v1/evals/runs/{run_id} -> 200 { status, counts:{...} } GET /v1/evals/runs/{run_id}/cases -> 200 list, case-creation order (A5) Launch enqueues onto app.state.eval_run_queue when present and FAILS OPEN to an inline drive (the vector-store ingest idiom) — the durable Redis worker loop is a follow-up; the seam is in place and no CHECK depends on it. TDD: 13 red-first CHECKS (M1-M8, A1-A6, E1-E6, R:*) — governance-per-case, over-budget-no-dial, one-usage-record-per-dial, per-tenant-breaker isolation, ZDR-atomic-zero-results, resume-no-rebill, timeout-errored-run-continues, uniform-404, case-creation-order, empty-set-vacuous, cross-tenant-launch + billing-identity, snapshot-fixed-at-launch, bounded-per-tenant-concurrency. make ci green incl. migration-parity + guardrails no-new-tables gates; pyright + ruff clean. ADD: eval-run-executor frozen sha256:2abaa868 (reseal binding A1/A2/A3/A6/M8) -> gate PASS. R7 evals-regression-gate. author: Tin Dang --- .add/tasks/eval-run-executor.d/runs/1.md | 15 + .add/tasks/eval-run-executor.d/runs/2.md | 26 + .add/tasks/eval-run-executor.d/runs/3.md | 29 + .add/tasks/eval-run-executor.md | 121 ++- apps/gateway/migrations/env.py | 1 + .../f5b2d8c41a37_eval_run_executor.py | 98 +++ .../gateway/src/gateway/core/error_catalog.py | 9 + .../src/gateway/evals/runs/__init__.py | 0 .../src/gateway/evals/runs/api/__init__.py | 0 .../src/gateway/evals/runs/api/run_router.py | 212 ++++++ .../evals/runs/application/__init__.py | 0 .../evals/runs/application/run_executor.py | 252 +++++++ .../src/gateway/evals/runs/domain/__init__.py | 0 .../src/gateway/evals/runs/domain/entities.py | 35 + .../src/gateway/evals/runs/domain/errors.py | 15 + .../src/gateway/evals/runs/domain/ports.py | 87 +++ .../evals/runs/infrastructure/__init__.py | 0 .../gateway/evals/runs/infrastructure/orm.py | 115 +++ .../evals/runs/infrastructure/repository.py | 175 +++++ .../runs/infrastructure/upstream_adapter.py | 117 +++ apps/gateway/src/gateway/evals/wire_id.py | 12 +- apps/gateway/src/gateway/main.py | 33 + apps/gateway/tests/evals_runs/__init__.py | 0 .../evals_runs/test_eval_run_executor.py | 697 ++++++++++++++++++ .../tests/guardrails/test_guardrails_core.py | 7 +- .../tests/migrations/test_migrations.py | 2 + 26 files changed, 2039 insertions(+), 19 deletions(-) create mode 100644 .add/tasks/eval-run-executor.d/runs/1.md create mode 100644 .add/tasks/eval-run-executor.d/runs/2.md create mode 100644 .add/tasks/eval-run-executor.d/runs/3.md create mode 100644 apps/gateway/migrations/versions/f5b2d8c41a37_eval_run_executor.py create mode 100644 apps/gateway/src/gateway/evals/runs/__init__.py create mode 100644 apps/gateway/src/gateway/evals/runs/api/__init__.py create mode 100644 apps/gateway/src/gateway/evals/runs/api/run_router.py create mode 100644 apps/gateway/src/gateway/evals/runs/application/__init__.py create mode 100644 apps/gateway/src/gateway/evals/runs/application/run_executor.py create mode 100644 apps/gateway/src/gateway/evals/runs/domain/__init__.py create mode 100644 apps/gateway/src/gateway/evals/runs/domain/entities.py create mode 100644 apps/gateway/src/gateway/evals/runs/domain/errors.py create mode 100644 apps/gateway/src/gateway/evals/runs/domain/ports.py create mode 100644 apps/gateway/src/gateway/evals/runs/infrastructure/__init__.py create mode 100644 apps/gateway/src/gateway/evals/runs/infrastructure/orm.py create mode 100644 apps/gateway/src/gateway/evals/runs/infrastructure/repository.py create mode 100644 apps/gateway/src/gateway/evals/runs/infrastructure/upstream_adapter.py create mode 100644 apps/gateway/tests/evals_runs/__init__.py create mode 100644 apps/gateway/tests/evals_runs/test_eval_run_executor.py diff --git a/.add/tasks/eval-run-executor.d/runs/1.md b/.add/tasks/eval-run-executor.d/runs/1.md new file mode 100644 index 00000000..bc74cc10 --- /dev/null +++ b/.add/tasks/eval-run-executor.d/runs/1.md @@ -0,0 +1,15 @@ +--- +type: Run +runtime: process +task: /tasks/eval-run-executor.md +computation: "uv run pytest tests/evals_runs/test_eval_run_executor.py -q --no-cov -p no:randomly --junitxml tests/evals_runs/_add_junit.xml" +receipt: + kind: command-exit + ids: unknown + exit: 0 + freshness: mtime + at: 2026-08-13 + stdout: '10 passed in 5.76s' + note: '' +generated: { by: process:run, at: 2026-08-13 } +--- diff --git a/.add/tasks/eval-run-executor.d/runs/2.md b/.add/tasks/eval-run-executor.d/runs/2.md new file mode 100644 index 00000000..c9bbca88 --- /dev/null +++ b/.add/tasks/eval-run-executor.d/runs/2.md @@ -0,0 +1,26 @@ +--- +type: Run +runtime: process +task: /tasks/eval-run-executor.md +computation: "uv run pytest tests/evals_runs/test_eval_run_executor.py -q --no-cov -p no:randomly --junitxml=/private/tmp/claude-501/-Users-tindang-workspaces-tind-repo-ai-proxy/88454d31-bfec-4421-87a2-2d625a0ab229/scratchpad/evals_runs_add.xml" +receipt: + kind: test-ids + ids: 10/10 reported + exit: 0 + freshness: mtime + at: 2026-08-13 + stdout: '10 passed in 5.65s' + note: '' + passed: + - tests.evals_runs.test_eval_run_executor::test_cross_tenant_and_absent_run_uniform_404 + - tests.evals_runs.test_eval_run_executor::test_empty_set_run_completes_vacuously + - tests.evals_runs.test_eval_run_executor::test_one_usage_record_per_dialed_case + - tests.evals_runs.test_eval_run_executor::test_over_budget_run_refuses_every_case_no_dial + - tests.evals_runs.test_eval_run_executor::test_per_tenant_breaker_isolation + - tests.evals_runs.test_eval_run_executor::test_results_aligned_to_case_creation_order + - tests.evals_runs.test_eval_run_executor::test_resume_does_not_rebill_terminal_cases + - tests.evals_runs.test_eval_run_executor::test_run_enters_governance_per_case + - tests.evals_runs.test_eval_run_executor::test_timeout_case_errored_run_continues + - tests.evals_runs.test_eval_run_executor::test_zdr_run_refused_atomically_zero_results +generated: { by: process:run, at: 2026-08-13 } +--- diff --git a/.add/tasks/eval-run-executor.d/runs/3.md b/.add/tasks/eval-run-executor.d/runs/3.md new file mode 100644 index 00000000..487a7ddf --- /dev/null +++ b/.add/tasks/eval-run-executor.d/runs/3.md @@ -0,0 +1,29 @@ +--- +type: Run +runtime: process +task: /tasks/eval-run-executor.md +computation: "uv run pytest tests/evals_runs/test_eval_run_executor.py -q --no-cov -p no:randomly --junitxml=/private/tmp/claude-501/-Users-tindang-workspaces-tind-repo-ai-proxy/88454d31-bfec-4421-87a2-2d625a0ab229/scratchpad/evals_runs_final.xml" +receipt: + kind: test-ids + ids: 13/13 reported + exit: 0 + freshness: mtime + at: 2026-08-13 + stdout: '13 passed in 7.19s' + note: '' + passed: + - tests.evals_runs.test_eval_run_executor::test_bounded_per_tenant_concurrency + - tests.evals_runs.test_eval_run_executor::test_cross_tenant_and_absent_run_uniform_404 + - tests.evals_runs.test_eval_run_executor::test_cross_tenant_launch_and_billing_identity + - tests.evals_runs.test_eval_run_executor::test_empty_set_run_completes_vacuously + - tests.evals_runs.test_eval_run_executor::test_one_usage_record_per_dialed_case + - tests.evals_runs.test_eval_run_executor::test_over_budget_run_refuses_every_case_no_dial + - tests.evals_runs.test_eval_run_executor::test_per_tenant_breaker_isolation + - tests.evals_runs.test_eval_run_executor::test_results_aligned_to_case_creation_order + - tests.evals_runs.test_eval_run_executor::test_resume_does_not_rebill_terminal_cases + - tests.evals_runs.test_eval_run_executor::test_run_enters_governance_per_case + - tests.evals_runs.test_eval_run_executor::test_run_snapshot_fixed_at_launch + - tests.evals_runs.test_eval_run_executor::test_timeout_case_errored_run_continues + - tests.evals_runs.test_eval_run_executor::test_zdr_run_refused_atomically_zero_results +generated: { by: process:run, at: 2026-08-13 } +--- diff --git a/.add/tasks/eval-run-executor.md b/.add/tasks/eval-run-executor.md index 921c63a1..3f477c02 100644 --- a/.add/tasks/eval-run-executor.md +++ b/.add/tasks/eval-run-executor.md @@ -1,47 +1,134 @@ --- type: Task title: eval-run-executor -status: direction +status: done milestone: evals-regression-gate needs: - eval-set-store.md gives: - S1 the run-execution port — tenant-keyed breaker, bounded concurrency, timeouts, one usage_record per case through the governance path generated: { by: add/3.2.0, at: 2026-08-12 } -verified: [] +verified: + - { by: "Tin Dang", at: 2026-08-13, act: freeze, authority: process, direction: "sha256:1353206db5fb398c" } + - { by: "cli", at: 2026-08-13, act: brief, authority: process, brief: "sha256:1181f2f2c2711ebd" } + - { by: "process:run", at: 2026-08-13, act: run, authority: process, outcome: PASS, receipt: /tasks/eval-run-executor.d/runs/1.md } + - { by: "process:run", at: 2026-08-13, act: run, authority: process, outcome: PASS, receipt: /tasks/eval-run-executor.d/runs/2.md } + - { by: "Tin Dang", at: 2026-08-13, act: refreeze, authority: process, direction: "sha256:2abaa8687a94642b" } + - { by: "cli", at: 2026-08-13, act: brief, authority: process, brief: "sha256:c478d7392f02d99d" } + - { by: "process:run", at: 2026-08-13, act: run, authority: process, outcome: PASS, receipt: /tasks/eval-run-executor.d/runs/3.md } + - { by: "add-worker", at: 2026-08-13, act: gate, authority: process, outcome: PASS, receipt: /tasks/eval-run-executor.d/runs/3.md, brief: "sha256:c478d7392f02d99d", reason: "13/13 CHECKS bound green (M1-M8, A1-A6, E1-E6, R:*); migration-parity + guardrails manifests + pyright + ruff all green; per-tenant breaker isolation, ZDR-atomic launch/mid-run refusal, auth-scoped resume (no key material at rest), bounded per-tenant concurrency all verified" } +advised_by: appsec-engineer --- ## CARD goal: run a set against a model through the governance path — per-tenant breaker, bounded concurrency, timeouts, partial-run resumability why: a run is a burst of billed upstream calls; it must enter through the budget/credit/tier guards and never trip a shared breaker that degrades another tenant -beat: direction · next: add freeze eval-run-executor +beat: done · next: add status ## RULES -- M1 +- M1 A run replays each case's stored `request_body` against the run's named model THROUGH the existing governance path (`NonChatGovernance.authorize` / the completion use-case), never around it — every ordered guard (auth · expiry · allowlist · catalog/model_checker tri-state · per-key/team/tenant budget · credit · tier · rate-limit) runs per case exactly as for a live request. A run enters no back door. +- M2 A run produces EXACTLY ONE `usage_records` row per case that actually dials upstream, with the SAME shape a live request of that `request_body` would produce (billed on the SERVED model id, not `request_body["model"]`) — no unmetered case, no double-bill, no eval side-channel. A case refused before the dial (M3) writes NO usage row. +- M3 A case whose per-case governance check fails (budget/credit/tier/rate/allowlist/disabled model) is REFUSED with NO upstream call and NO usage row for that case; the run records the case as `refused` with the governance reason. A tenant at their credit/budget limit cannot spend through a run. (Assert NO provider call was made.) +- M4 The per-case upstream dial goes through a breaker + concurrency limiter keyed by `tenant_id` — the [[per-tenant-breaker-recurring-defect]] anchor: NEVER a process-global breaker/semaphore. Tenant A's failing burst opens A's breaker only; tenant B's request in the same process is unaffected. Every outbound dial has a per-call TIMEOUT; breaker-open and timeout both fail the case CLOSED (recorded `errored`), never hang the run. +- M5 A run persists per-case RESULTS including the model's `response_text` (the payload the console diff + re-scoring read). This is a payload-at-rest surface, so ZDR is enforced: a ZDR tenant's run is REFUSED outright at launch with 403 `ERR_ZDR_PAYLOAD_BLOCKED`, and the check is ATOMIC with the first result write (`raise_if_zdr_locked`, SELECT … FOR UPDATE) — a flip landing mid-run persists NOTHING further. (Same disposition as [[eval-set-store]] M3; assertion-only/redacted-run mode is a later-milestone follow-up, not R7.) +- M6 Reads/writes are tenant-scoped in the SAME query that resolves the row: a run/set/case owned by another tenant is a uniform 404 `ERR_EVAL_SET_NOT_FOUND` / `ERR_EVAL_RUN_NOT_FOUND`, never a distinguishable error (carried isolation invariant, #84). A run grants NO model visibility a normal request lacks — the candidate model resolves under `tenant_id IS NULL OR = :tenant`, including finetuned models. +- M7 A run is DURABLE and resumable: each case result is committed as it completes, and a run interrupted (crash/redeploy) mid-flight can be resumed to drive only the cases still `pending` — no case is dialed (and billed) twice on resume (a completed/errored/refused case is terminal; re-drive is a no-op). The run's terminal status is derived from its cases, not written speculatively. +- M8 The execution engine is reached through a `typing.Protocol` port with a zero-network fake usable from `app.state` (a fake upstream that returns canned responses / raises / times out on command) — the use-case never imports the concrete provider client; the breaker/governance/timeout behavior is provable without a network (backend-architect lens). -- R: -> "" +- R:GOVERNANCE_BYPASS a case dials upstream without passing the same governance guards a live request would -> "a run must enter through governance, never around it" +- R:ZDR_BLOCKED a ZDR tenant's run persists a response payload at rest -> "ERR_ZDR_PAYLOAD_BLOCKED" +- R:GLOBAL_BREAKER one tenant's run opens a breaker/exhausts a limiter that degrades another tenant -> "the breaker + concurrency limiter are per-tenant, never global" +- R:DOUBLE_BILL a case produces zero or more-than-one usage_record, or a resumed run re-dials a terminal case -> "exactly one usage_record per dialed case; resume never re-bills" +- R:RUN_NOT_FOUND a run/set/case owned by another tenant is distinguishable from an absent one -> "ERR_EVAL_RUN_NOT_FOUND" ## ASSUMPTIONS -- A1 [who] covers: · the request does not say ; taking -> -- A2 [which] covers: · the request does not say ; taking -> -- A3 [when] covers: · the request does not say ; taking -> -- A4 [absent] covers: · the request does not say ; taking -> -- A5 [order] covers: · the request does not say ; taking -> -- A6 [experience] covers: · the request does not say ; taking -> -every `gives:` surface is swept on every dimension; `[] n/a · ` retires one. one line, one silence — split, never bundle. `· probe: ` declares a reading checkable: cite its A id from CHECKS and the gate holds the PASS to it. +- A1 [who] covers: S1 · the request does not say who may launch a run or whose key it bills; taking "any authenticated member of the owning tenant launches a run on THAT tenant's set, billed to the launching key exactly as that key's live traffic; superadmin only via impersonation, never cross-tenant" -> if wrong, a run bills the wrong key or leaks a cross-tenant set. · probe: tenant B launching a run on tenant A's set gets 404; a run's usage_records carry the launching tenant/key. +- A2 [which] covers: S1 · the request does not say which cases a run covers; taking "a run covers ALL cases in the set AT LAUNCH TIME (a snapshot by created_at); cases added to the set after launch are NOT in the in-flight run — a later run picks them up" -> if wrong, a run's case set shifts under it and the verdict has no fixed denominator. · probe: adding a case after launch does not change the running run's case count. +- A3 [when] covers: S1 · the request does not say the concurrency/timeout bounds; taking "a bounded per-tenant concurrency (a small fixed ceiling, config-driven) and a per-call timeout from the same settings the live proxy uses — a run is throttled to protect shared upstream capacity, not fanned out unbounded" -> if wrong, a run either crawls or stampedes the provider and other tenants. · probe: a run of N>ceiling cases never has more than the ceiling in flight for that tenant at once. +- A4 [absent] covers: S1 · the request does not say what a run of an EMPTY set (or a set whose every case is refused) means; taking "an empty/all-refused run completes with status `completed`, zero dialed cases, a 0/0 result — a valid (if vacuous) run, never an error; the verdict layer decides what 0/0 means" -> if wrong, a legitimate empty set crashes the runner. · probe: running an empty set returns a completed run with an empty result list, no exception. +- A5 [order] covers: S1 · the request does not say the result/case ordering or how a tie in completion is broken; taking "results are aligned to the case's stable creation order (the eval-set-store A5 order), NOT wall-clock completion order — so a baseline and a candidate run of the same set align case-for-case regardless of which finished first" -> if wrong, the per-case diff misaligns and the verdict compares apples to oranges. · probe: two runs of the same set expose results in identical case order. +- A6 [experience] covers: S1 · the request does not say what a partially-failed run shows; taking "a run surfaces a per-case status (`completed` scored-later · `refused` with the governance reason · `errored` with a breaker/timeout reason) plus a run-level rollup, so an operator sees WHICH cases ran, which were refused for spend, which errored — never a bare 'run failed'" -> if wrong, a half-run is an unactionable red bar. · probe: a run mixing dialed + budget-refused + timed-out cases exposes a distinct, actionable status per case. ## PLAN -contract: -scope: +contract: +``` +POST /v1/evals/sets/{set_id}/runs body: { model } # launch a run of the set against a model + 201 -> { id:"er_<32hex>", eval_set_id, model, status:"pending", case_count, created_at } + 403 -> ERR_ZDR_PAYLOAD_BLOCKED # ZDR tenant, atomic with the first result write (M5) + 404 -> ERR_EVAL_SET_NOT_FOUND # absent OR cross-tenant set (M6) + 402/403/429 (per governance) at PER-CASE dial time -> the case is `refused`, not the whole run +GET /v1/evals/runs/{run_id} -> { id, eval_set_id, model, status, case_count, + counts:{completed,refused,errored,pending}, created_at } +GET /v1/evals/runs/{run_id}/cases -> [ { eval_case_id, status, response_text?, usage_record_id?, + reason? } ] # tenant-scoped, CASE creation order (A5) + # response_text present only for `completed` cases; a ZDR tenant never reaches this (M5) + +Schema (new tables, tenant-scoped; response_text is the ZDR-gated payload-at-rest surface): + eval_runs(id uuid pk, tenant_id uuid fk, eval_set_id uuid fk, model text, status text, + created_at timestamptz) — index (tenant_id, eval_set_id, created_at) + eval_case_results(id uuid pk, tenant_id uuid fk, eval_run_id uuid fk ON DELETE CASCADE, + eval_case_id uuid fk, status text, response_text text null, reason text null, + usage_record_id uuid null, created_at timestamptz) + — unique (eval_run_id, eval_case_id) # idempotent per case on resume (M7, R:DOUBLE_BILL) + — index (tenant_id, eval_run_id, created_at) + +Execution (per case, bounded per-tenant concurrency, each committed as it lands): + authorize via governance (M1) -> refused? record `refused`+reason, NO dial, NO usage (M3) + -> else dial upstream through the PER-TENANT breaker + timeout (M4) + -> breaker-open/timeout -> record `errored`+reason (fail closed) + -> ok -> raise_if_zdr_locked(s, tenant) then persist `completed`+response_text, one usage row (M2/M5) + status derived from case counts (M7); resume drives only `pending` cases, terminal cases are no-ops. +``` + +BUILD NOTE (reuse seam — implementation guidance, PLAN not sealed): + - Per case, REUSE `CompletionUseCase.complete(raw_key=, body=, upstream=, usage_recorder=, + ...)`. This is what makes M1 (enter through governance) + M2 (one usage_record, SAME shape as + a live request, billed on the SERVED model) hold BY CONSTRUCTION rather than by re-implementation. + A governance refusal surfaces as the ProblemError `complete()` already raises (budget/credit/ + tier/rate) -> map to `refused`+reason, NO dial (M3). An upstream 5xx / CircuitOpenError -> + `errored`. + - ⚠ M4 HARD-STOP realization: the LIVE completion path's breaker (`BoundCircuitBreakerUpstream` + over `app.state.circuit_breaker`) is GLOBAL — one breaker per app instance, NOT per tenant + ([[per-tenant-breaker-recurring-defect]]). Reusing it for an eval BURST would let tenant A's + run open the SHARED breaker and DoS tenant B's LIVE traffic (R:GLOBAL_BREAKER). So the executor + MUST pass its OWN `upstream` into `complete()`: a per-tenant `CircuitBreaker` (from a registry + keyed by tenant_id) wrapping the raw `app.state.completion_upstream` delegate — the `upstream=` + parameter is the injection seam. Eval dials then trip only that tenant's eval-breaker; the live + global breaker is never touched by an eval run, and the per-tenant concurrency semaphore + (also a tenant-keyed registry, OWNED by the executor) is separate from the breaker. + - Response extraction (A2 handoff to scorers): pull the assistant text from the `complete()` + json_body (choices[0].message.content) and persist it as `response_text`; scoring is a LATER + pure pass (deterministic-scorers), not this task. +scope (may touch): `apps/gateway/src/gateway/evals/` (NEW: `runs/domain/{ports,entities,errors}.py` · `runs/application/run_executor.py` [the use-case + resume driver] · `runs/infrastructure/{orm,repository,upstream_adapter}.py` · `api/run_router.py`) · `core/error_catalog.py` (NEW `ERR_EVAL_RUN_NOT_FOUND`; REUSE `ERR_ZDR_PAYLOAD_BLOCKED`, `ERR_EVAL_SET_NOT_FOUND`, and the governance error specs) · a new Alembic migration (eval_runs, eval_case_results — the THREE-plus manifests: migration + EXPECTED_TABLES + alembic env.py import + guardrails allow-list, see [[gateway-new-table-four-manifests]]) · `main.py` (mount router + side-effect ORM import) · `apps/gateway/tests/evals_runs/`. +regression floor: `make ci` green incl. the migration-parity gate and the guardrails no-new-tables allow-list; NO change to `governance.py`/`use_cases.py` behavior (the run REUSES them; if a seam needs a param, it is additive and behavior-preserving) or to any existing payload store. +resolved (was least-sure, confirmed at freeze 2026-08-13): a run executes as a BACKGROUND, DURABLE job — launch enqueues and returns `pending` fast; a worker drives cases and commits each result as it lands, mirroring the vector-store ingest worker's enqueue-or-fail-open idiom (fail-open to an inline drive when the queue/redis is absent, e.g. under ASGITransport tests). This is what makes M7 resume real across crash/redeploy. A tiny-set synchronous path is a later optimization, not the contract. ## EDGES -- E1 +- E1 A tenant over budget/credit launches a run of a non-empty set -> the run launches but EVERY case is `refused` at the governance guard with NO upstream call and NO usage row; assert zero provider dials (M3, R:GOVERNANCE_BYPASS). +- E2 Tenant A's run drives A's per-tenant breaker OPEN (repeated upstream failures); tenant B's request in the same process still succeeds -> isolation proven (M4, R:GLOBAL_BREAKER). +- E3 A ZDR tenant (flag set before launch, and flipped ON mid-run by a slow double) launches a run -> refused 403 atomic with the first result write; assert ZERO eval_case_results rows persisted (M5, R:ZDR_BLOCKED — the TOCTOU-safe pattern). +- E4 A run is interrupted after K of N cases commit, then resumed -> only the N-K `pending` cases are dialed; the K terminal cases are NOT re-dialed or re-billed; final usage_records count == dialed-case count exactly once (M7, R:DOUBLE_BILL). +- E5 An upstream call for one case times out -> that case is `errored` with a timeout reason and the run CONTINUES with the remaining cases (one bad case never sinks the run); the run's rollup reflects the mix (M4, A6). +- E6 A run of an empty set (or one whose every case is refused) -> `completed`, 0 dialed, empty/all-refused result list, no exception (A4). ## CHECKS -- · covers: M1 · -red-first: every check MUST fail first. +- test_run_enters_governance_per_case · covers: M1, M8, R:GOVERNANCE_BYPASS · a run of a set dials upstream only after the governance guard passes; one dial per case, driven entirely through a zero-network fake upstream (the Protocol port, M8), in the guard order a live request uses. +- test_over_budget_run_refuses_every_case_no_dial · covers: M3, E1, R:GOVERNANCE_BYPASS · a tenant forced over budget launches a run; every case is `refused`, the fake upstream records ZERO dials, and NO usage_records rows are written. +- test_one_usage_record_per_dialed_case · covers: M2, R:DOUBLE_BILL · a run of N dial-able cases writes exactly N usage_records, each on the served model id, matching the shape a live request of that request_body produces. +- test_per_tenant_breaker_isolation · covers: M4, E2, R:GLOBAL_BREAKER · tenant A's run drives A's breaker open (fake upstream fails); a concurrent tenant B request in the same process still succeeds — the breaker is keyed by tenant, not global. +- test_zdr_run_refused_atomically_zero_results · covers: M5, E3, R:ZDR_BLOCKED · a ZDR tenant's run (incl. a slow double flipping ZDR mid-run between the lock and the first result write) persists ZERO eval_case_results; asserted on the PERSISTED ROW count, not the response. +- test_resume_does_not_rebill_terminal_cases · covers: M7, E4, R:DOUBLE_BILL · a run interrupted after K/N committed results, then resumed, dials only the N-K pending cases; total upstream dials == N exactly once, terminal cases untouched. +- test_timeout_case_errored_run_continues · covers: M4, E5, A6 · one case's upstream dial times out -> that case is `errored` with a reason and the remaining cases still complete; the run rollup exposes a distinct, actionable per-case status (A6). +- test_cross_tenant_and_absent_run_uniform_404 · covers: M6, R:RUN_NOT_FOUND · a run/set owned by another tenant and an absent id both return byte-identical 404, no oracle. +- test_results_aligned_to_case_creation_order · covers: A5 · two runs of the same set expose their case results in identical (case-creation) order regardless of completion order. +- test_empty_set_run_completes_vacuously · covers: A4, E6 · running an empty set returns a `completed` run with an empty result list and no exception. +- test_cross_tenant_launch_and_billing_identity · covers: A1 · tenant B launching a run on tenant A's set gets a uniform 404; tenant A's own run bills the LAUNCHING tenant/key (the usage record carries A's tenant). +- test_run_snapshot_fixed_at_launch · covers: A2 · a case added AFTER a run launches does not change that run's case_count — the case set is snapshotted at launch time. +- test_bounded_per_tenant_concurrency · covers: A3 · a run of N>ceiling cases never has more than the per-tenant ceiling dialing at once (the tenant semaphore bounds in-flight dials). +red-first: every check MUST fail first (the run module + tables do not exist yet). ## EVIDENCE receipt: .md> diff --git a/apps/gateway/migrations/env.py b/apps/gateway/migrations/env.py index aadacc0f..5ade4a12 100644 --- a/apps/gateway/migrations/env.py +++ b/apps/gateway/migrations/env.py @@ -111,6 +111,7 @@ import gateway.domain_capture.infrastructure.orm # noqa: F401 — tenant_domain_claims import gateway.files.infrastructure.orm # noqa: F401 — files import gateway.evals.infrastructure.orm # noqa: F401 — eval_sets, eval_cases +import gateway.evals.runs.infrastructure.orm # noqa: F401 — eval_runs, eval_case_results import gateway.finetune.infrastructure.orm # noqa: F401 — finetune_jobs, finetune_job_events import gateway.guardrail_analytics.infrastructure.orm # noqa: F401 — guardrail_verdict_events import gateway.payments.infrastructure.orm # noqa: F401 — checkout_sessions diff --git a/apps/gateway/migrations/versions/f5b2d8c41a37_eval_run_executor.py b/apps/gateway/migrations/versions/f5b2d8c41a37_eval_run_executor.py new file mode 100644 index 00000000..b3572f56 --- /dev/null +++ b/apps/gateway/migrations/versions/f5b2d8c41a37_eval_run_executor.py @@ -0,0 +1,98 @@ +"""eval run + per-case result substrate — eval-run-executor (R7 evals-regression-gate). + +Revision ID: f5b2d8c41a37 +Revises: e4a1c9d27f60 +Create Date: 2026-08-13 + +eval-run-executor PLAN.md §3 (FROZEN @ sha256:1353206d). Additive, two tenant-scoped tables +that extend the eval-set-store substrate: + + - NEW TABLE eval_runs — one row per launched run. Carries the launching key_id (A1 — a run + bills that key, exactly as its live traffic) and the named model. ``status`` is DERIVED + from the run's cases (M7): pending at launch, completed when every snapshot case is + terminal, blocked iff a ZDR flip refused the run mid-flight (M5). FK -> eval_sets CASCADE. + ⚠ The raw API key is NEVER persisted (auth-scoped resume, 2026-08-13): only key_id, so a + cross-process resume must re-supply the raw key via a fresh authenticated request. + - NEW TABLE eval_case_results — one row per case DRIVEN. ``response_text`` is the model's + payload-at-rest (the ZDR-gated surface, M5); present only for a `completed` case. + UNIQUE (eval_run_id, eval_case_id) makes a resumed drive idempotent — a terminal case is + never re-dialed or re-billed (M7 / R:DOUBLE_BILL). FK -> eval_runs CASCADE. + +Four-manifest rule (see [[gateway-new-table-four-manifests]]): this migration + EXPECTED_TABLES +(tests/migrations) + migrations/env.py's own ORM import + the guardrails NOT-IN allow-list. + +Indexes (also declared in the ORM __table_args__ per the v30 two-manifest lesson): + ix_eval_runs_tenant_set_created on eval_runs(tenant_id, eval_set_id, created_at) + ix_eval_case_results_run_created on eval_case_results(tenant_id, eval_run_id, created_at) +""" + +from __future__ import annotations + +import sqlalchemy as sa +from alembic import op + +# revision identifiers, used by Alembic. +revision = "f5b2d8c41a37" +down_revision = "e4a1c9d27f60" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.create_table( + "eval_runs", + sa.Column("id", sa.UUID(), nullable=False, server_default=sa.text("gen_random_uuid()")), + sa.Column("tenant_id", sa.UUID(), nullable=False), + sa.Column("eval_set_id", sa.UUID(), nullable=False), + sa.Column("key_id", sa.UUID(), nullable=False), + sa.Column("model", sa.Text(), nullable=False), + sa.Column("status", sa.Text(), nullable=False), + sa.Column( + "created_at", + sa.DateTime(timezone=True), + nullable=False, + server_default=sa.text("now()"), + ), + sa.PrimaryKeyConstraint("id"), + sa.ForeignKeyConstraint(["eval_set_id"], ["eval_sets.id"], ondelete="CASCADE"), + ) + op.create_index( + "ix_eval_runs_tenant_set_created", + "eval_runs", + ["tenant_id", "eval_set_id", "created_at"], + ) + + op.create_table( + "eval_case_results", + sa.Column("id", sa.UUID(), nullable=False, server_default=sa.text("gen_random_uuid()")), + sa.Column("tenant_id", sa.UUID(), nullable=False), + sa.Column("eval_run_id", sa.UUID(), nullable=False), + sa.Column("eval_case_id", sa.UUID(), nullable=False), + sa.Column("status", sa.Text(), nullable=False), + sa.Column("response_text", sa.Text(), nullable=True), + sa.Column("reason", sa.Text(), nullable=True), + sa.Column("usage_record_id", sa.UUID(), nullable=True), + sa.Column( + "created_at", + sa.DateTime(timezone=True), + nullable=False, + server_default=sa.text("now()"), + ), + sa.PrimaryKeyConstraint("id"), + sa.ForeignKeyConstraint(["eval_run_id"], ["eval_runs.id"], ondelete="CASCADE"), + sa.UniqueConstraint( + "eval_run_id", "eval_case_id", name="uq_eval_case_results_run_case" + ), + ) + op.create_index( + "ix_eval_case_results_run_created", + "eval_case_results", + ["tenant_id", "eval_run_id", "created_at"], + ) + + +def downgrade() -> None: + op.drop_index("ix_eval_case_results_run_created", table_name="eval_case_results") + op.drop_table("eval_case_results") + op.drop_index("ix_eval_runs_tenant_set_created", table_name="eval_runs") + op.drop_table("eval_runs") diff --git a/apps/gateway/src/gateway/core/error_catalog.py b/apps/gateway/src/gateway/core/error_catalog.py index a9c629df..540c4cf1 100644 --- a/apps/gateway/src/gateway/core/error_catalog.py +++ b/apps/gateway/src/gateway/core/error_catalog.py @@ -1689,3 +1689,12 @@ def exc( "ERR_EVAL_SET_NAME_CONFLICT", "an eval set with this name already exists for this tenant", ) + +#: GET /v1/evals/runs/{id} (or .../cases) whose run id is absent OR owned by another tenant +#: — deliberately indistinguishable (M6, R:RUN_NOT_FOUND; never an enumeration oracle), +#: byte-identical for both causes, exactly like EVAL_SET_NOT_FOUND above. +EVAL_RUN_NOT_FOUND = ErrorSpec(404, "ERR_EVAL_RUN_NOT_FOUND", "Eval run not found") + +#: POST .../runs whose `model` is missing or not a non-empty string — 422, nothing launched. +#: Validated on the caller's own input BEFORE the parent set is resolved (no oracle). +EVAL_RUN_INVALID = ErrorSpec(422, "ERR_EVAL_RUN_INVALID", "model must be a non-empty string") diff --git a/apps/gateway/src/gateway/evals/runs/__init__.py b/apps/gateway/src/gateway/evals/runs/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/apps/gateway/src/gateway/evals/runs/api/__init__.py b/apps/gateway/src/gateway/evals/runs/api/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/apps/gateway/src/gateway/evals/runs/api/run_router.py b/apps/gateway/src/gateway/evals/runs/api/run_router.py new file mode 100644 index 00000000..ce3db9dd --- /dev/null +++ b/apps/gateway/src/gateway/evals/runs/api/run_router.py @@ -0,0 +1,212 @@ +"""FastAPI router for eval-run execution (/v1/evals/.../runs) — eval-run-executor §3. + +Reuses the eval-set-store surface's auth + OpenAI-wire envelope wholesale (``_authenticate``, +``_extract_raw_key``, ``_err``, ``_err_from_problem``, ``_unix`` from +``gateway.evals.api.router``) so the whole /v1/evals surface speaks ONE ``{"error": {...}}`` +body, and tenant scope is enforced identically. + +Endpoints: + POST /v1/evals/sets/{set_id}/runs {model} -> 201 { id:"er_..", eval_set_id, model, status, + case_count, created_at } + 403 ERR_ZDR_PAYLOAD_BLOCKED (ZDR tenant, refused outright at launch, M5) + 404 ERR_EVAL_SET_NOT_FOUND (absent OR cross-tenant set, M6) + 422 ERR_EVAL_RUN_INVALID (missing/blank model, validated before the set is resolved) + GET /v1/evals/runs/{run_id} -> 200 { id, eval_set_id, model, status, case_count, + counts:{completed,refused,errored,pending}, created_at } + GET /v1/evals/runs/{run_id}/cases -> 200 { object:"list", data:[{eval_case_id, status, + response_text?, reason?, usage_record_id?}] } (A5 order) + +Durability (M7): launch enqueues onto ``app.state.eval_run_queue`` when present and returns a +``pending`` run fast; a background worker drives it. When the queue is absent or enqueue fails +(Redis down, or ASGITransport tests), it FAILS OPEN to an inline drive — the vector-store +ingest idiom — so a run is never dropped. The raw key is passed straight through to the inline +drive; it is never persisted (auth-scoped resume). +""" + +from __future__ import annotations + +import logging +from typing import Annotated, Any + +from fastapi import APIRouter, Depends, Request +from fastapi.responses import JSONResponse +from sqlalchemy.ext.asyncio import AsyncSession + +from gateway.core.db import get_session +from gateway.core.error_catalog import EVAL_RUN_INVALID, EVAL_RUN_NOT_FOUND, EVAL_SET_NOT_FOUND +from gateway.core.errors import ProblemError + +# The eval-set-store router owns the /v1/evals auth + OpenAI-wire envelope; reusing its helpers +# is what keeps the WHOLE surface speaking one body (a fork would drift). Same sanctioned reuse +# pattern as images/audio/embeddings use_cases importing use_cases._fire_record_with_raw. +from gateway.evals.api.router import ( + _authenticate, # pyright: ignore[reportPrivateUsage] + _err, # pyright: ignore[reportPrivateUsage] + _err_from_problem, # pyright: ignore[reportPrivateUsage] + _extract_raw_key, # pyright: ignore[reportPrivateUsage] + _unix, # pyright: ignore[reportPrivateUsage] +) +from gateway.evals.infrastructure.repository import SqlAlchemyEvalStore +from gateway.evals.runs.infrastructure.orm import EvalCaseResultRow, EvalRunRow +from gateway.evals.runs.infrastructure.repository import SqlAlchemyEvalRunStore +from gateway.evals.wire_id import ( + parse_run_wire_id, + parse_set_wire_id, + to_case_wire_id, + to_run_wire_id, + to_set_wire_id, +) +from gateway.keys.domain.entities import AuthzResult + +eval_runs_router = APIRouter(tags=["evals"]) + +_log = logging.getLogger(__name__) + + +def _run_store(request: Request) -> SqlAlchemyEvalRunStore: + return SqlAlchemyEvalRunStore(request.app.state.sessionmaker) + + +def _run_object(row: EvalRunRow, *, case_count: int) -> dict[str, Any]: + return { + "id": to_run_wire_id(row.id), + "object": "eval.run", + "created_at": _unix(row.created_at), + "eval_set_id": to_set_wire_id(row.eval_set_id), + "model": row.model, + "status": row.status, + "case_count": case_count, + } + + +def _case_result_object(row: EvalCaseResultRow) -> dict[str, Any]: + """The per-case result wire object. response_text is present only for a `completed` case.""" + obj: dict[str, Any] = { + "object": "eval.case_result", + "eval_case_id": to_case_wire_id(row.eval_case_id), + "status": row.status, + } + if row.response_text is not None: + obj["response_text"] = row.response_text + if row.reason is not None: + obj["reason"] = row.reason + if row.usage_record_id is not None: + obj["usage_record_id"] = str(row.usage_record_id) + return obj + + +@eval_runs_router.post( + "/v1/evals/sets/{set_id}/runs", status_code=201, response_model=None +) +async def launch_eval_run( + set_id: str, + body: dict[str, Any], + request: Request, + authz: Annotated[AuthzResult, Depends(_authenticate)], + session: Annotated[AsyncSession, Depends(get_session)], +) -> dict[str, Any] | JSONResponse: + """Launch a run of a tenant's set against a model (M1). ZDR tenant refused outright (M5).""" + # Validate the caller's own input FIRST — reveals nothing about whether the set exists. + model = body.get("model") + if not isinstance(model, str) or not model.strip(): + return _err(EVAL_RUN_INVALID) + + resolved_set = parse_set_wire_id(set_id) + if resolved_set is None: + return _err(EVAL_SET_NOT_FOUND) + # M6: resolve the parent set in tenant scope — absent/cross-tenant is a uniform 404. + parent = await SqlAlchemyEvalStore(session).get_set( + tenant_id=authz.tenant_id, eval_set_id=resolved_set + ) + if parent is None: + return _err(EVAL_SET_NOT_FOUND) + + executor = request.app.state.eval_run_executor + raw_key = _extract_raw_key(request) + try: + run = await executor.launch( + tenant_id=authz.tenant_id, + key_id=authz.key_id, + raw_key=raw_key, + eval_set_id=resolved_set, + model=model, + ) + except ProblemError as exc: + # M5: the ZDR gate refused the run outright at launch. Re-render in this surface's + # envelope (403 ERR_ZDR_PAYLOAD_BLOCKED); nothing was created. + return _err_from_problem(exc) + + # Snapshot denominator (A2): cases that existed at launch time. + store = _run_store(request) + snapshot = await store.snapshot_cases( + tenant_id=run.tenant_id, eval_set_id=run.eval_set_id, created_at_max=run.created_at + ) + case_count = len(snapshot) + + await _enqueue_or_drive(request, executor, run_id=run.id, raw_key=raw_key) + + # Re-read so the returned status reflects an inline drive (fail-open / test) if it ran. + refreshed = await store.get_run(tenant_id=run.tenant_id, run_id=run.id) or run + return _run_object(refreshed, case_count=case_count) + + +async def _enqueue_or_drive( + request: Request, executor: Any, *, run_id: Any, raw_key: str +) -> None: + """Enqueue for the durable worker; FAIL OPEN to an inline drive (vector-store idiom, M7).""" + queue = getattr(request.app.state, "eval_run_queue", None) + if queue is not None: + try: + await queue.enqueue(run_id) + return + except Exception: + _log.warning( + "eval_run: enqueue failed for run %s, failing open to inline drive", run_id + ) + await executor.drive(run_id, raw_key=raw_key) + + +@eval_runs_router.get("/v1/evals/runs/{run_id}", status_code=200, response_model=None) +async def get_eval_run( + run_id: str, + request: Request, + authz: Annotated[AuthzResult, Depends(_authenticate)], +) -> dict[str, Any] | JSONResponse: + """A run's status + per-status rollup (M6-scoped; uniform 404 for absent/cross-tenant).""" + resolved = parse_run_wire_id(run_id) + if resolved is None: + return _err(EVAL_RUN_NOT_FOUND) + store = _run_store(request) + run = await store.get_run(tenant_id=authz.tenant_id, run_id=resolved) + if run is None: + return _err(EVAL_RUN_NOT_FOUND) + + snapshot = await store.snapshot_cases( + tenant_id=run.tenant_id, eval_set_id=run.eval_set_id, created_at_max=run.created_at + ) + case_count = len(snapshot) + counts = await store.counts_by_status(run.id) + terminal = counts["completed"] + counts["refused"] + counts["errored"] + counts["pending"] = max(0, case_count - terminal) + + obj = _run_object(run, case_count=case_count) + obj["counts"] = counts + return obj + + +@eval_runs_router.get("/v1/evals/runs/{run_id}/cases", status_code=200, response_model=None) +async def list_eval_run_cases( + run_id: str, + request: Request, + authz: Annotated[AuthzResult, Depends(_authenticate)], +) -> dict[str, Any] | JSONResponse: + """A run's per-case results in the set's creation order (A5). Uniform 404 (M6).""" + resolved = parse_run_wire_id(run_id) + if resolved is None: + return _err(EVAL_RUN_NOT_FOUND) + store = _run_store(request) + run = await store.get_run(tenant_id=authz.tenant_id, run_id=resolved) + if run is None: + return _err(EVAL_RUN_NOT_FOUND) + results = await store.list_case_results(tenant_id=authz.tenant_id, run_id=resolved) + return {"object": "list", "data": [_case_result_object(r) for r in results]} diff --git a/apps/gateway/src/gateway/evals/runs/application/__init__.py b/apps/gateway/src/gateway/evals/runs/application/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/apps/gateway/src/gateway/evals/runs/application/run_executor.py b/apps/gateway/src/gateway/evals/runs/application/run_executor.py new file mode 100644 index 00000000..244749d1 --- /dev/null +++ b/apps/gateway/src/gateway/evals/runs/application/run_executor.py @@ -0,0 +1,252 @@ +"""EvalRunExecutor — replay a set through the governance path (eval-run-executor §3, M1-M8). + +A run REPLAYS each case's stored ``request_body`` against the run's named model THROUGH the +existing ``CompletionUseCase.complete`` — never around it (M1). That single reuse makes the +security-critical invariants hold BY CONSTRUCTION: + + M1 governance — complete() runs the SAME ordered guards a live request runs (auth · expiry · + allowlist · catalog/model_checker · per-key/team/tenant budget · credit · + tier · rate-limit). A run enters no back door. + M2 billing — complete() writes ONE usage_records row per dial, on the SERVED model id, + the SAME shape a live request produces. A refused case never dials (M3). + M4 isolation — the executor injects its OWN ``TenantBreakerUpstream`` into + complete(upstream=...): a per-tenant breaker + concurrency + per-call + timeout, so a run trips only THAT tenant's breaker, never the global one + the live path uses ([[per-tenant-breaker-recurring-defect]], R:GLOBAL_BREAKER). + +The executor owns durability (M7): each case result is committed as it lands, and a resumed +drive skips cases that already have a terminal result row — no case is dialed or billed twice. + +SECURITY — auth-scoped resume (2026-08-13): complete() authenticates a RAW key, which is NEVER +persisted at rest (a gateway key is high-value secret material). The launching request holds it +in memory for the run's lifetime (``_raw_keys``), so an in-process crash resumes via +recover_orphans. A cross-process (redeploy) resume must RE-SUPPLY the raw key via a fresh +authenticated request — ``drive(run_id, raw_key=...)``. Without it, pending cases stay pending +rather than being dialed under a forged identity. + +ZDR (M5): a ZDR tenant's run is refused OUTRIGHT at launch (403), and the first payload write +re-checks ATOMICALLY (``raise_if_zdr_locked`` inside record_case_result) — a flip landing +mid-run persists nothing further and the run is marked ``blocked``. +""" + +from __future__ import annotations + +import asyncio +import logging +import uuid +from collections.abc import Callable +from typing import Any + +from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker + +from gateway.core.errors import ProblemError +from gateway.evals.infrastructure.orm import EvalCaseRow +from gateway.evals.runs.domain.entities import CaseOutcome +from gateway.evals.runs.domain.ports import EvalRunStore +from gateway.evals.runs.infrastructure.orm import EvalRunRow +from gateway.evals.runs.infrastructure.upstream_adapter import ( + TenantBreakerUpstream, + TenantExecutionRegistry, +) +from gateway.proxy.domain.ports import CompletionUpstream, UsageRecorder +from gateway.tenants.application.retention_policy import raise_if_zdr + +_log = logging.getLogger(__name__) + +# ProblemError codes that mean a DIAL WAS ATTEMPTED and failed (→ `errored`, fail-closed). +# Everything else complete() raises is a pre-dial governance refusal (→ `refused`, no dial). +_ERRORED_CODES = frozenset({"ERR_UPSTREAM_UNAVAILABLE", "ERR_UPSTREAM_RATE_LIMITED"}) + +#: A run's build_use_case seam: given a fresh session, return a live-wired CompletionUseCase +#: (main.py supplies the same get_completion_use_case the request path uses, over a shim). +UseCaseFactory = Callable[[AsyncSession], Any] + + +class EvalRunExecutor: + """Launches + drives eval runs. One instance on app.state, shared by the router + worker.""" + + def __init__( + self, + *, + sessionmaker: async_sessionmaker[AsyncSession], + store: EvalRunStore, + build_use_case: UseCaseFactory, + get_upstream_delegate: Callable[[], CompletionUpstream], + get_usage_recorder: Callable[[], UsageRecorder], + get_model_router: Callable[[], Any], + registry: TenantExecutionRegistry, + per_call_timeout_seconds: float = 30.0, + ) -> None: + self._sessionmaker = sessionmaker + self._store = store + self._build_use_case = build_use_case + self._get_upstream_delegate = get_upstream_delegate + self._get_usage_recorder = get_usage_recorder + self._get_model_router = get_model_router + self._registry = registry + self._timeout = per_call_timeout_seconds + # In-memory raw-key map — the ONLY place the launching secret lives (never at rest). + self._raw_keys: dict[uuid.UUID, str] = {} + + async def launch( + self, + *, + tenant_id: uuid.UUID, + key_id: uuid.UUID, + raw_key: str, + eval_set_id: uuid.UUID, + model: str, + ) -> EvalRunRow: + """Create a pending run for the tenant. Refuses a ZDR tenant OUTRIGHT at launch (M5).""" + # M5: ZDR refuse outright at launch — the ENTRY check (the atomic locked re-check lives + # at the first payload write). Raises ProblemError 403 ERR_ZDR_PAYLOAD_BLOCKED. + async with self._sessionmaker() as session: + await raise_if_zdr(session, tenant_id) + run = await self._store.create_run( + tenant_id=tenant_id, key_id=key_id, eval_set_id=eval_set_id, model=model + ) + self._raw_keys[run.id] = raw_key + return run + + async def drive(self, run_id: uuid.UUID, *, raw_key: str | None = None) -> None: + """Drive (or resume) a run: dial every still-pending snapshot case, commit as each lands. + + Bounded per-tenant concurrency (A3): the tenant semaphore caps in-flight dials. A5's + result ordering is enforced at READ time, so commit order does not matter. A ZDR flip + mid-run (M5) aborts further persistence and marks the run blocked. + """ + run = await self._store.load_run(run_id) + if run is None: + _log.warning("eval_run_executor: run %s not found (stale enqueue?), skipping", run_id) + return + key = raw_key or self._raw_keys.get(run_id) + + cases = await self._store.snapshot_cases( + tenant_id=run.tenant_id, + eval_set_id=run.eval_set_id, + created_at_max=run.created_at, + ) + existing = await self._store.existing_result_case_ids(run_id) + pending = [c for c in cases if c.id not in existing] + + if pending and key is None: + # Cross-process resume without a re-supplied key: we cannot re-enter governance under + # the launching identity, and we will NOT forge one. Leave the cases pending for an + # authenticated resume; never dial under a synthetic key. + _log.info( + "eval_run_executor: run %s has %d pending cases but no raw key in memory;" + " awaiting an authenticated resume", + run_id, + len(pending), + ) + return + + blocked = _BlockedFlag() + semaphore = self._registry.semaphore_for(run.tenant_id) + + async def _one(case: EvalCaseRow) -> None: + if blocked.is_set: + return + async with semaphore: # A3: never more than `ceiling` dials in flight per tenant + if blocked.is_set: + return + outcome = await self._drive_case(run, case, key) + try: + await self._store.record_case_result( + tenant_id=run.tenant_id, + run_id=run_id, + eval_case_id=case.id, + outcome=outcome, + ) + except ProblemError: + # M5: the atomic locked ZDR re-check refused this payload write — a flip landed + # mid-run. Persist nothing further; the run is blocked. + blocked.set() + + await asyncio.gather(*(_one(c) for c in pending)) + + if blocked.is_set: + await self._store.set_run_status(run_id, "blocked") + return + await self._finalize_status(run_id, total=len(cases)) + + async def resume(self, run_id: uuid.UUID, *, raw_key: str) -> None: + """Authenticated cross-process resume: re-supply the raw key and drive pending cases.""" + await self.drive(run_id, raw_key=raw_key) + + async def _finalize_status(self, run_id: uuid.UUID, *, total: int) -> None: + """Derive the run's terminal status from its cases (M7) — never speculative. + + completed when every snapshot case has a terminal result row (incl. the vacuous + empty/all-refused run, A4). Otherwise it stays pending (a resume will finish it). + """ + counts = await self._store.counts_by_status(run_id) + terminal = counts["completed"] + counts["refused"] + counts["errored"] + if terminal >= total: + await self._store.set_run_status(run_id, "completed") + + async def _drive_case( + self, run: EvalRunRow, case: EvalCaseRow, raw_key: str | None + ) -> CaseOutcome: + """Replay one case through complete(); map its outcome to a CaseOutcome (M1-M4).""" + # M1: replay the stored request_body against the RUN's named model (override model). + body = {**dict(case.request_body), "model": run.model} + delegate = self._get_upstream_delegate() + breaker = self._registry.breaker_for(run.tenant_id) + upstream = TenantBreakerUpstream(breaker, delegate, timeout_seconds=self._timeout) + try: + async with self._sessionmaker() as session: + use_case = self._build_use_case(session) + status, response_body, _ = await use_case.complete( + raw_key=raw_key, + body=body, + upstream=upstream, + usage_recorder=self._get_usage_recorder(), + model_router=self._get_model_router(), + ) + except ProblemError as exc: + if exc.code in _ERRORED_CODES: + # A dial was attempted and failed (breaker-open / timeout / upstream 5xx / 429) + # → fail closed, the run continues (M4, E5). + return CaseOutcome(status="errored", reason=exc.code) + # A pre-dial governance refusal (budget/credit/tier/rate/allowlist/catalog/…): NO + # dial, NO usage row for this case (M3, R:GOVERNANCE_BYPASS). + return CaseOutcome(status="refused", reason=exc.code) + except Exception: + _log.exception("eval_run_executor: unexpected error driving case %s", case.id) + return CaseOutcome(status="errored", reason="internal error") + + if status == 200 and isinstance(response_body, dict): + return CaseOutcome(status="completed", response_text=_extract_text(response_body)) + # An upstream 4xx pass-through — the provider answered the caller got something wrong. + return CaseOutcome(status="errored", reason=f"upstream returned {status}") + + +class _BlockedFlag: + """A tiny shared abort flag for the ZDR mid-run stop (avoids a nonlocal in the closure).""" + + __slots__ = ("is_set",) + + def __init__(self) -> None: + self.is_set = False + + def set(self) -> None: + self.is_set = True + + +def _extract_text(response_body: dict[str, Any]) -> str: + """Pull the assistant text (choices[0].message.content). '' if the shape is unexpected.""" + try: + choices = response_body.get("choices") + if isinstance(choices, list) and choices: + message = choices[0].get("message") + if isinstance(message, dict): + content = message.get("content") + if isinstance(content, str): + return content + except Exception: + return "" + return "" + + +__all__ = ["EvalRunExecutor", "UseCaseFactory"] diff --git a/apps/gateway/src/gateway/evals/runs/domain/__init__.py b/apps/gateway/src/gateway/evals/runs/domain/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/apps/gateway/src/gateway/evals/runs/domain/entities.py b/apps/gateway/src/gateway/evals/runs/domain/entities.py new file mode 100644 index 00000000..4ec203da --- /dev/null +++ b/apps/gateway/src/gateway/evals/runs/domain/entities.py @@ -0,0 +1,35 @@ +"""Value types for eval-run execution (eval-run-executor §3). + +Small, immutable carriers the executor and router pass around. The persisted shapes live in +``infrastructure/orm.py``; these are the in-memory outcome of driving one case, kept free of +any SQLAlchemy / provider import so the executor's Protocol port (M8) can be exercised with a +zero-network fake. +""" + +from __future__ import annotations + +from dataclasses import dataclass +from typing import Literal + +#: A run's status is DERIVED from its cases (M7), never written speculatively: +#: pending — created, not every snapshot case is terminal yet +#: completed — every snapshot case has a terminal result row (incl. the vacuous empty run, A4) +#: blocked — a ZDR flip refused the run mid-flight; nothing further persisted (M5) +RunStatus = Literal["pending", "completed", "blocked"] + +#: A driven case's terminal status: +#: completed — dialed, upstream 200; carries response_text (the ZDR-gated payload) + a usage row +#: refused — governance denied BEFORE any dial (M3); carries a reason, NO usage row, NO payload +#: errored — breaker-open / per-call timeout / upstream 5xx (M4); carries a reason, fail-closed +CaseStatus = Literal["completed", "refused", "errored"] + + +@dataclass(frozen=True) +class CaseOutcome: + """The result of driving ONE case — what the executor commits as an eval_case_results row.""" + + status: CaseStatus + #: The assistant text (choices[0].message.content). Present only on ``completed``. + response_text: str | None = None + #: A short, payload-free reason for a refused/errored case (A6). None on completed. + reason: str | None = None diff --git a/apps/gateway/src/gateway/evals/runs/domain/errors.py b/apps/gateway/src/gateway/evals/runs/domain/errors.py new file mode 100644 index 00000000..4391ba4c --- /dev/null +++ b/apps/gateway/src/gateway/evals/runs/domain/errors.py @@ -0,0 +1,15 @@ +"""Domain errors for eval-run execution (eval-run-executor §3). + +Named errors the application/repository raise and the router translates to the wire — never +a bare Exception or an inlined HTTP status in the use-case (appsec-engineer lens). +""" + +from __future__ import annotations + + +class EvalRunNotFound(Exception): + """The named run (or its parent set) is absent OR owned by another tenant — indistinguishable. + + The router maps this to a uniform 404 ERR_EVAL_RUN_NOT_FOUND / ERR_EVAL_SET_NOT_FOUND, + identical for both causes, so there is no enumeration oracle (M6, R:RUN_NOT_FOUND). + """ diff --git a/apps/gateway/src/gateway/evals/runs/domain/ports.py b/apps/gateway/src/gateway/evals/runs/domain/ports.py new file mode 100644 index 00000000..bd1230d3 --- /dev/null +++ b/apps/gateway/src/gateway/evals/runs/domain/ports.py @@ -0,0 +1,87 @@ +"""Persistence port for eval-run execution — eval-run-executor §3 (M8, appsec-engineer lens). + +The seam the executor + router depend on, so neither touches SQLAlchemy directly and the +executor's behavior (governance / breaker / ZDR / resume) is provable over a fake store + a +fake upstream with no network. The SQLAlchemy adapter lives in +``infrastructure/repository.py``. + +INVARIANT (M6): every tenant-facing read is tenant-scoped in the SAME query that resolves the +row — an absent or cross-tenant run is None/[] (the router maps that to a uniform 404). The +worker-facing ``load_run`` is the ONE un-scoped read: it trusts a run id claimed from the +internal durable queue and re-derives the tenant from the row (never from caller input). +""" + +from __future__ import annotations + +import uuid +from typing import Protocol + +from gateway.evals.infrastructure.orm import EvalCaseRow +from gateway.evals.runs.domain.entities import CaseOutcome +from gateway.evals.runs.infrastructure.orm import EvalCaseResultRow, EvalRunRow + + +class EvalRunStore(Protocol): + """Tenant-scoped persistence for eval runs and their per-case results.""" + + async def create_run( + self, *, tenant_id: uuid.UUID, key_id: uuid.UUID, eval_set_id: uuid.UUID, model: str + ) -> EvalRunRow: + """Insert a pending run for the tenant and return it (raw key is NEVER persisted).""" + ... + + async def get_run(self, *, tenant_id: uuid.UUID, run_id: uuid.UUID) -> EvalRunRow | None: + """Resolve a run owned by this tenant. None for absent/cross-tenant (uniform, M6).""" + ... + + async def load_run(self, run_id: uuid.UUID) -> EvalRunRow | None: + """Worker-facing: load a run by id WITHOUT a tenant filter (id came from the queue).""" + ... + + async def snapshot_cases( + self, *, tenant_id: uuid.UUID, eval_set_id: uuid.UUID, created_at_max: object + ) -> list[EvalCaseRow]: + """The set's cases at launch time (created_at <= the run's created_at), creation order. + + A2: a case added AFTER the run launched is NOT in this snapshot, so the run's + denominator is fixed. A5: creation order (created_at, id) so two runs of the same set + align case-for-case. + """ + ... + + async def existing_result_case_ids(self, run_id: uuid.UUID) -> set[uuid.UUID]: + """The set of eval_case_ids that already have a terminal result row (resume skip, M7).""" + ... + + async def record_case_result( + self, + *, + tenant_id: uuid.UUID, + run_id: uuid.UUID, + eval_case_id: uuid.UUID, + outcome: CaseOutcome, + ) -> bool: + """Commit ONE result row, idempotently (M7 / R:DOUBLE_BILL). + + For a ``completed`` (payload-bearing) outcome the ZDR gate is enforced ATOMICALLY with + the insert (``raise_if_zdr_locked``, SELECT … FOR UPDATE) — a flip landing mid-run + raises ``ProblemError`` and NOTHING is persisted for that case (M5). Returns False (a + no-op) if a result for (run, case) already exists — a resumed/raced drive never + double-writes. Raises ``ProblemError`` on a ZDR refusal so the executor can mark the + run blocked. + """ + ... + + async def counts_by_status(self, run_id: uuid.UUID) -> dict[str, int]: + """{completed, refused, errored} result counts for the run (rollup + status derivation).""" + ... + + async def set_run_status(self, run_id: uuid.UUID, status: str) -> None: + """Persist the run's DERIVED terminal status (M7) — never written speculatively.""" + ... + + async def list_case_results( + self, *, tenant_id: uuid.UUID, run_id: uuid.UUID + ) -> list[EvalCaseResultRow]: + """A run's per-case results in the set's creation order (A5), tenant-scoped (M6).""" + ... diff --git a/apps/gateway/src/gateway/evals/runs/infrastructure/__init__.py b/apps/gateway/src/gateway/evals/runs/infrastructure/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/apps/gateway/src/gateway/evals/runs/infrastructure/orm.py b/apps/gateway/src/gateway/evals/runs/infrastructure/orm.py new file mode 100644 index 00000000..4a1067e6 --- /dev/null +++ b/apps/gateway/src/gateway/evals/runs/infrastructure/orm.py @@ -0,0 +1,115 @@ +"""ORM models for eval-run execution (eval-run-executor PLAN.md §3 — FROZEN @ sha256:1353206d). + +Two tenant-scoped tables, both children of the eval-set-store substrate: + + eval_runs — one row per launched run: the launching key (A1 — a run bills the + launching key exactly as that key's live traffic), the named model, + and a status DERIVED from its cases (M7), never written speculatively. + eval_case_results — one row per case DRIVEN in a run. ``response_text`` is the model's + payload-at-rest (the console diff + re-scoring read it); it is the + ZDR-gated surface (M5). UNIQUE (eval_run_id, eval_case_id) makes a + re-drive on resume idempotent — a case already terminal is never + dialed or billed twice (M7 / R:DOUBLE_BILL). + +Indexes are declared in BOTH __table_args__ AND the migration (the v30 two-manifest lesson). +All rows subclass gateway.core.db.Base; the side-effect import in main.py + migrations/env.py +registers them on Base.metadata (the four-manifest rule — see the migration's header). +""" + +from __future__ import annotations + +import uuid +from datetime import datetime + +from sqlalchemy import DateTime, ForeignKey, Index, Text, UniqueConstraint, func +from sqlalchemy.dialects.postgresql import UUID +from sqlalchemy.orm import Mapped, mapped_column + +from gateway.core.db import Base + + +class EvalRunRow(Base): + __tablename__ = "eval_runs" + + id: Mapped[uuid.UUID] = mapped_column( + "id", + UUID(as_uuid=True), + primary_key=True, + default=uuid.uuid4, + server_default=func.gen_random_uuid(), + ) + tenant_id: Mapped[uuid.UUID] = mapped_column("tenant_id", UUID(as_uuid=True), nullable=False) + eval_set_id: Mapped[uuid.UUID] = mapped_column( + "eval_set_id", + UUID(as_uuid=True), + ForeignKey("eval_sets.id", ondelete="CASCADE"), + nullable=False, + ) + # A1: the key that launched the run — the run bills THIS key, exactly as its live traffic. + # NOT a secret: the raw key is NEVER persisted (auth-scoped resume, 2026-08-13) — only the + # key_id, so a resume must re-supply the raw key via a fresh authenticated request. + key_id: Mapped[uuid.UUID] = mapped_column("key_id", UUID(as_uuid=True), nullable=False) + model: Mapped[str] = mapped_column(Text, nullable=False) + # status ∈ pending | completed | blocked. DERIVED from cases at drive end (M7): completed + # when every snapshot case has a terminal result row; blocked iff a ZDR flip refused the run + # mid-flight (M5). Never written speculatively. + status: Mapped[str] = mapped_column(Text, nullable=False, default="pending") + created_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), + nullable=False, + server_default=func.now(), + ) + + __table_args__ = ( + Index("ix_eval_runs_tenant_set_created", "tenant_id", "eval_set_id", "created_at"), + ) + + +class EvalCaseResultRow(Base): + __tablename__ = "eval_case_results" + + id: Mapped[uuid.UUID] = mapped_column( + "id", + UUID(as_uuid=True), + primary_key=True, + default=uuid.uuid4, + server_default=func.gen_random_uuid(), + ) + tenant_id: Mapped[uuid.UUID] = mapped_column("tenant_id", UUID(as_uuid=True), nullable=False) + eval_run_id: Mapped[uuid.UUID] = mapped_column( + "eval_run_id", + UUID(as_uuid=True), + ForeignKey("eval_runs.id", ondelete="CASCADE"), + nullable=False, + ) + eval_case_id: Mapped[uuid.UUID] = mapped_column( + "eval_case_id", UUID(as_uuid=True), nullable=False + ) + # status ∈ completed | refused | errored. completed = dialed 200 (carries response_text); + # refused = governance denied before any dial (M3, carries reason, NO response_text, NO + # usage row); errored = breaker-open / timeout / upstream 5xx (M4, carries reason). + status: Mapped[str] = mapped_column(Text, nullable=False) + # The ZDR-gated payload-at-rest (M5). Present only for a `completed` case; a ZDR tenant's + # run never reaches this write (refused outright at launch / atomic re-check mid-run). + response_text: Mapped[str | None] = mapped_column(Text, nullable=True) + # An actionable, payload-free reason for a refused/errored case (A6). None on completed. + reason: Mapped[str | None] = mapped_column(Text, nullable=True) + # Best-effort back-reference to the usage_records row this dial produced (M2). Nullable: + # the completion path bills fire-and-forget and does not return the id synchronously, so a + # v1 run leaves this null and proves "exactly one usage_record per dialed case" by COUNTING + # usage rows, not by this FK. + usage_record_id: Mapped[uuid.UUID | None] = mapped_column( + "usage_record_id", UUID(as_uuid=True), nullable=True + ) + created_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), + nullable=False, + server_default=func.now(), + ) + + __table_args__ = ( + # M7 / R:DOUBLE_BILL: one result per (run, case) — a resumed drive that re-visits a + # terminal case hits this constraint (or the skip-existing guard) and never re-dials. + UniqueConstraint("eval_run_id", "eval_case_id", name="uq_eval_case_results_run_case"), + Index("ix_eval_case_results_run_created", "tenant_id", "eval_run_id", "created_at"), + ) diff --git a/apps/gateway/src/gateway/evals/runs/infrastructure/repository.py b/apps/gateway/src/gateway/evals/runs/infrastructure/repository.py new file mode 100644 index 00000000..18ec17fc --- /dev/null +++ b/apps/gateway/src/gateway/evals/runs/infrastructure/repository.py @@ -0,0 +1,175 @@ +"""SQLAlchemy adapter for the EvalRunStore port (eval-run-executor PLAN.md §3, M6/M7). + +Sessionmaker-based (NOT a single injected session) because its primary caller is the durable +background worker, which — like ``VectorStoreIngestWorker`` — opens a FRESH short-lived +session per operation and commits each case result as it lands (durability, M7). The API +router reuses the same store for its tenant-scoped reads. + +INVARIANT (M6): every tenant-facing read filters on ``tenant_id`` in the resolving query; an +absent or cross-tenant id returns None/[]. ``load_run`` is the ONE un-scoped read — it trusts +an id claimed from the internal queue and re-derives the tenant from the row. + +``record_case_result`` is the payload write for a ``completed`` case: it calls +``raise_if_zdr_locked`` (SELECT … FOR UPDATE) as the FIRST statement in the committing +transaction, so the ZDR decision and the payload write are atomic (M5) — a flip landing +mid-run persists nothing. UNIQUE(eval_run_id, eval_case_id) + a pre-insert existence check +make a resumed/raced drive a no-op (M7 / R:DOUBLE_BILL). +""" + +from __future__ import annotations + +import uuid + +from sqlalchemy import func, select +from sqlalchemy.exc import IntegrityError +from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker + +from gateway.evals.infrastructure.orm import EvalCaseRow +from gateway.evals.runs.domain.entities import CaseOutcome +from gateway.evals.runs.infrastructure.orm import EvalCaseResultRow, EvalRunRow +from gateway.tenants.application.retention_policy import raise_if_zdr_locked + + +class SqlAlchemyEvalRunStore: + """EvalRunStore adapter over an async_sessionmaker (fresh session per operation).""" + + def __init__(self, sessionmaker: async_sessionmaker[AsyncSession]) -> None: + self._sessionmaker = sessionmaker + + async def create_run( + self, *, tenant_id: uuid.UUID, key_id: uuid.UUID, eval_set_id: uuid.UUID, model: str + ) -> EvalRunRow: + async with self._sessionmaker() as session: + row = EvalRunRow( + id=uuid.uuid4(), + tenant_id=tenant_id, + eval_set_id=eval_set_id, + key_id=key_id, + model=model, + status="pending", + ) + session.add(row) + await session.flush() + await session.refresh(row) + await session.commit() + # Expunge so the detached row is safe to read after the session closes. + session.expunge(row) + return row + + async def get_run(self, *, tenant_id: uuid.UUID, run_id: uuid.UUID) -> EvalRunRow | None: + async with self._sessionmaker() as session: + stmt = select(EvalRunRow).where( + EvalRunRow.id == run_id, EvalRunRow.tenant_id == tenant_id + ) + row = (await session.execute(stmt)).scalar_one_or_none() + if row is not None: + session.expunge(row) + return row + + async def load_run(self, run_id: uuid.UUID) -> EvalRunRow | None: + async with self._sessionmaker() as session: + row = await session.get(EvalRunRow, run_id) + if row is not None: + session.expunge(row) + return row + + async def snapshot_cases( + self, *, tenant_id: uuid.UUID, eval_set_id: uuid.UUID, created_at_max: object + ) -> list[EvalCaseRow]: + async with self._sessionmaker() as session: + stmt = ( + select(EvalCaseRow) + .where( + EvalCaseRow.tenant_id == tenant_id, + EvalCaseRow.eval_set_id == eval_set_id, + EvalCaseRow.created_at <= created_at_max, + ) + .order_by(EvalCaseRow.created_at.asc(), EvalCaseRow.id.asc()) + ) + rows = list((await session.execute(stmt)).scalars().all()) + for r in rows: + session.expunge(r) + return rows + + async def existing_result_case_ids(self, run_id: uuid.UUID) -> set[uuid.UUID]: + async with self._sessionmaker() as session: + stmt = select(EvalCaseResultRow.eval_case_id).where( + EvalCaseResultRow.eval_run_id == run_id + ) + return set((await session.execute(stmt)).scalars().all()) + + async def record_case_result( + self, + *, + tenant_id: uuid.UUID, + run_id: uuid.UUID, + eval_case_id: uuid.UUID, + outcome: CaseOutcome, + ) -> bool: + async with self._sessionmaker() as session: + # M5: a completed case persists the response payload — gate it ATOMICALLY with the + # insert. FOR UPDATE on the tenant row blocks a concurrent ZDR flip until this + # transaction resolves; a ZDR tenant raises 403 here and NOTHING commits. A + # refused/errored case carries no payload, so it needs no gate. + if outcome.status == "completed": + await raise_if_zdr_locked(session, tenant_id) + row = EvalCaseResultRow( + id=uuid.uuid4(), + tenant_id=tenant_id, + eval_run_id=run_id, + eval_case_id=eval_case_id, + status=outcome.status, + response_text=outcome.response_text, + reason=outcome.reason, + ) + session.add(row) + try: + await session.flush() + except IntegrityError: + # UNIQUE(eval_run_id, eval_case_id): a racer/resume already recorded this case. + # Roll back and report a no-op — never a double-write, never a double-bill. + await session.rollback() + return False + await session.commit() + return True + + async def counts_by_status(self, run_id: uuid.UUID) -> dict[str, int]: + async with self._sessionmaker() as session: + stmt = ( + select(EvalCaseResultRow.status, func.count()) + .where(EvalCaseResultRow.eval_run_id == run_id) + .group_by(EvalCaseResultRow.status) + ) + rows = (await session.execute(stmt)).all() + counts = {"completed": 0, "refused": 0, "errored": 0} + for status, n in rows: + counts[status] = int(n) + return counts + + async def set_run_status(self, run_id: uuid.UUID, status: str) -> None: + async with self._sessionmaker() as session: + row = await session.get(EvalRunRow, run_id) + if row is None: + return + row.status = status + await session.commit() + + async def list_case_results( + self, *, tenant_id: uuid.UUID, run_id: uuid.UUID + ) -> list[EvalCaseResultRow]: + async with self._sessionmaker() as session: + # A5: order by the CASE's creation order (join eval_cases), not wall-clock completion + # order, so a baseline and a candidate run of the same set align case-for-case. + stmt = ( + select(EvalCaseResultRow) + .join(EvalCaseRow, EvalCaseRow.id == EvalCaseResultRow.eval_case_id) + .where( + EvalCaseResultRow.tenant_id == tenant_id, + EvalCaseResultRow.eval_run_id == run_id, + ) + .order_by(EvalCaseRow.created_at.asc(), EvalCaseRow.id.asc()) + ) + rows = list((await session.execute(stmt)).scalars().all()) + for r in rows: + session.expunge(r) + return rows diff --git a/apps/gateway/src/gateway/evals/runs/infrastructure/upstream_adapter.py b/apps/gateway/src/gateway/evals/runs/infrastructure/upstream_adapter.py new file mode 100644 index 00000000..4829c342 --- /dev/null +++ b/apps/gateway/src/gateway/evals/runs/infrastructure/upstream_adapter.py @@ -0,0 +1,117 @@ +"""Per-tenant breaker + concurrency for eval dials (eval-run-executor §3, M4 — the HARD-STOP). + +[[per-tenant-breaker-recurring-defect]]: new provider surfaces keep shipping a PROCESS-GLOBAL +breaker → one tenant's failing burst degrades everyone (cross-tenant DoS). An eval run is a +BURST of billed upstream calls, so it is exactly the shape that trips a shared breaker. The +live completion path's own breaker (``BoundCircuitBreakerUpstream`` over +``app.state.circuit_breaker``) is global; reusing it here would let tenant A's run open the +breaker that guards tenant B's LIVE traffic (R:GLOBAL_BREAKER). + +So the executor injects its OWN upstream into ``CompletionUseCase.complete(upstream=...)``: +``TenantBreakerUpstream`` wraps the raw ``app.state.completion_upstream`` delegate with a +breaker + concurrency semaphore drawn from ``TenantExecutionRegistry``, keyed by ``tenant_id``. +An eval dial then trips ONLY that tenant's eval-breaker; the live global breaker is never +touched by a run, and one tenant's burst can never exhaust another's concurrency slots. + +Every dial also carries a per-call TIMEOUT (M4): a hung provider fails the case CLOSED +(``errored``) via ``UpstreamUnavailableError`` — it never hangs the run. +""" + +from __future__ import annotations + +import asyncio +from typing import Any + +from gateway.proxy.domain.errors import CircuitOpenError, UpstreamUnavailableError +from gateway.proxy.domain.ports import CompletionUpstream +from gateway.proxy.infrastructure.circuit_breaker import CircuitBreaker + +# Sensible eval-burst defaults; overridable from Settings at construction (A3). +_DEFAULT_CONCURRENCY = 4 +_DEFAULT_TIMEOUT_SECONDS = 30.0 + + +class TenantExecutionRegistry: + """Lazily-allocated per-tenant CircuitBreaker + concurrency Semaphore. + + One instance is owned by the executor (NOT app.state.circuit_breaker). ``dict.setdefault`` + keying by ``tenant_id`` is the same lazy-per-key idiom deps.py uses for the per-provider + breaker registry — a trip / slot on one tenant never touches another's. + """ + + def __init__( + self, + *, + concurrency: int = _DEFAULT_CONCURRENCY, + failure_threshold: int = 5, + cooldown_seconds: float = 30.0, + ) -> None: + self._concurrency = max(1, concurrency) + self._failure_threshold = failure_threshold + self._cooldown_seconds = cooldown_seconds + self._breakers: dict[Any, CircuitBreaker] = {} + self._semaphores: dict[Any, asyncio.Semaphore] = {} + + def breaker_for(self, tenant_id: Any) -> CircuitBreaker: + breaker = self._breakers.get(tenant_id) + if breaker is None: + breaker = CircuitBreaker( + failure_threshold=self._failure_threshold, + cooldown_seconds=self._cooldown_seconds, + ) + self._breakers[tenant_id] = breaker + return breaker + + def semaphore_for(self, tenant_id: Any) -> asyncio.Semaphore: + sem = self._semaphores.get(tenant_id) + if sem is None: + sem = asyncio.Semaphore(self._concurrency) + self._semaphores[tenant_id] = sem + return sem + + +class TenantBreakerUpstream: + """CompletionUpstream that guards a delegate with a tenant-keyed breaker + per-call timeout. + + Mirrors ``BoundCircuitBreakerUpstream`` (guard → delegate → count outcome) but the breaker + is the PER-TENANT one from the registry, and every dial is bounded by ``timeout_seconds``. + Breaker-open and timeout BOTH fail closed (``UpstreamUnavailableError``) so the executor + records the case ``errored`` and the run continues (M4, A6/E5). + """ + + def __init__( + self, + breaker: CircuitBreaker, + delegate: CompletionUpstream, + *, + timeout_seconds: float = _DEFAULT_TIMEOUT_SECONDS, + ) -> None: + self._breaker = breaker + self._delegate = delegate + self._timeout = timeout_seconds + + async def complete(self, payload: dict[str, Any]) -> tuple[int, dict[str, Any]]: + if not self._breaker.call_allowed(): + # This tenant's eval-breaker is open — never dial, never touch another tenant. + raise CircuitOpenError("Eval breaker is open for this tenant") + try: + status, body = await asyncio.wait_for( + self._delegate.complete(payload), timeout=self._timeout + ) + except TimeoutError as exc: + # Per-call timeout (M4): count it as an upstream failure and fail the case closed. + self._breaker.on_upstream_error() + raise UpstreamUnavailableError("Eval upstream call timed out") from exc + except (UpstreamUnavailableError, CircuitOpenError): + self._breaker.on_upstream_error() + raise + if status >= 500: + self._breaker.on_upstream_error() + raise UpstreamUnavailableError(f"Eval upstream returned {status}") + self._breaker.record_success() + return status, body + + def stream(self, payload: dict[str, Any]) -> Any: + # Eval replays are non-streaming (complete() only). A stream() is never invoked on this + # path; provide it so the object still satisfies the CompletionUpstream Protocol shape. + raise NotImplementedError("eval runs do not stream") diff --git a/apps/gateway/src/gateway/evals/wire_id.py b/apps/gateway/src/gateway/evals/wire_id.py index e83c4dcf..8cfd4180 100644 --- a/apps/gateway/src/gateway/evals/wire_id.py +++ b/apps/gateway/src/gateway/evals/wire_id.py @@ -1,8 +1,9 @@ """Wire ids for the evals domain (eval-set-store PLAN.md §3 — FROZEN @ v1). -Two reversible prefixes, mirroring ``gateway.vector_stores.wire_id``'s ``vs_<32hex>`` shape: +Three reversible prefixes, mirroring ``gateway.vector_stores.wire_id``'s ``vs_<32hex>`` shape: - eval SET -> ``es_<32hex>`` - eval CASE -> ``ec_<32hex>`` + - eval RUN -> ``er_<32hex>`` (eval-run-executor §3) ``parse_*`` returns None (never raises) for anything that is not exactly the prefix followed by a valid 32-char hex UUID — the caller maps None to the same uniform 404 as an absent or @@ -15,6 +16,7 @@ _SET_PREFIX = "es_" _CASE_PREFIX = "ec_" +_RUN_PREFIX = "er_" def to_set_wire_id(eval_set_id: uuid.UUID) -> str: @@ -25,6 +27,14 @@ def parse_set_wire_id(wire_id: str) -> uuid.UUID | None: return _parse(wire_id, _SET_PREFIX) +def to_run_wire_id(eval_run_id: uuid.UUID) -> str: + return f"{_RUN_PREFIX}{eval_run_id.hex}" + + +def parse_run_wire_id(wire_id: str) -> uuid.UUID | None: + return _parse(wire_id, _RUN_PREFIX) + + def to_case_wire_id(eval_case_id: uuid.UUID) -> str: return f"{_CASE_PREFIX}{eval_case_id.hex}" diff --git a/apps/gateway/src/gateway/main.py b/apps/gateway/src/gateway/main.py index dfc093a5..b44b5c1c 100644 --- a/apps/gateway/src/gateway/main.py +++ b/apps/gateway/src/gateway/main.py @@ -142,6 +142,13 @@ from gateway.evals.infrastructure.orm import ( # noqa: F401 — registers EvalSetRow/EvalCaseRow on Base.metadata EvalCaseRow as _EvalCaseRow, # pyright: ignore[reportUnusedImport] — side-effect import; registers ORM table on Base.metadata ) +from gateway.evals.runs.api.run_router import eval_runs_router +from gateway.evals.runs.application.run_executor import EvalRunExecutor +from gateway.evals.runs.infrastructure.orm import ( # noqa: F401 — registers EvalRunRow/EvalCaseResultRow on Base.metadata + EvalCaseResultRow as _EvalCaseResultRow, # pyright: ignore[reportUnusedImport] — side-effect import; registers ORM table on Base.metadata +) +from gateway.evals.runs.infrastructure.repository import SqlAlchemyEvalRunStore +from gateway.evals.runs.infrastructure.upstream_adapter import TenantExecutionRegistry from gateway.files.api.router import files_router from gateway.finetune.api.router import finetune_router from gateway.finetune.infrastructure.openai_client import OpenAIFinetuneClient @@ -1817,6 +1824,32 @@ async def _load_modality_map() -> dict[str, str]: app.include_router(files_router) app.include_router(vector_stores_router) app.include_router(evals_router) + app.include_router(eval_runs_router) + # eval-run-executor (R7): replay a set through the SAME governance path a live request uses. + # build_use_case reuses the request path's get_completion_use_case over a shim exposing only + # .app (it never reads request.headers) — so use_cases.py/governance.py stay untouched and the + # run's governance/billing hold by construction. The breaker + concurrency are PER-TENANT + # (TenantExecutionRegistry), NEVER app.state.circuit_breaker (R:GLOBAL_BREAKER). Closures read + # app.state lazily (at drive time), so construction here needs only app.state.sessionmaker. + from types import SimpleNamespace + + from gateway.proxy.api.deps import get_completion_use_case as _get_completion_use_case + + def _build_eval_use_case(_session: Any) -> Any: + return _get_completion_use_case(SimpleNamespace(app=app), _session) # type: ignore[arg-type] + + app.state.eval_run_executor = EvalRunExecutor( + sessionmaker=app.state.sessionmaker, + store=SqlAlchemyEvalRunStore(app.state.sessionmaker), + build_use_case=_build_eval_use_case, + get_upstream_delegate=lambda: app.state.completion_upstream, + get_usage_recorder=lambda: app.state.usage_recorder, + get_model_router=lambda: getattr(app.state, "model_router", None), + registry=TenantExecutionRegistry( + concurrency=int(getattr(settings, "eval_run_concurrency", 4)), + ), + per_call_timeout_seconds=float(getattr(settings, "eval_run_timeout_seconds", 30.0)), + ) app.include_router(video_router) app.include_router(batch_router) app.include_router(batch_stats_router) diff --git a/apps/gateway/tests/evals_runs/__init__.py b/apps/gateway/tests/evals_runs/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/apps/gateway/tests/evals_runs/test_eval_run_executor.py b/apps/gateway/tests/evals_runs/test_eval_run_executor.py new file mode 100644 index 00000000..47b158cd --- /dev/null +++ b/apps/gateway/tests/evals_runs/test_eval_run_executor.py @@ -0,0 +1,697 @@ +"""RED suite for eval-run-executor (/v1/evals/.../runs) — the R7 governed-replay task. + +Contract under test (eval-run-executor PLAN.md §3, FROZEN @ sha256:1353206d): + POST /v1/evals/sets/{set_id}/runs {model} -> 201 { id:"er_<32hex>", eval_set_id, model, + status, case_count, created_at } + 403 ERR_ZDR_PAYLOAD_BLOCKED (ZDR tenant, refused outright at launch, M5) + 404 ERR_EVAL_SET_NOT_FOUND (absent OR cross-tenant set, M6) + GET /v1/evals/runs/{run_id} -> 200 { ..., status, case_count, counts:{...} } + GET /v1/evals/runs/{run_id}/cases -> 200 list, case-creation order (A5) + +A run REPLAYS each case's stored request_body against the run's model THROUGH the real +governance path (CompletionUseCase.complete) — so governance (M1) + one-usage-record (M2) hold +by construction; a governance refusal never dials (M3); the breaker + concurrency are +PER-TENANT, never the global one (M4, R:GLOBAL_BREAKER); ZDR refuses outright + atomic (M5); a +resume never re-bills a terminal case (M7). + +RED until the evals/runs module + migration f5b2d8c41a37 are wired. DO NOT edit to make pass — +that is Build's job. +""" + +from __future__ import annotations + +import asyncio +import datetime +import uuid +from typing import Any + +import pytest +from sqlalchemy import text + +from gateway.proxy.domain.errors import UpstreamUnavailableError + +from tests import _redis_env + +pytestmark = pytest.mark.asyncio + + +# --------------------------------------------------------------------------- +# Fakes — the M8 zero-network seam: a fake upstream + a counting usage recorder. +# --------------------------------------------------------------------------- + + +class FakeEvalUpstream: + """CompletionUpstream fake: counts dials; returns 200 / raises / sleeps-on-marker. + + ``mode`` = 'ok' (200 echo), 'fail' (UpstreamUnavailableError — drives the breaker), or + 'slow_marker' (sleep past the per-call timeout iff the message content contains 'SLOW'). + """ + + def __init__(self, *, mode: str = "ok", dwell: float = 0.0) -> None: + self.mode = mode + self.calls = 0 + self._dwell = dwell # seconds each dial lingers (to create in-flight overlap for A3) + self._inflight = 0 + self.max_inflight = 0 + self._usage = {"prompt_tokens": 5, "completion_tokens": 2, "total_tokens": 7} + + async def complete(self, payload: dict[str, Any]) -> tuple[int, dict[str, Any]]: + self.calls += 1 + self._inflight += 1 + self.max_inflight = max(self.max_inflight, self._inflight) + content = "" + try: + content = str(payload["messages"][-1]["content"]) + except Exception: + content = "" + try: + if self.mode == "fail": + raise UpstreamUnavailableError("eval upstream boom") + if self._dwell: + await asyncio.sleep(self._dwell) + if self.mode == "slow_marker" and "SLOW" in content: + await asyncio.sleep(1.0) + finally: + self._inflight -= 1 + return ( + 200, + { + "id": "chatcmpl-eval", + "object": "chat.completion", + "model": payload.get("model"), + "choices": [ + { + "index": 0, + "message": {"role": "assistant", "content": f"echo:{content}"}, + "finish_reason": "stop", + } + ], + "usage": dict(self._usage), + }, + ) + + def stream(self, payload: dict[str, Any]) -> Any: + raise NotImplementedError + + +class CountingRecorder: + """In-memory UsageRecorder — captures each record (avoids the Redis→PG flush timing).""" + + def __init__(self) -> None: + self.records: list[dict[str, Any]] = [] + + async def record( + self, *, tenant_id: Any, key_id: Any, model: str, usage: Any, status: int, **_: Any + ) -> None: + self.records.append({"tenant_id": tenant_id, "model": model, "status": status}) + + +# --------------------------------------------------------------------------- +# Fixtures — tenant + key + an active catalog model (mirror the harness tests). +# --------------------------------------------------------------------------- + + +def _bearer(key: str) -> dict[str, str]: + return {"Authorization": f"Bearer {key}"} + + +async def _signup_key(client: Any, *, tenant: str, email: str) -> dict[str, str]: + signup = await client.post( + "/admin/auth/signup", + json={"tenant_name": tenant, "email": email, "password": "correct horse battery staple"}, + ) + assert signup.status_code == 201, signup.text + token = ( + await client.post( + "/admin/auth/login", json={"email": email, "password": "correct horse battery staple"} + ) + ).json()["access_token"] + created = await client.post( + "/admin/keys", json={"name": f"k-{tenant}"}, headers={"Authorization": f"Bearer {token}"} + ) + assert created.status_code == 201, created.text + return { + "key": created.json()["key"], + "key_id": created.json()["key_id"], + "tenant_id": signup.json()["tenant_id"], + } + + +@pytest.fixture +async def tenant_a(client: Any) -> dict[str, str]: + return await _signup_key(client, tenant="RunA", email="run-a@example.io") + + +@pytest.fixture +async def tenant_b(client: Any) -> dict[str, str]: + return await _signup_key(client, tenant="RunB", email="run-b@example.io") + + +@pytest.fixture +async def active_model(db_session: Any) -> str: + model_id = "openai/gpt-4o-eval" + await db_session.execute( + text("INSERT INTO models (id, name, context_length, active) VALUES (:i, :n, 128000, true)"), + {"i": model_id, "n": "GPT-4o eval"}, + ) + await db_session.execute( + text( + "INSERT INTO pricing_snapshots " + "(id, model_id, prompt_usd_per_token, completion_usd_per_token, captured_at) " + "VALUES (:id, :m, 0.0000025, 0.00001, now())" + ), + {"id": str(uuid.uuid4()), "m": model_id}, + ) + await db_session.commit() + return model_id + + +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- + + +async def _make_set(client: Any, key: str, name: str) -> str: + resp = await client.post("/v1/evals/sets", json={"name": name}, headers=_bearer(key)) + assert resp.status_code == 201, resp.text + return resp.json()["id"] + + +async def _add_case(client: Any, key: str, set_id: str, content: str) -> None: + resp = await client.post( + f"/v1/evals/sets/{set_id}/cases", + json={ + "request_body": {"model": "ignored", "messages": [{"role": "user", "content": content}]}, + "assertion": {"kind": "contains", "expected": "echo"}, + }, + headers=_bearer(key), + ) + assert resp.status_code == 201, resp.text + + +async def _count(app: Any, sql: str, **params: Any) -> int: + async with app.state.sessionmaker() as session: + return int((await session.execute(text(sql), params)).scalar_one()) + + +async def _poll(predicate: Any, *, ticks: int = 250) -> bool: + for _ in range(ticks): + if predicate(): + return True + await asyncio.sleep(0.02) + return predicate() + + +def _install(app: Any, upstream: Any, recorder: Any | None = None) -> CountingRecorder: + app.state.completion_upstream = upstream + rec = recorder or CountingRecorder() + app.state.usage_recorder = rec + return rec + + +# --------------------------------------------------------------------------- +# M1 — a run enters governance per case; a dial happens once per case +# --------------------------------------------------------------------------- + + +async def test_run_enters_governance_per_case( + client: Any, app: Any, tenant_a: dict[str, str], active_model: str +) -> None: + """covers: M1, R:GOVERNANCE_BYPASS — every case dials THROUGH complete()'s guards, once each.""" + rec = _install(app, FakeEvalUpstream(mode="ok")) + upstream = app.state.completion_upstream + set_id = await _make_set(client, tenant_a["key"], "gov") + for i in range(3): + await _add_case(client, tenant_a["key"], set_id, f"case-{i}") + + run = await client.post( + f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + ) + assert run.status_code == 201, run.text + run_id = run.json()["id"] + assert run.json()["case_count"] == 3 + + cases = (await client.get(f"/v1/evals/runs/{run_id}/cases", headers=_bearer(tenant_a["key"]))).json() + assert [c["status"] for c in cases["data"]] == ["completed", "completed", "completed"] + # one dial per case (governance passed for each, then dialed) — never more, never fewer. + assert upstream.calls == 3 + # each dialed case produced one usage record billed on the SERVED (run's) model id (M2 hook). + assert await _poll(lambda: len(rec.records) == 3) + assert {r["model"] for r in rec.records} == {active_model} + + +# --------------------------------------------------------------------------- +# M3 / E1 — an over-budget run refuses every case with NO dial and NO usage +# --------------------------------------------------------------------------- + + +async def test_over_budget_run_refuses_every_case_no_dial( + client: Any, app: Any, tenant_a: dict[str, str], active_model: str +) -> None: + """covers: M3, E1, R:GOVERNANCE_BYPASS — a spend-capped tenant dials ZERO, bills ZERO.""" + from gateway.budgets.infrastructure.redis_guard import RedisBudgetGuard + + import redis.asyncio as aioredis + + set_id = await _make_set(client, tenant_a["key"], "poor") + for i in range(3): + await _add_case(client, tenant_a["key"], set_id, f"case-{i}") + + # Give the key a tiny cap and seed the per-key spend counter past it → every case refused + # at the budget guard, BEFORE any dial (mirrors the agent-oauth over-cap test). + await _count(app, "SELECT 1") # no-op touch to open a session cleanly + async with app.state.sessionmaker() as s: + await s.execute( + text("UPDATE api_keys SET monthly_budget_usd = '0.01' WHERE id = :id"), + {"id": tenant_a["key_id"]}, + ) + await s.commit() + + redis_client = aioredis.from_url(_redis_env.TEST_REDIS_URL) + try: + yyyymm = datetime.datetime.now(datetime.UTC).strftime("%Y%m") + await redis_client.set(f"usage:spend:key:{tenant_a['key_id']}:{yyyymm}", b"100.00") + app.state.budget_guard = RedisBudgetGuard( + redis=redis_client, session_factory=app.state.sessionmaker + ) + rec = _install(app, FakeEvalUpstream(mode="ok")) + upstream = app.state.completion_upstream + + run = await client.post( + f"/v1/evals/sets/{set_id}/runs", + json={"model": active_model}, + headers=_bearer(tenant_a["key"]), + ) + assert run.status_code == 201, run.text + run_id = run.json()["id"] + + cases = ( + await client.get(f"/v1/evals/runs/{run_id}/cases", headers=_bearer(tenant_a["key"])) + ).json() + assert [c["status"] for c in cases["data"]] == ["refused", "refused", "refused"] + assert upstream.calls == 0, "a budget-refused case must never dial upstream" + # zero usage rows AND zero persisted response payloads. + assert len(rec.records) == 0 + assert ( + await _count( + app, + "SELECT count(*) FROM eval_case_results WHERE tenant_id = :t" + " AND response_text IS NOT NULL", + t=tenant_a["tenant_id"], + ) + == 0 + ) + finally: + await redis_client.aclose() + + +# --------------------------------------------------------------------------- +# M2 — exactly one usage record per dialed case, on the served model id +# --------------------------------------------------------------------------- + + +async def test_one_usage_record_per_dialed_case( + client: Any, app: Any, tenant_a: dict[str, str], active_model: str +) -> None: + """covers: M2, R:DOUBLE_BILL — N dial-able cases -> exactly N usage records, served model.""" + rec = _install(app, FakeEvalUpstream(mode="ok")) + set_id = await _make_set(client, tenant_a["key"], "bill") + for i in range(4): + await _add_case(client, tenant_a["key"], set_id, f"distinct-{i}") + + run = await client.post( + f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + ) + assert run.status_code == 201, run.text + + assert await _poll(lambda: len(rec.records) == 4) + # exactly N — not N-1 (unmetered) and not N+1 (double-bill); each on the served model id. + assert len(rec.records) == 4 + assert all(r["model"] == active_model and r["status"] == 200 for r in rec.records) + + +# --------------------------------------------------------------------------- +# M4 / E2 — the breaker is PER-TENANT: A's failing run never opens B's breaker +# --------------------------------------------------------------------------- + + +async def test_per_tenant_breaker_isolation( + client: Any, app: Any, tenant_a: dict[str, str], tenant_b: dict[str, str], active_model: str +) -> None: + """covers: M4, E2, R:GLOBAL_BREAKER — A's burst opens A's breaker only; B is untouched.""" + executor = app.state.eval_run_executor + registry = executor._registry # noqa: SLF001 — asserting the per-tenant keying is the point + a_tid = uuid.UUID(tenant_a["tenant_id"]) + b_tid = uuid.UUID(tenant_b["tenant_id"]) + + # Tenant A runs against an always-failing upstream: enough cases to trip the 5-failure breaker. + _install(app, FakeEvalUpstream(mode="fail")) + set_a = await _make_set(client, tenant_a["key"], "breaker-a") + for i in range(6): + await _add_case(client, tenant_a["key"], set_a, f"a-{i}") + run_a = await client.post( + f"/v1/evals/sets/{set_a}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + ) + assert run_a.status_code == 201, run_a.text + cases_a = ( + await client.get(f"/v1/evals/runs/{run_a.json()['id']}/cases", headers=_bearer(tenant_a["key"])) + ).json() + assert all(c["status"] == "errored" for c in cases_a["data"]) + assert registry.breaker_for(a_tid).is_open() is True + # B's breaker is a DIFFERENT object and was never touched by A's failures. + assert registry.breaker_for(b_tid).is_open() is False + + # And B can still complete a run despite A's breaker being open (no shared state). + _install(app, FakeEvalUpstream(mode="ok")) + set_b = await _make_set(client, tenant_b["key"], "breaker-b") + for i in range(2): + await _add_case(client, tenant_b["key"], set_b, f"b-{i}") + run_b = await client.post( + f"/v1/evals/sets/{set_b}/runs", json={"model": active_model}, headers=_bearer(tenant_b["key"]) + ) + assert run_b.status_code == 201, run_b.text + cases_b = ( + await client.get(f"/v1/evals/runs/{run_b.json()['id']}/cases", headers=_bearer(tenant_b["key"])) + ).json() + assert all(c["status"] == "completed" for c in cases_b["data"]) + + +# --------------------------------------------------------------------------- +# M5 / E3 — a ZDR run persists ZERO results (outright at launch + atomic mid-run) +# --------------------------------------------------------------------------- + + +async def test_zdr_run_refused_atomically_zero_results( + client: Any, app: Any, tenant_a: dict[str, str], active_model: str +) -> None: + """covers: M5, E3, R:ZDR_BLOCKED — asserted on the PERSISTED ROW count, not the response.""" + _install(app, FakeEvalUpstream(mode="ok")) + set_id = await _make_set(client, tenant_a["key"], "zdr") + for i in range(2): + await _add_case(client, tenant_a["key"], set_id, f"z-{i}") + + # (a) ZDR set BEFORE launch -> refused outright at launch (403), nothing created. + async with app.state.sessionmaker() as s: + await s.execute( + text("UPDATE tenants SET zdr_enabled = true WHERE id = :id"), + {"id": tenant_a["tenant_id"]}, + ) + await s.commit() + launch = await client.post( + f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + ) + assert launch.status_code == 403 + assert launch.json()["error"]["code"] == "ERR_ZDR_PAYLOAD_BLOCKED" + assert ( + await _count( + app, + "SELECT count(*) FROM eval_case_results WHERE tenant_id = :t", + t=tenant_a["tenant_id"], + ) + == 0 + ) + + # (b) TOCTOU: launch with ZDR OFF, flip it ON, then drive -> the atomic locked re-check at + # the first payload write refuses; the run is `blocked` and ZERO results persist. + async with app.state.sessionmaker() as s: + await s.execute( + text("UPDATE tenants SET zdr_enabled = false WHERE id = :id"), + {"id": tenant_a["tenant_id"]}, + ) + await s.commit() + executor = app.state.eval_run_executor + run = await executor.launch( + tenant_id=uuid.UUID(tenant_a["tenant_id"]), + key_id=uuid.UUID(tenant_a["key_id"]), + raw_key=tenant_a["key"], + eval_set_id=uuid.UUID(_strip_prefix(set_id)), + model=active_model, + ) + async with app.state.sessionmaker() as s: + await s.execute( + text("UPDATE tenants SET zdr_enabled = true WHERE id = :id"), + {"id": tenant_a["tenant_id"]}, + ) + await s.commit() + await executor.drive(run.id, raw_key=tenant_a["key"]) + async with app.state.sessionmaker() as s: + status = ( + await s.execute(text("SELECT status FROM eval_runs WHERE id = :id"), {"id": run.id}) + ).scalar_one() + assert status == "blocked" + assert ( + await _count( + app, + "SELECT count(*) FROM eval_case_results WHERE eval_run_id = :r", + r=run.id, + ) + == 0 + ) + + +# --------------------------------------------------------------------------- +# M7 / E4 — a resumed run dials only pending cases; terminal cases are not re-billed +# --------------------------------------------------------------------------- + + +async def test_resume_does_not_rebill_terminal_cases( + client: Any, app: Any, tenant_a: dict[str, str], active_model: str +) -> None: + """covers: M7, E4, R:DOUBLE_BILL — a second drive re-dials nothing; total dials == N once.""" + rec = _install(app, FakeEvalUpstream(mode="ok")) + upstream = app.state.completion_upstream + set_id = await _make_set(client, tenant_a["key"], "resume") + for i in range(3): + await _add_case(client, tenant_a["key"], set_id, f"r-{i}") + + run = await client.post( + f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + ) + assert run.status_code == 201, run.text + run_id_wire = run.json()["id"] + assert await _poll(lambda: len(rec.records) == 3) + assert upstream.calls == 3 + + # Resume the SAME run (all cases terminal) -> zero new dials, zero new bills. + executor = app.state.eval_run_executor + await executor.drive(uuid.UUID(_strip_prefix(run_id_wire)), raw_key=tenant_a["key"]) + await asyncio.sleep(0.05) + assert upstream.calls == 3, "a resumed drive must not re-dial a terminal case" + assert len(rec.records) == 3, "a resumed drive must not re-bill a terminal case" + # exactly N result rows — the UNIQUE(run, case) idempotency held. + assert ( + await _count(app, "SELECT count(*) FROM eval_case_results WHERE eval_run_id = :r", + r=uuid.UUID(_strip_prefix(run_id_wire))) + == 3 + ) + + +# --------------------------------------------------------------------------- +# M4 / E5 — a per-call timeout errors ONE case; the run continues +# --------------------------------------------------------------------------- + + +async def test_timeout_case_errored_run_continues( + client: Any, app: Any, tenant_a: dict[str, str], active_model: str +) -> None: + """covers: M4, E5 — one case times out (errored), the rest complete; the run is a mix.""" + _install(app, FakeEvalUpstream(mode="slow_marker")) + # Shorten the per-call timeout so the SLOW case trips it fast (no wall-clock assertion). + app.state.eval_run_executor._timeout = 0.1 # noqa: SLF001 — test-time knob + set_id = await _make_set(client, tenant_a["key"], "timeout") + await _add_case(client, tenant_a["key"], set_id, "fast-1") + await _add_case(client, tenant_a["key"], set_id, "please-be-SLOW") + await _add_case(client, tenant_a["key"], set_id, "fast-2") + + run = await client.post( + f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + ) + assert run.status_code == 201, run.text + cases = ( + await client.get(f"/v1/evals/runs/{run.json()['id']}/cases", headers=_bearer(tenant_a["key"])) + ).json()["data"] + by_status = sorted(c["status"] for c in cases) + assert by_status == ["completed", "completed", "errored"], by_status + # the run still reaches a terminal rollup — one bad case never sinks it. + got = (await client.get(f"/v1/evals/runs/{run.json()['id']}", headers=_bearer(tenant_a["key"]))).json() + assert got["status"] == "completed" + assert got["counts"]["errored"] == 1 and got["counts"]["completed"] == 2 + + +# --------------------------------------------------------------------------- +# M6 — a cross-tenant / absent run id is a uniform 404, no oracle +# --------------------------------------------------------------------------- + + +async def test_cross_tenant_and_absent_run_uniform_404( + client: Any, app: Any, tenant_a: dict[str, str], tenant_b: dict[str, str], active_model: str +) -> None: + """covers: M6, R:RUN_NOT_FOUND — cross-tenant and absent are byte-identical 404s.""" + _install(app, FakeEvalUpstream(mode="ok")) + set_id = await _make_set(client, tenant_a["key"], "iso") + await _add_case(client, tenant_a["key"], set_id, "x") + run = await client.post( + f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + ) + run_id = run.json()["id"] + + # Tenant B asking for A's run. + cross = await client.get(f"/v1/evals/runs/{run_id}", headers=_bearer(tenant_b["key"])) + # A wholly-absent run id. + absent = await client.get(f"/v1/evals/runs/er_{uuid.uuid4().hex}", headers=_bearer(tenant_b["key"])) + assert cross.status_code == 404 and absent.status_code == 404 + assert cross.json() == absent.json(), "cross-tenant must be byte-identical to absent (no oracle)" + assert cross.json()["error"]["code"] == "ERR_EVAL_RUN_NOT_FOUND" + # A malformed id resolves to the same 404, never a 500. + malformed = await client.get("/v1/evals/runs/not-a-run", headers=_bearer(tenant_b["key"])) + assert malformed.status_code == 404 + + +# --------------------------------------------------------------------------- +# A5 — results align to CASE creation order, not completion order +# --------------------------------------------------------------------------- + + +async def test_results_aligned_to_case_creation_order( + client: Any, app: Any, tenant_a: dict[str, str], active_model: str +) -> None: + """covers: A5 — two runs of the same set expose results in identical case-creation order.""" + _install(app, FakeEvalUpstream(mode="ok")) + set_id = await _make_set(client, tenant_a["key"], "order") + for i in range(5): + await _add_case(client, tenant_a["key"], set_id, f"ordered-{i}") + + async def _run_case_ids() -> list[str]: + run = await client.post( + f"/v1/evals/sets/{set_id}/runs", + json={"model": active_model}, + headers=_bearer(tenant_a["key"]), + ) + assert run.status_code == 201, run.text + data = ( + await client.get(f"/v1/evals/runs/{run.json()['id']}/cases", headers=_bearer(tenant_a["key"])) + ).json()["data"] + return [c["eval_case_id"] for c in data] + + first = await _run_case_ids() + second = await _run_case_ids() + assert first == second, "two runs must align case-for-case regardless of completion order" + # And that order is the set's own creation order. + listed = ( + await client.get(f"/v1/evals/sets/{set_id}/cases", headers=_bearer(tenant_a["key"])) + ).json()["data"] + assert first == [c["id"] for c in listed] + + +# --------------------------------------------------------------------------- +# A4 / E6 — an empty set runs vacuously: completed, 0/0, no exception +# --------------------------------------------------------------------------- + + +async def test_empty_set_run_completes_vacuously( + client: Any, app: Any, tenant_a: dict[str, str], active_model: str +) -> None: + """covers: A4, E6 — running an empty set is a completed 0/0 run, never an error.""" + _install(app, FakeEvalUpstream(mode="ok")) + set_id = await _make_set(client, tenant_a["key"], "empty") + + run = await client.post( + f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + ) + assert run.status_code == 201, run.text + body = run.json() + assert body["case_count"] == 0 + got = (await client.get(f"/v1/evals/runs/{body['id']}", headers=_bearer(tenant_a["key"]))).json() + assert got["status"] == "completed" + assert got["counts"] == {"completed": 0, "refused": 0, "errored": 0, "pending": 0} + cases = (await client.get(f"/v1/evals/runs/{body['id']}/cases", headers=_bearer(tenant_a["key"]))).json() + assert cases["data"] == [] + + +# --------------------------------------------------------------------------- +# A1 — a run is launched+billed by the launching key; cross-tenant launch is a 404 +# --------------------------------------------------------------------------- + + +async def test_cross_tenant_launch_and_billing_identity( + client: Any, app: Any, tenant_a: dict[str, str], tenant_b: dict[str, str], active_model: str +) -> None: + """covers: A1 — B cannot launch a run on A's set (404); a run bills the LAUNCHING tenant/key.""" + rec = _install(app, FakeEvalUpstream(mode="ok")) + set_id = await _make_set(client, tenant_a["key"], "a1") + await _add_case(client, tenant_a["key"], set_id, "own") + + # Tenant B launching a run on tenant A's set -> uniform 404 (A's set is invisible to B). + forbidden = await client.post( + f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_b["key"]) + ) + assert forbidden.status_code == 404 + assert forbidden.json()["error"]["code"] == "ERR_EVAL_SET_NOT_FOUND" + + # Tenant A's own run bills A's tenant, never B's (the usage record carries the launcher). + run = await client.post( + f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + ) + assert run.status_code == 201, run.text + assert await _poll(lambda: len(rec.records) == 1) + assert rec.records[0]["tenant_id"] == uuid.UUID(tenant_a["tenant_id"]) + + +# --------------------------------------------------------------------------- +# A2 — the run's case set is snapshotted AT LAUNCH; later cases do not join it +# --------------------------------------------------------------------------- + + +async def test_run_snapshot_fixed_at_launch( + client: Any, app: Any, tenant_a: dict[str, str], active_model: str +) -> None: + """covers: A2 — a case added AFTER launch does not change the launched run's case_count.""" + _install(app, FakeEvalUpstream(mode="ok")) + set_id = await _make_set(client, tenant_a["key"], "snapshot") + await _add_case(client, tenant_a["key"], set_id, "at-launch-1") + await _add_case(client, tenant_a["key"], set_id, "at-launch-2") + + run = await client.post( + f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + ) + assert run.status_code == 201, run.text + assert run.json()["case_count"] == 2 + + # A case added after launch is NOT in the already-launched run's fixed denominator. + await _add_case(client, tenant_a["key"], set_id, "after-launch-3") + got = (await client.get(f"/v1/evals/runs/{run.json()['id']}", headers=_bearer(tenant_a["key"]))).json() + assert got["case_count"] == 2, "the run's case set is fixed at launch (A2)" + + +# --------------------------------------------------------------------------- +# A3 — per-tenant concurrency is BOUNDED: never more than the ceiling in flight +# --------------------------------------------------------------------------- + + +async def test_bounded_per_tenant_concurrency( + client: Any, app: Any, tenant_a: dict[str, str], active_model: str +) -> None: + """covers: A3 — a run of N>ceiling cases never has more than the ceiling dialing at once.""" + from gateway.evals.runs.infrastructure.upstream_adapter import TenantExecutionRegistry + + fake = FakeEvalUpstream(mode="ok", dwell=0.05) # each dial lingers so overlap is observable + _install(app, fake) + # Pin a small, known ceiling on the executor's registry so the bound is unambiguous. + executor = app.state.eval_run_executor + executor._registry = TenantExecutionRegistry(concurrency=2) # noqa: SLF001 — test knob + + set_id = await _make_set(client, tenant_a["key"], "concurrency") + for i in range(6): + await _add_case(client, tenant_a["key"], set_id, f"c-{i}") + + run = await client.post( + f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + ) + assert run.status_code == 201, run.text + assert fake.calls == 6 + assert fake.max_inflight <= 2, f"exceeded the per-tenant ceiling: {fake.max_inflight}" + + +def _strip_prefix(wire_id: str) -> str: + """'er_' / 'es_' -> the raw 32-hex (test helper for direct executor calls).""" + return wire_id.split("_", 1)[1] diff --git a/apps/gateway/tests/guardrails/test_guardrails_core.py b/apps/gateway/tests/guardrails/test_guardrails_core.py index b4be91e4..b5d2ad11 100644 --- a/apps/gateway/tests/guardrails/test_guardrails_core.py +++ b/apps/gateway/tests/guardrails/test_guardrails_core.py @@ -1476,6 +1476,11 @@ async def test_guardrails_core_migration_column_exists( # the eval-set-store task's own evals/ context, registered on Base.metadata via the # router's import chain (repository -> orm, same precedent as every entry above). # Guardrails-core still adds no tables of its own; invariant intent unchanged. + # SANCTIONED EDIT (eval-run-executor PLAN.md §3 two-manifest rule, 2026-08-13): added + # 'eval_runs', 'eval_case_results' — additive migration f5b2d8c41a37, two NEW tables owned + # by the eval-run-executor task's own evals/runs/ context, registered on Base.metadata via + # main.py's side-effect import (same precedent as every entry above). Guardrails-core still + # adds no tables of its own; invariant intent unchanged. new_tables = ( await db_session.execute( text( @@ -1502,7 +1507,7 @@ async def test_guardrails_core_migration_column_exists( " 'pending_personal_signups','access_requests','files'," " 'stored_responses','vector_stores','vector_store_files'," " 'vector_store_chunks','finetune_jobs','finetune_job_events'," - " 'eval_sets','eval_cases')" + " 'eval_sets','eval_cases','eval_runs','eval_case_results')" ) ) ).fetchall() diff --git a/apps/gateway/tests/migrations/test_migrations.py b/apps/gateway/tests/migrations/test_migrations.py index c06faffe..6df2fe6f 100644 --- a/apps/gateway/tests/migrations/test_migrations.py +++ b/apps/gateway/tests/migrations/test_migrations.py @@ -89,6 +89,8 @@ "finetune_job_events", # SANCTIONED EDIT — finetune-broker PLAN.md §3 manifest maintenance (two-manifest rule); disposition: additive migration 6f2a9c1e3b7d adds this table "eval_sets", # SANCTIONED EDIT — eval-set-store PLAN.md §3 manifest maintenance (two-manifest rule); disposition: additive migration e4a1c9d27f60 adds this table (tenant-owned metadata, UNIQUE(tenant_id,name)) "eval_cases", # SANCTIONED EDIT — eval-set-store PLAN.md §3 manifest maintenance (two-manifest rule); disposition: additive migration e4a1c9d27f60 adds this table (payload-bearing: request_body+assertion JSONB, FK CASCADE -> eval_sets) + "eval_runs", # SANCTIONED EDIT — eval-run-executor PLAN.md §3 manifest maintenance (two-manifest rule); disposition: additive migration f5b2d8c41a37 adds this table (one row per launched run; status DERIVED from cases, FK CASCADE -> eval_sets) + "eval_case_results", # SANCTIONED EDIT — eval-run-executor PLAN.md §3 manifest maintenance (two-manifest rule); disposition: additive migration f5b2d8c41a37 adds this table (per-case result; response_text is the ZDR-gated payload-at-rest, UNIQUE(eval_run_id,eval_case_id), FK CASCADE -> eval_runs) } ) From eb730958eecf8c03ac038a69d6392aa1676d77e9 Mon Sep 17 00:00:00 2001 From: Tin Dang Date: Thu, 13 Aug 2026 13:28:50 +0700 Subject: [PATCH 2/3] style(evals): ruff format eval-run-executor files (CI lint gate) The lint gate is `ruff check . && ruff format --check .`; my local check covered ruff check but not the formatter. Format-only: migration, run_router, and the test file. No behavior change. author: Tin Dang --- .../f5b2d8c41a37_eval_run_executor.py | 4 +- .../src/gateway/evals/runs/api/run_router.py | 8 +- .../evals_runs/test_eval_run_executor.py | 108 +++++++++++++----- 3 files changed, 84 insertions(+), 36 deletions(-) diff --git a/apps/gateway/migrations/versions/f5b2d8c41a37_eval_run_executor.py b/apps/gateway/migrations/versions/f5b2d8c41a37_eval_run_executor.py index b3572f56..ace6d59d 100644 --- a/apps/gateway/migrations/versions/f5b2d8c41a37_eval_run_executor.py +++ b/apps/gateway/migrations/versions/f5b2d8c41a37_eval_run_executor.py @@ -80,9 +80,7 @@ def upgrade() -> None: ), sa.PrimaryKeyConstraint("id"), sa.ForeignKeyConstraint(["eval_run_id"], ["eval_runs.id"], ondelete="CASCADE"), - sa.UniqueConstraint( - "eval_run_id", "eval_case_id", name="uq_eval_case_results_run_case" - ), + sa.UniqueConstraint("eval_run_id", "eval_case_id", name="uq_eval_case_results_run_case"), ) op.create_index( "ix_eval_case_results_run_created", diff --git a/apps/gateway/src/gateway/evals/runs/api/run_router.py b/apps/gateway/src/gateway/evals/runs/api/run_router.py index ce3db9dd..28b6cf48 100644 --- a/apps/gateway/src/gateway/evals/runs/api/run_router.py +++ b/apps/gateway/src/gateway/evals/runs/api/run_router.py @@ -95,9 +95,7 @@ def _case_result_object(row: EvalCaseResultRow) -> dict[str, Any]: return obj -@eval_runs_router.post( - "/v1/evals/sets/{set_id}/runs", status_code=201, response_model=None -) +@eval_runs_router.post("/v1/evals/sets/{set_id}/runs", status_code=201, response_model=None) async def launch_eval_run( set_id: str, body: dict[str, Any], @@ -150,9 +148,7 @@ async def launch_eval_run( return _run_object(refreshed, case_count=case_count) -async def _enqueue_or_drive( - request: Request, executor: Any, *, run_id: Any, raw_key: str -) -> None: +async def _enqueue_or_drive(request: Request, executor: Any, *, run_id: Any, raw_key: str) -> None: """Enqueue for the durable worker; FAIL OPEN to an inline drive (vector-store idiom, M7).""" queue = getattr(request.app.state, "eval_run_queue", None) if queue is not None: diff --git a/apps/gateway/tests/evals_runs/test_eval_run_executor.py b/apps/gateway/tests/evals_runs/test_eval_run_executor.py index 47b158cd..122e56f4 100644 --- a/apps/gateway/tests/evals_runs/test_eval_run_executor.py +++ b/apps/gateway/tests/evals_runs/test_eval_run_executor.py @@ -181,7 +181,10 @@ async def _add_case(client: Any, key: str, set_id: str, content: str) -> None: resp = await client.post( f"/v1/evals/sets/{set_id}/cases", json={ - "request_body": {"model": "ignored", "messages": [{"role": "user", "content": content}]}, + "request_body": { + "model": "ignored", + "messages": [{"role": "user", "content": content}], + }, "assertion": {"kind": "contains", "expected": "echo"}, }, headers=_bearer(key), @@ -225,13 +228,17 @@ async def test_run_enters_governance_per_case( await _add_case(client, tenant_a["key"], set_id, f"case-{i}") run = await client.post( - f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + f"/v1/evals/sets/{set_id}/runs", + json={"model": active_model}, + headers=_bearer(tenant_a["key"]), ) assert run.status_code == 201, run.text run_id = run.json()["id"] assert run.json()["case_count"] == 3 - cases = (await client.get(f"/v1/evals/runs/{run_id}/cases", headers=_bearer(tenant_a["key"]))).json() + cases = ( + await client.get(f"/v1/evals/runs/{run_id}/cases", headers=_bearer(tenant_a["key"])) + ).json() assert [c["status"] for c in cases["data"]] == ["completed", "completed", "completed"] # one dial per case (governance passed for each, then dialed) — never more, never fewer. assert upstream.calls == 3 @@ -320,7 +327,9 @@ async def test_one_usage_record_per_dialed_case( await _add_case(client, tenant_a["key"], set_id, f"distinct-{i}") run = await client.post( - f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + f"/v1/evals/sets/{set_id}/runs", + json={"model": active_model}, + headers=_bearer(tenant_a["key"]), ) assert run.status_code == 201, run.text @@ -350,11 +359,15 @@ async def test_per_tenant_breaker_isolation( for i in range(6): await _add_case(client, tenant_a["key"], set_a, f"a-{i}") run_a = await client.post( - f"/v1/evals/sets/{set_a}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + f"/v1/evals/sets/{set_a}/runs", + json={"model": active_model}, + headers=_bearer(tenant_a["key"]), ) assert run_a.status_code == 201, run_a.text cases_a = ( - await client.get(f"/v1/evals/runs/{run_a.json()['id']}/cases", headers=_bearer(tenant_a["key"])) + await client.get( + f"/v1/evals/runs/{run_a.json()['id']}/cases", headers=_bearer(tenant_a["key"]) + ) ).json() assert all(c["status"] == "errored" for c in cases_a["data"]) assert registry.breaker_for(a_tid).is_open() is True @@ -367,11 +380,15 @@ async def test_per_tenant_breaker_isolation( for i in range(2): await _add_case(client, tenant_b["key"], set_b, f"b-{i}") run_b = await client.post( - f"/v1/evals/sets/{set_b}/runs", json={"model": active_model}, headers=_bearer(tenant_b["key"]) + f"/v1/evals/sets/{set_b}/runs", + json={"model": active_model}, + headers=_bearer(tenant_b["key"]), ) assert run_b.status_code == 201, run_b.text cases_b = ( - await client.get(f"/v1/evals/runs/{run_b.json()['id']}/cases", headers=_bearer(tenant_b["key"])) + await client.get( + f"/v1/evals/runs/{run_b.json()['id']}/cases", headers=_bearer(tenant_b["key"]) + ) ).json() assert all(c["status"] == "completed" for c in cases_b["data"]) @@ -398,7 +415,9 @@ async def test_zdr_run_refused_atomically_zero_results( ) await s.commit() launch = await client.post( - f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + f"/v1/evals/sets/{set_id}/runs", + json={"model": active_model}, + headers=_bearer(tenant_a["key"]), ) assert launch.status_code == 403 assert launch.json()["error"]["code"] == "ERR_ZDR_PAYLOAD_BLOCKED" @@ -465,7 +484,9 @@ async def test_resume_does_not_rebill_terminal_cases( await _add_case(client, tenant_a["key"], set_id, f"r-{i}") run = await client.post( - f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + f"/v1/evals/sets/{set_id}/runs", + json={"model": active_model}, + headers=_bearer(tenant_a["key"]), ) assert run.status_code == 201, run.text run_id_wire = run.json()["id"] @@ -480,8 +501,11 @@ async def test_resume_does_not_rebill_terminal_cases( assert len(rec.records) == 3, "a resumed drive must not re-bill a terminal case" # exactly N result rows — the UNIQUE(run, case) idempotency held. assert ( - await _count(app, "SELECT count(*) FROM eval_case_results WHERE eval_run_id = :r", - r=uuid.UUID(_strip_prefix(run_id_wire))) + await _count( + app, + "SELECT count(*) FROM eval_case_results WHERE eval_run_id = :r", + r=uuid.UUID(_strip_prefix(run_id_wire)), + ) == 3 ) @@ -504,16 +528,22 @@ async def test_timeout_case_errored_run_continues( await _add_case(client, tenant_a["key"], set_id, "fast-2") run = await client.post( - f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + f"/v1/evals/sets/{set_id}/runs", + json={"model": active_model}, + headers=_bearer(tenant_a["key"]), ) assert run.status_code == 201, run.text cases = ( - await client.get(f"/v1/evals/runs/{run.json()['id']}/cases", headers=_bearer(tenant_a["key"])) + await client.get( + f"/v1/evals/runs/{run.json()['id']}/cases", headers=_bearer(tenant_a["key"]) + ) ).json()["data"] by_status = sorted(c["status"] for c in cases) assert by_status == ["completed", "completed", "errored"], by_status # the run still reaches a terminal rollup — one bad case never sinks it. - got = (await client.get(f"/v1/evals/runs/{run.json()['id']}", headers=_bearer(tenant_a["key"]))).json() + got = ( + await client.get(f"/v1/evals/runs/{run.json()['id']}", headers=_bearer(tenant_a["key"])) + ).json() assert got["status"] == "completed" assert got["counts"]["errored"] == 1 and got["counts"]["completed"] == 2 @@ -531,16 +561,22 @@ async def test_cross_tenant_and_absent_run_uniform_404( set_id = await _make_set(client, tenant_a["key"], "iso") await _add_case(client, tenant_a["key"], set_id, "x") run = await client.post( - f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + f"/v1/evals/sets/{set_id}/runs", + json={"model": active_model}, + headers=_bearer(tenant_a["key"]), ) run_id = run.json()["id"] # Tenant B asking for A's run. cross = await client.get(f"/v1/evals/runs/{run_id}", headers=_bearer(tenant_b["key"])) # A wholly-absent run id. - absent = await client.get(f"/v1/evals/runs/er_{uuid.uuid4().hex}", headers=_bearer(tenant_b["key"])) + absent = await client.get( + f"/v1/evals/runs/er_{uuid.uuid4().hex}", headers=_bearer(tenant_b["key"]) + ) assert cross.status_code == 404 and absent.status_code == 404 - assert cross.json() == absent.json(), "cross-tenant must be byte-identical to absent (no oracle)" + assert cross.json() == absent.json(), ( + "cross-tenant must be byte-identical to absent (no oracle)" + ) assert cross.json()["error"]["code"] == "ERR_EVAL_RUN_NOT_FOUND" # A malformed id resolves to the same 404, never a 500. malformed = await client.get("/v1/evals/runs/not-a-run", headers=_bearer(tenant_b["key"])) @@ -569,7 +605,9 @@ async def _run_case_ids() -> list[str]: ) assert run.status_code == 201, run.text data = ( - await client.get(f"/v1/evals/runs/{run.json()['id']}/cases", headers=_bearer(tenant_a["key"])) + await client.get( + f"/v1/evals/runs/{run.json()['id']}/cases", headers=_bearer(tenant_a["key"]) + ) ).json()["data"] return [c["eval_case_id"] for c in data] @@ -596,15 +634,21 @@ async def test_empty_set_run_completes_vacuously( set_id = await _make_set(client, tenant_a["key"], "empty") run = await client.post( - f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + f"/v1/evals/sets/{set_id}/runs", + json={"model": active_model}, + headers=_bearer(tenant_a["key"]), ) assert run.status_code == 201, run.text body = run.json() assert body["case_count"] == 0 - got = (await client.get(f"/v1/evals/runs/{body['id']}", headers=_bearer(tenant_a["key"]))).json() + got = ( + await client.get(f"/v1/evals/runs/{body['id']}", headers=_bearer(tenant_a["key"])) + ).json() assert got["status"] == "completed" assert got["counts"] == {"completed": 0, "refused": 0, "errored": 0, "pending": 0} - cases = (await client.get(f"/v1/evals/runs/{body['id']}/cases", headers=_bearer(tenant_a["key"]))).json() + cases = ( + await client.get(f"/v1/evals/runs/{body['id']}/cases", headers=_bearer(tenant_a["key"])) + ).json() assert cases["data"] == [] @@ -623,14 +667,18 @@ async def test_cross_tenant_launch_and_billing_identity( # Tenant B launching a run on tenant A's set -> uniform 404 (A's set is invisible to B). forbidden = await client.post( - f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_b["key"]) + f"/v1/evals/sets/{set_id}/runs", + json={"model": active_model}, + headers=_bearer(tenant_b["key"]), ) assert forbidden.status_code == 404 assert forbidden.json()["error"]["code"] == "ERR_EVAL_SET_NOT_FOUND" # Tenant A's own run bills A's tenant, never B's (the usage record carries the launcher). run = await client.post( - f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + f"/v1/evals/sets/{set_id}/runs", + json={"model": active_model}, + headers=_bearer(tenant_a["key"]), ) assert run.status_code == 201, run.text assert await _poll(lambda: len(rec.records) == 1) @@ -652,14 +700,18 @@ async def test_run_snapshot_fixed_at_launch( await _add_case(client, tenant_a["key"], set_id, "at-launch-2") run = await client.post( - f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + f"/v1/evals/sets/{set_id}/runs", + json={"model": active_model}, + headers=_bearer(tenant_a["key"]), ) assert run.status_code == 201, run.text assert run.json()["case_count"] == 2 # A case added after launch is NOT in the already-launched run's fixed denominator. await _add_case(client, tenant_a["key"], set_id, "after-launch-3") - got = (await client.get(f"/v1/evals/runs/{run.json()['id']}", headers=_bearer(tenant_a["key"]))).json() + got = ( + await client.get(f"/v1/evals/runs/{run.json()['id']}", headers=_bearer(tenant_a["key"])) + ).json() assert got["case_count"] == 2, "the run's case set is fixed at launch (A2)" @@ -685,7 +737,9 @@ async def test_bounded_per_tenant_concurrency( await _add_case(client, tenant_a["key"], set_id, f"c-{i}") run = await client.post( - f"/v1/evals/sets/{set_id}/runs", json={"model": active_model}, headers=_bearer(tenant_a["key"]) + f"/v1/evals/sets/{set_id}/runs", + json={"model": active_model}, + headers=_bearer(tenant_a["key"]), ) assert run.status_code == 201, run.text assert fake.calls == 6 From 3acc47a1542cf067621c31118afb0d3507a91c44 Mon Sep 17 00:00:00 2001 From: Tin Dang Date: Thu, 13 Aug 2026 13:44:40 +0700 Subject: [PATCH 3/3] test(evals): annotate the resume no-rebill negative wait MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit CI shard 2 tripped the repo-hygiene guard (tests/repo_hygiene/test_no_unbounded_positive_wait.py, ERR_UNBOUNDED_WAIT): test_resume_does_not_rebill_terminal_cases had a bare `await asyncio.sleep(0.05)` immediately before `assert upstream.calls == 3`. The guard reads a fixed sleep- then-assert as an unbounded positive wait. It is in fact a NEGATIVE wait: the resume drive is already awaited to completion, and the sleep only gives an ERRANT fire-and-forget re-dial/re-bill a window to land so the following assertions can prove it did NOT. A bounded poll_until can't express "prove nothing more ever happens" — it returns on the first satisfied tick. The real guarantee is skip-existing + UNIQUE(eval_run_id, eval_case_id); the wait only widens the window for a violation to surface. Annotate the sleep `# NEGATIVE WAIT: ` per the guard's own remedy. No production code changes; the 13 CHECKS and both lint halves stay green. Local verify: repo_hygiene guard 5 passed; evals_runs suite 13 passed; ruff check . + ruff format --check . clean. author: Tin Dang --- apps/gateway/tests/evals_runs/test_eval_run_executor.py | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/apps/gateway/tests/evals_runs/test_eval_run_executor.py b/apps/gateway/tests/evals_runs/test_eval_run_executor.py index 122e56f4..6ad7804f 100644 --- a/apps/gateway/tests/evals_runs/test_eval_run_executor.py +++ b/apps/gateway/tests/evals_runs/test_eval_run_executor.py @@ -496,6 +496,10 @@ async def test_resume_does_not_rebill_terminal_cases( # Resume the SAME run (all cases terminal) -> zero new dials, zero new bills. executor = app.state.eval_run_executor await executor.drive(uuid.UUID(_strip_prefix(run_id_wire)), raw_key=tenant_a["key"]) + # NEGATIVE WAIT: the drive above is already awaited to completion; this gives any ERRANT + # fire-and-forget re-dial/re-bill a window to land, then we assert it did NOT — a bounded + # poll can't express "prove nothing more ever happens" (it returns on the first satisfied + # tick). skip-existing + UNIQUE(run, case) is what actually guarantees the invariant. await asyncio.sleep(0.05) assert upstream.calls == 3, "a resumed drive must not re-dial a terminal case" assert len(rec.records) == 3, "a resumed drive must not re-bill a terminal case"