budget/cache accounting: cold split after route rebind, complete failed-send rows, one prepared wrap-up candidate

Fifth authoritative review: an A→B→A route rebind reprojected the transcript
without invalidating the task's cache split — it does now; the failed-send
observer emitted earlier dispatched attempts only when the terminal exception
carried a positive capture — it now queries the final ledger state per collected
attempt (released reservations excluded), so a terminal BudgetExceeded or
attempt-limit still yields one row per earlier physical send; the wrap-up
candidate is prepared ONCE after the forced augmentations through the same
main-message preparation as dispatch (vision caption routing included) and that
exact candidate is the initial forced dispatch, fenced by the physical-attempt
predicate; the cost-stop rail stamps are written only beside the actual
fallback/forced-return sites. loop.py 283 460 bytes (−154). Manifest untouched.

Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
This commit is contained in:
Ouroboros 2026-09-04 03:37:34 +03:00
parent 5dc32a8afc
commit 18f78bcb87
7 changed files with 298 additions and 51 deletions

View file

@ -247,13 +247,6 @@ def _check_budget_limits(
budget_remaining_usd: Optional[float],
cost_ceiling: Optional["task_pacing.CostCeiling"] = None,
) -> Optional[Tuple[str, Dict[str, Any], Dict[str, Any]]]:
"""Return a final response when a budget axis requires stopping.
``cost_ceiling`` is resolved once at loop start. Root-capped tasks decide
on ledger-accounted tree spend; own cost is the disclosed fallback and
unknown spend is never $0. The axes are independent: a
``budget_remaining_usd`` of None means the owner explicitly chose no
finite global budget, and must not silence a live per-task root cap."""
accumulated_usage = ctx.accumulated_usage
raw_task_cost = accumulated_usage.get("cost")
task_cost = float(raw_task_cost) if raw_task_cost is not None else None
@ -272,7 +265,6 @@ def _check_budget_limits(
_force_plan_disclosure(tool_ctx, trace, forced_reason="budget_exhausted")
if tool_ctx is not None else ""
)
# Record the owed panel even when rejection happens before work.
_record_forced_finalization(
ctx,
trace,
@ -308,18 +300,10 @@ def _check_budget_limits(
forced_prompt = f"[BUDGET LIMIT] {finish_reason} {_FORCED_BEST_EFFORT_TAIL}"
prospective_messages = [dict(message) for message in ctx.messages]
_append_or_merge_user_message(prospective_messages, forced_prompt)
request_args = dict(
llm=ctx.llm, model=ctx.active_model, reasoning_effort=ctx.active_effort,
tools=ctx.tool_schemas, allow_server_web_search=_server_web_allowed_by_task(
getattr(getattr(ctx, "tools", None), "_ctx", None)
), prompt_tokens=prompt_estimate,
)
wrapup_request = task_pacing.prospective_wrapup_attempt_request(
messages=prospective_messages, **request_args,
) if not ctx.active_use_local else None
request_args = dict(model=ctx.active_model, prompt_tokens=prompt_estimate,
use_local=ctx.active_use_local)
wrapup_args = dict(
request=wrapup_request, root_cap_usd=cost_ceiling.root_cap_usd,
deciding_usd=deciding,
**request_args, root_cap_usd=cost_ceiling.root_cap_usd, deciding_usd=deciding,
)
wrapup_fits = task_pacing.wrapup_reservation_fits(**wrapup_args)
if wrapup_fits is True and task_pacing.wrapup_reservation_fits(
@ -329,25 +313,34 @@ def _check_budget_limits(
priced_prompt = _prepare_forced_prompt(ctx, forced_prompt, trace)
prospective_messages = [dict(message) for message in ctx.messages]
_append_or_merge_user_message(prospective_messages, priced_prompt)
wrapup_args["request"] = task_pacing.prospective_wrapup_attempt_request(
messages=prospective_messages, **request_args,
wrapup_request, send_messages = task_pacing.prepared_wrapup_candidate(
ctx, prospective_messages,
allow_server_web_search=_server_web_allowed_by_task(
getattr(getattr(ctx, "tools", None), "_ctx", None)),
)
wrapup_args = dict(
request=wrapup_request, root_cap_usd=cost_ceiling.root_cap_usd,
deciding_usd=deciding,
)
wrapup_fits = task_pacing.wrapup_reservation_fits(**wrapup_args)
two_fit = task_pacing.wrapup_reservation_fits(
**wrapup_args, reservation_count=2,
)
accumulated_usage["cost_stop_spend_basis"] = spend_basis
accumulated_usage["cost_stop_rail"] = "wrapup_reservation_last_fit"
if wrapup_fits is False:
accumulated_usage["cost_stop_spend_basis"] = spend_basis
accumulated_usage["cost_stop_rail"] = "wrapup_reservation_last_fit"
return _forced_fallback_result(
ctx, trace, finish_reason, "budget_exhausted",
source="budget_wrapup_unaffordable",
)
if wrapup_fits is not True or two_fit is not False:
return None
accumulated_usage["cost_stop_spend_basis"] = spend_basis
accumulated_usage["cost_stop_rail"] = "wrapup_reservation_last_fit"
return _forced_final_answer(
ctx, prompt=priced_prompt, _prompt_prepared=True,
fallback_text=finish_reason, reason_code="budget_exhausted",
_initial_messages=send_messages, _admitted_request=wrapup_request,
)
if deciding is not None and ceiling_usd is not None and deciding > ceiling_usd:
if wrapup_fits is False:
@ -360,8 +353,6 @@ def _check_budget_limits(
else f"Task tree spent ${deciding:.3f} (ledger-accounted incl. in-flight holds)"
)
elif spend_basis == task_pacing.SPEND_BASIS_OWN_TREE_UNKNOWN:
# Stopping on a disclosed lower bound beats not stopping at all, but
# the substitution is stated, never silent (BIBLE P1).
spent_text = (
f"Task spent ${deciding:.3f} on its OWN calls (the tree-accounted total "
"is unavailable right now, so subagent spend is not included — this is a "
@ -377,8 +368,6 @@ def _check_budget_limits(
f"{spent_text}, over the in-task cost ceiling ${ceiling_usd:.2f}{cap_text}. "
"Budget exhausted."
)
# The basis rides the usage record too, so a later reader can tell a
# tree-decided stop from an own-cost stand-in without parsing prose.
accumulated_usage["cost_stop_spend_basis"] = spend_basis
return _forced_final_answer(
ctx,
@ -4321,8 +4310,17 @@ def _drain_forced_owner_directives(
return True
def _call_forced_model_once(ctx: _RoundLimitContext) -> str:
def _call_forced_model_once(
ctx: _RoundLimitContext, *, initial_messages: Any = None, admitted_request: Any = None,
) -> str:
response_meta: Dict[str, Any] = {}
identity = (
"model", "provider", "candidate_raw_sha256", "candidate_raw_size_bytes",
)
candidate_predicate = (
lambda actual: all(getattr(actual, key, None) == getattr(admitted_request, key, None) for key in identity)
if admitted_request is not None else None
)
final_msg, _final_cost = call_llm_with_retry(
ctx.llm,
ctx.messages,
@ -4343,6 +4341,8 @@ def _call_forced_model_once(ctx: _RoundLimitContext) -> str:
allow_server_web_search=_server_web_allowed_by_task(
getattr(getattr(ctx, "tools", None), "_ctx", None)
),
initial_messages=initial_messages,
candidate_predicate=candidate_predicate,
)
ctx.accumulated_usage["_forced_response_meta"] = response_meta
return str((final_msg or {}).get("content") or "").strip()
@ -4683,6 +4683,8 @@ def _forced_final_answer(
single_semantic_turn: bool = False,
provider_terminal: bool = False,
_prompt_prepared: bool = False,
_initial_messages: Any = None,
_admitted_request: Any = None,
) -> Tuple[str, Dict[str, Any], Dict[str, Any]]:
"""Forced rail."""
live_trace = getattr(ctx, "llm_trace", None)
@ -4705,9 +4707,12 @@ def _forced_final_answer(
for attempt in range(1 if single_semantic_turn else 2):
try:
ctx.accumulated_usage.pop("_forced_response_meta", None)
extracted, response_meta = forced_response_parts(
_call_forced_model_once(ctx), ctx.accumulated_usage,
)
if attempt == 0 and _admitted_request is not None:
forced = _call_forced_model_once(
ctx, initial_messages=_initial_messages, admitted_request=_admitted_request)
else:
forced = _call_forced_model_once(ctx)
extracted, response_meta = forced_response_parts(forced, ctx.accumulated_usage)
except BudgetExceeded:
_drain_forced_owner_directives(ctx, llm_trace)
raise
@ -4836,12 +4841,6 @@ def _rebind_context_fit_plan(
preferred_mode: str,
tool_schemas: List[Dict[str, Any]],
) -> Tuple[Any, str]:
"""Recalibrate the captured immutable core for one new exact route.
Route switches reuse the plan's already-rendered Low/Max projections; only
exact-route evidence, calibration, and fit are rebound. This avoids both a
stale initial-route retry plan and a second context-builder/intent corpus.
"""
if plan is None or not all(
hasattr(plan, name) for name in ("max_projection", "low_projection", "core_sha256")
):
@ -4907,6 +4906,7 @@ def _rebind_context_fit_plan(
mode = initial_mode
projected_prompt_tokens = rebound.projected_tokens_with_tools(mode, tool_schemas)
messages[:] = rebound.reproject_transcript(messages, mode)
invalidate_task_cache_splits(getattr(tools._ctx, "task_id", ""))
tools._ctx.context_fit_plan = rebound
tools._ctx.messages = messages
tools._ctx.active_context_mode = mode

View file

@ -1319,7 +1319,7 @@ def call_llm_with_retry(
candidate_predicate: Optional[Callable[[Any], Any]] = None,
task_attempt: Any = None,
response_meta_out: Optional[Dict[str, Any]] = None,
transport_reserve_sec: Optional[float] = None,
transport_reserve_sec: Optional[float] = None, initial_messages: Optional[List[Dict[str, Any]]] = None,
) -> Tuple[Optional[Dict[str, Any]], Optional[float]]:
"""Call one model with bounded retries and deadline-aware transport."""
msg = None
@ -1346,7 +1346,7 @@ def call_llm_with_retry(
request_ref: Dict[str, Any] = {}
try:
_emit_llm_operation(event_queue, task_id, llm_call_id, "started", task_attempt, execution_id, round_id)
send_messages = _prepare_main_messages(
send_messages = initial_messages if attempt == 0 and initial_messages is not None else _prepare_main_messages(
messages, model=model, llm=llm, accumulated_usage=accumulated_usage,
drive_root=drive_root, task_id=task_id, event_queue=event_queue,
use_local=use_local, task_attempt=task_attempt, deadline_ts=deadline_ts,
@ -1421,7 +1421,7 @@ def call_llm_with_retry(
_emit_main_llm_call_state(event_queue, call_identity, "started")
resp_msg, usage = _send_main_candidate(
llm, kwargs, model=model, use_local=use_local, deadline_ts=deadline_ts,
physical_context=physical_context, candidate_predicate=candidate_predicate,
physical_context=physical_context, candidate_predicate=candidate_predicate if attempt == 0 else None,
)
msg = resp_msg
_take_custom_receipts(usage, msg, accumulated_usage)

View file

@ -28,7 +28,10 @@ from ouroboros.review_slot_cancel import ( # noqa: F401 — re-exported seam su
_slot_cancel_outcome,
)
from ouroboros.review_dispatch import bind_api_review_paid_stamp, invoke_review_paid_stamp
from ouroboros.usage_accounting import POSITIVE_PHYSICAL_ATTEMPT_STATES
from ouroboros.usage_accounting import (
POSITIVE_PHYSICAL_ATTEMPT_STATES, _drive_root, _final_rows, _locked,
_read_records_locked_cached, current_usage_scope,
)
from ouroboros.triad_review import (
ACCEPTANCE_SURFACE_RULES,
REVIEW_JSON_ARRAY_CONTRACT,
@ -324,12 +327,31 @@ class ReviewSlotExecutor:
def _observe_failed_send(self, exc: BaseException) -> None:
capture = getattr(exc, "physical_attempt_capture", None)
if str(getattr(capture, "state", "") or "") in POSITIVE_PHYSICAL_ATTEMPT_STATES:
attempt_ids = [str(value) for value in (getattr(exc, "ledger_attempt_ids", None) or []) if value]
attempt_ids = [str(value) for value in (getattr(exc, "ledger_attempt_ids", None) or []) if value]
capture_id = str(getattr(capture, "attempt_id", "") or "")
if capture_id and capture_id not in attempt_ids:
attempt_ids.append(capture_id)
rows: Dict[str, Dict[str, Any]] = {}
try:
scope = current_usage_scope()
root = _drive_root(getattr(scope, "drive_root", None))
with _locked(root):
finals = _final_rows(_read_records_locked_cached(root))
rows = {attempt_id: finals[attempt_id] for attempt_id in attempt_ids if attempt_id in finals}
except Exception:
log.debug("failed to resolve review attempt states", exc_info=True)
capture_state = str(getattr(capture, "state", "") or "")
for attempt_id in attempt_ids:
row = rows.get(attempt_id, {})
state = str(row.get("state") or (capture_state if attempt_id == capture_id else ""))
if state not in POSITIVE_PHYSICAL_ATTEMPT_STATES and not (
not rows and attempt_id != capture_id
):
continue
self._observe_usage({
"resolved_model": str(getattr(capture, "model", "") or ""),
"provider": str(getattr(capture, "provider", "") or ""),
"ledger_attempt_ids": attempt_ids or [str(getattr(capture, "attempt_id", "") or "")],
"resolved_model": str(row.get("model") or getattr(capture, "model", "") or ""),
"provider": str(row.get("provider") or getattr(capture, "provider", "") or ""),
"ledger_attempt_ids": [attempt_id],
})
def prompt_payload(self) -> Dict[str, Any]:

View file

@ -669,6 +669,30 @@ def prospective_wrapup_attempt_request(
return _merge_scope(_attempt_request(target, candidate))[0]
def prepared_wrapup_candidate(
ctx: Any, messages: list[Dict[str, Any]], *, allow_server_web_search: bool,
) -> Tuple[Any, list[Dict[str, Any]]]:
"""Prepare the exact first-send transcript and price that same payload."""
from ouroboros.loop_llm_call import _prepare_main_messages
send_messages = _prepare_main_messages(
messages, model=ctx.active_model, llm=ctx.llm,
accumulated_usage=ctx.accumulated_usage,
drive_root=ctx.drive_root or pathlib.Path(ctx.drive_logs or ".").parent,
task_id=ctx.task_id, event_queue=ctx.event_queue,
use_local=ctx.active_use_local,
task_attempt=ctx.accumulated_usage.get("_task_attempt"),
deadline_ts=ctx.deadline_ts,
)
request = prospective_wrapup_attempt_request(
llm=ctx.llm, messages=send_messages, model=ctx.active_model,
reasoning_effort=ctx.active_effort, tools=ctx.tool_schemas,
allow_server_web_search=allow_server_web_search,
prompt_tokens=int(ctx.accumulated_usage.get("_context_prompt_estimate") or 0),
)
return request, send_messages
def wrapup_reservation_fits(
*,
model: str = "",

View file

@ -240,6 +240,54 @@ def test_route_rebind_keeps_owner_projection_on_small_confirmed_route(monkeypatc
assert rebound.route_fp == "small-route"
def test_route_rebind_a_b_a_forgets_the_old_a_cache_split(monkeypatch, tmp_path):
from ouroboros import context, loop, usage_accounting
from ouroboros.tools.registry import ToolRegistry
plan = _plan()
registry = ToolRegistry(repo_dir=tmp_path, drive_root=tmp_path)
registry._ctx.task_id = "switch-back"
registry._ctx.event_queue = None
monkeypatch.setattr(
context,
"_context_fit_route",
lambda task, **_kw: (
{"model": task["model"], "provider": "openrouter", "use_local": False},
SimpleNamespace(
status="confirmed", stale=False, window_tokens=500_000,
route_fp=f"route-{task['model']}",
),
),
)
monkeypatch.setattr(
usage_accounting,
"estimate_cost_optional",
lambda _model, _prompt, _completion, *, cache_usage, **_kw: float(
cache_usage["cache_write_tokens"]
),
)
request = usage_accounting.AttemptRequest(
model="anthropic/model-a", provider="openrouter", task_id="switch-back",
prompt_tokens_estimate=1_000,
)
usage_accounting.stash_task_cache_split(
"switch-back", request.model, 800, provider=request.provider, ttl_seconds=300,
)
assert usage_accounting._reservation_cost(request) == 200
messages = plan.messages_for("max")
plan, _ = loop._rebind_context_fit_plan(
plan, registry, messages, model="anthropic/model-b", use_local=False,
preferred_mode="max", tool_schemas=[],
)
loop._rebind_context_fit_plan(
plan, registry, messages, model=request.model, use_local=False,
preferred_mode="max", tool_schemas=[],
)
assert usage_accounting._reservation_cost(request) == 1_000
def test_route_switch_without_immutable_core_fails_loudly(tmp_path):
from ouroboros import loop
from ouroboros.tools.registry import ToolRegistry

View file

@ -3092,6 +3092,70 @@ def test_terminal_failed_reviewer_retry_emits_one_row_per_dispatched_attempt(tmp
assert rows[0]["ledger_attempt_ids"] != rows[1]["ledger_attempt_ids"]
def test_terminal_budget_refusal_keeps_prior_dispatched_retry_usage(tmp_path):
from ouroboros.usage_accounting import (
AttemptRequest, capture_attempt_ids, execute_physical_attempt,
)
class BudgetStopsRetryLLM:
def chat(self, **_kwargs):
request = AttemptRequest(
model="same/model", provider="openrouter",
reservation_usd=1.0, global_limit_usd=1.0,
)
with capture_attempt_ids():
try:
execute_physical_attempt(
request, lambda: (_ for _ in ()).throw(RuntimeError("dispatched")),
)
except RuntimeError:
pass
execute_physical_attempt(request, lambda: None)
ctx = SimpleNamespace(task_id="budget-refusal-usage", event_queue=None, pending_events=[])
run_review_request(
ReviewRequest(surface="task_acceptance", goal="review", task_id=ctx.task_id),
slots=[ReviewSlot(slot_id="slot_a", model="same/model")],
drive_root=tmp_path, llm=BudgetStopsRetryLLM(), usage_ctx=ctx,
)
rows = [event for event in ctx.pending_events if event.get("type") == "llm_usage"]
assert len(rows) == 1
assert len(rows[0]["ledger_attempt_ids"]) == 1
def test_terminal_attempt_limit_keeps_prior_send_but_excludes_released_hold(tmp_path):
from ouroboros.usage_accounting import (
AttemptRequest, capture_attempt_ids, execute_physical_attempt,
physical_attempt_limit,
)
class RailStopsRetryLLM:
def chat(self, **_kwargs):
request = AttemptRequest(
model="same/model", provider="openrouter", reservation_usd=0.0,
)
with capture_attempt_ids(), physical_attempt_limit(1):
try:
execute_physical_attempt(
request, lambda: (_ for _ in ()).throw(RuntimeError("dispatched")),
)
except RuntimeError:
pass
execute_physical_attempt(request, lambda: None)
ctx = SimpleNamespace(task_id="rail-refusal-usage", event_queue=None, pending_events=[])
run_review_request(
ReviewRequest(surface="task_acceptance", goal="review", task_id=ctx.task_id),
slots=[ReviewSlot(slot_id="slot_a", model="same/model")],
drive_root=tmp_path, llm=RailStopsRetryLLM(), usage_ctx=ctx,
)
rows = [event for event in ctx.pending_events if event.get("type") == "llm_usage"]
assert len(rows) == 1
assert len(rows[0]["ledger_attempt_ids"]) == 1
def test_internal_reviewer_transport_attempts_each_get_one_usage_row(tmp_path):
class RetriedLLM:
def chat(self, **_kwargs):

View file

@ -491,11 +491,100 @@ class TestWrapupAffordabilityRail:
assert result[2]["kwargs"]["reason_code"] == "budget_exhausted"
assert result[2]["kwargs"]["_prompt_prepared"] is True
assert "forced delegation note" in result[2]["kwargs"]["prompt"]
assert "[BUDGET LIMIT]" in builds[1]["messages"][-1]["content"]
assert "forced delegation note" in builds[1]["messages"][-1]["content"]
assert "service evidence" in str(builds[1]["messages"])
assert len(builds) == 2
assert [call["request"] for call in calls] == [request] * 4
assert "[BUDGET LIMIT]" in builds[0]["messages"][-1]["content"]
assert "forced delegation note" in builds[0]["messages"][-1]["content"]
assert "service evidence" in str(builds[0]["messages"])
assert len(builds) == 1
assert [call.get("request") for call in calls] == [None, None, request, request]
def test_repriced_nondecision_does_not_stamp_a_cost_stop(self, monkeypatch):
ctx = _ctx()
monkeypatch.setattr(
"ouroboros.loop._loop_tree_accounting", lambda **_k: {"accounted_usd": 20.0},
)
answers = iter((True, False, None, False))
monkeypatch.setattr(
task_pacing, "wrapup_reservation_fits", lambda **_kwargs: next(answers),
)
monkeypatch.setattr(
task_pacing, "prospective_wrapup_attempt_request", lambda **_kwargs: object(),
)
assert _check_budget_limits(ctx, None, self._ceiling(50.0)) is None
assert "cost_stop_spend_basis" not in ctx.accumulated_usage
assert "cost_stop_rail" not in ctx.accumulated_usage
def test_captioned_wrapup_candidate_is_the_initial_forced_dispatch(
self, monkeypatch, tmp_path,
):
import ouroboros.loop as loop_module
from ouroboros.llm import LLMClient
ctx = _ctx(drive_logs=tmp_path / "logs", llm=LLMClient(api_key="unused"))
ctx.messages = [{
"role": "user",
"content": [
{"type": "text", "text": "inspect"},
{
"type": "image_url",
"image_url": {"url": "data:image/png;base64,aaa"},
"_caption": "wire caption",
},
],
}]
monkeypatch.setenv("OUROBOROS_IMAGE_INPUT_MODE", "caption")
monkeypatch.setattr(
"ouroboros.loop._loop_tree_accounting", lambda **_k: {"accounted_usd": 20.0},
)
answers = iter((True, False, True, False))
monkeypatch.setattr(
task_pacing, "wrapup_reservation_fits", lambda **_kwargs: next(answers),
)
built = {}
real_build = task_pacing.prospective_wrapup_attempt_request
def build(**kwargs):
built["messages"] = kwargs["messages"]
built["request"] = real_build(**kwargs)
return built["request"]
monkeypatch.setattr(task_pacing, "prospective_wrapup_attempt_request", build)
dispatched = {}
def call(*_args, **kwargs):
dispatched["messages"] = kwargs["initial_messages"]
from ouroboros.llm import _attempt_request, _finalized_physical_candidate
from ouroboros.request_wire_recovery import request_wire_call_scope
target = ctx.llm._resolve_remote_target(ctx.active_model)
with request_wire_call_scope():
candidate = ctx.llm._build_remote_candidate(
target, kwargs["initial_messages"], ctx.active_effort,
built["request"].max_completion_tokens, "auto", None, ctx.tool_schemas,
skip_capability_fetch=True,
)
ctx.llm._normalize_payload_cache_ttl(target, candidate)
candidate = _finalized_physical_candidate(
target, candidate,
"messages" if target.get("provider") == "anthropic" else "chat.completions",
)
actual = _attempt_request(target, candidate)
dispatched["accepted"] = kwargs["candidate_predicate"](actual)
dispatched["sha256"] = actual.candidate_raw_sha256
return {"content": "wrapped up"}, 0.0
monkeypatch.setattr(loop_module, "call_llm_with_retry", call)
result = _check_budget_limits(ctx, None, self._ceiling(50.0))
assert result is not None
assert dispatched["messages"] is built["messages"]
assert dispatched["accepted"] is True
assert dispatched["sha256"] == built["request"].candidate_raw_sha256
assert built["messages"][0]["content"][1] == {
"type": "text", "text": "[image caption: wire caption]",
}
assert ctx.messages[0]["content"][1]["type"] == "image_url"
def test_a_missing_prompt_estimate_keeps_the_rail_silent(self, monkeypatch):
ctx = _ctx(accumulated_usage={"cost": 1.0})