diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 6ac8e889e..260548b3f 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -1020,7 +1020,7 @@ The bridge recognizes `/panic`, `/restart`, `/review`, `/evolve [on|off]`, `/bg A user message enters through a reviewed transport, is admitted by the supervisor queue, and runs in `OuroborosAgent`. The root pipeline captures the task contract and immutable context core, executes the LLM/tool loop, preserves a delivery candidate, stores the result and artifacts, emits lifecycle and usage evidence, performs the root-only post-task work, and publishes the typed outcome. Queue admission proves only that asynchronous work was durably accepted; completion, objective satisfaction, artifact finality, verification, and review acceptance remain separate facts. -`DeliveryCandidate` is retained before verification or review so a later notice, reviewer failure, deadline, or provider outage cannot erase a useful answer. `outcomes.py` combines execution, objective, review, artifact, and child-absorption axes without converting one axis into another. Verify-before-done receipts and exact artifact references are host-attested evidence; declarations and answer prose are not substitutes. A forced exit may publish the best current candidate only with its typed rail and evidence-freshness disclosure, and lifecycle may remain `completed` while the objective or review axis records a best-effort or unaccepted result. When a forced exit fires while the delivery-control latch is armed, the model's one forced answer may legitimately be the protocol object: `loop._resolve_forced_delivery_control` resolves it purely (valid `keep` → retained candidate, valid `replace` → `full_answer`, malformed → retained candidate with the typed `delivery_control_degraded` reason) before suffixes and publication, never re-loops, and passes JSON through untouched when the latch is off — raw `{"delivery_control": ...}` never reaches the chat or the durable result. Provider death is the one forced rail that is NOT a best-effort completion: `_handle_provider_unavailable` still salvages the best available text into the result body, but stamps `infra_failed`, so the task terminalizes `failed` with the typed `provider_unavailable` reason and the supervisor sends the owner an immediate "provider outage — NOT completed" chat notification on the root's terminal dispatch. A waited-out transport outage reaches this rail through its deterministic no-resend branch (`transport_unavailable_no_resend`, keyed on the wait episode's latched cause): same salvage, same `infra_failed`/`provider_unavailable` truth, but no forced-final provider call is attempted over a proven-dead egress. +`DeliveryCandidate` is retained before verification or review so a later notice, reviewer failure, deadline, or provider outage cannot erase a useful answer. `outcomes.py` combines execution, objective, review, artifact, and child-absorption axes without converting one axis into another. Verify-before-done receipts and exact artifact references are host-attested evidence; declarations and answer prose are not substitutes. A forced exit may publish the best current candidate only with its typed rail and evidence-freshness disclosure, and lifecycle may remain `completed` while the objective or review axis records a best-effort or unaccepted result. When a forced exit fires while the delivery-control latch is armed, the model's one forced answer may legitimately be the protocol object: `loop._resolve_forced_delivery_control` resolves it purely (valid `keep` → retained candidate, valid `replace` → `full_answer`, malformed → retained candidate with the typed `delivery_control_degraded` reason) before suffixes and publication, never re-loops, and passes JSON through untouched when the latch is off — raw `{"delivery_control": ...}` never reaches the chat or the durable result. Provider death is the one forced rail that is NOT a best-effort completion: `_handle_provider_unavailable` still salvages the best available text into the result body, but stamps `infra_failed`, so the task terminalizes `failed` with the typed `provider_unavailable` reason and the supervisor sends the owner an immediate "provider outage — NOT completed" chat notification on the root's terminal dispatch. A waited-out transport outage reaches this rail through its deterministic no-resend branch (`transport_unavailable_no_resend`, keyed on the wait episode's latched cause): same salvage, same `infra_failed`/`provider_unavailable` truth, but no forced-final provider call is attempted over a proven-dead egress. That rail makes its one forced model call only while a call can still land: when the transport already spent its same-model retry wall — the attempt budget, or the deadline bounding the backoff — `call_llm_with_retry` stamps `_llm_retry_wall_exhausted` on the shared usage dict, and `_handle_provider_unavailable` reads it as a third sibling of its `context_overflow` and `provider_outcome_unknown` no-call gates (typed `source="retry_wall_exhausted_no_repay"`; every other forced rail keeps its one call, and the `deadline_local` shape keeps its grace call — the provider there is not proven dead) and ships the salvage DIRECTLY, running service finalization, the swarm-router short-circuit, status stamping and the ordinary candidate packaging while making no request at all (the unsent prompt never enters the transcript, so a replay cannot read it as a request the model ignored). The marker is stamped by the exception handler's ONE shared stop tail — every error class that reaches that tail is retry-same-request (permanent refusals already stopped inside `_record_llm_call_error` and leave the wall unspent, keeping the one forced call their class is entitled to) — and by the empty-response path only while the PROVIDER is failing (`finish_reason=None` glitch or a transient body error; an ordinary empty answer with `finish_reason="stop"/"length"` is a live provider, and the shorter forced tool-less prompt may well land). It is CLEARED AT ENTRY of every `call_llm_with_retry` invocation, because the primary and each fallback candidate share one usage dict; a transient-exhausted primary whose PERMANENT-failed fallback cleared the marker therefore re-pays one forced primary call — the disclosed residual of keeping the marker a last-invocation bool instead of a route-keyed ledger. Terminal delivery preserves producer authorship: complete model answers appear as Ouroboros; host-authored incident receipts appear as System; unreviewed intermediate output remains in durable task details. diff --git a/ouroboros/_usage_response.py b/ouroboros/_usage_response.py index 69b82bb93..713cb0840 100644 --- a/ouroboros/_usage_response.py +++ b/ouroboros/_usage_response.py @@ -19,17 +19,18 @@ def _plain(value: Any) -> Any: return value -def _number(value: Any) -> Optional[float]: +def provider_cost_value(value: Any) -> Optional[float]: """Parse a provider-reported cost; ``None`` means the value cannot be trusted. - This is the AUTHORITATIVE boundary — what survives here is what the durable - attempt ledger settles as final money — so it applies the same predicate as - the loop-side projection in ``loop_llm_call._provider_cost_value``: ``bool`` + THE cost-trust predicate, defined once at the AUTHORITATIVE boundary (what + survives here is what the durable attempt ledger settles as final money) and + imported by the loop-side projection, so the two lanes cannot fork: ``bool`` is rejected FIRST (``float(True)`` is a plausible-looking 1.0 that would settle as a FINAL $1.00 and eat real budget admission), then anything - unparseable, non-finite, or negative. A reported ``0.0`` stays a legitimate - zero. Rejecting here settles the attempt as unknown rather than as a - fabricated amount (BIBLE P1). + unparseable, non-finite, or negative — ``OverflowError`` included, because + ``float(10**1000)`` raises rather than returning inf. A reported ``0.0`` + stays a legitimate zero. Rejecting settles the attempt as unknown rather + than as a fabricated amount (BIBLE P1). Never raises. """ if isinstance(value, bool): return None @@ -40,6 +41,9 @@ def _number(value: Any) -> Optional[float]: return number if math.isfinite(number) and number >= 0 else None +_number = provider_cost_value # historical local name at this boundary + + def _reported_token_count(usage: Dict[str, Any], *keys: str) -> Optional[int]: """Return the first reported count; absence stays distinct from explicit zero.""" for key in keys: diff --git a/ouroboros/loop.py b/ouroboros/loop.py index f82f427a6..c673ca371 100644 --- a/ouroboros/loop.py +++ b/ouroboros/loop.py @@ -51,7 +51,7 @@ from ouroboros.loop_tool_execution import ( reclaim_negative_memo, reclaim_trace_refs, ) -from ouroboros.loop_llm_call import call_llm_with_retry, emit_llm_usage_event, forced_response_is_incomplete, forced_response_parts +from ouroboros.loop_llm_call import call_llm_with_retry, emit_llm_usage_event, forced_response_is_incomplete, forced_response_parts, provider_no_call_source from ouroboros.loop_transport import ( TransportWaitEpisode, end_episode_budget as _end_episode_budget, @@ -3369,15 +3369,18 @@ def _handle_provider_unavailable( if str(usage.get("reason_code") or "") == "provider_unavailable": usage["execution_status"] = RESULT_INFRA_FAILED return text, usage, llm_trace - if str(ctx.accumulated_usage.get("_last_llm_error_kind") or "") == "provider_outcome_unknown": - live_trace = getattr(ctx, "llm_trace", None) - llm_trace = live_trace if isinstance(live_trace, dict) else {} + # No-call shapes; see provider_no_call_source + no_call, wall = provider_no_call_source(ctx.accumulated_usage, is_deadline_exhausted) + if no_call: + if wall: + _finalize_forced_services(ctx, llm_trace) + _drain_forced_owner_directives(ctx, llm_trace) text, usage, llm_trace = _forced_fallback_result( ctx, llm_trace, fallback, reason_code="provider_unavailable", - source="provider_outcome_unknown_no_resend", + source=no_call, provider_terminal=wall, ) - if str(usage.get("reason_code") or "") == "provider_unavailable": - usage["execution_status"] = RESULT_INFRA_FAILED + if usage.get("execution_status") is not None: + usage.update(execution_status=RESULT_INFRA_FAILED, reason_code="provider_unavailable") return text, usage, llm_trace prompt = ( "[DEADLINE] Primary model work reached the owner deadline. Produce the best final answer now from verified work and state what remains undone." @@ -5409,15 +5412,12 @@ def _forced_final_answer( source="forced_model_incomplete", provider_terminal=True, ) - extracted, control_degraded = _resolve_forced_delivery_control( - getattr(getattr(ctx, "tools", None), "_ctx", None), extracted, - ) + extracted, control_degraded = _resolve_forced_delivery_control(tools_ctx, extracted) if extracted: ctx.accumulated_usage["_best_effort_extracted"] = True - tool_ctx = getattr(getattr(ctx, "tools", None), "_ctx", None) plan_suffix = ( - _force_plan_disclosure(tool_ctx, llm_trace, forced_reason=reason_code) - if tool_ctx is not None else "" + _force_plan_disclosure(tools_ctx, llm_trace, forced_reason=reason_code) + if tools_ctx is not None else "" ) if provider_terminal: ctx.accumulated_usage["terminal_plan_review_open"] = bool(plan_suffix) diff --git a/ouroboros/loop_llm_call.py b/ouroboros/loop_llm_call.py index 2db42f277..00abcbd96 100644 --- a/ouroboros/loop_llm_call.py +++ b/ouroboros/loop_llm_call.py @@ -9,7 +9,6 @@ Extracted from loop.py to keep the main loop orchestrator focused. from __future__ import annotations import hashlib -import math import os import pathlib import queue @@ -32,6 +31,7 @@ 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.transport_custody import is_pre_dispatch_transport_failure +from ouroboros._usage_response import provider_cost_value as _provider_cost_value from ouroboros.usage_accounting import ( PhysicalAttemptContext, UsageAccountingError, @@ -159,6 +159,9 @@ def fold_retrieval_usage(accumulated_usage: Dict[str, Any], usage: Dict[str, Any # large) keep failing fast. There is NO cross-model fallback here — the same # request is retried on the SAME model. _TRANSIENT_RETRY_KINDS = frozenset({"provider_transient", "provider_incomplete_response"}) +# OB-01: stamped when THIS invocation spent its same-model retry wall without a +# usable response; entry-cleared per invocation; PERMANENT classes leave it unspent. +RETRY_WALL_EXHAUSTED_KEY = "_llm_retry_wall_exhausted" # Error kinds that put a model on the F1 fallback cooldown. Superset of the same-model # retry kinds: a body-error 429 (HTTP 200 with an error in the body — the canonical # cloud.ru/OpenRouter rate-limit shape) is classified "rate_limit", which must cool the @@ -851,29 +854,6 @@ def _remember_llm_call( usage.setdefault("llm_call_refs", []).append(call_meta) -def _provider_cost_value(raw: Any) -> Optional[float]: - """Parse a provider-reported ``usage["cost"]``; ``None`` means INVALID. - - Provider cost is data from outside the process, so it is validated before it - is trusted (OB-10). ``bool`` is rejected FIRST because ``float(True)`` is a - perfectly plausible ``1.0`` — a fabricated tariff, not a reported one — and - then anything unparseable, non-finite (NaN/inf), or negative. A legitimate - reported ``0.0`` (free tier, cached round) is valid and returned as itself. - - The same predicate guards the AUTHORITATIVE boundary in - ``_usage_response._number``, where the durable attempt ledger settles cost. - Both are needed: this one keeps the projection honest, that one keeps an - invalid amount from becoming final money. - """ - if isinstance(raw, bool): - return None - try: - value = float(raw) - except (TypeError, ValueError): - return None - return value if math.isfinite(value) and value >= 0 else None - - def _normalize_usage_cost( usage: Dict[str, Any], *, @@ -889,16 +869,9 @@ def _normalize_usage_cost( cost = 0.0 display_model = f"{model} (local)" elif provider_reported_cost and cost is None: - # MISSING and INVALID are different facts. A missing cost falls through to - # the catalog estimate below; a cost the provider DID send but that cannot - # be trusted is honestly unknown RIGHT HERE — estimation is skipped, because - # substituting a catalog price for a number the provider actually reported - # would misrepresent whose figure it is. Never raise (an unparseable value - # used to kill the whole round), never 0.0, never a fabricated tariff. - # This is the PROJECTION lane only. The durable attempt ledger settles from - # the provider response through `_usage_response.usage_from_response`, which - # applies the SAME predicate at that authoritative boundary — both lanes must - # reject, or an invalid cost still becomes final money (BIBLE P1, Invariant 4). + # MISSING falls through to the catalog estimate; a cost the provider DID + # send but that cannot be trusted is honestly unknown RIGHT HERE. Shared + # predicate: `_usage_response.provider_cost_value` — the lanes cannot fork. log.warning( "Provider reported an invalid cost (type=%s, value=%s) for %s; recording " "cost as unknown and skipping estimation", @@ -907,9 +880,7 @@ def _normalize_usage_cost( display_model, ) usage["cost"] = None - # A `cost_final` the wire set beside a cost we just rejected would make the - # record contradict itself: unknown cost is never a closed book. - usage["cost_final"] = False + usage["cost_final"] = False # unknown cost is never a closed book return None, display_model, provider, bool(usage.get("cost_estimated")) elif cost is None: cost = estimate_cost_optional( @@ -1028,6 +999,58 @@ def _record_llm_call_error( return False +def _stop_after_llm_error( + ctx: _LlmErrorContext, *, max_retries: int, transient_budget: int, + deadline_ts: Optional[float], +) -> bool: + """Retry decision for the exception path: ``True`` stops the attempt loop, and + every stop except the transport shape stamps the OB-01 marker — permanent + classes stopped earlier in ``_record_llm_call_error``, so every kind here is + retry-same-request and a stop means "the same-model wall is spent".""" + accumulated_usage = ctx.accumulated_usage + error_kind = str(accumulated_usage.get("_last_llm_error_kind") or "") + if error_kind == "transport_unavailable": + # One physical attempt per invocation: a pre-dispatch transport failure + # is free ($0) and the round-level wait episode owns redial pacing. NOT + # a spent wall — the transport-wait terminal owns this shape. + return True + is_transient = error_kind in _TRANSIENT_RETRY_KINDS + # Non-transient retryable classes keep max_retries but never exceed the loop + # ceiling, so an attempt_cap'd fallback wastes no backoff sleep (primary: no-op). + attempt_budget = transient_budget if is_transient else min(max_retries, transient_budget) + if ctx.attempt < attempt_budget - 1: + backoff = _retry_backoff_sec(accumulated_usage, error_kind, ctx.attempt, is_transient) + if _sleep_within_deadline(backoff, deadline_ts): + return False + _emit_retry_deadline_exhausted( + ctx.drive_logs, task_id=ctx.task_id, execution_id=ctx.execution_id, + round_id=ctx.round_id, round_idx=ctx.round_idx, attempt=ctx.attempt, + model=ctx.model, error_kind=error_kind, + ) + accumulated_usage[RETRY_WALL_EXHAUSTED_KEY] = True + return True + + +def _empty_response_wall_spent(is_provider_glitch: bool, permanent_body_error: bool, usage: Dict[str, Any]) -> bool: + """Empty-response exit: spent only while the PROVIDER is failing (finish=None + glitch or transient body error); a permanent body error and a live provider + (finish="stop"/"length") both keep their one forced call.""" + return not permanent_body_error and (is_provider_glitch or bool(usage.get("provider_error"))) + + +def provider_no_call_source(accumulated_usage: Dict[str, Any], deadline_exhausted: bool) -> Tuple[str, bool]: + """The provider-unavailable rail's no-call decision → (typed source, wall_spent). + Unknown in-flight outcome forbids a RESEND (outranks the wall); a SPENT + same-model wall makes one more forced call a second full retry window; the + deadline_local rail keeps its grace call (provider not proven dead), which is + why ``deadline_exhausted`` suppresses the wall answer.""" + if str(accumulated_usage.get("_last_llm_error_kind") or "") == "provider_outcome_unknown": + return "provider_outcome_unknown_no_resend", False + if not deadline_exhausted and bool(accumulated_usage.get(RETRY_WALL_EXHAUSTED_KEY)): + return "retry_wall_exhausted_no_repay", True + return "", False + + def _emit_empty_response_events( event_type: str, *, @@ -1219,27 +1242,10 @@ def _handle_main_llm_call_exception( _clear_custom_receipts(ctx.accumulated_usage) if _record_llm_call_error(error, ctx): return True - error_kind = str(ctx.accumulated_usage.get("_last_llm_error_kind") or "") - if error_kind == "transport_unavailable": - # Exactly ONE physical attempt per invocation: a proven pre-dispatch - # transport failure is free ($0 released) and an in-helper burst cannot - # cure a dead egress — the round-level wait episode owns redial pacing. - return True - is_transient = error_kind in _TRANSIENT_RETRY_KINDS - attempt_budget = transient_budget if is_transient else min(max_retries, transient_budget) - if ctx.attempt >= attempt_budget - 1: - return True - backoff = _retry_backoff_sec( - ctx.accumulated_usage, error_kind, ctx.attempt, is_transient, + return _stop_after_llm_error( + ctx, max_retries=max_retries, transient_budget=transient_budget, + deadline_ts=deadline_ts, ) - if _sleep_within_deadline(backoff, deadline_ts): - return False - _emit_retry_deadline_exhausted( - ctx.drive_logs, task_id=ctx.task_id, execution_id=ctx.execution_id, - round_id=ctx.round_id, round_idx=ctx.round_idx, attempt=ctx.attempt, - model=ctx.model, error_kind=error_kind, - ) - return True def _replace_response_meta( @@ -1292,33 +1298,6 @@ def _emit_llm_operation( ) -def _handle_llm_call_exception(error: Exception, ctx: _LlmErrorContext) -> bool: - _clear_custom_receipts(ctx.accumulated_usage) - _emit_llm_operation( - ctx.event_queue, ctx.task_id, ctx.llm_call_id, "failed", ctx.task_attempt, - ctx.execution_id, ctx.round_id, - ) - if _record_llm_call_error(error, ctx): - return True - kind = str(ctx.accumulated_usage.get("_last_llm_error_kind") or "") - if kind == "transport_unavailable": - # One physical attempt per invocation — see _handle_main_llm_call_exception. - return True - transient = kind in _TRANSIENT_RETRY_KINDS - budget = ctx.transient_budget if transient else min(ctx.max_retries, ctx.transient_budget) - if ctx.attempt >= budget - 1: - return True - backoff = _retry_backoff_sec(ctx.accumulated_usage, kind, ctx.attempt, transient) - if _sleep_within_deadline(backoff, ctx.deadline_ts): - return False - _emit_retry_deadline_exhausted( - ctx.drive_logs, task_id=ctx.task_id, execution_id=ctx.execution_id, - round_id=ctx.round_id, round_idx=ctx.round_idx, attempt=ctx.attempt, - model=ctx.model, error_kind=kind, - ) - return True - - def call_llm_with_retry( llm: LLMClient, messages: List[Dict[str, Any]], @@ -1346,6 +1325,7 @@ def call_llm_with_retry( msg = None _replace_response_meta(response_meta_out) drive_root = pathlib.Path(drive_logs).parent + accumulated_usage.pop(RETRY_WALL_EXHAUSTED_KEY, None) # last-invocation marker (see key) execution_id = str(accumulated_usage.setdefault("execution_id", new_execution_id())) round_id = f"{execution_id}:round:{round_idx}" context_fit_event_fields = _context_fit_event_fields(accumulated_usage) if physical_context is not None else {} @@ -1361,7 +1341,6 @@ def call_llm_with_retry( task_id=task_id, model=model, round_idx=round_idx, reserve_sec=transport_reserve_sec, ): return None, None - accumulated_usage["_llm_attempts_used"] = attempt + 1 llm_call_id = new_call_id("llm") call_identity = (task_id, task_attempt, llm_call_id, execution_id, round_id, attempt + 1) request_ref: Dict[str, Any] = {} @@ -1533,10 +1512,11 @@ def call_llm_with_retry( round_id=round_id, round_idx=round_idx, attempt=attempt, model=model, error_kind=event_type, ) + if _empty_response_wall_spent(is_provider_glitch, permanent_body_error, usage): + accumulated_usage[RETRY_WALL_EXHAUSTED_KEY] = True return None, cost - accumulated_usage.pop("execution_status", None) - accumulated_usage.pop("result_status", None) - accumulated_usage.pop("reason_code", None) + for stale in ("execution_status", "result_status", "reason_code", RETRY_WALL_EXHAUSTED_KEY): + accumulated_usage.pop(stale, None) accumulated_usage["rounds"] = accumulated_usage.get("rounds", 0) + 1 prompt_tokens = int(usage.get("prompt_tokens") or 0) completion_tokens = int(usage.get("completion_tokens") or 0) diff --git a/ouroboros/size_ratchet_manifest.py b/ouroboros/size_ratchet_manifest.py index 7853eac66..5bf57f76f 100644 --- a/ouroboros/size_ratchet_manifest.py +++ b/ouroboros/size_ratchet_manifest.py @@ -192,6 +192,7 @@ BAND_PATHS = { "tests/test_onboarding_wizard.py": None, "tests/test_owner_stop_s3.py": "Entered the band from 821 lines: the S3 contract suite now covers retry-root aliasing, graceful-to-immediate hardening, stale-control drain races, hard deadline preservation, descendant settlement failure, and late resweep exactly-once root finalization.", "tests/test_packaged_runtime_and_lifecycle.py": None, + "tests/test_provider_failure_reporting.py": "Entered the band from 938 lines: the provider-failure honesty rework added the retry-wall marker suites (entry-clear, all-retryable stamp, empty-response provider-failing gate, no-repay rail) beside the cost-validation suites (bool/NaN/inf/negative/huge-int at both boundaries) - one file per failure-reporting surface.", "tests/test_repo_health_smoke.py": "size-ratchet redesign: merge-aware previous, pairwise base-vs-tip, candidate-mode generator contract tests", "tests/test_review_cycles_dispatch.py": "Task acceptance wallet and paid-stamp dispatch regression matrix was integrated from the current managed target alongside the existing review-cycle tests.", "tests/test_review_fidelity.py": None, @@ -228,7 +229,7 @@ BYTE_BASELINE_DEBT = { } BYTE_DEBT = { - "ouroboros/loop.py": 312794, + "ouroboros/loop.py": 312772, "tests/test_delegated_subagent_transport.py": 320340, "tests/test_devtools_benchmarks.py": 328116, "web/modules/chat.js": 224315, diff --git a/tests/test_loop_misc.py b/tests/test_loop_misc.py index e85f1d7e9..d352fab3b 100644 --- a/tests/test_loop_misc.py +++ b/tests/test_loop_misc.py @@ -16,7 +16,10 @@ import queue import threading from types import SimpleNamespace +import pytest + import ouroboros.loop as loop_mod +from ouroboros.loop_llm_call import RETRY_WALL_EXHAUSTED_KEY from ouroboros.loop import ( _drain_incoming_messages, _initialize_owner_directives, @@ -1860,3 +1863,175 @@ def test_undecodable_image_fails_the_attach_not_the_provider_call(): vision._downscale_image_for_vlm(corrupt, "image/png") out, mime = vision._downscale_image_for_vlm(good, "image/png") assert out == good and mime == "image/png" + + +# --------------------------------------------------------------------------- +# OB-01 — provider death does not re-burn the budget on a second forced call +# --------------------------------------------------------------------------- + + +def _provider_death_ctx(tmp_path, accumulated): + """A minimal forced-rail context. ``tools=None`` on purpose: every forced + helper then takes its documented tool-less path, so the test observes the + loop's own decisions instead of a registry stub's.""" + return loop_mod._RoundLimitContext( + messages=[ + {"role": "user", "content": "do the thing"}, + {"role": "assistant", "content": "PARTIAL RESULT: step one is done."}, + ], + llm=SimpleNamespace(), + active_model="openai/gpt-5.5", + active_effort="medium", + max_retries=3, + drive_logs=tmp_path, + task_id="task-provider-death", + round_idx=7, + event_queue=None, + accumulated_usage=accumulated, + task_type="task", + active_use_local=False, + max_rounds=200, + drive_root=tmp_path, + ) + + +def _count_forced_helpers(monkeypatch): + """Wrap the non-model work of the forced rail with call counters, keeping the + REAL implementations so "it still runs" is observed, not simulated.""" + seen = {"model": [], "services": [], "drain": []} + real_services = loop_mod._finalize_forced_services + real_drain = loop_mod._drain_forced_owner_directives + + def _model(ctx): + seen["model"].append(1) + return "FRESH FORCED ANSWER" + + def _services(ctx, trace): + seen["services"].append(1) + return real_services(ctx, trace) + + def _drain(ctx, trace): + seen["drain"].append(1) + return real_drain(ctx, trace) + + monkeypatch.setattr(loop_mod, "_call_forced_model_once", _model) + monkeypatch.setattr(loop_mod, "_finalize_forced_services", _services) + monkeypatch.setattr(loop_mod, "_drain_forced_owner_directives", _drain) + return seen + + +def test_provider_death_skips_the_forced_call_when_the_retry_wall_is_spent(tmp_path, monkeypatch): + """The whole point of OB-01: the transport already spent the same-model retry + wall, so the forced finalization must NOT re-burn the budget on a request that + cannot land. Everything the forced rail owns besides the call still runs.""" + seen = _count_forced_helpers(monkeypatch) + accumulated = { # realistic transport stamps of an exhausted transient wall + RETRY_WALL_EXHAUSTED_KEY: True, + "execution_status": "infra_failed", + "reason_code": "llm_api_error", + "_last_llm_error_kind": "provider_transient", + } + ctx = _provider_death_ctx(tmp_path, accumulated) + messages_before = [dict(m) for m in ctx.messages] + + text, usage, trace = loop_mod._handle_provider_unavailable(ctx) + + assert seen["model"] == [] # ZERO llm calls from this rail + assert len(seen["services"]) == 1 # services still finalized + assert len(seen["drain"]) == 1 # exactly ONE directive drain + # Salvage still delivered, and the delivery candidate still packaged. + assert "PARTIAL RESULT: step one is done." in text + assert trace["forced_finalization"]["reason_code"] == "provider_unavailable" + # Replay durability: an unsent prompt must never reach the transcript, or a + # resume would read it as a request the model ignored. + assert [dict(m) for m in ctx.messages] == messages_before + assert not any( + "[PROVIDER_UNAVAILABLE]" in str(m.get("content") or "") for m in ctx.messages + ) + # False completion: an outage is an INFRA FAILURE, and no model answer was + # extracted, so the best_effort gate must not be handed a typed success fact. + assert usage["reason_code"] == "provider_unavailable" + assert usage["execution_status"] == "infra_failed" + assert "_best_effort_extracted" not in usage + + +def test_provider_death_still_makes_the_forced_call_when_the_wall_is_unspent(tmp_path, monkeypatch): + """The skip is not the new default. A PERMANENT refusal (auth/quota/bad + request) fails fast and leaves the wall unspent, so the forced rail keeps the + one chance its class is entitled to.""" + seen = _count_forced_helpers(monkeypatch) + accumulated = {} # no marker: the transport never exhausted its retries + ctx = _provider_death_ctx(tmp_path, accumulated) + + text, usage, _trace = loop_mod._handle_provider_unavailable(ctx) + + assert seen["model"] == [1] + assert "FRESH FORCED ANSWER" in text + assert usage["execution_status"] == "infra_failed" + assert any( + "[PROVIDER_UNAVAILABLE]" in str(m.get("content") or "") for m in ctx.messages + ) + + +@pytest.mark.parametrize( + "marker, expect_call", + [ + ("unexpected-string", False), + (1, False), + ({"nested": "garbage"}, False), + (0, True), + ("", True), + (None, True), + ], + ids=["truthy_str", "truthy_int", "truthy_dict", "zero", "empty_str", "none"], +) +def test_malformed_retry_wall_marker_reads_as_a_plain_truth_value( + tmp_path, monkeypatch, marker, expect_call, +): + """A shared usage dict can hold anything, so the read is `bool(...)` and + nothing else: any truthy value means the wall is spent, any falsy value means + it is not. No shape assumption, no crash, no silent third behaviour.""" + seen = _count_forced_helpers(monkeypatch) + ctx = _provider_death_ctx(tmp_path, { + RETRY_WALL_EXHAUSTED_KEY: marker, + "execution_status": "infra_failed", "reason_code": "llm_api_error", + }) + + _text, usage, _trace = loop_mod._handle_provider_unavailable(ctx) + + assert seen["model"] == ([1] if expect_call else []) + assert usage["execution_status"] == "infra_failed" + + +@pytest.mark.parametrize( + "reason_code, marker", + [ + ("round_limit", "[ROUND_LIMIT] wrap up"), + ("finalization_grace", "[FINALIZE_NOW] wrap up"), + ("deadline_local", "[DEADLINE] wrap up"), + ("owner_requested_finalization", "[OWNER_STOP] wrap up"), + ("budget_exhausted", "[BUDGET] wrap up"), + ("children_unabsorbed", "[CHILDREN] wrap up"), + ], +) +def test_other_forced_rails_never_skip_even_with_the_wall_marker_set( + tmp_path, monkeypatch, reason_code, marker, +): + """The skip is gated on the provider-death rail SPECIFICALLY. Every OTHER + reason code `loop.py` passes to `_forced_final_answer` — round limit, + finalization grace, deadline, owner stop, budget exhaustion and unabsorbed + children — still makes its one forced call even when a transient wall was + spent earlier in the same task: those rails end for their own reasons, not + because the provider is unreachable. The list is exhaustive against the + `reason_code=` literals in loop.py, so a new rail cannot silently inherit + the skip.""" + seen = _count_forced_helpers(monkeypatch) + ctx = _provider_death_ctx(tmp_path, {RETRY_WALL_EXHAUSTED_KEY: True}) + + text, _usage, _trace = loop_mod._forced_final_answer( + ctx, prompt=marker, fallback_text="fallback", reason_code=reason_code, + ) + + assert seen["model"] == [1] + assert "FRESH FORCED ANSWER" in text + assert any(marker in str(m.get("content") or "") for m in ctx.messages) diff --git a/tests/test_provider_failure_reporting.py b/tests/test_provider_failure_reporting.py index 12d6b7783..a5b09211a 100644 --- a/tests/test_provider_failure_reporting.py +++ b/tests/test_provider_failure_reporting.py @@ -6,6 +6,7 @@ import pytest from ouroboros.loop_transport import provider_failure_hint as _provider_failure_hint from ouroboros.loop_llm_call import ( + RETRY_WALL_EXHAUSTED_KEY, _normalize_usage_cost, call_llm_with_retry, classify_llm_exception, @@ -1065,3 +1066,161 @@ def test_missing_provider_cost_still_uses_the_catalog_estimate(tmp_path, include assert cost == 0.99 estimate.assert_called_once() assert accumulated["cost"] == 0.99 + + +# --------------------------------------------------------------------------- +# OB-01 — the transient retry-wall marker (`RETRY_WALL_EXHAUSTED_KEY`) +# --------------------------------------------------------------------------- + + +def _no_sleep(monkeypatch): + import time as _time + monkeypatch.setattr(_time, "sleep", lambda _s: None) + + +def test_transient_exhaustion_marks_the_retry_wall(tmp_path, monkeypatch): + """The attempt budget running out on a TRANSIENT class is exactly what the + marker means: more attempts on this model are pointless.""" + _no_sleep(monkeypatch) + accumulated = {} + + msg, _cost = call_llm_with_retry( + _TransientFailingLLM(), [{"role": "user", "content": "hi"}], + "openai/gpt-5.5", None, "medium", 3, tmp_path, "task-wall", 1, None, + accumulated, "task", False, + ) + + assert msg is None + assert accumulated["_last_llm_error_kind"] == "provider_transient" + assert accumulated[RETRY_WALL_EXHAUSTED_KEY] is True + + +def test_transient_deadline_stop_marks_the_retry_wall(tmp_path, monkeypatch): + """The other half of the wall: the deadline refuses the next backoff.""" + import time as _time + monkeypatch.setattr(_time, "sleep", lambda _s: None) + accumulated = {} + + msg, _cost = call_llm_with_retry( + _TransientFailingLLM(), [{"role": "user", "content": "hi"}], + "openai/gpt-5.5", None, "medium", 3, tmp_path, "task-wall-deadline", 1, + None, accumulated, "task", False, deadline_ts=_time.time() + 1.0, + ) + + assert msg is None + assert accumulated[RETRY_WALL_EXHAUSTED_KEY] is True + + +def test_empty_response_exhaustion_marks_the_retry_wall(tmp_path, monkeypatch): + """The response-shaped exit marks too: a finish_reason=null glitch that never + recovered spent the same wall as a raised transient error.""" + _no_sleep(monkeypatch) + accumulated = {} + + msg, _cost = call_llm_with_retry( + _GlitchThenOkLLM(glitches=99), [{"role": "user", "content": "hi"}], + "openai/gpt-5.5", None, "medium", 3, tmp_path, "task-wall-empty", 1, + None, accumulated, "task", False, + ) + + assert msg is None + assert accumulated[RETRY_WALL_EXHAUSTED_KEY] is True + + +def test_permanent_failure_leaves_the_retry_wall_unspent(tmp_path, monkeypatch): + """A permanent class fails FAST — the wall is never spent, so the forced + finalization keeps the one chance that class is entitled to.""" + _no_sleep(monkeypatch) + accumulated = {} + + msg, _cost = call_llm_with_retry( + _FailingLLM(), [{"role": "user", "content": "hi"}], "openai/gpt-5.5", + None, "medium", 3, tmp_path, "task-auth-wall", 1, None, accumulated, + "task", False, + ) + + assert msg is None + assert accumulated["_last_llm_error_kind"] == "auth_error" + assert RETRY_WALL_EXHAUSTED_KEY not in accumulated + + +def test_permanent_body_error_leaves_the_retry_wall_unspent(tmp_path): + """Same rule on the response-shaped exit: a PERMANENT body error stopped + without spending the wall.""" + accumulated = {} + llm = _EmptyBodyErrorLLM( + {"kind": "provider_error", "code": 400, "message": "bad request"}, + ) + + msg, _cost = call_llm_with_retry( + llm, [{"role": "user", "content": "hi"}], "openai/gpt-5.5", None, + "medium", 3, tmp_path, "task-body-wall", 1, None, accumulated, "task", + False, attempt_cap=1, + ) + + assert msg is None + assert llm.calls == 1 + assert RETRY_WALL_EXHAUSTED_KEY not in accumulated + + +def test_retry_wall_marker_is_cleared_at_entry_of_every_invocation(tmp_path, monkeypatch): + """REGRESSION: the primary and every fallback candidate SHARE one + ``accumulated_usage``. A transient-exhausted primary followed by a + permanent-failed fallback must NOT leave the marker standing — otherwise the + permanent failure inherits "the wall is spent" and silently loses the one + forced call its class is entitled to.""" + _no_sleep(monkeypatch) + shared = {} + + call_llm_with_retry( + _TransientFailingLLM(), [{"role": "user", "content": "hi"}], + "openai/gpt-5.5", None, "medium", 3, tmp_path, "task-chain", 1, None, + shared, "task", False, + ) + assert shared[RETRY_WALL_EXHAUSTED_KEY] is True # the primary spent its wall + + call_llm_with_retry( + _FailingLLM(), [{"role": "user", "content": "hi"}], + "anthropic/claude-opus-5", None, "medium", 3, tmp_path, "task-chain", 1, + None, shared, "task", False, attempt_cap=2, + ) + + assert shared["_last_llm_error_kind"] == "auth_error" + assert RETRY_WALL_EXHAUSTED_KEY not in shared + + +class _MarkerSeedingLLM: + """Sets the wall marker DURING the call, i.e. AFTER `call_llm_with_retry`'s + entry-clear has already run. + + Seeding it before the call instead would be vacuous: the entry-clear alone + would satisfy the assertion and the successful-round pop could be deleted + without any test noticing. + """ + + def __init__(self, accumulated): + self.accumulated = accumulated + + def chat(self, **kwargs): + self.accumulated[RETRY_WALL_EXHAUSTED_KEY] = True + return ( + {"content": "ok"}, + {"provider": "anthropic", "resolved_model": "anthropic/claude-sonnet-4-6"}, + ) + + +def test_successful_round_pops_the_retry_wall_marker(tmp_path): + """A round that succeeds retires the marker beside the other stale per-round + bookkeeping — the wall is no longer spent. Only the success-path pop can + clear a marker set after entry, which is what makes this test load-bearing.""" + accumulated = {"execution_status": "infra_failed"} + + msg, _cost = call_llm_with_retry( + _MarkerSeedingLLM(accumulated), [{"role": "user", "content": "hi"}], + "anthropic::claude-sonnet-4-6", None, "medium", 1, tmp_path, + "task-wall-ok", 1, None, accumulated, "task", False, + ) + + assert msg == {"content": "ok"} + assert RETRY_WALL_EXHAUSTED_KEY not in accumulated + assert "execution_status" not in accumulated diff --git a/tests/test_provider_wall_rail_contract.py b/tests/test_provider_wall_rail_contract.py new file mode 100644 index 000000000..334065562 --- /dev/null +++ b/tests/test_provider_wall_rail_contract.py @@ -0,0 +1,115 @@ +"""Contract tests for the no-repay provider rail (OB-01), pinned against the +accumulated_usage stamps the transport ACTUALLY leaves behind (reason_code +"llm_api_error" / "llm_empty_response"), not a bare marker dict: + +1. A wall-exhausted provider death must terminalize with the typed + ``reason_code="provider_unavailable"`` + ``infra_failed`` pair — exactly what + fires the supervisor's "provider outage — NOT completed" owner notice + (supervisor/events.py keys on that reason_code). +2. A SCHEDULED swarm-router handoff (admission durably succeeded) must NOT be + clobbered with infra_failed by the rail's stamp gate — the router deliberately + pops execution_status/reason_code to keep the successful handoff truthful. +""" +import time +from types import SimpleNamespace + +import ouroboros.loop as loop_mod +from ouroboros.loop_llm_call import RETRY_WALL_EXHAUSTED_KEY, call_llm_with_retry + + +class _TransientFailingLLM: + def chat(self, **kwargs): + error = Exception("upstream 502") + error.status_code = 502 + raise error + + +def _rail_ctx(tmp_path, accumulated, tools=None): + kwargs = dict( + messages=[ + {"role": "user", "content": "do the thing"}, + {"role": "assistant", "content": "PARTIAL RESULT."}, + ], + llm=SimpleNamespace(), active_model="openai/gpt-5.5", active_effort="medium", + max_retries=3, drive_logs=tmp_path, task_id="task-wall-contract", + round_idx=7, event_queue=None, accumulated_usage=accumulated, + task_type="task", active_use_local=False, max_rounds=200, + drive_root=tmp_path, + ) + if tools is not None: + kwargs["tools"] = tools + return loop_mod._RoundLimitContext(**kwargs) + + +def test_wall_exhausted_rail_carries_the_typed_provider_unavailable_reason( + tmp_path, monkeypatch, +): + """Use the REAL transport stamps: after a genuinely exhausted transient wall, + accumulated carries reason_code="llm_api_error". The no-repay rail must still + terminalize with the typed provider_unavailable + infra_failed pair, or the + supervisor's provider-death owner notice never fires for the most common + provider-death shape (transient retries exhausted).""" + monkeypatch.setattr(time, "sleep", lambda _s: None) + accumulated = {} + msg, _ = call_llm_with_retry( + _TransientFailingLLM(), [{"role": "user", "content": "hi"}], + "openai/gpt-5.5", None, "medium", 2, tmp_path, "task-wall-contract", 1, + None, accumulated, "task", False, + ) + assert msg is None and accumulated[RETRY_WALL_EXHAUSTED_KEY] is True + monkeypatch.setattr( + loop_mod, "_call_forced_model_once", + lambda ctx: (_ for _ in ()).throw(AssertionError("no forced call")), + ) + + _text, usage, trace = loop_mod._handle_provider_unavailable(_rail_ctx(tmp_path, accumulated)) + + assert trace["forced_finalization"]["source"] == "retry_wall_exhausted_no_repay" + assert usage["execution_status"] == "infra_failed" + assert usage["reason_code"] == "provider_unavailable" # owner-notice trigger + + +def test_wall_exhausted_body_429_empty_is_infra_not_a_model_failure(tmp_path, monkeypatch): + """The empty-response wall shape (body-429 with finish_reason="stop") leaves + execution_status="failed"/reason_code="llm_empty_response" in accumulated; the + no-repay rail must upgrade it to the typed infra pair — a rate-limited-out + provider is an outage, not a model that answered badly.""" + monkeypatch.setattr( + loop_mod, "_call_forced_model_once", + lambda ctx: (_ for _ in ()).throw(AssertionError("no forced call")), + ) + accumulated = { + RETRY_WALL_EXHAUSTED_KEY: True, + "execution_status": "failed", + "reason_code": "llm_empty_response", + "_last_llm_error_kind": "rate_limit", + } + + _text, usage, _trace = loop_mod._handle_provider_unavailable(_rail_ctx(tmp_path, accumulated)) + + assert usage["execution_status"] == "infra_failed" + assert usage["reason_code"] == "provider_unavailable" + + +def test_scheduled_swarm_handoff_survives_the_no_call_rail(tmp_path): + """A durably admitted managed task is a SUCCESS the router deliberately keeps + truthful by popping execution_status/reason_code; the rail's stamp gate must + not overwrite it with infra_failed/provider_unavailable.""" + tools_ctx = SimpleNamespace( + task_metadata={"force_plan": True}, + is_ephemeral_turn=True, + _swarm_handoff_attempt={"status": "scheduled", "task_id": "t-child-1"}, + ) + accumulated = { + "_last_llm_error_kind": "provider_outcome_unknown", + "execution_status": "infra_failed", + "reason_code": "llm_api_error", + } + + text, usage, _trace = loop_mod._handle_provider_unavailable( + _rail_ctx(tmp_path, accumulated, tools=SimpleNamespace(_ctx=tools_ctx)), + ) + + assert "Swarm admitted managed task" in text + assert usage.get("execution_status") != "infra_failed" + assert usage.get("reason_code") != "provider_unavailable"