Skip to content
Open
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
89 changes: 82 additions & 7 deletions s16code/events/engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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:
Expand All @@ -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
Expand Down Expand Up @@ -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", "")))
Expand Down Expand Up @@ -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
Expand Down
5 changes: 5 additions & 0 deletions s16code/runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -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", []))}
Expand Down
102 changes: 100 additions & 2 deletions tests/test_autonomy_governor.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,8 @@

import pytest

from s16code.economics.budget import RunBudget
from s16code.economics.config import EconomicsConfig
from s16code.events import (
AutonomousEventEngine,
AutonomyGovernor,
Expand Down Expand Up @@ -41,15 +43,25 @@ 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] = []
self.spend = spend

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):
Expand Down Expand Up @@ -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


Expand Down