mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-04 12:56:34 +00:00
Convergence fixes A7-A13: finalize_now honest no-resend during outage, typed proxy chains, loopback route fact, terminal wordings, Q14 margin, contract tests
- A7: finalize_now (deadline/ceiling/owner stop) arriving during an active transport-wait episode no longer dispatches a forced provider call over the proven-dead egress: _maybe_early_finalize routes the exit through the transport no-resend terminal (finalize_now_transport_terminal helper closes the episode's durable evidence with detail=finalize_now); mailbox-drain control bookkeeping is unchanged and the task terminalizes promptly. - A8: is_pre_dispatch_transport_failure types the standard unreachable-proxy chain (requests ProxyError -> MaxRetryError -> urllib3 ProxyError -> NewConnectionError/ConnectTimeoutError); proxy HTTP responses and post-dispatch read failures stay untyped. - A9: additive route_is_loopback fact on AttemptRequest/PhysicalAttemptCapture (set from the target base_url host in _attempt_request via is_loopback_base_url); classify_llm_exception excludes loopback routes from transport_unavailable alongside provider=='local', so a stopped Ollama/LM Studio/vLLM server fails fast instead of waiting out a phantom outage. - A10: provider_terminal_fallback_text keys fail-fast wording on NOT wait_eligible, adds an eligible-but-zero-wait wording, and replaces INTERRUPTED with outage phrasing that cannot collide with the supervisor's STATUS_INTERRUPTED lifecycle term. - A11: the Q14 final free redial reserves a named 3s margin (_FINAL_REDIAL_MARGIN_SEC) so round-top overhead cannot get the granted redial refused by the admission gate. - A12: contract tests — error_kind_changed episode exit resumes ordinary fallback; unknown-outcome redial takes provider_outcome_unknown_no_resend with zero further dials; episode window carries zero dispatched paid attempts; nonstandard OUROBOROS_TRANSIENT_RETRY_MAX keeps the one-attempt transport contract; review actors' physical_attempt_limit(2) rail pinned. - A13: _WAIT_BACKOFF_START_SEC promoted to config.py as NETWORK_WAIT_BACKOFF_START_SEC (config stays 1600 lines); generalized _self_check_round comment; dropped redundant capture-None guard in the classify branch; CLB RUNBOOK notes headless tasks keep the idle rail as an additional bound. Size-ratchet manifest regenerated (loop.py byte debt shrinks 312847 -> 312831).
This commit is contained in:
parent
f360d47822
commit
e9bf6f148e
11 changed files with 523 additions and 87 deletions
|
|
@ -30,7 +30,9 @@ Field-tested configuration and operational hazards from the 2026-07-20 full 1-se
|
|||
supervisor's absolute per-attempt ceiling (`OUROBOROS_TASK_ABS_CEILING_SEC`, default 6h),
|
||||
not a deadline or budget rail: a dead egress holds the task up to that ceiling instead of
|
||||
failing it after the burst. The wait is visible as durable `network_wait` events in the
|
||||
isolated server's `events.jsonl`.
|
||||
isolated server's `events.jsonl`. Note: idle-rail survival via waiting progress notes
|
||||
requires a real chat thread; headless tasks without a `chat_id` keep the idle rail
|
||||
(reaper) as an additional bound on the wait.
|
||||
- `OUROBOROS_TOTAL_BUDGET=200` per domain-seed (measured 1-seed domain costs:
|
||||
poker $117 · bsm $60 · cohort $58 · code $33 · sales $26 · db $25 — a $60 cap silently
|
||||
truncates poker mid-rollout).
|
||||
|
|
|
|||
|
|
@ -46,11 +46,11 @@ OWNER_STOP_OUTER_CAP_SEC = 600
|
|||
NESTED_SETTLEMENT_MARGIN_SEC = 30 # Structural ordering margin, not a cognition timeout.
|
||||
# Owner-note cadence while a task waits out a provider-connection outage; the effective interval is min(this, idle_timeout/2) so the notes also keep the idle rail alive.
|
||||
NETWORK_WAIT_NOTE_INTERVAL_SEC = 300
|
||||
# Cadence for intrinsic self-pacing checkpoints when a task has NO deadline_at (headless
|
||||
# benchmark runs). Advisory only — surfaces elapsed/rounds/cost for self-pacing; 0 disables.
|
||||
# First free-redial pause of a transport-wait episode; doubles per wait iteration up to the existing 60s transient backoff cap (Q10: an existing bound, not a new knob).
|
||||
NETWORK_WAIT_BACKOFF_START_SEC = 4.0
|
||||
# Cadence for intrinsic self-pacing checkpoints when a task has NO deadline_at (headless benchmark runs). Advisory only — surfaces elapsed/rounds/cost for self-pacing; 0 disables.
|
||||
PACING_INTERVAL_DEFAULT_SEC = 600
|
||||
# Supervisor-loop liveness deadline (WS3, v6.34.0): a watchdog thread flags the main supervisor
|
||||
# loop STALLED if it has not ticked within this many seconds (healthy tick ~0.5s, real wedges only). 0 disables.
|
||||
# Supervisor-loop liveness deadline (WS3, v6.34.0): a watchdog thread flags the main supervisor loop STALLED if it has not ticked within this many seconds (healthy tick ~0.5s, real wedges only). 0 disables.
|
||||
SUPERVISOR_LIVENESS_DEADLINE_DEFAULT_SEC = 90
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -49,6 +49,7 @@ from ouroboros.usage_accounting import (
|
|||
last_physical_attempt_capture,
|
||||
usage_scope,
|
||||
)
|
||||
from ouroboros.transport_custody import is_loopback_base_url
|
||||
from ouroboros.utils import in_worker_process, sanitize_tool_result_for_log
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
|
@ -273,6 +274,7 @@ def _attempt_request(
|
|||
candidate_context_size_bytes=len(context),
|
||||
candidate_measurement_kind="canonical_json_v1",
|
||||
physical_context=current_physical_attempt_context(),
|
||||
route_is_loopback=is_loopback_base_url(target.get("base_url")),
|
||||
)
|
||||
|
||||
|
||||
|
|
@ -689,8 +691,8 @@ class LLMClient:
|
|||
_SUPPORTED_PARAMS_CACHE: Dict[str, set] = {}
|
||||
_SUPPORTED_PARAMS_FETCHED: bool = False
|
||||
# Did the one-shot /models fetch actually reach OpenRouter (HTTP 200 + parse)?
|
||||
# Distinguishes a provider OUTAGE from a route with no metadata, so Capability
|
||||
# Evidence can mark STATUS_FAILED (transient) vs STATUS_UNPROBEABLE (v6.33.0 P4).
|
||||
# Splits provider OUTAGE from a route with no metadata, so Capability Evidence
|
||||
# can mark STATUS_FAILED (transient) vs STATUS_UNPROBEABLE (v6.33.0 P4).
|
||||
_CAPABILITIES_FETCH_OK: bool = False
|
||||
# OpenRouter-reported context window per model id (provider_metadata evidence).
|
||||
_CONTEXT_LENGTH_CACHE: Dict[str, int] = {}
|
||||
|
|
@ -721,10 +723,9 @@ class LLMClient:
|
|||
cls._CAPABILITIES_FETCH_OK = False # set True only on a clean 200 + parse
|
||||
try:
|
||||
import requests
|
||||
# 5s, not 15s: this fetch is on the synchronous capability-probe path
|
||||
# behind the max-context-mode gate (settings save / max toggle). A slow
|
||||
# probe must fail-closed quickly (-> window unknown -> max blocked with
|
||||
# the owner-ack escape), never hang the save (v6.33.0 WS4 timing budget).
|
||||
# 5s, not 15s: this fetch sits on the synchronous capability-probe path
|
||||
# behind the max-context-mode gate (settings save / max toggle); a slow
|
||||
# probe must fail-closed quickly, never hang the save (v6.33.0 WS4).
|
||||
resp = requests.get(
|
||||
"https://openrouter.ai/api/v1/models",
|
||||
timeout=5,
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@
|
|||
|
||||
from __future__ import annotations
|
||||
|
||||
import functools
|
||||
import json
|
||||
import hashlib
|
||||
import os
|
||||
|
|
@ -52,7 +53,9 @@ from ouroboros.loop_tool_execution import (
|
|||
)
|
||||
from ouroboros.loop_llm_call import call_llm_with_retry, emit_llm_usage_event, forced_response_is_incomplete, forced_response_parts
|
||||
from ouroboros.loop_transport import (
|
||||
TransportWaitEpisode,
|
||||
fallback_chain_allowed as _fallback_chain_allowed,
|
||||
finalize_now_transport_terminal as _finalize_now_transport_terminal,
|
||||
last_assistant_text as _last_assistant_text,
|
||||
provider_terminal_fallback_text as _provider_terminal_fallback_text,
|
||||
reconcile_transport_wait as _reconcile_transport_wait,
|
||||
|
|
@ -361,10 +364,9 @@ def _resolve_task_cost_ceiling(
|
|||
|
||||
|
||||
# Bounded staleness for the two DECIDING cost surfaces (ceiling check and
|
||||
# milestone note). The free stash is refreshed by every dispatch under this
|
||||
# root — at most one round old, zero reads — but ONE round can block 900s in
|
||||
# wait_tasks while children spend (the shape both dead waves had), and the
|
||||
# pacing refresh only covers deadline-less tasks, so a round outliving this
|
||||
# milestone note). The free stash refreshes on every dispatch under this root,
|
||||
# but ONE round can block 900s in wait_tasks while children spend, and the
|
||||
# pacing refresh only covers deadline-less tasks — so a round outliving this
|
||||
# bound pays for exactly one real projection read. Never per-round (see the
|
||||
# usage_accounting telemetry note and the e4a87344 contention class).
|
||||
_TREE_ACCOUNTING_MAX_STALE_SEC = 120.0
|
||||
|
|
@ -1591,8 +1593,7 @@ def _execute_task_acceptance_panel(ctx: _TaskAcceptanceContext) -> Any:
|
|||
# cannot fit the remaining root budget is declined up front as a terminal
|
||||
# DEGRADED (no-quorum semantics) instead of dying mid-wave. The estimate
|
||||
# renders the REAL per-slot message pair; the rare second physical attempt
|
||||
# is deliberately not multiplied in — a fail-open coarse filter, not a
|
||||
# hard reservation.
|
||||
# is not multiplied in — a fail-open coarse filter, not a reservation.
|
||||
from ouroboros.tools.review_helpers import review_wave_budget_gate
|
||||
|
||||
try:
|
||||
|
|
@ -1795,8 +1796,7 @@ def _apply_task_acceptance_result(
|
|||
if dialogue_terminal:
|
||||
# v6.74.0 (A5): a reviewer quorum judged the dialogue no longer
|
||||
# actionable (unreachable_here / stable_disagreement). Finalize through
|
||||
# the EXISTING honest path, recording BOTH positions (findings in the
|
||||
# run record, dispositions on the obligation rows) with one owner-
|
||||
# the EXISTING honest path, recording BOTH positions with one owner-
|
||||
# visible line. Reviewer authorship — not a host timer or a
|
||||
# unilateral agent give-up.
|
||||
ctx.tools._ctx._task_acceptance_reviewed = True
|
||||
|
|
@ -2576,9 +2576,8 @@ def _compute_subagent_handoff(tools: Any, drive_root: Any, task_id: str, content
|
|||
# P5: the reminder is suppressed ONLY by structured signals — a child
|
||||
# discarded/cancelled (filtered above) or absorbed (unchanged
|
||||
# signature). NEVER by parsing final PROSE for status words. Fires once
|
||||
# per CHANGE, not every round; if the agent still finalizes with
|
||||
# unhandled children, the no-tool / forced finalization paths append a
|
||||
# loud orphan note via _forced_orphan_note (P1).
|
||||
# per CHANGE, not every round; finalizing with unhandled children still
|
||||
# appends a loud orphan note via _forced_orphan_note (P1).
|
||||
_ = nonterminal_children # (kept for readability; trigger is change-based)
|
||||
if children and signature and signature != previous:
|
||||
tools._ctx._subagent_handoff_signature = signature
|
||||
|
|
@ -2607,7 +2606,7 @@ def _maybe_inject_self_check(
|
|||
REMINDER_INTERVAL = 15
|
||||
if round_idx <= 1 or round_idx % REMINDER_INTERVAL != 0 or round_idx >= max_rounds:
|
||||
return False
|
||||
# Free transport-wait redials re-enter the SAME round: one self-check per round.
|
||||
# Non-incrementing round re-entries (e.g. free redials): one self-check per round.
|
||||
if accumulated_usage.get("_self_check_round") == round_idx:
|
||||
return False
|
||||
accumulated_usage["_self_check_round"] = round_idx
|
||||
|
|
@ -3315,14 +3314,15 @@ def _handle_owner_stop_finalization(
|
|||
|
||||
def _handle_provider_unavailable(
|
||||
ctx: _RoundLimitContext, *, error_kind: str = "provider_unavailable",
|
||||
wait_cause: str = "", waited: bool = False,
|
||||
wait_cause: str = "", waited: bool = False, wait_eligible: bool = True,
|
||||
) -> Tuple[str, Dict[str, Any], Dict[str, Any]]:
|
||||
"""Salvage provider failure without an unsafe retry.
|
||||
|
||||
``wait_cause`` is the transport-wait episode's latched cause: it selects the
|
||||
deterministic no-resend terminal even when a later refusal (e.g. the
|
||||
deadline admission gate) overwrote the mutable ``_last_llm_error_kind``.
|
||||
``waited`` keeps the terminal text honest for zero-wait turns.
|
||||
``waited``/``wait_eligible`` keep the terminal text honest for zero-wait
|
||||
turns (fail-fast vs a managed task whose window left no time to wait).
|
||||
"""
|
||||
kind = str(error_kind or "")
|
||||
is_context_overflow = kind == "context_overflow"
|
||||
|
|
@ -3345,6 +3345,7 @@ def _handle_provider_unavailable(
|
|||
fallback = _provider_terminal_fallback_text(
|
||||
ctx.accumulated_usage, is_context_overflow=is_context_overflow,
|
||||
is_transport_wait=is_transport_wait, waited=waited,
|
||||
wait_eligible=wait_eligible,
|
||||
is_deadline_exhausted=is_deadline_exhausted,
|
||||
)
|
||||
if is_context_overflow:
|
||||
|
|
@ -3360,9 +3361,8 @@ def _handle_provider_unavailable(
|
|||
return text, usage, llm_trace
|
||||
if is_transport_wait:
|
||||
# No-resend terminal: salvage, no forced-final call over a dead egress.
|
||||
# Stamp BEFORE the composer (_handle_owner_stop_finalization pattern):
|
||||
# a SCHEDULED swarm-router handoff deliberately clears it and stays
|
||||
# truthful; the guard mirrors the sibling no-resend branch below.
|
||||
# Stamp BEFORE the composer (owner-stop pattern): a SCHEDULED swarm
|
||||
# handoff deliberately clears it; guard mirrors the sibling below.
|
||||
live_trace = getattr(ctx, "llm_trace", None)
|
||||
llm_trace = live_trace if isinstance(live_trace, dict) else {}
|
||||
ctx.accumulated_usage["execution_status"] = RESULT_INFRA_FAILED
|
||||
|
|
@ -3434,20 +3434,28 @@ def _maybe_deadline_local_finalize(
|
|||
|
||||
def _maybe_early_finalize(
|
||||
limit_ctx: _RoundLimitContext, tools: ToolRegistry, controls: Dict[str, Any],
|
||||
*, allow_deadline_local: bool = True,
|
||||
*, transport_episode: Optional[TransportWaitEpisode] = None,
|
||||
) -> Optional[Tuple[str, Dict[str, Any], Dict[str, Any]]]:
|
||||
"""Consume supervisor grace first, then a local deadline."""
|
||||
if controls.get("finalize_now"):
|
||||
if transport_episode is not None:
|
||||
# Every finalize_now flavor during an active outage takes the
|
||||
# honest no-resend terminal (rationale in the helper); control
|
||||
# bookkeeping was done by the drain, terminalization is prompt.
|
||||
return _finalize_now_transport_terminal(
|
||||
transport_episode, drive_logs=limit_ctx.drive_logs,
|
||||
task_id=limit_ctx.task_id, model=limit_ctx.active_model,
|
||||
handle_provider_unavailable=functools.partial(
|
||||
_handle_provider_unavailable, limit_ctx),
|
||||
)
|
||||
if controls.get("finalize_deadline_ts") is not None:
|
||||
_narrow_round_deadline(
|
||||
limit_ctx, controls["finalize_deadline_ts"],
|
||||
)
|
||||
return _handle_forced_finalization(limit_ctx, str(controls["finalize_now"]))
|
||||
if not allow_deadline_local:
|
||||
# An active transport-wait episode owns the deadline sliver: its last
|
||||
# free redial + no-resend terminal replace the paid graceful finalize
|
||||
# call, which cannot succeed over a proven-dead egress and would fork
|
||||
# the terminal story (deadline_local vs transport no-resend).
|
||||
if transport_episode is not None:
|
||||
# An active episode owns the deadline sliver: its last free redial +
|
||||
# no-resend terminal replace the paid deadline_local finalize call.
|
||||
return None
|
||||
return _maybe_deadline_local_finalize(limit_ctx, tools)
|
||||
|
||||
|
|
@ -5965,18 +5973,14 @@ def _maybe_inject_finalization_nudges(
|
|||
emit_progress("No-op attempt nudge injected before final response.")
|
||||
llm_trace["reasoning_notes"].append("No-op attempt nudge injected before final response.")
|
||||
return True
|
||||
# P2 one-shot final-answer-marker nudge: the turn produced REAL work AND
|
||||
# visible prose but no FINAL ANSWER marker — the typed extractor would drop
|
||||
# it and a forced/deadline finalization would score empty. Strengthen the
|
||||
# BEHAVIOR (ask the agent to mark its OWN answer), never mine prose into a
|
||||
# claimed answer (Bible P5). Own latch, ordered AFTER verify/red/A3
|
||||
# (grounding outranks formatting); mutually exclusive with the A3 no-op
|
||||
# nudge; forced paths return earlier. Structural facts only. The protocol
|
||||
# gate alone suffices: answer_protocol="final_answer_line" itself declares
|
||||
# a machine-extracted deliverable, so the nudge must not ALSO require a
|
||||
# declared expected_output — GAIA-shaped contracts keep expected_output
|
||||
# empty, and that extra gate once suppressed the only salvage surface
|
||||
# (a v6.56.0 run finalized a last-round refusal empty despite 24 calls).
|
||||
# P2 one-shot final-answer-marker nudge: REAL work + visible prose but no
|
||||
# FINAL ANSWER marker — the typed extractor would drop it and a forced
|
||||
# finalization would score empty. Strengthen the BEHAVIOR (agent marks its
|
||||
# OWN answer), never mine prose into a claimed answer (Bible P5). Own
|
||||
# latch, ordered AFTER verify/red/A3; forced paths return earlier. The
|
||||
# protocol gate alone suffices: it must not ALSO require expected_output —
|
||||
# GAIA-shaped contracts keep it empty, and that extra gate once suppressed
|
||||
# the only salvage surface (v6.56.0: last-round refusal finalized empty).
|
||||
if (
|
||||
not getattr(tools._ctx, "_final_marker_nudged", False)
|
||||
and _answer_protocol_active(tools._ctx) # v6.60.0: marker nudge is protocol-gated
|
||||
|
|
@ -6761,8 +6765,7 @@ def run_llm_loop(
|
|||
try:
|
||||
while True:
|
||||
if free_redial:
|
||||
# A transport-wait redial re-enters the SAME logical round: the
|
||||
# outage must not burn the round budget.
|
||||
# A transport-wait redial re-enters the SAME logical round.
|
||||
free_redial = False
|
||||
else:
|
||||
round_idx += 1
|
||||
|
|
@ -6817,11 +6820,10 @@ def run_llm_loop(
|
|||
_owner_msg_seen,
|
||||
owner_ctx=ctx,
|
||||
)
|
||||
# Early-exit per round: supervisor finalize_now, else loop-local real-
|
||||
# deadline finalize (headless runs that get no finalize_now) — finalize
|
||||
# best-effort rather than be killed mid-step with nothing.
|
||||
# Early-exit per round: supervisor finalize_now, else loop-local
|
||||
# real-deadline finalize (headless runs with no finalize_now).
|
||||
_early_final = _maybe_early_finalize(
|
||||
limit_ctx, tools, _controls, allow_deadline_local=transport_wait is None)
|
||||
limit_ctx, tools, _controls, transport_episode=transport_wait)
|
||||
if _early_final is not None:
|
||||
text, accumulated_usage, forced_trace = _early_final
|
||||
_merge_finalization_trace(llm_trace, forced_trace)
|
||||
|
|
@ -6904,8 +6906,7 @@ def run_llm_loop(
|
|||
emit_progress=emit_progress, context_fit_plan=context_fit_plan,
|
||||
active_context_mode=active_context_mode)
|
||||
# Post-chain reconcile with the FRESH kind: a MID-chain outage
|
||||
# latches an episode too, so no forced-final call dials a dead
|
||||
# egress (rationale in reconcile_transport_wait's docstring).
|
||||
# latches too (see reconcile_transport_wait's docstring).
|
||||
transport_wait = _reconcile_transport_wait(
|
||||
transport_wait, ctx, msg_present=msg is not None,
|
||||
error_kind=str(accumulated_usage.get("_last_llm_error_kind") or ""),
|
||||
|
|
@ -6928,6 +6929,8 @@ def run_llm_loop(
|
|||
wait_cause=transport_wait.wait_cause if transport_wait is not None else "",
|
||||
waited=transport_wait is not None and (
|
||||
transport_wait.wait_iterations > 0 or transport_wait.redials > 0),
|
||||
wait_eligible=(
|
||||
transport_wait.wait_eligible if transport_wait is not None else True),
|
||||
)
|
||||
_merge_finalization_trace(llm_trace, forced_trace)
|
||||
return text, accumulated_usage, llm_trace
|
||||
|
|
|
|||
|
|
@ -799,9 +799,9 @@ def classify_llm_exception(exc: Exception, safe_error: str = "") -> LlmErrorClas
|
|||
"provider_outcome_unknown", False, status_code, provider_code,
|
||||
)
|
||||
if (
|
||||
capture is not None
|
||||
and str(getattr(capture, "state", "") or "") == "released"
|
||||
str(getattr(capture, "state", "") or "") == "released"
|
||||
and str(getattr(capture, "provider", "") or "") != "local"
|
||||
and not bool(getattr(capture, "route_is_loopback", False))
|
||||
and is_pre_dispatch_transport_failure(exc)
|
||||
):
|
||||
# Typed $0 fact: the request never left this host toward a REMOTE
|
||||
|
|
@ -809,7 +809,9 @@ def classify_llm_exception(exc: Exception, safe_error: str = "") -> LlmErrorClas
|
|||
# Retrying the same request is free and safe; pacing/waiting is owned
|
||||
# by the round-level episode in loop.py, never by this helper. A local
|
||||
# provider's connect failure stays on the generic path — a stopped
|
||||
# local server is not a network outage worth waiting out.
|
||||
# local server is not a network outage worth waiting out — and so does
|
||||
# a loopback OpenAI-compatible route (Ollama / LM Studio / vLLM): its
|
||||
# provider stamp is remote-shaped but the server is on this host.
|
||||
return LlmErrorClassification(
|
||||
"transport_unavailable", True, status_code, provider_code,
|
||||
)
|
||||
|
|
|
|||
|
|
@ -27,6 +27,7 @@ from dataclasses import dataclass
|
|||
from typing import Any, Callable, Dict, List, Optional
|
||||
|
||||
from ouroboros.config import (
|
||||
NETWORK_WAIT_BACKOFF_START_SEC,
|
||||
NETWORK_WAIT_NOTE_INTERVAL_SEC,
|
||||
get_finalization_grace_sec,
|
||||
get_task_idle_timeout_sec,
|
||||
|
|
@ -37,9 +38,11 @@ from ouroboros.utils import append_jsonl, utc_now_iso
|
|||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
# First free-redial pause; doubles per wait iteration up to the existing
|
||||
# transient backoff cap (Q10: the 60s cap is an existing bound, not a new knob).
|
||||
_WAIT_BACKOFF_START_SEC = 4.0
|
||||
# Reserve for the final free redial near the owner deadline (Q14): round-top
|
||||
# overhead (message drain, checkpoints, transcript seal, token measurement)
|
||||
# routinely eats about a second on long transcripts, and a granted redial that
|
||||
# the admission gate then refuses is a wasted grant.
|
||||
_FINAL_REDIAL_MARGIN_SEC = 3.0
|
||||
|
||||
|
||||
@dataclass
|
||||
|
|
@ -289,7 +292,7 @@ def transport_wait_step(
|
|||
if remaining is not None and remaining <= 0:
|
||||
return _ended("deadline_exhausted")
|
||||
backoff = min(
|
||||
_WAIT_BACKOFF_START_SEC * (2.0 ** min(episode.wait_iterations, 4)),
|
||||
NETWORK_WAIT_BACKOFF_START_SEC * (2.0 ** min(episode.wait_iterations, 4)),
|
||||
_TRANSIENT_BACKOFF_CAP_SEC,
|
||||
)
|
||||
note_interval = max(
|
||||
|
|
@ -299,9 +302,9 @@ def transport_wait_step(
|
|||
# The sleep never exceeds the note interval, so waiting notes keep the idle
|
||||
# rail alive even on owner-lowered idle timeouts.
|
||||
sleep_sec = min(backoff, note_interval)
|
||||
if remaining is not None and remaining < sleep_sec + 1.0:
|
||||
if remaining is not None and remaining < sleep_sec + _FINAL_REDIAL_MARGIN_SEC:
|
||||
# One last free redial just before the admission window closes (Q14).
|
||||
sleep_sec = max(0.0, remaining - 1.0)
|
||||
sleep_sec = max(0.0, remaining - _FINAL_REDIAL_MARGIN_SEC)
|
||||
episode.final_redial_done = True
|
||||
if time.monotonic() - episode.last_note_monotonic >= note_interval:
|
||||
episode.last_note_monotonic = time.monotonic()
|
||||
|
|
@ -327,6 +330,37 @@ def transport_wait_step(
|
|||
return True
|
||||
|
||||
|
||||
def finalize_now_transport_terminal(
|
||||
episode: TransportWaitEpisode,
|
||||
*,
|
||||
drive_logs: pathlib.Path,
|
||||
task_id: str,
|
||||
model: str,
|
||||
handle_provider_unavailable: Callable[..., Any],
|
||||
) -> Any:
|
||||
"""Route a finalize_now that lands during an active episode to the honest
|
||||
transport no-resend terminal.
|
||||
|
||||
Every finalize_now flavor (supervisor deadline, cost ceiling, owner stop)
|
||||
normally dispatches one forced summarize call — but over a proven-dead
|
||||
egress that paid path can only fail at $0 with identical salvage, so the
|
||||
deterministic no-resend terminal wins. The episode's durable evidence is
|
||||
closed with an ``ended`` row first; the caller passes a partial of its
|
||||
``_handle_provider_unavailable`` so terminal composition stays in loop.py.
|
||||
"""
|
||||
emit_network_wait_event(
|
||||
drive_logs, task_id=task_id, phase="ended",
|
||||
elapsed_sec=time.monotonic() - episode.started_monotonic,
|
||||
redials=episode.redials, model=model, detail="finalize_now",
|
||||
)
|
||||
return handle_provider_unavailable(
|
||||
error_kind="transport_unavailable",
|
||||
wait_cause=episode.wait_cause,
|
||||
waited=episode.wait_iterations > 0 or episode.redials > 0,
|
||||
wait_eligible=episode.wait_eligible,
|
||||
)
|
||||
|
||||
|
||||
def task_deadline_epoch(tools: Any) -> Optional[float]:
|
||||
"""Return the task deadline for retry backoff."""
|
||||
meta = getattr(tools._ctx, "task_metadata", {})
|
||||
|
|
@ -355,13 +389,17 @@ def provider_terminal_fallback_text(
|
|||
is_context_overflow: bool,
|
||||
is_transport_wait: bool,
|
||||
waited: bool,
|
||||
wait_eligible: bool = True,
|
||||
is_deadline_exhausted: bool,
|
||||
) -> str:
|
||||
"""Owner-facing terminal text when provider death left nothing to salvage.
|
||||
|
||||
``waited`` is the episode's wait fact (iterations or redials happened): a
|
||||
wait-ineligible direct/ephemeral turn fails fast and must not claim it
|
||||
"waited and redialed until its own limits ran out".
|
||||
``wait_eligible`` is the episode's turn-class fact and ``waited`` its wait
|
||||
fact (iterations or redials happened). The fast-fail wording is keyed on
|
||||
NOT wait_eligible: a managed task whose admission window was already spent
|
||||
before the first wait iteration is not "this interactive turn". The
|
||||
waited-out wording deliberately avoids the supervisor's lifecycle term
|
||||
INTERRUPTED (STATUS_INTERRUPTED means pre-requeue, not terminal).
|
||||
"""
|
||||
if is_context_overflow:
|
||||
return (
|
||||
|
|
@ -369,16 +407,23 @@ def provider_terminal_fallback_text(
|
|||
"Any files written so far are preserved in the workspace."
|
||||
)
|
||||
if is_transport_wait:
|
||||
if not wait_eligible:
|
||||
return (
|
||||
"⚠️ Could not establish a provider connection; this interactive turn fails fast "
|
||||
"— retry when connectivity returns. Any files written so far are preserved in "
|
||||
"the workspace."
|
||||
)
|
||||
if waited:
|
||||
return (
|
||||
"⚠️ Could not establish a provider connection; the task waited and redialed "
|
||||
"until its own limits ran out and was INTERRUPTED, not completed. Any files "
|
||||
"written so far are preserved in the workspace. Retry when connectivity returns."
|
||||
"until its own limits ran out and ended as a provider outage, not completed. "
|
||||
"Any files written so far are preserved in the workspace. Retry when "
|
||||
"connectivity returns."
|
||||
)
|
||||
return (
|
||||
"⚠️ Could not establish a provider connection; this interactive turn fails fast "
|
||||
"— retry when connectivity returns. Any files written so far are preserved in "
|
||||
"the workspace."
|
||||
"⚠️ Could not establish a provider connection, and the owner deadline left no "
|
||||
"time to wait; the task ended as a provider outage, not completed. Any files "
|
||||
"written so far are preserved in the workspace. Retry when connectivity returns."
|
||||
)
|
||||
if is_deadline_exhausted:
|
||||
return "⚠️ The owner deadline ended primary model work; any files written so far are preserved."
|
||||
|
|
|
|||
|
|
@ -221,7 +221,7 @@ BYTE_BASELINE_DEBT = {
|
|||
}
|
||||
|
||||
BYTE_DEBT = {
|
||||
"ouroboros/loop.py": 312847,
|
||||
"ouroboros/loop.py": 312831,
|
||||
"tests/test_delegated_subagent_transport.py": 320568,
|
||||
"tests/test_devtools_benchmarks.py": 328116,
|
||||
"web/modules/chat.js": 224315,
|
||||
|
|
|
|||
|
|
@ -4,9 +4,31 @@ from __future__ import annotations
|
|||
|
||||
import logging
|
||||
from typing import Any
|
||||
from urllib.parse import urlsplit
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
_LOOPBACK_HOSTS = frozenset({"localhost", "127.0.0.1", "::1"})
|
||||
|
||||
|
||||
def is_loopback_base_url(base_url: Any) -> bool:
|
||||
"""True when the configured route targets this very host.
|
||||
|
||||
A loopback OPENAI_COMPATIBLE_BASE_URL (the documented Ollama / LM Studio /
|
||||
vLLM setups) is a LOCAL server even though its provider name is not
|
||||
"local": its connect failure means that server is down, not that the
|
||||
network egress is — so such routes must never classify as a remote
|
||||
transport outage worth waiting out.
|
||||
"""
|
||||
text = str(base_url or "").strip()
|
||||
if not text:
|
||||
return False
|
||||
try:
|
||||
host = urlsplit(text).hostname or ""
|
||||
except ValueError:
|
||||
return False
|
||||
return host.lower() in _LOOPBACK_HOSTS
|
||||
|
||||
|
||||
def is_pre_dispatch_transport_failure(exc: BaseException) -> bool:
|
||||
"""Return true only for exceptions raised before request bytes can be sent."""
|
||||
|
|
@ -35,11 +57,24 @@ def is_pre_dispatch_transport_failure(exc: BaseException) -> bool:
|
|||
return True
|
||||
if not isinstance(exc, requests.exceptions.ConnectionError):
|
||||
return False
|
||||
# requests.exceptions.ProxyError subclasses ConnectionError; both the
|
||||
# direct and the proxied connect failure arrive as MaxRetryError args.
|
||||
for value in getattr(exc, "args", ()):
|
||||
if isinstance(value, urllib3.exceptions.MaxRetryError):
|
||||
reason = getattr(value, "reason", None)
|
||||
if isinstance(reason, urllib3.exceptions.ConnectTimeoutError):
|
||||
return True
|
||||
if isinstance(reason, urllib3.exceptions.ProxyError):
|
||||
# An unreachable proxy is a pre-dispatch fact only with
|
||||
# nested connect-time evidence (NewConnectionError is a
|
||||
# ConnectTimeoutError subclass); a proxy HTTP response or
|
||||
# a post-dispatch read failure never matches.
|
||||
nested = getattr(reason, "original_error", None)
|
||||
if isinstance(nested, (
|
||||
urllib3.exceptions.ConnectTimeoutError,
|
||||
urllib3.exceptions.NewConnectionError,
|
||||
)):
|
||||
return True
|
||||
except Exception: # pragma: no cover - optional transport dependency
|
||||
pass
|
||||
return False
|
||||
|
|
|
|||
|
|
@ -235,6 +235,9 @@ class AttemptRequest:
|
|||
candidate_context_size_bytes: Optional[int] = None
|
||||
candidate_measurement_kind: Literal["canonical_json_v1", "opaque"] = "opaque"
|
||||
physical_context: Optional[PhysicalAttemptContext] = None
|
||||
# Route-locality fact (additive): the base_url host is localhost/127.0.0.1/::1
|
||||
# (loopback OpenAI-compatible installs — Ollama / LM Studio / vLLM).
|
||||
route_is_loopback: bool = False
|
||||
@dataclass(frozen=True)
|
||||
class AttemptReservation:
|
||||
attempt_id: str
|
||||
|
|
@ -263,6 +266,7 @@ class PhysicalAttemptCapture:
|
|||
provider_code: str = ""
|
||||
provider_error_type: str = ""
|
||||
provider_error: str = ""
|
||||
route_is_loopback: bool = False # see AttemptRequest.route_is_loopback
|
||||
|
||||
|
||||
@contextlib.contextmanager
|
||||
|
|
@ -481,10 +485,9 @@ def usage_breakdown(
|
|||
"by_category": by_category,
|
||||
"by_task": by_task,
|
||||
"by_root": by_root,
|
||||
# Execution-axis filter (v6.91): the delegated (subscription-harness) rows
|
||||
# only — a VIEW over the same rows for "where did the money go" readers,
|
||||
# never a third monetary sum or authority. Disclosed-free sessions settle
|
||||
# at $0 here; undisclosed spend stays in `unknown`.
|
||||
# Execution-axis filter (v6.91): delegated (subscription-harness) rows only —
|
||||
# a VIEW over the same rows for "where did the money go" readers, never a third
|
||||
# monetary sum or authority. Disclosed-free settles $0; undisclosed stays `unknown`.
|
||||
"delegated": _with_integrity(
|
||||
_breakdown_bucket([row for row in rows if str(row.get("kind") or "") == "subscription_session"]),
|
||||
integrity_degraded,
|
||||
|
|
@ -539,10 +542,9 @@ def _reservation_cost(request: AttemptRequest) -> Optional[float]:
|
|||
|
||||
prompt_cache_ttl = str(request.prompt_cache_ttl or "").strip().lower()
|
||||
if prompt_cache_ttl not in PROMPT_CACHE_TTL_SCALE:
|
||||
# An inspectable marker-free candidate writes no cache at all. Keep
|
||||
# the historical conservative base-tier reservation without
|
||||
# misreporting a TTL on the eventual physical settlement. Opaque
|
||||
# construction sites still fall back to the owner setting.
|
||||
# An inspectable marker-free candidate writes no cache at all. Keep the
|
||||
# historical conservative base-tier reservation without misreporting a TTL on
|
||||
# the physical settlement; opaque sites still fall back to the owner setting.
|
||||
prompt_cache_ttl = (
|
||||
"default"
|
||||
if request.candidate_measurement_kind == "canonical_json_v1"
|
||||
|
|
@ -685,8 +687,7 @@ def reserve_attempt(request: AttemptRequest) -> AttemptReservation:
|
|||
root_task_id=scope.root_task_id,
|
||||
)
|
||||
ensure_legacy_imported(root)
|
||||
# IMPORTANT: live catalog I/O belongs before ``with _locked(root)`` below.
|
||||
# The lock protects only the atomic budget read/check/append transaction.
|
||||
# IMPORTANT: live catalog I/O belongs before ``with _locked(root)`` below — the lock protects only the atomic budget read/check/append transaction.
|
||||
bound = _reservation_cost(request)
|
||||
pricing_known = bound is not None
|
||||
attempt_id = uuid.uuid4().hex
|
||||
|
|
@ -824,9 +825,8 @@ def _append_single_settled_row(
|
|||
existing = _final_rows(records).get(attempt_id)
|
||||
if existing is not None:
|
||||
def identity_value(source: Dict[str, Any], key: str) -> Any:
|
||||
# Rows written before physical_attempt_v1 omitted these optional
|
||||
# keys. Missing and explicit empty are the same legacy identity;
|
||||
# a non-empty wave/slot still conflicts with either one.
|
||||
# Rows written before physical_attempt_v1 omitted these optional keys:
|
||||
# missing == explicit empty; a non-empty wave/slot conflicts with either.
|
||||
return str(source.get(key) or "") if key in REVIEW_ATTRIBUTION_KEYS else source.get(key)
|
||||
|
||||
if any(identity_value(existing, key) != identity_value(row, key) for key in comparable):
|
||||
|
|
@ -1199,6 +1199,7 @@ def _record_attempt_capture(
|
|||
provider_code=code,
|
||||
provider_error_type=error_type,
|
||||
provider_error=error,
|
||||
route_is_loopback=bool(request.route_is_loopback),
|
||||
)
|
||||
_LAST_PHYSICAL_ATTEMPT.set(capture)
|
||||
if exc is not None:
|
||||
|
|
@ -1384,8 +1385,7 @@ def _legacy_snapshot(root: pathlib.Path) -> Tuple[list[Dict[str, Any]], Dict[str
|
|||
except OSError as exc:
|
||||
raise UsageAccountingError(f"cannot snapshot legacy usage source {path}: {exc}") from exc
|
||||
hashes = {name: hashlib.sha256(snapshots[name]).hexdigest() if name in snapshots else "" for name in sources}
|
||||
# Settings are owner-secret state: prove non-mutation by hash, but never copy
|
||||
# their contents into the usage archive.
|
||||
# Settings are owner-secret state: prove non-mutation by hash, never copy contents.
|
||||
try:
|
||||
hashes["settings.json"] = hashlib.sha256(settings_path.read_bytes()).hexdigest()
|
||||
except FileNotFoundError:
|
||||
|
|
|
|||
|
|
@ -670,3 +670,288 @@ def test_wait_step_caps_sleep_at_note_interval_for_low_idle_timeouts(tmp_path, m
|
|||
assert redial is True
|
||||
assert sleeps == [30.0] # min(60s backoff cap, 60/2 note interval)
|
||||
assert len(notes) == 1 # the periodic note fired for the lowered interval
|
||||
|
||||
|
||||
# ------------------------------------------------- route locality (A9, loopback)
|
||||
|
||||
def test_released_loopback_route_failure_stays_generic():
|
||||
"""A loopback OPENAI_COMPATIBLE_BASE_URL install (Ollama / LM Studio /
|
||||
vLLM) stamps a remote-shaped provider name, but a stopped LOCAL server is
|
||||
not a network outage worth waiting out."""
|
||||
exc = httpx.ConnectError("connection refused")
|
||||
exc.physical_attempt_capture = ua.PhysicalAttemptCapture(
|
||||
attempt_id="pa-lb", model="m", provider="openai-compatible",
|
||||
state="released", candidate_measurement_kind="opaque",
|
||||
route_is_loopback=True,
|
||||
)
|
||||
assert classify_llm_exception(exc).kind != "transport_unavailable"
|
||||
|
||||
|
||||
def test_released_remote_compatible_route_still_classifies_transport_unavailable():
|
||||
"""The additive default keeps every remote route on the wait path."""
|
||||
exc = httpx.ConnectError("connection refused")
|
||||
exc.physical_attempt_capture = ua.PhysicalAttemptCapture(
|
||||
attempt_id="pa-rc", model="m", provider="openai-compatible",
|
||||
state="released", candidate_measurement_kind="opaque",
|
||||
)
|
||||
assert classify_llm_exception(exc).kind == "transport_unavailable"
|
||||
|
||||
|
||||
def test_attempt_request_carries_route_locality_from_target():
|
||||
from ouroboros.llm import _attempt_request
|
||||
|
||||
loopback = _attempt_request(
|
||||
{"provider": "openai-compatible", "usage_model": "m",
|
||||
"base_url": "http://localhost:11434/v1"},
|
||||
{"model": "m", "messages": []},
|
||||
)
|
||||
assert loopback.route_is_loopback is True
|
||||
remote = _attempt_request(
|
||||
{"provider": "openrouter", "usage_model": "m",
|
||||
"base_url": "https://openrouter.ai/api/v1"},
|
||||
{"model": "m", "messages": []},
|
||||
)
|
||||
assert remote.route_is_loopback is False
|
||||
|
||||
|
||||
def test_attempt_capture_propagates_route_locality(tmp_path):
|
||||
reservation = ua.AttemptReservation(
|
||||
attempt_id="pa-cap", drive_root=tmp_path, model="m",
|
||||
provider="openai-compatible", reservation_upper_bound_usd=None,
|
||||
)
|
||||
request = ua.AttemptRequest(
|
||||
model="m", provider="openai-compatible", route_is_loopback=True,
|
||||
)
|
||||
capture = ua._record_attempt_capture(reservation, request, "released")
|
||||
assert capture.route_is_loopback is True
|
||||
|
||||
|
||||
# ------------------------------------------ finalize_now during an episode (A7)
|
||||
|
||||
def test_finalize_now_during_episode_takes_no_resend_terminal_via_mailbox(tmp_path, monkeypatch):
|
||||
"""finalize_now (deadline / ceiling / owner stop flavors share this exit)
|
||||
arriving MID-SLEEP through the REAL owner mailbox: the interruptible sleep
|
||||
wakes within a slice and the terminal is the transport no-resend — zero
|
||||
further provider dials, never a forced-final paid call over the dead
|
||||
egress."""
|
||||
from ouroboros.owner_mailbox import KIND_FINALIZE_NOW, write_owner_message
|
||||
|
||||
fake_call, calls = _transport_failing_call(fail_times=99)
|
||||
monkeypatch.setattr(loop_mod, "call_llm_with_retry", fake_call)
|
||||
monkeypatch.setenv("OUROBOROS_TASK_REVIEW_MODE", "off")
|
||||
monkeypatch.delenv("USE_LOCAL_FALLBACK", raising=False)
|
||||
registry = ToolRegistry(repo_dir=tmp_path, drive_root=tmp_path)
|
||||
notes = []
|
||||
timer = threading.Timer(0.3, lambda: write_owner_message(
|
||||
tmp_path, "budget ceiling reached", "t-wait", kind=KIND_FINALIZE_NOW,
|
||||
))
|
||||
timer.start()
|
||||
start = time.monotonic()
|
||||
try:
|
||||
result, usage, trace = run_llm_loop(**_loop_kwargs(tmp_path, registry, notes))
|
||||
finally:
|
||||
timer.cancel()
|
||||
elapsed = time.monotonic() - start
|
||||
|
||||
assert elapsed < 10.0 # woke within a sleep slice, not the full backoff ladder
|
||||
assert calls["n"] == 1 # the woken redial exited at the round top: zero further dials
|
||||
assert usage.get("execution_status") == "infra_failed"
|
||||
assert usage.get("reason_code") == "provider_unavailable"
|
||||
assert trace.get("forced_finalization", {}).get("source") == "transport_unavailable_no_resend"
|
||||
assert "ended as a provider outage" in result
|
||||
events = _read_network_wait_events(tmp_path)
|
||||
assert events[-1]["phase"] == "ended"
|
||||
assert events[-1].get("detail") == "finalize_now"
|
||||
|
||||
|
||||
# --------------------------------------------------- episode exit contracts (A12)
|
||||
|
||||
def test_redial_failing_with_different_kind_ends_episode_and_resumes_fallback(tmp_path, monkeypatch):
|
||||
"""A redial that gets past the connect phase but fails differently proves
|
||||
the transport passable: the episode ends and the ORDINARY fallback chain
|
||||
resumes for the fresh kind."""
|
||||
calls = {"n": 0}
|
||||
|
||||
def fake_call(_llm, _messages, _model, _tools, _effort, _max_retries, _drive_logs,
|
||||
_task_id, _round_idx, _event_queue, accumulated_usage, *_a, **_k):
|
||||
calls["n"] += 1
|
||||
accumulated_usage["_last_llm_error_kind"] = (
|
||||
"transport_unavailable" if calls["n"] == 1 else "provider_transient"
|
||||
)
|
||||
return None, 0.0
|
||||
|
||||
chain_calls = {"n": 0}
|
||||
|
||||
def fake_chain(**kwargs):
|
||||
chain_calls["n"] += 1
|
||||
kwargs["accumulated_usage"].pop("_last_llm_error_kind", None)
|
||||
return (
|
||||
{"role": "assistant", "content": "fallback-ok"}, "other/model", False,
|
||||
kwargs["context_fit_plan"], kwargs["active_context_mode"],
|
||||
)
|
||||
|
||||
monkeypatch.setattr(loop_transport, "interruptible_wait_sleep", lambda _sec, _wake: False)
|
||||
monkeypatch.setattr(loop_mod, "call_llm_with_retry", fake_call)
|
||||
monkeypatch.setattr(loop_mod, "_run_cross_model_fallback_chain", fake_chain)
|
||||
monkeypatch.setenv("OUROBOROS_TASK_REVIEW_MODE", "off")
|
||||
monkeypatch.delenv("USE_LOCAL_FALLBACK", raising=False)
|
||||
notes = []
|
||||
result, usage, _trace = run_llm_loop(**_loop_kwargs(tmp_path, ToolRegistry(repo_dir=tmp_path, drive_root=tmp_path), notes))
|
||||
|
||||
assert result == "fallback-ok"
|
||||
assert calls["n"] == 2 # primary, then the one free redial
|
||||
assert chain_calls["n"] == 1 # ordinary policy resumed after the episode ended
|
||||
assert usage.get("reason_code") is None
|
||||
events = _read_network_wait_events(tmp_path)
|
||||
assert events[-1]["phase"] == "ended"
|
||||
assert events[-1].get("detail") == "error_kind_changed:provider_transient"
|
||||
|
||||
|
||||
def test_redial_with_unknown_outcome_takes_unknown_no_resend_terminal(tmp_path, monkeypatch):
|
||||
"""A main-dispatch redial whose socket outcome is unknown ends the episode
|
||||
AND the round: provider_outcome_unknown keeps its own no-resend terminal —
|
||||
zero further dials of any kind."""
|
||||
calls = {"n": 0}
|
||||
|
||||
def fake_call(_llm, _messages, _model, _tools, _effort, _max_retries, _drive_logs,
|
||||
_task_id, _round_idx, _event_queue, accumulated_usage, *_a, **_k):
|
||||
calls["n"] += 1
|
||||
accumulated_usage["_last_llm_error_kind"] = (
|
||||
"transport_unavailable" if calls["n"] == 1 else "provider_outcome_unknown"
|
||||
)
|
||||
# Mirror _record_llm_call_error's stamps (the real failure path).
|
||||
accumulated_usage.update(execution_status="infra_failed", reason_code="llm_api_error")
|
||||
return None, 0.0
|
||||
|
||||
def _chain_must_not_run(**_kwargs):
|
||||
raise AssertionError("no fallback chain after an unknown-outcome redial")
|
||||
|
||||
monkeypatch.setattr(loop_transport, "interruptible_wait_sleep", lambda _sec, _wake: False)
|
||||
monkeypatch.setattr(loop_mod, "call_llm_with_retry", fake_call)
|
||||
monkeypatch.setattr(loop_mod, "_run_cross_model_fallback_chain", _chain_must_not_run)
|
||||
monkeypatch.setenv("OUROBOROS_TASK_REVIEW_MODE", "off")
|
||||
monkeypatch.setenv("OUROBOROS_MODEL_FALLBACKS", "other/model")
|
||||
monkeypatch.delenv("USE_LOCAL_FALLBACK", raising=False)
|
||||
notes = []
|
||||
_result, usage, trace = run_llm_loop(**_loop_kwargs(tmp_path, ToolRegistry(repo_dir=tmp_path, drive_root=tmp_path), notes))
|
||||
|
||||
assert calls["n"] == 2 # zero dials after the unknown outcome
|
||||
assert usage.get("execution_status") == "infra_failed"
|
||||
assert trace.get("forced_finalization", {}).get("source") == "provider_outcome_unknown_no_resend"
|
||||
events = _read_network_wait_events(tmp_path)
|
||||
assert events[-1]["phase"] == "ended"
|
||||
assert events[-1].get("detail") == "error_kind_changed:provider_outcome_unknown"
|
||||
|
||||
|
||||
def test_episode_ledger_evidence_has_zero_dispatched_paid_attempts(tmp_path, monkeypatch):
|
||||
"""Between episode entry and the last wait, every durable failure row is
|
||||
the typed $0 transport kind and no usage row records a dispatched paid
|
||||
attempt (released rows are the legitimate evidence)."""
|
||||
sleeps = []
|
||||
monkeypatch.setattr(loop_transport, "interruptible_wait_sleep",
|
||||
lambda sec, _wake: (sleeps.append(sec), False)[1])
|
||||
monkeypatch.setenv("OUROBOROS_TASK_REVIEW_MODE", "off")
|
||||
monkeypatch.delenv("USE_LOCAL_FALLBACK", raising=False)
|
||||
llm = _FlakyChatLLM(fail_times=3)
|
||||
notes = []
|
||||
result, _usage, _trace = run_llm_loop(
|
||||
**_loop_kwargs(tmp_path, ToolRegistry(repo_dir=tmp_path, drive_root=tmp_path), notes, llm=llm))
|
||||
|
||||
assert result == "done"
|
||||
rows = [json.loads(line) for line in (tmp_path / "events.jsonl").read_text().splitlines() if line.strip()]
|
||||
entered_idx = next(i for i, r in enumerate(rows) if r.get("type") == "network_wait" and r.get("phase") == "entered")
|
||||
last_wait_idx = max(i for i, r in enumerate(rows) if r.get("type") == "network_wait" and r.get("phase") == "waiting")
|
||||
window = rows[entered_idx:last_wait_idx + 1]
|
||||
api_errors = [r for r in window if r.get("type") == "llm_api_error"]
|
||||
assert api_errors, "the episode's failures must leave durable evidence"
|
||||
assert all(r.get("error_kind") == "transport_unavailable" for r in api_errors)
|
||||
assert not any(r.get("type") == "llm_usage" for r in window)
|
||||
|
||||
|
||||
def test_nonstandard_transient_retry_max_keeps_one_attempt_transport_contract(tmp_path, monkeypatch):
|
||||
"""OUROBOROS_TRANSIENT_RETRY_MAX tunes the transient burst, never the
|
||||
one-physical-attempt transport contract."""
|
||||
monkeypatch.setenv("OUROBOROS_TRANSIENT_RETRY_MAX", "12")
|
||||
llm = _RaisingLLM(_typed_transport_exc)
|
||||
usage = {}
|
||||
msg, _cost = call_llm_with_retry(
|
||||
llm, [{"role": "user", "content": "hi"}], "test-model", None, "low", 3,
|
||||
tmp_path, "t-max", 1, None, usage,
|
||||
)
|
||||
assert msg is None
|
||||
assert llm.calls == 1
|
||||
assert usage.get("_last_llm_error_kind") == "transport_unavailable"
|
||||
|
||||
|
||||
def test_review_actor_physical_send_rail_stays_bounded_at_two():
|
||||
"""Review actors (P3 / scope / acceptance) dispatch under
|
||||
physical_attempt_limit(2) — the review_substrate contract — so a transport
|
||||
outage can never turn a review slot into an unbounded redial loop."""
|
||||
from ouroboros.usage_accounting import (
|
||||
PhysicalAttemptLimitExceeded,
|
||||
_claim_physical_dispatch,
|
||||
physical_attempt_limit,
|
||||
)
|
||||
|
||||
with physical_attempt_limit(2):
|
||||
_claim_physical_dispatch()
|
||||
_claim_physical_dispatch()
|
||||
with pytest.raises(PhysicalAttemptLimitExceeded):
|
||||
_claim_physical_dispatch()
|
||||
|
||||
|
||||
# ------------------------------------------------ terminal wordings + Q14 margin
|
||||
|
||||
def test_terminal_wordings_cover_three_wait_outcomes():
|
||||
"""Three honest terminal texts: fail-fast is keyed on NOT wait_eligible; a
|
||||
managed zero-wait task names the spent window; the waited-out wording says
|
||||
outage without the supervisor's lifecycle term INTERRUPTED."""
|
||||
kwargs = dict(is_context_overflow=False, is_transport_wait=True, is_deadline_exhausted=False)
|
||||
fast = loop_transport.provider_terminal_fallback_text(
|
||||
{}, waited=False, wait_eligible=False, **kwargs)
|
||||
assert "fails fast" in fast
|
||||
waited = loop_transport.provider_terminal_fallback_text(
|
||||
{}, waited=True, wait_eligible=True, **kwargs)
|
||||
assert "ended as a provider outage, not completed" in waited
|
||||
assert "INTERRUPTED" not in waited
|
||||
zero_wait = loop_transport.provider_terminal_fallback_text(
|
||||
{}, waited=False, wait_eligible=True, **kwargs)
|
||||
assert "left no time to wait" in zero_wait
|
||||
assert "fails fast" not in zero_wait
|
||||
|
||||
|
||||
def test_final_redial_reserves_named_margin_before_admission_close(tmp_path, monkeypatch):
|
||||
"""Q14/A11: the last free redial sleeps to remaining minus the named 3s
|
||||
margin (round-top overhead routinely eats ~1s), then the next step
|
||||
terminalizes deterministically."""
|
||||
from ouroboros.config import get_finalization_grace_sec
|
||||
from ouroboros.deadline_utils import dispatch_window_remaining_sec
|
||||
|
||||
assert loop_transport._FINAL_REDIAL_MARGIN_SEC == 3.0
|
||||
sleeps = []
|
||||
monkeypatch.setattr(loop_transport, "interruptible_wait_sleep",
|
||||
lambda sec, _wake: (sleeps.append(sec), False)[1])
|
||||
deadline = datetime.now(timezone.utc) + timedelta(seconds=get_finalization_grace_sec() + 40)
|
||||
tools = SimpleNamespace(_ctx=SimpleNamespace(
|
||||
task_metadata={"deadline_at": deadline.isoformat()}, task_attempt=None,
|
||||
))
|
||||
episode = loop_transport.TransportWaitEpisode(
|
||||
started_monotonic=time.monotonic(), wait_iterations=10, # backoff at the 60s cap
|
||||
)
|
||||
remaining = dispatch_window_remaining_sec(
|
||||
deadline_ts=deadline.timestamp(), reserve_sec=get_finalization_grace_sec(),
|
||||
)
|
||||
redial = loop_transport.transport_wait_step(
|
||||
episode, tools=tools, error_kind="transport_unavailable",
|
||||
drive_root=None, drive_logs=tmp_path, task_id="t-margin", model="m",
|
||||
emit_progress=lambda _n: None, incoming_messages=None, owner_msg_seen=set(),
|
||||
)
|
||||
assert redial is True
|
||||
assert episode.final_redial_done is True
|
||||
assert sleeps[0] == pytest.approx(
|
||||
remaining - loop_transport._FINAL_REDIAL_MARGIN_SEC, abs=1.0)
|
||||
assert loop_transport.transport_wait_step(
|
||||
episode, tools=tools, error_kind="transport_unavailable",
|
||||
drive_root=None, drive_logs=tmp_path, task_id="t-margin", model="m",
|
||||
emit_progress=lambda _n: None, incoming_messages=None, owner_msg_seen=set(),
|
||||
) is False # deadline_after_final_redial
|
||||
|
|
|
|||
|
|
@ -80,3 +80,66 @@ def test_requests_read_timeout_connection_error_does_not_prove_pre_dispatch(data
|
|||
urllib3.exceptions.MaxRetryError(None, "/messages", reason=reason)
|
||||
)
|
||||
assert not is_pre_dispatch_transport_failure(wrapped)
|
||||
|
||||
|
||||
def test_requests_proxy_error_with_nested_connect_evidence_proves_pre_dispatch(data_root):
|
||||
"""The standard unreachable-proxy chain (native Anthropic behind a dead
|
||||
proxy): requests.exceptions.ProxyError -> MaxRetryError -> urllib3
|
||||
ProxyError -> NewConnectionError is typed pre-dispatch evidence."""
|
||||
import requests
|
||||
import urllib3
|
||||
from ouroboros.transport_custody import is_pre_dispatch_transport_failure
|
||||
|
||||
nested = urllib3.exceptions.NewConnectionError(
|
||||
None, "Failed to establish a new connection: [Errno 111] Connection refused"
|
||||
)
|
||||
proxy = urllib3.exceptions.ProxyError("Cannot connect to proxy.", nested)
|
||||
wrapped = requests.exceptions.ProxyError(
|
||||
urllib3.exceptions.MaxRetryError(None, "/messages", reason=proxy)
|
||||
)
|
||||
assert is_pre_dispatch_transport_failure(wrapped)
|
||||
|
||||
|
||||
def test_requests_proxy_error_with_connect_timeout_evidence_proves_pre_dispatch(data_root):
|
||||
import requests
|
||||
import urllib3
|
||||
from ouroboros.transport_custody import is_pre_dispatch_transport_failure
|
||||
|
||||
nested = urllib3.exceptions.ConnectTimeoutError("timed out connecting to proxy")
|
||||
proxy = urllib3.exceptions.ProxyError("Cannot connect to proxy.", nested)
|
||||
wrapped = requests.exceptions.ProxyError(
|
||||
urllib3.exceptions.MaxRetryError(None, "/messages", reason=proxy)
|
||||
)
|
||||
assert is_pre_dispatch_transport_failure(wrapped)
|
||||
|
||||
|
||||
def test_requests_proxy_error_without_connect_evidence_stays_untyped(data_root):
|
||||
"""A proxy failure that is NOT connect-time (a proxy HTTP response, a
|
||||
post-dispatch read failure) must never release custody."""
|
||||
import requests
|
||||
import urllib3
|
||||
from ouroboros.transport_custody import is_pre_dispatch_transport_failure
|
||||
|
||||
proxy = urllib3.exceptions.ProxyError(
|
||||
"Your proxy appears to only use HTTP and not HTTPS",
|
||||
urllib3.exceptions.HTTPError("bad proxy response"),
|
||||
)
|
||||
wrapped = requests.exceptions.ProxyError(
|
||||
urllib3.exceptions.MaxRetryError(None, "/messages", reason=proxy)
|
||||
)
|
||||
assert not is_pre_dispatch_transport_failure(wrapped)
|
||||
|
||||
|
||||
@pytest.mark.parametrize("url,expected", [
|
||||
("http://localhost:11434/v1", True),
|
||||
("http://127.0.0.1:1234/v1", True),
|
||||
("http://[::1]:8000/v1", True),
|
||||
("https://openrouter.ai/api/v1", False),
|
||||
("https://api.anthropic.com/v1", False),
|
||||
("", False),
|
||||
("not a url", False),
|
||||
])
|
||||
def test_is_loopback_base_url(url, expected):
|
||||
from ouroboros.transport_custody import is_loopback_base_url
|
||||
|
||||
assert is_loopback_base_url(url) is expected
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue