From 8ece77d930d38f28dabc4b6655f5c3109f521951 Mon Sep 17 00:00:00 2001 From: ranjitha13g Date: Thu, 13 Aug 2026 19:33:57 +0530 Subject: [PATCH] Meter what a run and a gate actually cost, so the window ceilings hold daily_budget and daily_triage_budget could never fire: both accumulated 0.00 forever, and the morning report's two cost lines were permanently zero. The governor was correct; the engine fed it zero, because it read money from fields only its own test doubles wrote. Three facets, one root cause. engine.py read result['spend_usd'], but AgentRuntime.run() reports budget.spent and never emitted that key. _spend_of read reply['metered_calls'], but GatewayClient.complete -- the callable the HTTP event route passes in -- carries no meter at all. And it looked for cost_usd/usd, while BudgetedGateway writes Charge.as_dict(), whose money field is cost. Fix: runtime.run() publishes spend_usd at the top level; the engine falls back to budget.spent; _spend_of also reads cost, and prices an unmetered reply's own tokens against config/pricing.yaml via the existing Pricing.cost(). No new price appears in Python. The runtime double now builds its return from a real RunBudget.snapshot() instead of a hand-written dict, so this class of drift fails the suite rather than passing it. Reverting the two source files fails 5 tests, two of which already shipped green. 357 passed, up from 353. Co-Authored-By: Claude Opus 5 (1M context) --- s16code/events/engine.py | 89 +++++++++++++++++++++++++--- s16code/runtime.py | 5 ++ tests/test_autonomy_governor.py | 102 +++++++++++++++++++++++++++++++- 3 files changed, 187 insertions(+), 9 deletions(-) diff --git a/s16code/events/engine.py b/s16code/events/engine.py index ab32753..2e1b176 100644 --- a/s16code/events/engine.py +++ b/s16code/events/engine.py @@ -3,10 +3,13 @@ import asyncio import fnmatch import json +import os from datetime import UTC, datetime +from functools import lru_cache from typing import Any, Awaitable, Callable from s16code.core.memory import MemoryScope +from s16code.economics.pricing import Pricing from .governor import AutonomyGovernor from .models import EventEnvelope, Subscription @@ -25,17 +28,84 @@ def _json_object(text: str) -> dict[str, Any]: return value -def _spend_of(reply: dict[str, Any]) -> float: - """What the gate itself cost. Deciding not to act is not free.""" +def _spend_of(reply: dict[str, Any], pricing: Pricing | None = None) -> float: + """What the gate itself cost. Deciding not to act is not free. + + Two shapes reach here. A reply that came through ``BudgetedGateway`` carries + per-call meter records under ``metered_calls``; the field they price under is + ``cost`` (``Charge.as_dict``), so that name has to be read or the controller's + own records total zero. A reply that came straight off ``GatewayClient`` + carries no meter at all — only the token counts — so it is priced here + against the same table the rest of the system bills with. Returning 0.0 + because nobody attached a meter is how ``daily_triage_budget`` silently stops + being a ceiling. + """ + metered = False total = 0.0 for call in reply.get("metered_calls", []) or []: if not isinstance(call, dict): continue + for key in ("cost_usd", "usd", "cost"): + if key in call: + try: + total += float(call[key] or 0.0) + except (TypeError, ValueError): + continue + metered = True + break + if metered: + return total + if pricing is None: + return 0.0 + return pricing.cost( + reply.get("model"), + input_tokens=int(reply.get("input_tokens") or 0), + output_tokens=int(reply.get("output_tokens") or 0), + cache_read_tokens=int(reply.get("cache_read_input_tokens") or 0), + cache_write_tokens=int(reply.get("cache_creation_input_tokens") or 0), + ) + + +def _run_spend_of(result: dict[str, Any]) -> float: + """What a dispatched run cost, from whichever field the runtime reported. + + ``AgentRuntime.run`` publishes ``spend_usd`` directly; it also returns the + whole ``RunBudget.snapshot()``, whose ``spent`` is the same number. Reading + only one of them is what made ``daily_budget`` unenforceable: the engine + asked for a key no production runtime emitted, got ``None`` every time, and + recorded a window spend of zero forever. + """ + for candidate in (result.get("spend_usd"), (result.get("budget") or {}).get("spent")): + if candidate is None: + continue try: - total += float(call.get("cost_usd") or call.get("usd") or 0.0) + return float(candidate) except (TypeError, ValueError): continue - return total + return 0.0 + + +@lru_cache(maxsize=4) +def _pricing_for(config_dir: str | None) -> Pricing | None: + """The price table for one config directory, read once. + + Keyed on the directory rather than cached outright, because ``S16_CONFIG_DIR`` + picks the table and a cache that ignored it would serve one profile's prices + to another. + """ + try: + from s16code.economics.config import EconomicsConfig + + return EconomicsConfig.load(config_dir).pricing + except Exception: # pragma: no cover - config ships with every checkout + # A price table we cannot read must not stop the agent noticing events. + # It costs the dollar ceilings, not the loop, and _spend_of falls back + # to whatever meter the reply carried. + return None + + +def _default_pricing() -> Pricing | None: + return _pricing_for(os.getenv("S16_CONFIG_DIR") or None) class AutonomousEventEngine: @@ -46,10 +116,15 @@ class AutonomousEventEngine: """ def __init__(self, store: EventStore, runtime: Any, *, max_concurrent: int = 4, - governor: AutonomyGovernor | None = None) -> None: + governor: AutonomyGovernor | None = None, + pricing: Pricing | None = None) -> None: self.store, self.runtime = store, runtime self.max_concurrent = max_concurrent self.governor = governor or AutonomyGovernor(store) + # The gate is a model call on every matching event. Pricing it needs the + # same table the runs are billed against, or "cost of watching" is a + # number nobody computes and daily_triage_budget bounds nothing. + self.pricing = pricing if pricing is not None else _default_pricing() # One semaphore for the engine, not one per call. Constructing it inside # process() bounded fan-out across subscriptions for a single event and # bounded nothing at all across concurrent events, which is exactly the @@ -110,7 +185,7 @@ async def decide(subscription: Subscription) -> dict[str, Any]: ) reply = await llm(json.dumps({"subscription": subscription.model_dump(mode="json"), "event": event.model_dump(mode="json")}), system) - triage_cost = _spend_of(reply) + triage_cost = _spend_of(reply, self.pricing) self.governor.record(subscription.id, kind="triage", usd=triage_cost, now=now) raw = _json_object(str(reply.get("text", ""))) @@ -145,7 +220,7 @@ async def decide(subscription: Subscription) -> dict[str, Any]: initial_evidence={"event": event.model_dump(mode="json")}, ) self.governor.record(subscription.id, kind="run", - usd=float(result.get("spend_usd") or 0.0), now=now) + usd=_run_spend_of(result), now=now) decision.update(run_id=result["run_id"], run_status=result["status"], acted=True) self.store.add_decision(event.source, event.id, decision) return decision diff --git a/s16code/runtime.py b/s16code/runtime.py index b8fcef7..32615d0 100644 --- a/s16code/runtime.py +++ b/s16code/runtime.py @@ -963,6 +963,11 @@ async def execute_once(task: TaskSpec) -> dict[str, Any] | Deferred: "edges": snapshot.edges}, "trace": {"planner": getattr(planner, "last_selection", {"mode": "deterministic"}), "agents": trace}, "events": [event.__dict__ for event in self.graph.events(run_id)], "principal": who, + # What this run cost, at the top level. The windowed governor in + # s16code.events meters against this, and a number reachable only + # by digging into a nested snapshot is a number a caller forgets + # to read. + "spend_usd": run_budget.spent if run_budget is not None else 0.0, "budget": run_budget.snapshot() if run_budget is not None else None, "economics": economics_config.describe() if economics_config is not None else None, "allocations": list(getattr(planner, "allocations", []))} diff --git a/tests/test_autonomy_governor.py b/tests/test_autonomy_governor.py index 3b8a241..0b7866c 100644 --- a/tests/test_autonomy_governor.py +++ b/tests/test_autonomy_governor.py @@ -13,6 +13,8 @@ import pytest +from s16code.economics.budget import RunBudget +from s16code.economics.config import EconomicsConfig from s16code.events import ( AutonomousEventEngine, AutonomyGovernor, @@ -41,7 +43,14 @@ def _subscription(**changes) -> Subscription: class _Runtime: - """A runtime that records whether it was ever asked to do work.""" + """A runtime that records whether it was ever asked to do work. + + The returned dict is built from a real :class:`RunBudget` snapshot rather + than hand-written, because a double that invents its own shape cannot fail + when the seam it stands in for drifts. The engine has to be able to read + what ``AgentRuntime.run`` genuinely emits (runtime.py, the ``budget`` key), + not what a test wishes it emitted. + """ def __init__(self, spend: float = 0.001) -> None: self.runs: list[str] = [] @@ -49,7 +58,10 @@ def __init__(self, spend: float = 0.001) -> None: async def run(self, *, prompt: str, **_: object) -> dict[str, object]: self.runs.append(prompt) - return {"run_id": f"run-{len(self.runs)}", "status": "completed", "spend_usd": self.spend} + budget = RunBudget(total=max(self.spend, 0.01), run_id=f"run-{len(self.runs)}") + budget.spent = self.spend + return {"run_id": f"run-{len(self.runs)}", "status": "completed", "answer": "", + "budget": budget.snapshot()} def _relevance_llm(relevant: bool = True, cost: float = 0.00002): @@ -188,6 +200,92 @@ async def test_a_daily_budget_caps_the_window_and_shrinks_the_last_run(tmp_path) assert any(item["control"] == "daily_budget" for item in store.refusals()) +async def test_the_daily_budget_holds_when_the_runtime_reports_only_a_snapshot(tmp_path) -> None: + """A window ceiling must read the field the runtime actually publishes. + + ``AgentRuntime.run`` returns the whole ``RunBudget.snapshot()``. A governor + that meters some other key records ``$0.00`` for every run, and then + ``spent >= daily_budget`` is a comparison that can never be true no matter + how much the agent spends. + """ + store = EventStore(tmp_path) + + class _SnapshotOnlyRuntime: + def __init__(self) -> None: + self.runs: list[str] = [] + + async def run(self, *, prompt: str, **_: object) -> dict[str, object]: + self.runs.append(prompt) + budget = RunBudget(total=0.01, run_id=f"run-{len(self.runs)}") + budget.spent = 0.004 + return {"run_id": budget.run_id, "status": "completed", + "budget": budget.snapshot()} + + runtime = _SnapshotOnlyRuntime() + engine = AutonomousEventEngine(store, runtime) + store.put_subscription(_subscription(budget=0.01, daily_budget=0.01)) + + for index in range(6): + await engine.process(_event(id=f"snap-{index}"), llm=_relevance_llm()) + + day = datetime.now(UTC).date().isoformat() + assert store.window_spend("watch", day, kind="run") == pytest.approx(0.012) + assert len(runtime.runs) == 3 + assert any(item["control"] == "daily_budget" for item in store.refusals()) + + +# --------------------------------------------------------- the cost of watching + + +async def test_the_triage_ceiling_holds_when_the_gateway_reports_only_tokens(tmp_path) -> None: + """Being awake is billed even when nobody attached a meter to the reply. + + ``GatewayClient.complete`` — the callable the HTTP event route hands the + engine — returns text, provider, model and token counts. It carries no + ``metered_calls``. If that shape prices at zero then ``daily_triage_budget`` + never trips, and the one control the session says must survive a + denial-of-wallet attack is the unmetered surface the attack aims at. + """ + store = EventStore(tmp_path) + engine = AutonomousEventEngine(store, _Runtime()) + # gemini-3.1-flash-lite is $0.25/$1.50 per Mtok in config/pricing.yaml, so + # each gate call below prices at $0.0013 and the fifth exhausts the ceiling. + store.put_subscription(_subscription(daily_triage_budget=0.005)) + + calls: list[str] = [] + + async def gateway_shaped(prompt: str, system: str): # noqa: ARG001 + calls.append(prompt) + return {"text": json.dumps({"relevant": False, "reason": "routine", "goal": ""}), + "provider": "gemini", "model": "gemini-3.1-flash-lite", + "input_tokens": 4_000, "output_tokens": 200} + + for index in range(40): + await engine.process(_event(id=f"awake-{index}"), llm=gateway_shaped) + + day = datetime.now(UTC).date().isoformat() + assert store.window_spend("watch", day, kind="triage") > 0.0 + assert len(calls) < 40, "the triage ceiling never stopped a single gate call" + assert any(item["control"] == "daily_triage_budget" for item in store.refusals()) + + +def test_a_controller_meter_record_is_priced_under_its_own_key() -> None: + """``BudgetedGateway`` writes ``Charge.as_dict()``, whose money field is + ``cost``. Reading only ``cost_usd``/``usd`` totals a real meter at zero.""" + from s16code.events.engine import _spend_of + + assert _spend_of({"metered_calls": [{"cost": 0.0025}]}) == pytest.approx(0.0025) + assert _spend_of({"metered_calls": [{"cost_usd": 0.0025}]}) == pytest.approx(0.0025) + + +def test_the_shipped_price_table_prices_a_gateway_reply() -> None: + """The engine's default pricing has to be able to price a real reply, or the + fallback above is decorative.""" + pricing = EconomicsConfig.load().pricing + + assert pricing.cost("gemini-3.1-flash-lite", input_tokens=4_000, output_tokens=200) > 0.0 + + # -------------------------------------------------------------------- liveness