Deliver the budget wrap-up answer even when its admitted candidate drifted

The graceful budget stop admits one wrap-up call against a candidate priced
before the send, and an identity predicate refuses the send when the real
payload differs. The real main-loop send had grown `stream` (+`stream_options`)
and the priced copy had not -- 54 bytes on OpenAI-compatible routes, 14 on direct
Anthropic, where the lane added `stream` after its builder -- so every API-lane
budget wrap-up was refused before dispatch, the forced path swallowed the typed
refusal as a generic error, and the task ended Failed with an empty final answer
after finalizing its services. The parity pins never saw it: each rebuilt the
"actual" candidate with the prospective builder's own arguments.

- task_pacing.main_loop_wire_options is the one owner of the payload-shaping
  options of a main-loop send; call_llm_with_retry and the prospective wrap-up
  builder spread the same dict. `stream` is a builder input on every lane; the
  direct Anthropic lane no longer mutates the payload after its builder.
- A refused admitted candidate is recorded as a typed `forced_candidate_drift`
  checkpoint and the answer is asked for once more without the predicate: the
  refusal happens before dispatch and releases its reservation, and the ledger
  fence prices the send it actually sees. A closed dispatch window stays a
  deadline and is never retried.
- tests/test_wrapup_real_send_parity.py drives call_llm_with_retry itself on four
  routes, the whole rail from the last-fit decision to a delivered best-effort
  answer ("Done with warnings", not "Failed"), the typed retry and the deadline
  exception. The hand-built mirrors take the shared options.
- The stop sentence opens with the bound that binds and, when that is the
  shared wallet, says how to lift it; an owner whose wallet ran dry at $125 of a
  $400 cap read the old sentence as a broken per-task cap.

Disclosed residual: in the server process (direct-chat turns) the real send
fetches model capabilities while the priced copy skips that fetch; a difference
there now costs one released reservation and a retry instead of the answer.
This commit is contained in:
Ouroboros 2026-09-20 00:18:19 +03:00
parent 4d046a289c
commit 2c2151125b
10 changed files with 318 additions and 32 deletions

View file

@ -592,7 +592,7 @@ Before summary, reflection or consolidation starts, the root freezes one shared
The terminal checkpoint also retains `root_phase_checkpoint.accounting`: one cumulative root-tree observation from the existing ledger breakdown, with root identity, accounted/reserved/unresolved amounts, non-final and unresolved row counts, unknown exposure and integrity. It stays separate from own-task fields and from the read-time `TaskCostBreakdown`; public detail and saved-result recovery retain the same observation. The existing checkpoint owner refreshes it, including explicit unavailable/null evidence after a failed read, and canonical replica protection preserves it. These row facts establish neither an exact invoice nor closure of all local paid work; phase and execution ownership remain independent.
The in-task pacing stop is an explicitly unreserved planning threshold in the shared global pool, resolved once per task as a typed `CostCeiling` (`disabled`, `active`, `exhausted_soft_land`, `unknown`; `task_pacing.py`), published in the start-of-task runtime budget block as `in_task_cost_ceiling` and consumed by the loop as the SAME object, so the number the mind is shown and the number that stops the task cannot differ. The root of a tree resolves the minimum of the configured percentage of global remaining budget and the root-tree cap, minus a small planning margin; enabled descendants keep that original number instead of taking a percentage of a later wallet. Explicitly disabled profiles keep the early axis disabled; a root without a task cap resolves from the starting wallet while actual global/root monetary admission stays independent; every host cost surface prints the bound that binds first and names it. The loop decides against subtree-accounted spend including in-flight holds (an own-cost fallback is disclosed as a lower bound), for every task that has a root, with or without a per-task cap — a global-only ceiling decides on the subtree, not the parent alone — and checks the axis only after tool-call rounds; graceful finalization runs before the ledger fence and never weakens it: unknown is not zero, an unknown price fails open, a disabled ceiling never arms it. The margin pulls the stop earlier so a post-round affordability crossing can borrow the fence's OWN per-attempt reservation — cache-aware from the task's last settled split for the same provider, normalized route and review surface (`_usage_cache_splits.py`, process-local; a lost entry only re-prices a full write) — and soft-land while one wrap-up call is still admitted instead of dying answerless at the fence. The comparison reserves nothing (a competing task can consume it first; a cache expiry may under-reserve by one write) — the ledger fence at the full cap still binds. The bounded context proxy only pre-screens: a proxy stop, and any prompt the proxy can understate (native image parts), is decided by exact pricing on a copy of the transcript before service finalization. A task that cannot reserve even one wrap-up says so in its own words; the exhausted-ceiling soft landing prices the same prepared candidate and ends as `budget_wrapup_unaffordable` rather than a fence pause when it cannot fit. A reply that still asks for a tool is incomplete on every rail, even beside a `replace` delivery control.
The in-task pacing stop is an explicitly unreserved planning threshold in the shared global pool, resolved once per task as a typed `CostCeiling` (`disabled`, `active`, `exhausted_soft_land`, `unknown`; `task_pacing.py`), published in the start-of-task runtime budget block as `in_task_cost_ceiling` and consumed by the loop as the SAME object, so the number the mind is shown and the number that stops the task cannot differ. The root of a tree resolves the minimum of the configured percentage of global remaining budget and the root-tree cap, minus a small planning margin; enabled descendants keep that original number instead of taking a percentage of a later wallet. Explicitly disabled profiles keep the early axis disabled; a root without a task cap resolves from the starting wallet while actual global/root monetary admission stays independent; every host cost surface prints the bound that binds first and names it. The loop decides against subtree-accounted spend including in-flight holds (an own-cost fallback is disclosed as a lower bound), for every task that has a root, with or without a per-task cap — a global-only ceiling decides on the subtree, not the parent alone — and checks the axis only after tool-call rounds; graceful finalization runs before the ledger fence and never weakens it: unknown is not zero, an unknown price fails open, a disabled ceiling never arms it. The margin pulls the stop earlier so a post-round affordability crossing can borrow the fence's OWN per-attempt reservation — cache-aware from the task's last settled split for the same provider, normalized route and review surface (`_usage_cache_splits.py`, process-local; a lost entry only re-prices a full write) — and soft-land while one wrap-up call is still admitted instead of dying answerless at the fence. The admitted candidate takes the send's wire options from their one owner (`task_pacing.main_loop_wire_options`; a lane never adds a payload key after its builder), because a candidate that differs from the real send is refused by the identity predicate before dispatch; such a refusal is recorded (`forced_candidate_drift`) and the answer is asked for once more unpredicated — the fence prices the send it sees, and drift must never cost the owner the final answer. The comparison reserves nothing (a competing task can consume it first; a cache expiry may under-reserve by one write) — the ledger fence at the full cap still binds. The bounded context proxy only pre-screens: a proxy stop, and any prompt the proxy can understate (native image parts), is decided by exact pricing on a copy of the transcript before service finalization. A task that cannot reserve even one wrap-up says so in its own words; the exhausted-ceiling soft landing prices the same prepared candidate and ends as `budget_wrapup_unaffordable` rather than a fence pause when it cannot fit. A reply that still asks for a tool is incomplete on every rail, even beside a `replace` delivery control.
One resolver answers the configured global budget for every agent-side reader (the supervisor's own startup and settings-reload carriers still parse the raw setting and map absence to no limit — a pre-existing, tracked gap), and it answers LIVE: a worker re-projects settings into its environment only at task start, so the resolver reads the saved document first (unlocked, re-parsed only when the file changed; any read failure, a refused benchmark pin included, leaves the environment answering and is never remembered) and no task captures the limit — the fence and every wallet projection follow the owner's current `TOTAL_BUDGET`, raised or lowered, at the next reservation, while the in-task ceiling stays the number resolved at start. An absent setting is silence, resolving to the shipped default, while an explicitly non-positive value means no finite global budget and silences the loop-side global axis. Pre-dispatch pricing (`pricing.py`) is an exact-route, bounded, best-effort lookup from the provider's current catalog — only the normalized exact model id and provider-supplied fields count; no manual price table, prefix inheritance, numeric fallback or admission allowlist disguised as pricing, which would silently invent authority and grow stale. Unknown price is nullable and fail-open for admission while known spend stays below its limits: it reserves `None` and settles from provider-reported cost or a later exact price, else cost stays `None` with `cost_final=false`. A rejection settles at confirmed zero only when structural provider evidence proves it happened before upstream generation with zero usage, releasing the reservation so a provider storm cannot manufacture phantom budget exhaustion (`_usage_response.py` normalizes that block for accounting; adapters still read the raw `usage` dict themselves). `review_wave_admission` applies the same per-attempt math against the tighter of the global and root remainders — the two fences `reserve_attempt` enforces, the binding axis named — before skill, plan, task-acceptance or P3 commit-gate reviewers launch, and the managed-update assisted-apply floor reuses the same estimator before any destructive merge step. An unpayable reviewer row is bypassed, never swapped to another model, so the audit stays honest about which model reviewed; an unpriced slot is disclosed and contributes no invented price while priced siblings still bind — one unknown route cannot disable admission control for a paid wave.

View file

@ -704,7 +704,8 @@ and what enforces each.
provider hints and recovery; do not add a generic cache/retry framework.
Wrap-up calls keep schemas, server-web flag and `tool_choice` unchanged and
instruct in text, because removing tools or changing tool choice rebuilds
cached input. Preserve `context_fit.seal_task_transcript`'s single message
cached input; a main-loop payload option lives in `main_loop_wire_options`, never
in one lane after its builder (`tests/test_wrapup_real_send_parity.py`). Preserve `context_fit.seal_task_transcript`'s single message
marker as it moves between task and tool result; direct Anthropic and
OpenRouter keep their supported wire markers. OpenRouter's derived identity
excludes cache/host metadata, preserving real task/model differences and

View file

@ -469,6 +469,11 @@ class _AnthropicLaneMixin:
choice = self._build_anthropic_tool_choice(tool_choice)
if choice:
payload["tool_choice"] = choice
# A builder input, never a later mutation: whoever prices this payload
# before the send (the budget wrap-up) builds it here too, and a key added
# after the builder is a key the priced copy never has.
if remote_kwargs.get("stream"):
payload["stream"] = True
apply_processing_preference(target, payload)
return payload
@ -492,9 +497,8 @@ class _AnthropicLaneMixin:
payload = self._build_remote_candidate(
target, messages, reasoning_effort, max_tokens, tool_choice, temperature, tools,
stream=stream,
)
if stream:
payload["stream"] = True
prompt_cache_ttl = self._normalize_payload_cache_ttl(target, payload)
url = f"{str(target.get('base_url') or '').rstrip('/')}/messages"

View file

@ -913,6 +913,40 @@ def _resolve_forced_delivery_control(
)
def _send_admitted_forced_candidate(
ctx: _RoundLimitContext, initial_messages: Any, admitted_request: Any, reason_code: str,
) -> str:
"""Send the admitted wrap-up; one that drifted from its pricing is sent once more, unpredicated.
The identity predicate refuses BEFORE a byte leaves and the refused
reservation is released, so nothing was paid and nothing is sent twice. Ending
the task there cost the owner the whole final answer three times in one night
(a wire key the send had grown and the priced copy had not), while the money
that predicate guards is guarded again by the ledger fence, which prices the
send it actually sees. So the drift is recorded as a typed fact — the refused
attempt's row and sealed candidate carry the actual identity — and the answer
is asked for once more the ordinary way. A closed dispatch window is a
deadline, not drift, and keeps its own rail."""
from ouroboros.llm_attempt import PhysicalDispatchInterrupted
from ouroboros.usage_accounting import PhysicalAttemptPreconditionFailed
try:
return _loop()._call_forced_model_once(
ctx, initial_messages=initial_messages, admitted_request=admitted_request)
except PhysicalDispatchInterrupted:
raise
except PhysicalAttemptPreconditionFailed as refusal:
log.warning("Admitted %s wrap-up candidate drifted from its pricing; sending it unpredicated", reason_code)
_loop()._emit_checkpoint_event(ctx.event_queue, ctx.task_id, ctx.drive_logs, {
"checkpoint_kind": "forced_candidate_drift",
"reason_code": reason_code,
"refused_attempt_id": str(getattr(refusal, "attempt_id", "") or ""),
"admitted": {key: getattr(admitted_request, key, None) for key in (
"model", "provider", "candidate_raw_sha256", "candidate_raw_size_bytes")},
})
return _loop()._call_forced_model_once(ctx)
def _forced_final_answer(
ctx: _RoundLimitContext,
*,
@ -944,8 +978,8 @@ def _forced_final_answer(
try:
ctx.accumulated_usage.pop("_forced_response_meta", None)
if attempt == 0 and _admitted_request is not None:
forced = _loop()._call_forced_model_once(
ctx, initial_messages=_initial_messages, admitted_request=_admitted_request)
forced = _send_admitted_forced_candidate(
ctx, _initial_messages, _admitted_request, reason_code)
else:
forced = _loop()._call_forced_model_once(ctx)
extracted, response_meta = forced_response_parts(forced, ctx.accumulated_usage)

View file

@ -23,13 +23,13 @@ from ouroboros.deadline_utils import (
main_transport_timeout_sec as _main_transport_timeout,
)
from ouroboros.llm import LLMClient, LocalContextTooLargeError, add_usage
from ouroboros.llm_claudexor import cache_key_for_model, propagate_model_error
from ouroboros.llm_claudexor import propagate_model_error
from ouroboros.model_wait import propagate_model_control
from ouroboros.openai_chat_dispatch import CUSTOM_RECEIPTS_USAGE_KEY, pop_custom_validation_receipts
from ouroboros.llm_attempt import PROVIDER_POLICY_REFUSAL, _is_provider_policy_refusal # typed-refusal contract owner
from ouroboros.observability import new_call_id, new_execution_id, persist_call
from ouroboros.pricing import emit_llm_usage_event, estimate_cost_optional, infer_model_category
from ouroboros.provider_models import provider_for_model
from ouroboros.task_pacing import main_loop_wire_options
from ouroboros.transport_custody import attempt_custody_event_fields, is_pre_dispatch_transport_failure, is_retryable_transport_death
from ouroboros._usage_response import provider_cost_value as _provider_cost_value
from ouroboros.usage_accounting import (
@ -1378,11 +1378,11 @@ def call_llm_with_retry(
"context_mode": getattr(physical_context, "rendered_mode", None),
"reasoning_effort": effort,
"max_tokens": MAIN_LOOP_MAX_TOKENS,
"stream": True, "caller_deadline_ts": (None if deadline_ts is None
**main_loop_wire_options(model, allow_server_web_search=allow_server_web_search,
bypass_response_cache=response_cache_bypass_requested),
"caller_deadline_ts": (None if deadline_ts is None
else float(deadline_ts) - float(transport_reserve_sec or 0.0)),
"use_local": use_local, "cache_affinity": cache_key_for_model(model),
"allow_server_web_search": bool(allow_server_web_search) and provider_for_model(model) != "claudexor",
"bypass_response_cache": response_cache_bypass_requested and provider_for_model(model) != "claudexor",
"use_local": use_local,
"timeout": _main_transport_timeout(model, deadline_ts, reserve_sec=transport_reserve_sec),
}
request_ref = persist_observed_call(

View file

@ -602,30 +602,78 @@ def tree_spend_line(tree_info: Any, ceiling: Optional[CostCeiling] = None) -> st
)
def _wrapup_stop_facts(
deciding_usd: Optional[float], ceiling: CostCeiling, global_remaining_usd: Optional[float],
) -> Tuple[str, str]:
"""The stop's money facts with the BINDING bound first, and the owner's way out when there is one.
The old sentence opened with the tree cap whatever had stopped the task, so an
owner whose shared wallet ran dry at $125 of a $400 cap read it as a broken
per-task cap. The bound that binds is the one with less room; the wallet is
also the only one the owner can lift while the task is still running."""
cap = ceiling.root_cap_usd
room = None if cap is None or deciding_usd is None else float(cap) - float(deciding_usd)
if deciding_usd is None:
spent = "This task's tree spend is unavailable"
elif cap is not None:
spent = f"This task's tree spent ${deciding_usd:.2f} of its own ${cap:.2f} cap"
else:
spent = f"This task's tree spent ${deciding_usd:.2f}"
if global_remaining_usd is not None and (room is None or float(global_remaining_usd) <= room):
cleared = ", so the task cap is not what stopped it" if room is not None else ""
return (
f"The shared Total budget is nearly used up: ${global_remaining_usd:.2f} left across all "
f"tasks. {spent}{cleared}.",
" Raise Total budget in Settings for more room.",
)
wallet = (
f"; the shared Total budget still has ${global_remaining_usd:.2f}"
if global_remaining_usd is not None else ""
)
return f"{spent}{wallet}.", ""
def wrapup_unaffordable_text(deciding_usd: Optional[float], ceiling: CostCeiling, global_remaining_usd: Optional[float] = None) -> str:
"""The owner-facing reason a task ends without even one affordable wrap-up send."""
cap = ceiling.root_cap_usd
cap_text = f" of the ${cap:.2f} hard tree cap" if cap is not None else ""
spent = f"Task tree spent ${deciding_usd:.3f}{cap_text}" if deciding_usd is not None else "Task-tree spend is unavailable"
wallet = f"; global model budget remaining is ${global_remaining_usd:.3f}" if global_remaining_usd is not None else ""
facts, way_out = _wrapup_stop_facts(deciding_usd, ceiling, global_remaining_usd)
return (
f"{spent}{wallet}; not even one wrap-up call can "
"be reserved, so the host delivers the retained evidence without a model synthesis."
f"{facts} Not even one wrap-up call can be reserved, so the host delivers the "
f"retained evidence without a model synthesis.{way_out}"
)
def wrapup_last_fit_text(deciding_usd: Optional[float], ceiling: CostCeiling, global_remaining_usd: Optional[float] = None) -> str:
"""The owner-facing reason a task claims the last affordable wrap-up send."""
cap = ceiling.root_cap_usd
cap_text = f" of the ${cap:.2f} hard tree cap" if cap is not None else ""
spent = f"Task tree spent ${deciding_usd:.3f}{cap_text}" if deciding_usd is not None else "Task-tree spend is unavailable"
wallet = f"; global model budget remaining is ${global_remaining_usd:.3f}" if global_remaining_usd is not None else ""
facts, way_out = _wrapup_stop_facts(deciding_usd, ceiling, global_remaining_usd)
return (
f"{spent}{wallet}; one wrap-up call is still "
"admissible, but another similarly reserved work call would consume that room."
f"{facts} One wrap-up call still fits and another work call of this size would not, "
f"so the task is finishing now with its best current answer.{way_out}"
)
def main_loop_wire_options(
model: str, *, allow_server_web_search: bool, bypass_response_cache: bool = False,
) -> Dict[str, Any]:
"""The payload-shaping options EVERY main-loop send declares — one owner for the send and for its pricing.
The budget wrap-up is admitted against a candidate built BEFORE the send (below)
that must equal it byte for byte. Each option the send grew on its own — cache
affinity, then ``stream`` — made the two payloads differ, and the admitted
answer was refused at dispatch. ``loop_llm_call.call_llm_with_retry`` spreads
this dict into its send and the prospective builder below spreads the same
one, so an option added here reaches both and one added elsewhere reaches one."""
from ouroboros.llm_claudexor import cache_key_for_model
from ouroboros.provider_models import provider_for_model
remote = provider_for_model(model) != "claudexor"
return {
"stream": True,
"cache_affinity": cache_key_for_model(model),
"allow_server_web_search": bool(allow_server_web_search) and remote,
"bypass_response_cache": bool(bypass_response_cache) and remote,
}
def prospective_wrapup_attempt_request(
*, llm: Any, messages: list[Dict[str, Any]], model: str,
reasoning_effort: str, tools: Optional[list[Dict[str, Any]]] = None,
@ -671,9 +719,12 @@ def prospective_wrapup_attempt_request(
return _merge_scope(replace(_attempt_request(target, candidate),
force_unknown_reservation=True, max_completion_tokens=MAIN_LOOP_MAX_TOKENS))[0]
with request_wire_call_scope():
# The send's own options, from their one owner: a key the send carries and
# this copy does not is a different payload, and the admitted answer is refused.
candidate = llm._build_remote_candidate(
target, messages, reasoning_effort, MAIN_LOOP_MAX_TOKENS, "auto", None, tools,
skip_capability_fetch=True, allow_server_web_search=allow_server_web_search,
skip_capability_fetch=True,
**main_loop_wire_options(model, allow_server_web_search=allow_server_web_search),
)
llm._normalize_payload_cache_ttl(target, candidate)
candidate = _finalized_physical_candidate(

View file

@ -347,6 +347,7 @@ LEAVES: dict[str, tuple[str, str, frozenset[str]]] = {
"_degrade_retained_delivery_candidate", "_delivery_evidence_state",
"_delivery_replace_required", "_direct_child_results",
"_drain_forced_owner_directives", "_drain_incoming_messages",
"_emit_checkpoint_event",
"_end_task_acceptance_fence", "_finalize_forced_services",
"_finalize_task_services", "_force_plan_decision",
"_force_plan_disclosure", "_force_plan_reminder",

View file

@ -85,7 +85,20 @@ def test_projection_failure_and_torn_ledger_are_not_a_known_zero(tmp_path, monke
def test_global_remaining_is_disclosed_without_fabricating_a_tree_amount():
ceiling = task_pacing.CostCeiling(state="active", ceiling_usd=5.0)
text = task_pacing.wrapup_last_fit_text(None, ceiling, 2.0)
assert "spend is unavailable" in text and "global model budget remaining is $2.000" in text
assert "spend is unavailable" in text and "$2.00 left across all tasks" in text
def test_the_stop_sentence_opens_with_the_bound_that_binds():
"""An owner whose WALLET ran dry at $125 of a $400 cap must not read a per-task-cap story."""
capped = task_pacing.CostCeiling(state="active", ceiling_usd=239.66, root_cap_usd=400.0)
wallet = task_pacing.wrapup_last_fit_text(125.661, capped, 21.834)
assert wallet.startswith("The shared Total budget is nearly used up: $21.83 left across all tasks.")
assert "$125.66 of its own $400.00 cap, so the task cap is not what stopped it" in wallet
assert wallet.endswith("Raise Total budget in Settings for more room.")
cap = task_pacing.wrapup_last_fit_text(390.0, capped, 500.0)
assert cap.startswith("This task's tree spent $390.00 of its own $400.00 cap; the shared Total budget still has $500.00.")
assert "Raise Total budget" not in cap
assert task_pacing.wrapup_unaffordable_text(125.661, capped, 0.4).startswith("The shared Total budget is nearly used up")
def test_local_final_call_path_does_not_request_an_unneeded_wallet_projection(tmp_path, monkeypatch):
@ -148,7 +161,7 @@ def test_live_wallet_triggers_the_existing_final_call_path_without_a_root_cap(tm
accounting.reserve_attempt(_request(2.0, "another-root"))
result = loop._check_budget_limits(ctx, 10.0, ceiling)
assert result[0] == "verified final" and len(admitted) == 1 and len(prepared) == 1
assert "global model budget remaining is $2.000" in prepared[0]
assert "$2.00 left across all tasks" in prepared[0]
assert ctx.accumulated_usage["cost_stop_rail"] == "wrapup_reservation_last_fit"
assert _wrapup_global_remaining() == 0.5
accounting.release_attempt(admitted[0], "controlled final callback did not send")

View file

@ -17,6 +17,7 @@ import pytest
from ouroboros import task_pacing, usage_accounting
from ouroboros.contracts.task_contract import normalize_budget_profile
from ouroboros.loop import _check_budget_limits, _RoundLimitContext
from ouroboros.task_pacing import main_loop_wire_options
@pytest.fixture(autouse=True)
@ -211,6 +212,12 @@ class TestCacheAwareReservation:
# These pins rebuild the "actual" candidate by hand, so they carry the send's own
# options from their one owner; tests/test_wrapup_real_send_parity.py drives the
# real send, which is where a key these mirrors forget would be caught.
_MAIN_LOOP_OPTIONS = main_loop_wire_options("openai::gpt-test", allow_server_web_search=False)
def _patch_execute_candidate(monkeypatch, llm_module, execute):
"""Patch the physical candidate executor where the v7 lanes BIND it.
@ -381,6 +388,7 @@ class TestWrapupAffordability:
)
client._chat_anthropic(
target, messages, tools, "high", prospective.max_completion_tokens, "auto",
stream=_MAIN_LOOP_OPTIONS["stream"],
)
actual = captured["request"]
@ -416,7 +424,7 @@ class TestWrapupAffordability:
)
candidate = client._build_remote_candidate(
target, messages, "high", prospective.max_completion_tokens, "auto", None, tools,
skip_capability_fetch=True,
skip_capability_fetch=True, **_MAIN_LOOP_OPTIONS,
)
client._normalize_payload_cache_ttl(target, candidate)
client._create_chat_completion_with_retries(lambda **_kwargs: None, candidate, target)
@ -462,7 +470,7 @@ class TestWrapupAffordability:
)
candidate = client._build_remote_candidate(
target, messages, "high", prospective.max_completion_tokens, "auto", None, None,
skip_capability_fetch=True,
skip_capability_fetch=True, **_MAIN_LOOP_OPTIONS,
)
client._normalize_payload_cache_ttl(target, candidate)
client._create_chat_completion_with_retries(lambda **_kwargs: None, candidate, target)
@ -963,7 +971,7 @@ class TestWrapupAffordabilityRail:
assert result is not None
assert seen["source"] == "budget_wrapup_unaffordable"
assert seen["reason"] == "budget_exhausted"
assert "not even one wrap-up call" in seen["text"]
assert "Not even one wrap-up call" in seen["text"]
assert ctx.accumulated_usage["cost_stop_rail"] == "wrapup_reservation_last_fit"
assert [call.get("request") for call in calls] == [None, request, request]
@ -1040,7 +1048,7 @@ class TestWrapupAffordabilityRail:
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,
skip_capability_fetch=True, **_MAIN_LOOP_OPTIONS,
)
ctx.llm._normalize_payload_cache_ttl(target, candidate)
candidate = _finalized_physical_candidate(
@ -1092,7 +1100,7 @@ class TestWrapupAffordabilityRail:
def test_the_stop_text_names_the_cap_and_the_reason(self):
text = task_pacing.wrapup_last_fit_text(49.9, self._ceiling(50.0))
assert "$49.900" in text and "$50.00" in text
assert "$49.90 of its own $50.00 cap" in text
assert "wrap-up call" in text

View file

@ -0,0 +1,174 @@
"""The admitted budget wrap-up against the send the main loop REALLY makes.
Every earlier parity pin rebuilt the "actual" candidate by calling the prospective
builder's own builder with the prospective builder's own arguments, so a wire key
only the real send carries (``stream``) could never fail it — and in production it
cost three tasks their final answer in one night. This module drives
``call_llm_with_retry`` itself and meets the physical request at the executor seam,
the only place where the two payloads can honestly be compared."""
from __future__ import annotations
import queue
from types import SimpleNamespace
import pytest
from ouroboros import llm as llm_module
from ouroboros import task_pacing, usage_accounting
from ouroboros.contracts.task_contract import normalize_budget_profile
from ouroboros.llm import LLMClient
from ouroboros.llm_claudexor import cache_key_for_model
from ouroboros.loop import _check_budget_limits
from ouroboros.loop_llm_call import call_llm_with_retry
from ouroboros.task_pacing import main_loop_wire_options
from tests.test_tree_cost_ceiling import _ctx, _patch_execute_candidate
_IDENTITY = ("model", "provider", "candidate_raw_sha256", "candidate_raw_size_bytes")
_MESSAGES = [{"role": "system", "content": "policy"}, {"role": "user", "content": "wrap up"}]
_TOOLS = [{"type": "function", "function": {
"name": "probe", "description": "probe", "parameters": {"type": "object", "properties": {}},
}}]
_ROUTES = [
("openai::gpt-test", {"OPENAI_API_KEY": "unused"}),
("openai/gpt-test", {"OPENROUTER_API_KEY": "unused"}),
("anthropic/claude-test", {"OPENROUTER_API_KEY": "unused"}),
("anthropic::claude-test", {"ANTHROPIC_API_KEY": "unused"}),
]
class _Captured(Exception):
"""Stops the real send once the physical request exists."""
@pytest.mark.parametrize("model,env", _ROUTES)
def test_prospective_wrapup_candidate_is_the_candidate_the_main_loop_sends(monkeypatch, tmp_path, model, env):
for key, value in env.items():
monkeypatch.setenv(key, value)
captured = {}
def execute(request, send, before_dispatch):
captured["request"] = request
raise _Captured()
_patch_execute_candidate(monkeypatch, llm_module, execute)
client = LLMClient(api_key="unused")
logs = tmp_path / "logs"
logs.mkdir()
with usage_accounting.usage_scope(usage_accounting.UsageScope(
drive_root=tmp_path, task_id="parity", root_task_id="parity",
)):
prospective = task_pacing.prospective_wrapup_attempt_request(
llm=client, messages=_MESSAGES, model=model, reasoning_effort="high",
tools=_TOOLS, cache_affinity=cache_key_for_model(model),
)
try:
call_llm_with_retry(client, _MESSAGES, model, _TOOLS, "high", 1, logs, "parity", 1,
queue.Queue(), {}, initial_messages=_MESSAGES)
except _Captured:
pass
assert "request" in captured, "the real send never reached the physical executor"
assert {key: getattr(captured["request"], key) for key in _IDENTITY} == {
key: getattr(prospective, key) for key in _IDENTITY}
def test_every_wire_option_of_the_send_reaches_the_priced_copy():
"""The one owner: what the send declares is what the prospective builder is handed."""
options = main_loop_wire_options("openai/gpt-test", allow_server_web_search=True)
assert options["stream"] is True
assert set(options) == {"stream", "cache_affinity", "allow_server_web_search", "bypass_response_cache"}
def _completion(text):
body = {"id": "c1", "object": "chat.completion", "model": "gpt-test", "choices": [{
"index": 0, "finish_reason": "stop", "message": {"role": "assistant", "content": text}}],
"usage": {"prompt_tokens": 10, "completion_tokens": 5, "total_tokens": 15}}
return SimpleNamespace(model_dump=lambda: body)
def _last_fit_rail(monkeypatch, tmp_path, execute):
"""Drive the whole rail: last-fit decision -> admitted candidate -> the real send."""
monkeypatch.setenv("OPENROUTER_API_KEY", "unused")
monkeypatch.setattr("ouroboros.loop._loop_tree_accounting", lambda **_k: {"accounted_usd": 20.0})
answers = iter((True, False, True, False, True, False)) # proxy, exact probe, prepared: one fits, two do not
monkeypatch.setattr(task_pacing, "wrapup_reservation_fits", lambda **_kwargs: next(answers))
_patch_execute_candidate(monkeypatch, llm_module, execute)
logs = tmp_path / "logs"
logs.mkdir()
ctx = _ctx(drive_logs=logs, llm=LLMClient(api_key="unused"), active_model="openai/gpt-test")
ctx.messages = [dict(message) for message in _MESSAGES]
ceiling = task_pacing.resolve_cost_ceiling(None, normalize_budget_profile(None), root_cap_usd=50.0)
with usage_accounting.usage_scope(usage_accounting.UsageScope(
drive_root=tmp_path, task_id="task1", root_task_id="task1",
)):
return ctx, _check_budget_limits(ctx, None, ceiling)
def test_the_budget_rail_delivers_the_models_wrapup_as_a_best_effort_answer(monkeypatch, tmp_path):
from ouroboros.outcomes import EXECUTION_BEST_EFFORT, derive_loop_outcome
from ouroboros.project_dialogue import OUTCOME_PHASE_HEADLINE, outcome_phase
attempts, sends = [], []
def execute(request, send, before_dispatch):
attempts.append(request)
before_dispatch(SimpleNamespace(attempt_id=f"a{len(attempts)}", drive_root=tmp_path)) # the identity predicate runs here
sends.append(request)
return _completion("Done so far: parts 1-3 verified. Not finished: part 4.")
ctx, result = _last_fit_rail(monkeypatch, tmp_path, execute)
assert result is not None
assert (len(attempts), len(sends)) == (1, 1), "admitted on the first physical attempt: no refusal, no second send"
text, usage, trace = result
assert text.startswith("Done so far")
assert usage["reason_code"] == "budget_exhausted" and usage["_best_effort_extracted"] is True
axes = derive_loop_outcome(text, usage, trace)["outcome_axes"]
assert axes["execution"]["status"] == EXECUTION_BEST_EFFORT
phase = outcome_phase({"status": "completed", "outcome_axes": axes, "reason_code": "budget_exhausted"}, {})
assert OUTCOME_PHASE_HEADLINE[phase] == "Done with warnings" # not "Failed": the answer was delivered
def test_a_drifted_admitted_candidate_is_sent_once_more_instead_of_losing_the_answer(monkeypatch, tmp_path):
"""The predicate refuses before dispatch; the refusal must not cost the owner the final answer."""
import ouroboros.loop as loop_module
calls, events = [], []
real_once = loop_module._call_forced_model_once
def once(ctx, *, initial_messages=None, admitted_request=None):
calls.append(admitted_request is not None)
if admitted_request is not None:
raise usage_accounting.PhysicalAttemptPreconditionFailed(
"physical candidate precondition rejected dispatch", attempt_id="refused-1")
return real_once(ctx)
monkeypatch.setattr(loop_module, "_call_forced_model_once", once)
monkeypatch.setattr(loop_module, "_emit_checkpoint_event",
lambda _queue, _task, _logs, data: events.append(data))
ctx, result = _last_fit_rail(
monkeypatch, tmp_path, lambda request, send, before_dispatch: _completion("Recovered wrap-up."))
assert calls == [True, False], "one admitted send, then exactly one ordinary send"
assert result[0].startswith("Recovered wrap-up") and result[1]["_best_effort_extracted"] is True
drift = [event for event in events if event.get("checkpoint_kind") == "forced_candidate_drift"]
assert len(drift) == 1 and drift[0]["refused_attempt_id"] == "refused-1"
assert drift[0]["admitted"]["model"] and drift[0]["admitted"]["candidate_raw_sha256"]
def test_a_closed_dispatch_window_is_a_deadline_not_drift(monkeypatch, tmp_path):
import ouroboros.loop as loop_module
from ouroboros.llm_attempt import PhysicalDispatchInterrupted
calls = []
def once(ctx, *, initial_messages=None, admitted_request=None):
calls.append(admitted_request is not None)
raise PhysicalDispatchInterrupted("dispatch window closed", attempt_id="late-1")
monkeypatch.setattr(loop_module, "_call_forced_model_once", once)
ctx, result = _last_fit_rail(
monkeypatch, tmp_path, lambda request, send, before_dispatch: _completion("never reached"))
assert calls == [True], "a deadline refusal is never retried"
assert "never reached" not in result[0]