Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 6 additions & 3 deletions .dockerignore
Original file line number Diff line number Diff line change
Expand Up @@ -35,9 +35,12 @@
!prompts/quality-criteria.md
!prompts/workflow-prompts/impl-repair-claude.md
!prompts/workflow-prompts/ai-quality-review.md
# Promoted feedback cases hold users' data; they are synced here for eval runs
# and are never committed, so they must never reach an image either.
agents/evals/.cases
# The eval harness (agents/evals/) is not service code: its fixture cases are
# about 3 MB, its local reports hold renders and galleries, and its promoted
# feedback cases (agents/evals/.cases) hold users' data, which must never reach
# an image. The development-only fixture seed (AGENT_DEV_FIXTURE) reads the
# cases from a checkout, never from an image.
agents/evals

# …and the noise that lives inside the trees above. `__pycache__` matters twice
# over now that the image compiles its own bytecode: a stale host `.pyc` copied
Expand Down
6 changes: 6 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -286,3 +286,9 @@ ds-bundle/
app/.designsync-fonts/*.ttf
app/.designsync-fonts/*.otf
app/.designsync-fonts/*.woff2

# Agent network regression harness (agents/evals/): run reports, renders and review
# galleries are local; promoted feedback cases hold users' data and are synced from the
# private bucket, never committed. Baselines in agents/evals/baselines/ ARE committed.
agents/evals/reports/
agents/evals/.cases/
150 changes: 147 additions & 3 deletions agents/README.md

Large diffs are not rendered by default.

36 changes: 31 additions & 5 deletions agents/anyplot/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -266,13 +266,20 @@ def make_content_config(kind: ModelKind, settings: AgentSettings | None = None)


class JudgeVerdict(BaseModel):
"""What the judge decided, the reply language it saw, and the tokens it cost."""
"""What the judge decided, the reply language it saw, and the tokens it cost.

`tokens` is the total the budgets count; `input_tokens` and `output_tokens` split
it for the attribution log, so the eval harness can price the judge at the input
and output rates (output includes thinking tokens on Gemini).
"""

model_config = ConfigDict(extra="ignore")

verdict: Verdict
lang: str = "en"
tokens: int = Field(default=0, ge=0)
input_tokens: int = Field(default=0, ge=0)
output_tokens: int = Field(default=0, ge=0)

@field_validator("lang", mode="before")
@classmethod
Expand Down Expand Up @@ -353,8 +360,17 @@ async def _attempt(self, text: str, rubric: str) -> JudgeVerdict:
if not isinstance(answer, dict):
raise ValueError("the judge answered without a verdict")
usage = message.usage
tokens = int(getattr(usage, "input_tokens", 0) or 0) + int(getattr(usage, "output_tokens", 0) or 0)
return JudgeVerdict(**_JudgeAnswer.model_validate(answer).model_dump(), tokens=tokens)
input_tokens = sum(
int(getattr(usage, name, 0) or 0)
for name in ("input_tokens", "cache_read_input_tokens", "cache_creation_input_tokens")
)
output_tokens = int(getattr(usage, "output_tokens", 0) or 0)
return JudgeVerdict(
**_JudgeAnswer.model_validate(answer).model_dump(),
tokens=input_tokens + output_tokens,
input_tokens=input_tokens,
output_tokens=output_tokens,
)

async def judge(self, text: str, rubric: str) -> JudgeVerdict:
return await _within_budget(self._attempt, text, rubric, self.timeout_s)
Expand Down Expand Up @@ -395,8 +411,18 @@ async def _attempt(self, text: str, rubric: str) -> JudgeVerdict:
)
answer = _JudgeAnswer.model_validate_json(response.text or "")
usage = response.usage_metadata
tokens = int(getattr(usage, "total_token_count", 0) or 0) if usage else 0
return JudgeVerdict(**answer.model_dump(), tokens=tokens)

def count(name: str) -> int:
return int(getattr(usage, name, 0) or 0) if usage else 0

input_tokens = count("prompt_token_count") + count("tool_use_prompt_token_count")
output_tokens = count("candidates_token_count") + count("thoughts_token_count")
return JudgeVerdict(
**answer.model_dump(),
tokens=count("total_token_count"),
input_tokens=input_tokens,
output_tokens=output_tokens,
)

async def judge(self, text: str, rubric: str) -> JudgeVerdict:
return await _within_budget(self._attempt, text, rubric, self.timeout_s)
Expand Down
133 changes: 117 additions & 16 deletions agents/anyplot/pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,16 @@
`failed` names its reason; `not_ready` comes before any model call. Progress goes out
as content-free events with `custom_metadata={"anyplot_status": {"step", "attempt"}}`,
which the stream translator turns into `status` events and no model ever reads.

Every step also writes one content-free attribution line (`plugins/ledger.attribution`):
`pipeline_adapt` (the answer's outcome and edit count), `pipeline_check` (edit-apply
failure count, blocking validator rule ids, ADAPTATION rule ids), `pipeline_render`
(render and wall time, the gate ids that failed or reported), `pipeline_review` (the
verdict and its defect ids) and `pipeline_result` (status, reason, attempts, whether
the shipped render was padded or carries ADAPTATION or probe findings, and the
exception class behind reason `error`). The
eval harness (`agents/evals/matrix.py`) reads them per case; they never carry code,
data, column names or model text.
"""

import json
Expand All @@ -66,7 +76,7 @@
from .code.validate import validate_adaptation, validate_security
from .data.bindings import check_bindings
from .data.store import StoredDataset
from .plugins.ledger import budget_allows, ledger_for
from .plugins.ledger import RequestLedger, attribution, budget_allows, ledger_for
from .policy import DATA_PREAMBLE, fence
from .render.contract import RendererUnavailable, RenderJob, Theme
from .render.gates import data_rows, evaluate
Expand Down Expand Up @@ -145,6 +155,10 @@ class Candidate:
advisory_lines: list[str] = field(default_factory=list)
reviewed_ok: bool | None = None
render_id: str | None = None
adaptation_rules: list[str] = field(default_factory=list)
"""The ADAPTATION validator rule ids behind `adaptation_lines`, for the attribution log."""
advisory_gates: list[str] = field(default_factory=list)
"""The advisory probe gate ids (`G3`, `G5`, `G7`, `G8`) behind `advisory_lines`, for the attribution log."""


@dataclass
Expand All @@ -162,6 +176,8 @@ class Run:
reason: FailureReason | None = None
rendered: bool = False
unreviewed_why: str | None = None
error_type: str | None = None
"""The exception class that ended the run with reason `error` (content-free, for the attribution log)."""


def _line(text: str) -> str:
Expand Down Expand Up @@ -264,31 +280,42 @@ def _clip_feedback(lines: list[str]) -> list[str]:
return unique[:MAX_FEEDBACK]


def _check(working: str, *, library: str, palette: list[str], base: str) -> tuple[bool, list[str], list[str]]:
"""Validate a working form against its base: (may render, blocking lines, adaptation defect lines)."""
@dataclass(frozen=True)
class CheckResult:
"""The validator verdict on a working form: its lines for the repair, and its rule ids for the attribution log."""

may_render: bool
blocking: list[str]
defects: list[str]
blocking_rules: list[str]
adaptation_rules: list[str]


def _check(working: str, *, library: str, palette: list[str], base: str) -> CheckResult:
"""Validate a working form against its base: may it render, the blocking lines, the adaptation defect lines."""
security = validate_security(working, library=library)
adaptation = validate_adaptation(working, original_palette=palette)
blocking = [f"validator {f.rule}" + (f" (line {f.line})" if f.line else "") + f": {f.message}" for f in security]
blocking += [
f"validator {f.rule}" + (f" (line {f.line})" if f.line else "") + f": {f.message}"
for f in adaptation
if f.rule in BLOCKING_RULES
blocking_findings = [*security, *(f for f in adaptation if f.rule in BLOCKING_RULES)]
blocking = [
f"validator {f.rule}" + (f" (line {f.line})" if f.line else "") + f": {f.message}" for f in blocking_findings
]
blocking_rules = [f.rule for f in blocking_findings]
added = new_literal_chars(base, working)
if added > MAX_NEW_LITERAL_CHARS:
blocking_rules.append("literal-budget")
blocking.append(
f"validator literal-budget: the plan adds {added} characters of new string literals, "
f"{added - MAX_NEW_LITERAL_CHARS} over the limit of {MAX_NEW_LITERAL_CHARS} → derive text and values "
"from df instead of writing them out. Likely cause: data or long text written into the code."
)
soft = [f for f in adaptation if f.rule not in BLOCKING_RULES]
defects = [
f"{ADAPTATION_IDS.get(f.rule, 'SC-03')} (code): {f.message}"
+ (f" at line {f.line}" if f.line else "")
+ f" → derive it from df and the Imprint palette. Likely cause: the adaptation ({f.rule})."
for f in adaptation
if f.rule not in BLOCKING_RULES
for f in soft
]
return not blocking, blocking, defects
return CheckResult(not blocking, blocking, defects, blocking_rules, [f.rule for f in soft])


def finish(run: Run) -> PlotResult:
Expand Down Expand Up @@ -420,6 +447,7 @@ async def run_pipeline(ctx: Context, node_input: PipelineArgs) -> AsyncGenerator
except Exception as exc: # every failure ends in a PlotResult; the type is logged, never the message
logger.warning("pipeline failed: %s", type(exc).__name__)
run.reason = "error"
run.error_type = type(exc).__name__
run.unreviewed_why = run.unreviewed_why or "an internal error stopped the run"
finally:
ledger.pipeline_active = False
Expand All @@ -431,10 +459,37 @@ async def run_pipeline(ctx: Context, node_input: PipelineArgs) -> AsyncGenerator
except Exception as exc: # a result whose artifacts cannot be stored is not shippable
logger.warning("storing the version failed: %s", type(exc).__name__)
result = PlotResult(status="failed", reason="error", attempts=run.attempts)
run.error_type = run.error_type or type(exc).__name__
_attribute_result(ledger, run, result)
yield close_scope(scope)
yield _result_event(result)


def _attribute_result(ledger: RequestLedger, run: Run, result: PlotResult) -> None:
"""The content-free outcome line of one pipeline run: what the eval harness's pass and gate counts read.

`padded`, `adaptation` (ADAPTATION rule ids) and `advisory` (probe gate ids) describe
the shipped render; a failed run has none. `error` is the exception class behind
reason `error` (for example `RendererUnavailable`, which the harness treats as an
outage rather than a model failure), never its message.
"""
shipped = (run.best or run.padded) if result.status in ("ok", "needs_attention") else None
attribution(
"pipeline_result",
ledger,
status=result.status,
reason=result.reason,
attempts=result.attempts,
theme=run.theme,
reviewed=run.reviewed,
padded=bool(shipped and shipped.padded),
adaptation=list(shipped.adaptation_rules) if shipped else [],
advisory=list(shipped.advisory_gates) if shipped else [],
residual=len(result.residual_defects),
error=run.error_type if result.reason == "error" else None,
)


def close_scope(scope: str) -> Event:
"""The event that ends the sub-agents' isolation scope.

Expand Down Expand Up @@ -494,23 +549,44 @@ async def _attempts(
ledger.adapter_allow_full = request.allow_full
answer = await _adapt(ctx, scope, render_adapt_request(request, view), view.library)
if isinstance(answer, str):
attribution("pipeline_adapt", ledger, attempt=attempt, outcome="schema")
feedback, previous_plan = [answer], None
continue
plan = answer
if plan.full_code is not None and not request.allow_full:
attribution("pipeline_adapt", ledger, attempt=attempt, outcome="full_code_refused")
feedback, previous_plan = ["full_code is allowed only on the second attempt; send edits instead"], None
continue
attribution(
"pipeline_adapt",
ledger,
attempt=attempt,
outcome="plan",
edits=len(plan.edits),
full_code=plan.full_code is not None,
)

yield _status("checking", attempt)
applied = apply_plan(working, plan)
if applied.code is None:
attribution(
"pipeline_check", ledger, attempt=attempt, outcome="edits_failed", edit_failures=len(applied.failures)
)
feedback, previous_plan = applied.failures, plan
continue
may_render, blocking, adaptation_lines = _check(
applied.code, library=view.library, palette=palette, base=working
check = _check(applied.code, library=view.library, palette=palette, base=working)
adaptation_lines = check.defects
attribution(
"pipeline_check",
ledger,
attempt=attempt,
outcome="ok" if check.may_render else "rejected",
edit_failures=0,
validator=check.blocking_rules,
adaptation=check.adaptation_rules,
)
if not may_render:
feedback, previous_plan = blocking, plan
if not check.may_render:
feedback, previous_plan = check.blocking, plan
continue
working, previous_plan = applied.code, None
try:
Expand All @@ -521,6 +597,7 @@ async def _attempts(
parse_dates=list(dataset.parse_dates),
)
except ValueError as exc:
attribution("pipeline_check", ledger, attempt=attempt, outcome="loader_failed", edit_failures=0)
feedback = [_line(f"the code could not take the data loader: {exc}")]
continue

Expand All @@ -534,7 +611,21 @@ async def _attempts(
themes=(run.theme,),
timeout_s=deadline.clamp(settings.render_timeout_s),
)
report = evaluate(await services.backend.render(job), themes=job.themes, library=view.library, rows=rows)
started = time.monotonic()
rendered = await services.backend.render(job)
render_s = time.monotonic() - started
report = evaluate(rendered, themes=job.themes, library=view.library, rows=rows)
attribution(
"pipeline_render",
ledger,
attempt=attempt,
theme=run.theme,
render_s=round(render_s, 3),
wall_s={theme: round(output.wall_s, 3) for theme, output in rendered.outputs.items()},
passed=report.passed_host_gates,
canvas_ok=report.canvas_ok,
gates=report.failed_gates,
)
run.rendered = True
if not report.passed_host_gates:
feedback = [*report.blocking, *adaptation_lines]
Expand All @@ -549,6 +640,8 @@ async def _attempts(
canvas_line=report.canvas_defects[0] if report.canvas_defects else None,
adaptation_lines=adaptation_lines,
advisory_lines=list(report.advisory),
adaptation_rules=check.adaptation_rules,
advisory_gates=[gate for gate in report.failed_gates if gate.startswith("G")],
)
if report.canvas_ok:
run.best = candidate
Expand Down Expand Up @@ -583,8 +676,16 @@ async def _attempts(
),
)
if verdict is None:
attribution("pipeline_review", ledger, attempt=attempt, verdict="unreadable")
run.unreviewed_why = "the review answer could not be read"
return
attribution(
"pipeline_review",
ledger,
attempt=attempt,
verdict="ok" if verdict.ok else "defects",
defects=[defect.id for defect in verdict.defects],
)
run.reviewed = True
candidate.reviewed_ok = verdict.ok
if verdict.ok:
Expand Down
7 changes: 5 additions & 2 deletions agents/anyplot/plugins/budget.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,9 @@
`PlotResult(failed, reason=budget)` instead.
* `after_model_callback` books every response's tokens (prompt + candidates +
thoughts + tool-use prompt; cached tokens are counted separately and never twice),
the call and the `model_version`, and writes one content-free attribution line.
the call and the `model_version`, and writes one content-free attribution line
with the token counts by kind (`ledger.usage_breakdown`), which the eval harness
prices (`agents/evals/pricing.py`).
* `before_tool_callback` on `plot_pipeline` counts the user's daily pipeline runs
and refuses the call with `{"status": "error", "code": "budget"}` past
`AGENT_DAILY_PIPELINE_RUNS`. A second call in the same invocation, and a call
Expand All @@ -32,7 +34,7 @@
from ..policy import refusal
from ..services import get_services
from ..settings import get_settings
from .ledger import attribution, budget_allows, ledger_for, usage_tokens
from .ledger import attribution, budget_allows, ledger_for, usage_breakdown, usage_tokens
from .tool_safety import ToolSafetyPlugin


Expand Down Expand Up @@ -92,6 +94,7 @@ async def after_model_callback(
agent=callback_context.agent_name,
billable=billable,
cached=cached,
**usage_breakdown(llm_response.usage_metadata),
model_version=llm_response.model_version,
finish_reason=str(llm_response.finish_reason) if llm_response.finish_reason else None,
)
Expand Down
23 changes: 23 additions & 0 deletions agents/anyplot/plugins/ledger.py
Original file line number Diff line number Diff line change
Expand Up @@ -176,6 +176,29 @@ def count(name: str) -> int:
return billable, count("cached_content_token_count")


USAGE_FIELDS: dict[str, str] = {
"prompt": "prompt_token_count",
"candidates": "candidates_token_count",
"thoughts": "thoughts_token_count",
"tool_use_prompt": "tool_use_prompt_token_count",
"cache_write": "cache_creation_input_tokens",
}
"""The `usage_metadata` counts the attribution line carries by kind (`cached` is logged on its own).

`prompt` includes the cached and cache-write tokens (ADK folds Anthropic's disjoint counts
into one prompt count); `cache_write` is the Anthropic cache-creation count ADK attaches to
the usage object. The eval harness prices each kind at its own rate."""


def usage_breakdown(usage: Any) -> dict[str, int]:
"""The token counts of `usage_metadata` by kind (`USAGE_FIELDS`); a missing count is 0."""
counts: dict[str, int] = {}
for key, name in USAGE_FIELDS.items():
value = getattr(usage, name, None) if usage is not None else None
counts[key] = int(value) if isinstance(value, int) else 0
return counts


def argument_hash(arguments: Any) -> str:
"""A short, stable hash of tool arguments, so logs can correlate calls without their content."""
try:
Expand Down
Loading
Loading