diff --git a/s16code/events/engine.py b/s16code/events/engine.py index ab32753..0e867a3 100644 --- a/s16code/events/engine.py +++ b/s16code/events/engine.py @@ -134,18 +134,25 @@ async def decide(subscription: Subscription) -> dict[str, Any]: self.store.add_decision(event.source, event.id, decision) return decision budget = admit.detail.get("effective_budget", subscription.budget) - result = await self.runtime.run( - prompt=goal, - scope=MemoryScope(subscription.tenant_id, subscription.project_id, - subscription.user_id, subscription.agent_id), - llm=llm, source_uri=f"event://{event.source}/{event.id}", - source_author=event.source, - allowed_side_effects=set(subscription.allowed_side_effects), - budget=budget, transport=transport, - 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) + reserved = float(admit.detail.get("reserved") or 0.0) + try: + result = await self.runtime.run( + prompt=goal, + scope=MemoryScope(subscription.tenant_id, subscription.project_id, + subscription.user_id, subscription.agent_id), + llm=llm, source_uri=f"event://{event.source}/{event.id}", + source_author=event.source, + allowed_side_effects=set(subscription.allowed_side_effects), + budget=budget, transport=transport, + initial_evidence={"event": event.model_dump(mode="json")}, + ) + except Exception: + # A failed run refunds its reservation but keeps the + # consumed slot; the refusal path records it below. + self.governor.settle_run(subscription.id, reserved=reserved, actual=0.0, now=now) + raise + self.governor.settle_run(subscription.id, reserved=reserved, + actual=float(result.get("spend_usd") or 0.0), 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/events/governor.py b/s16code/events/governor.py index 16f5e92..1e4f236 100644 --- a/s16code/events/governor.py +++ b/s16code/events/governor.py @@ -116,36 +116,34 @@ def admit_triage(self, subscription: Subscription, *, now: datetime | None = Non # ------------------------------------------------------------------- doing def admit_run(self, subscription: Subscription, *, now: datetime | None = None) -> Verdict: - """May a relevant decision actually become a run?""" + """May a relevant decision actually become a run? + + The slot is reserved atomically at admit time, so overlapping events + cannot all observe an empty ledger and all start runs. + """ moment = now or datetime.now(UTC) day = _day(moment) - - if subscription.max_runs_per_day is not None: - runs = self.store.window_count(subscription.id, day, kind="run") - if runs >= subscription.max_runs_per_day: - return Verdict( - False, - f"daily run ceiling reached: {runs} of {subscription.max_runs_per_day}", - "max_runs_per_day", - {"runs": runs, "limit": subscription.max_runs_per_day}, - ) - - if subscription.daily_budget is not None: - spent = self.store.window_spend(subscription.id, day, kind="run") - if spent >= subscription.daily_budget: - return Verdict( - False, - f"daily budget exhausted: ${spent:.8f} of ${subscription.daily_budget:.8f}", - "daily_budget", - {"spent_usd": spent, "limit_usd": subscription.daily_budget}, - ) - # Hand the run the smaller of its per-run ceiling and what remains - # in the window, so one run cannot consume the whole day. - remaining = subscription.daily_budget - spent - per_run = subscription.budget if subscription.budget is not None else remaining - return Verdict(True, "", "", {"effective_budget": min(per_run, remaining)}) - - return ADMITTED + ok, control, detail = self.store.window_try_admit_run( + subscription.id, day, + max_runs=subscription.max_runs_per_day, + budget=subscription.daily_budget, + per_run_budget=subscription.budget, + ) + if not ok: + if control == "max_runs_per_day": + return Verdict(False, + f"daily run ceiling reached: {detail['runs']} of {detail['limit']}", + control, detail) + return Verdict(False, + f"daily budget exhausted: ${detail['spent_usd']:.8f} of ${detail['limit_usd']:.8f}", + control, detail) + return Verdict(True, "", "", detail) + + def settle_run(self, subscription_id: str, *, reserved: float, actual: float = 0.0, + now: datetime | None = None) -> None: + """Replace an admitted run's reservation with its metered spend.""" + self.store.window_settle_run(subscription_id, _day(now or datetime.now(UTC)), + reserved=reserved, actual=actual) # ----------------------------------------------------------------- ledger diff --git a/s16code/events/routes.py b/s16code/events/routes.py index b5b292d..8bb5eb3 100644 --- a/s16code/events/routes.py +++ b/s16code/events/routes.py @@ -26,7 +26,7 @@ async def put_subscription(subscription_id: str, body: Subscription, request: Re return {"accepted": True, "subscription": stored} -@router.get("/subscriptions") +@router.get("/subscriptions", dependencies=[Depends(require_control)]) async def subscriptions(request: Request): return {"subscriptions": [item.model_dump(mode="json") for item in request.app.state.event_store.subscriptions()]} diff --git a/s16code/events/store.py b/s16code/events/store.py index 73b6558..103dfcb 100644 --- a/s16code/events/store.py +++ b/s16code/events/store.py @@ -217,6 +217,51 @@ def window_spend(self, subscription_id: str, day: str, *, kind: str) -> float: bucket = self._load()["windows"].get(f"{subscription_id}|{day}", {}) return float(bucket.get(kind, {}).get("usd", 0.0)) + def window_try_admit_run(self, subscription_id: str, day: str, *, + max_runs: int | None, budget: float | None, + per_run_budget: float | None) -> tuple[bool, str, dict[str, Any]]: + """Reserve one run slot atomically, or refuse. Returns (ok, control, detail). + + The check and the reservation share one lock section, so overlapping + events cannot all observe an empty ledger and all start runs. + """ + with self._lock: + state = self._load() + bucket = state["windows"].setdefault(f"{subscription_id}|{day}", {}) + entry = bucket.setdefault("run", {"count": 0, "usd": 0.0, "reserved": 0.0}) + if max_runs is not None and int(entry["count"]) >= max_runs: + return (False, "max_runs_per_day", + {"runs": int(entry["count"]), "limit": max_runs}) + committed = float(entry["usd"]) + float(entry.get("reserved", 0.0)) + if budget is not None and committed >= budget: + return (False, "daily_budget", + {"spent_usd": committed, "limit_usd": budget}) + entry["count"] = int(entry["count"]) + 1 + if budget is not None: + remaining = budget - committed + per_run = per_run_budget if per_run_budget is not None else remaining + effective = min(per_run, remaining) + entry["reserved"] = float(entry.get("reserved", 0.0)) + effective + else: + effective = per_run_budget + self._save(state) + return (True, "", {"effective_budget": effective, + "reserved": 0.0 if budget is None else effective}) + + def window_settle_run(self, subscription_id: str, day: str, *, + reserved: float, actual: float = 0.0) -> None: + """Replace an admitted run's reservation with its metered spend.""" + with self._lock: + state = self._load() + bucket = state["windows"].get(f"{subscription_id}|{day}") + if bucket is None or "run" not in bucket: + return + entry = bucket["run"] + entry["usd"] = round(entry["usd"] + max(0.0, float(actual)), 12) + entry["reserved"] = round( + max(0.0, float(entry.get("reserved", 0.0)) - max(0.0, float(reserved))), 12) + self._save(state) + # --------------------------------------------------------------- refusals def record_refusal(self, *, control: str, reason: str, event: dict[str, Any] | None = None, diff --git a/s16code/routes.py b/s16code/routes.py index ee09f81..9234723 100644 --- a/s16code/routes.py +++ b/s16code/routes.py @@ -339,19 +339,19 @@ async def complete_waiting_job(body: CompletionBody, request: Request): "run": result, "channel_delivery": channel_delivery} -@router.post("/facts") +@router.post("/facts", dependencies=[Depends(require_control)]) async def fact(body: FactBody, request: Request): return request.app.state.runtime.remember_fact(text=body.text, scope=body.scope(), source_uri=body.source_uri, source_author=body.source_author, principal=Principal("gateway", "gateway"), supersedes_id=body.supersedes_id) -@router.post("/documents") +@router.post("/documents", dependencies=[Depends(require_control)]) async def document(body: IndexBody, request: Request): return request.app.state.runtime.index_document(text=body.text, source_uri=body.source_uri, scope=body.scope(), source_author=body.source_author) -@router.post("/memory/search") +@router.post("/memory/search", dependencies=[Depends(require_control)]) async def memory_search(body: SearchBody, request: Request): hits = request.app.state.runtime.memory.recall(body.query, body.scope(), kinds=body.kinds, limit=body.limit) diff --git a/s16code/ui/routes.py b/s16code/ui/routes.py index 564f2f0..8dceb9a 100644 --- a/s16code/ui/routes.py +++ b/s16code/ui/routes.py @@ -20,10 +20,12 @@ import json from pathlib import Path -from fastapi import APIRouter, HTTPException, Request +from fastapi import APIRouter, Depends, HTTPException, Request from fastapi.responses import HTMLResponse, StreamingResponse from pydantic import BaseModel, Field +from s16code.auth import require_control + from .agui import run_data_model, state_snapshot, to_agui_event from .catalog import catalog_manifest from .hitl import PendingAction, decide_resume @@ -152,7 +154,7 @@ class ActionBody(BaseModel): pending_summary: str = "" -@router.post("/v1/action") +@router.post("/v1/action", dependencies=[Depends(require_control)]) async def action(body: ActionBody, request: Request): try: node = request.app.state.runtime.graph.snapshot(body.run_id).nodes[body.node_id] diff --git a/tests/test_autonomy_governor.py b/tests/test_autonomy_governor.py index 3b8a241..9255cf7 100644 --- a/tests/test_autonomy_governor.py +++ b/tests/test_autonomy_governor.py @@ -8,6 +8,7 @@ """ from __future__ import annotations +import asyncio import json from datetime import UTC, datetime, timedelta @@ -43,12 +44,15 @@ def _subscription(**changes) -> Subscription: class _Runtime: """A runtime that records whether it was ever asked to do work.""" - def __init__(self, spend: float = 0.001) -> None: + def __init__(self, spend: float = 0.001, delay: float = 0.0) -> None: self.runs: list[str] = [] self.spend = spend + self.delay = delay async def run(self, *, prompt: str, **_: object) -> dict[str, object]: self.runs.append(prompt) + if self.delay: + await asyncio.sleep(self.delay) return {"run_id": f"run-{len(self.runs)}", "status": "completed", "spend_usd": self.spend} @@ -188,6 +192,28 @@ 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_run_ceiling_holds_under_concurrent_events(tmp_path) -> None: + """Overlapping events must not all pass an admit_run that only records + after the run finishes. The slot is reserved atomically at admit time, so + with a ceiling of one, exactly one run starts and the rest are refused.""" + store = EventStore(tmp_path) + # The run must yield control between admit and record, or the coroutines + # never interleave and the race cannot show itself. + runtime = _Runtime(delay=0.01) + engine = AutonomousEventEngine(store, runtime) + store.put_subscription(_subscription(max_runs_per_day=1)) + llm = _relevance_llm() + + outcomes = await asyncio.gather(*(engine.process(_event(id=f"overlap-{i}"), llm=llm) + for i in range(4))) + + assert len(runtime.runs) == 1 + refusals = [item for item in store.refusals() if item["control"] == "max_runs_per_day"] + assert len(refusals) == 3 + acted = [o for o in outcomes if o.get("decisions")] + assert sum(1 for o in acted for d in o["decisions"] if d.get("acted")) == 1 + + # -------------------------------------------------------------------- liveness diff --git a/tests/test_control_plane_auth.py b/tests/test_control_plane_auth.py index cbbfcb0..7b5ed64 100644 --- a/tests/test_control_plane_auth.py +++ b/tests/test_control_plane_auth.py @@ -20,6 +20,10 @@ "occurred_at": "2026-08-05T09:00:00Z"}), ("post", "/v1/agent/runs", {"prompt": "hello", "tenant_id": "t"}), ("post", "/v1/agent/runs/run-1/resume", {}), + # Memory writes are write paths too: they inject evidence the agent will + # later treat as authorised. The README promises every write path is gated. + ("post", "/v1/agent/facts", {"tenant_id": "t", "text": "x", "source_uri": "u"}), + ("post", "/v1/agent/documents", {"tenant_id": "t", "text": "x", "source_uri": "u"}), ] @@ -84,6 +88,44 @@ def test_read_only_observability_is_available_to_an_operator(app_client) -> None assert "NOT ALIVE" in markdown.text +def test_the_subscription_read_path_is_gated_like_its_write_path( + app_client, monkeypatch +) -> None: + """GET /subscriptions returns the authority objects (side effects, budgets, + instructions). Reading them is as sensitive as writing them, so the read + path must fail closed exactly like PUT /subscriptions/{id} does. + """ + monkeypatch.delenv("S16_CONTROL_TOKEN", raising=False) + unset = app_client.get("/v1/agent/subscriptions") + assert unset.status_code == 503 + assert "S16_CONTROL_TOKEN" in unset.json()["detail"] + + monkeypatch.setenv("S16_CONTROL_TOKEN", conftest.CONTROL_TOKEN) + wrong = app_client.get("/v1/agent/subscriptions", + headers={"Authorization": "Bearer wrong"}) + assert wrong.status_code == 401 + + +def test_the_action_route_is_a_write_path_and_fails_closed(app_client, monkeypatch) -> None: + """POST /v1/action resumes a parked node and re-runs it, spending money. + + It is a control-plane write path (the same resume the protected + /v1/agent/completions and /v1/agent/runs/{id}/resume routes perform), so it + must fail closed like every other write path. + """ + body = {"run_id": "run-x", "node_id": "n", "action": "approve"} + + monkeypatch.delenv("S16_CONTROL_TOKEN", raising=False) + unset = app_client.post("/v1/action", json=body) + assert unset.status_code == 503 + assert "S16_CONTROL_TOKEN" in unset.json()["detail"] + + monkeypatch.setenv("S16_CONTROL_TOKEN", conftest.CONTROL_TOKEN) + wrong = app_client.post("/v1/action", json=body, + headers={"Authorization": "Bearer wrong"}) + assert wrong.status_code == 401 + + def test_the_operator_console_is_served_and_is_read_only(app_client) -> None: """ยง9 promises an operations page. It has to exist, and it has to be inert.