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
4 changes: 3 additions & 1 deletion control_plane/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -344,7 +344,9 @@ async def get_run_record(
rec = gov.get_run(run_id)
if rec is None:
return JSONResponse({"error": "not found"}, status_code=404)
return JSONResponse(run_to_dict(rec))
body = run_to_dict(rec)
body["policy_stats"] = gov.get_policy_stats(principal.tenant_id, run_id)
return JSONResponse(body)

@app.get("/v1/run-records")
async def list_run_records(
Expand Down
39 changes: 39 additions & 0 deletions control_plane/store.py
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,13 @@
applied_at REAL NOT NULL,
PRIMARY KEY (tenant_id, idempotency_key)
);
CREATE TABLE IF NOT EXISTS run_policy_stats (
tenant_id TEXT NOT NULL,
run_id TEXT NOT NULL,
policy TEXT NOT NULL,
stats_json TEXT NOT NULL DEFAULT '{}',
PRIMARY KEY (tenant_id, run_id, policy)
);
CREATE TABLE IF NOT EXISTS run_state (
tenant_id TEXT NOT NULL,
run_id TEXT NOT NULL,
Expand Down Expand Up @@ -751,6 +758,8 @@ def _apply_step(self, tenant_id: str, run_id: str, ev: dict) -> None:
}
window.append(step_entry)
window = window[-RUN_STATE_WINDOW:]
if isinstance(ev.get("compaction"), dict):
self._add_policy_stats(tenant_id, run_id, "context_compaction", ev["compaction"])
step_count = (row["step_count"] if row else 0) + 1
velocity = _velocity_from_window(window)
self._db.execute(
Expand All @@ -762,6 +771,36 @@ def _apply_step(self, tenant_id: str, run_id: str, ev: dict) -> None:
(tenant_id, run_id, step_count, json.dumps(window), velocity, ev.get("ts")),
)

def _add_policy_stats(self, tenant_id: str, run_id: str, policy: str, delta: dict) -> None:
"""Fold one call's numeric metrics into the run's per-policy aggregate (contract §5).

Stored as JSON so a policy can report new metrics without a schema change. Every
numeric key is summed and ``calls`` counts contributions. Runs inside the batch
transaction after the idempotency check, so a deduped replay cannot double count.
"""
row = self._db.execute(
"SELECT stats_json FROM run_policy_stats WHERE tenant_id=? AND run_id=? AND policy=?",
(tenant_id, run_id, policy),
).fetchone()
stats = json.loads(row["stats_json"]) if row else {}
for k, v in delta.items():
if isinstance(v, (int, float)) and not isinstance(v, bool):
stats[k] = stats.get(k, 0) + v
stats["calls"] = stats.get("calls", 0) + 1
self._db.execute(
"INSERT INTO run_policy_stats(tenant_id, run_id, policy, stats_json) VALUES (?,?,?,?) "
"ON CONFLICT(tenant_id, run_id, policy) DO UPDATE SET stats_json=excluded.stats_json",
(tenant_id, run_id, policy, json.dumps(stats)),
)

def get_policy_stats(self, tenant_id: str, run_id: str) -> dict[str, dict]:
"""``{policy: {metric: value}}`` aggregated over the run."""
rows = self._db.execute(
"SELECT policy, stats_json FROM run_policy_stats WHERE tenant_id=? AND run_id=?",
(tenant_id, run_id),
)
return {r["policy"]: json.loads(r["stats_json"]) for r in rows}

def precheck(
self,
tenant_id: str,
Expand Down
12 changes: 11 additions & 1 deletion docs/api-contract.md
Original file line number Diff line number Diff line change
Expand Up @@ -123,7 +123,8 @@ Mirrors the `envelopes:batch` shape. All ledger writes go through here.
"boundary_id": "researcher.chat", "cost_micros": 10500, "cum_spent_micros": 10500,
"tool_signature": null, "result_hash": null,
"usage": { "input": 1000, "output": 500 },
"tags": { "provider": "anthropic", "model": "claude-sonnet-4-6" } }
"tags": { "provider": "anthropic", "model": "claude-sonnet-4-6" },
"compaction": { "tokens_before": 12000, "tokens_after": 8000, "tokens_saved": 4000 } } // optional (0.2.2)

// halt_mark / halt_clear
{ "kind": "halt_mark", "reason": "step_cap: 20 steps", "detector": "step_cap" }
Expand All @@ -145,6 +146,11 @@ Mirrors the `envelopes:batch` shape. All ledger writes go through here.
- `spent_add` → upsert every `target` in `ledger_spent`.
- `step` → upsert the `(tenant_id, run_id)` row in `run_state` (`step_count`, bounded
`window_json` ring, velocity inputs, `last_ts`).
If the step carries the optional `compaction` object, also sum its numeric fields
(plus a `calls` count) into `run_policy_stats.stats_json` for policy `context_compaction`,
keyed `(tenant_id, run_id, policy)`, in the same transaction. `runs` is not touched.
Exposed as `policy_stats` on `GET /v1/run-records/{run_id}`. Deduped
replays never reach this, so they do not double count. Values are chars/4 estimates.
- `halt_mark` / `halt_clear` → `ledger_halt` (+ flip `runs.status` when the run row
exists).
- Zero-cost crossings emit **only** a `step` — no `spent_add` with `delta_micros: 0`.
Expand Down Expand Up @@ -285,6 +291,10 @@ forward-only.
**unchanged**. Opening a 0.1 DB just adds the new tables/column and bumps
`user_version` to 2 — **no data loss, no break.**

### v2.1 — 0.2.2 (additive, auto-applied)
- `run_policy_stats (tenant_id, run_id, policy, stats_json)` — new table, PK
`(tenant_id, run_id, policy)`. `user_version` stays 2.

### v3 — 0.3.0 (destructive, ships after `tokenops <next>` — issue #11)
- Fold `run_registrations` identity columns into `runs`; `DROP TABLE run_registrations`.
- Make `runs` identity columns write-once; `steps` / `cost_micros` derived only.
Expand Down
61 changes: 61 additions & 0 deletions tests/test_ledger_events.py
Original file line number Diff line number Diff line change
Expand Up @@ -155,3 +155,64 @@ def test_bad_durability_header_is_400(make_client):
headers={"Durability": "whenever"},
)
assert r.status_code == 400


def _compacted_step(seq: int, saved: int, *, run_id: str = "run_1") -> dict:
ev = _step(seq, seq * 10_500, run_id=run_id)
ev["compaction"] = {
"tokens_before": saved + 8_000,
"tokens_after": 8_000,
"tokens_saved": saved,
}
return ev


def _stats(c, run_id: str = "run_1") -> dict:
return c.app.state.store.get_policy_stats("local", run_id)


def test_compaction_is_aggregated_per_run_and_policy(make_client):
c = make_client()
c.post(
"/v1/ledger/events:batch",
json={"events": [_compacted_step(1, 4_000), _step(2, 21_000), _compacted_step(3, 1_000)]},
)
assert _stats(c)["context_compaction"] == {
"tokens_before": 12_000 + 9_000,
"tokens_after": 16_000,
"tokens_saved": 5_000,
"calls": 2,
}


def test_replayed_compacted_step_does_not_double_count(make_client):
c = make_client()
ev = _compacted_step(1, 4_000)
c.post("/v1/ledger/events:batch", json={"events": [ev]})
r = c.post("/v1/ledger/events:batch", json={"events": [ev]})
assert r.json()["deduped"] == 1
assert _stats(c)["context_compaction"]["tokens_saved"] == 4_000


def test_compaction_survives_run_state_ring_eviction(make_client):
c = make_client()
n = 70 # > RUN_STATE_WINDOW
c.post(
"/v1/ledger/events:batch",
json={"events": [_compacted_step(i, 100) for i in range(1, n + 1)]},
)
assert _stats(c)["context_compaction"]["tokens_saved"] == 100 * n


def test_run_record_read_exposes_policy_stats_and_needs_no_runs_row(make_client):
c = make_client()
c.put("/v1/run-records", json={"run_id": "run_1", "agent": "researcher"})
c.post("/v1/ledger/events:batch", json={"events": [_compacted_step(1, 4_000)]})
body = c.get("/v1/run-records/run_1").json()
assert body["policy_stats"]["context_compaction"]["tokens_saved"] == 4_000


def test_steps_without_compaction_write_no_stats(make_client):
c = make_client()
c.post("/v1/ledger/events:batch", json={"events": [_step(1, 10_500)]})
assert _stats(c) == {}
8 changes: 8 additions & 0 deletions ui/run_detail.py
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,14 @@ def render_run_detail(store: Store, run: RunRecord) -> None:
c3.metric("Duration (s)", _duration(run))
c4.metric("Governance", mode)

compaction = store.get_policy_stats("local", run.run_id).get("context_compaction")
if compaction:
st.metric(
"Tokens saved by compaction (est.)",
f"{int(compaction.get('tokens_saved', 0)):,}",
help=f"{int(compaction.get('calls', 0))} compacted call(s); chars/4 estimate.",
)

with st.expander("Run metadata", expanded=False):
st.json(
{
Expand Down
2 changes: 1 addition & 1 deletion ui/views/admin.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@
"tool_output_cap": '{"cap_tokens": 8000}',
"progress_guard": '{"window": 6, "repeats": 3, "max_corrections": 2}',
"cost_guard": '{"threshold": 0.8, "mode": "minimize"}',
"context_compaction": '{"ctx_max": 100000, "has_hook": false}',
"context_compaction": '{"ctx_max": 100000}',
"output_runaway": '{"repeats": 4, "max_retries": 2}',
}

Expand Down
Loading