diff --git a/agent/turn_api_call.py b/agent/turn_api_call.py index 6c7aa3744d..56557b17df 100644 --- a/agent/turn_api_call.py +++ b/agent/turn_api_call.py @@ -2,12 +2,13 @@ the attempt while another session's Nous Portal rate limit is active), ``perform_api_call`` (streaming decision, MoA prepared-request handshake, LLM execution middleware wrapper, the redirect ``_model_request_active`` bracket and the response-vs-redirect crossing check) and -``handle_api_interrupt`` (``InterruptedError`` mid-call). Extracted from -``run_conversation``; nothing here imports ``agent.conversation_loop`` at module level. +``handle_api_interrupt`` (``InterruptedError`` mid-call). Nothing here imports +``agent.conversation_loop`` at module level (cycle). """ from __future__ import annotations +from contextlib import nullcontext from dataclasses import dataclass import logging import time @@ -18,6 +19,16 @@ from agent.message_metadata import append_message logger = logging.getLogger("agent.conversation_loop") +def stop_thinking_spinner(agent: Any, thinking_spinner: Any) -> None: + """Stop the spinner silently and clear the thinking callback; returns ``None`` so + callers can rebind ``thinking_spinner = stop_thinking_spinner(agent, thinking_spinner)``.""" + if thinking_spinner: + thinking_spinner.stop("") + if agent.thinking_callback: + agent.thinking_callback("") + return None + + @dataclass class ApiCallVerdict: """``action``: ``"fallthrough"`` (``response`` is ready for verification) or ``"break"`` @@ -29,58 +40,44 @@ class ApiCallVerdict: interrupted: Any +def _should_stream(agent: Any) -> bool: + """Streaming is preferred even without consumers (stale-stream / read-timeout health + checks); disabled on provider signal, ACP schemes, MoA without a display consumer, or + Mock clients in tests (SimpleNamespace, not stream iterators).""" + if getattr(agent, "_disable_streaming", False): + return False + _base = str(agent.base_url or "").lower() + if agent.provider in {"copilot-acp"} or _base.startswith(("acp://", "acp+tcp://")): + return False + if not agent._has_stream_consumers(): + if agent.provider == "moa": + return False + from unittest.mock import Mock + if isinstance(getattr(agent, "client", None), Mock): + return False + return True + + def perform_api_call( agent: Any, *, api_kwargs: Any, _original_api_kwargs: Any, _llm_middleware_trace: Any, _moa_prepared_request: Any, _retry: Any, thinking_spinner: Any, retry_count: Any, api_call_count: Any, api_request_id: Any, effective_task_id: Any, turn_id: Any, interrupted: Any, ) -> ApiCallVerdict: - """Issue the request. Streaming is preferred even without consumers (stale-stream / - read-timeout health checks) and disabled per provider signal, ACP schemes, MoA without a - display consumer, or Mock clients in tests.""" + """Issue the request (see ``_should_stream`` for the streaming decision).""" response = None - def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> ApiCallVerdict: + def _verdict(action: str) -> ApiCallVerdict: return ApiCallVerdict( action=action, response=response, thinking_spinner=thinking_spinner, interrupted=interrupted, ) - # Always prefer streaming even without consumers: it gives stale- - # stream/read-timeout health checks that quiet callers otherwise lack. - # Falls back if unsupported. def _stop_spinner(): nonlocal thinking_spinner - if thinking_spinner: - thinking_spinner.stop("") - thinking_spinner = None - if agent.thinking_callback: - agent.thinking_callback("") + thinking_spinner = stop_thinking_spinner(agent, thinking_spinner) - _use_streaming = True - # Provider signaled "stream not supported": stay non-streaming for the - # session. - if getattr(agent, "_disable_streaming", False): - _use_streaming = False - # ACP clients (`acp://` scheme, any vendor) return a plain - # SimpleNamespace, not a stream; mirrors the Responses API exclusion. - elif ( - agent.provider in {"copilot-acp"} - or str(agent.base_url or "").lower().startswith("acp://") - or str(agent.base_url or "").lower().startswith("acp+tcp://") - ): - _use_streaming = False - # MoA streams only with a display/TTS consumer - # (MoAChatCompletions.create() honors stream=True); else complete- - # response path. - elif agent.provider == "moa" and not agent._has_stream_consumers(): - _use_streaming = False - elif not agent._has_stream_consumers(): - # No consumer: still stream for health checking, except Mock clients - # in tests (SimpleNamespace, not stream iterators). - from unittest.mock import Mock - if isinstance(getattr(agent, "client", None), Mock): - _use_streaming = False + _use_streaming = _should_stream(agent) def _perform_api_call(next_api_kwargs): if agent.api_mode == "codex_responses": @@ -117,15 +114,14 @@ def perform_api_call( from hermes_cli.middleware import run_llm_execution_middleware + # The ``_model_request_active`` bracket is taken under the redirect lock when one exists, + # so redirect() can't observe a half-toggled flag. _model_request_active = getattr(agent, "_model_request_active", None) _redirect_lock = getattr(agent, "_pending_redirect_lock", None) - if _redirect_lock is not None: - with _redirect_lock: - if _model_request_active is not None: - _model_request_active.set() - elif _model_request_active is not None: - _model_request_active.set() - _redirect_crossed_response = False + _bracket = nullcontext() if _redirect_lock is None else _redirect_lock + with _bracket: + if _model_request_active is not None: + _model_request_active.set() try: response = run_llm_execution_middleware( api_kwargs, _perform_api_call, original_request=_original_api_kwargs, @@ -135,25 +131,17 @@ def perform_api_call( api_call_count=api_call_count, middleware_trace=list(_llm_middleware_trace), ) finally: - if _redirect_lock is not None: - with _redirect_lock: - if _model_request_active is not None: - _model_request_active.clear() - _redirect_crossed_response = bool( - agent._pending_redirect - ) - else: + with _bracket: if _model_request_active is not None: _model_request_active.clear() - _redirect_crossed_response = agent._has_pending_redirect() + _redirect_crossed_response = ( + bool(agent._pending_redirect) if _redirect_lock is not None + else agent._has_pending_redirect() + ) if _redirect_crossed_response: # Response and redirect can cross threads: discard the now-stale # response and rebuild from the correction rather than lose it. - if thinking_spinner: - thinking_spinner.stop("") - thinking_spinner = None - if agent.thinking_callback: - agent.thinking_callback("") + thinking_spinner = stop_thinking_spinner(agent, thinking_spinner) if agent.clear_interrupt(preserve_redirect=True): _retry.restart_with_redirected_messages = True else: @@ -180,36 +168,18 @@ def handle_api_interrupt( """``InterruptedError`` during the provider call: a pending redirect keeps its correction queued for the outer-loop rebuild; otherwise keep any streamed partial text so the next turn has a record of the half-finished reply.""" - from agent.conversation_loop import ( - INTERRUPT_WAITING_FOR_MODEL_PREFIX, - ) + from agent.conversation_loop import INTERRUPT_WAITING_FOR_MODEL_PREFIX - def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> ApiInterruptVerdict: - return ApiInterruptVerdict( - action=action, - thinking_spinner=thinking_spinner, - interrupted=interrupted, - final_response=final_response, - - ) - - if thinking_spinner: - thinking_spinner.stop("") - thinking_spinner = None - if agent.thinking_callback: - agent.thinking_callback("") - if agent._has_pending_redirect(): - # redirect() cancelled only this request: keep the correction - # queued, clear the cancellation bit, let the outer loop rebuild. - # Never materialize incomplete signed/encrypted reasoning items. - if agent.clear_interrupt(preserve_redirect=True): - _retry.restart_with_redirected_messages = True - return _verdict("break") + thinking_spinner = stop_thinking_spinner(agent, thinking_spinner) + # redirect() cancelled only this request: keep the correction queued, clear the + # cancellation bit, let the outer loop rebuild. Never materialize incomplete + # signed/encrypted reasoning items. + if agent._has_pending_redirect() and agent.clear_interrupt(preserve_redirect=True): + _retry.restart_with_redirected_messages = True + return ApiInterruptVerdict("break", thinking_spinner, interrupted, final_response) api_elapsed = time.time() - api_start_time agent._vprint(f"{agent.log_prefix}⚡ Interrupted during API call.", force=True) interrupted = True - # Keep assistant text already streamed before the stop, else the next - # turn has no record of the half-finished reply. _partial = agent._strip_think_blocks( getattr(agent, "_current_streamed_assistant_text", "") or "" ).strip() @@ -219,8 +189,7 @@ def handle_api_interrupt( else: final_response = f"{INTERRUPT_WAITING_FOR_MODEL_PREFIX}{api_elapsed:.1f}s elapsed)." agent._persist_session(messages, conversation_history) - return _verdict("break") - return _verdict("fallthrough") + return ApiInterruptVerdict("break", thinking_spinner, interrupted, final_response) @dataclass @@ -241,9 +210,7 @@ def nous_rate_limit_guard( ) -> NousRateGuardVerdict: """Skip the call if another session recorded a Nous Portal rate limit: every attempt (incl. SDK retries) counts against RPH. Never lets the guard itself break the agent loop.""" - from agent.conversation_loop import ( - _arm_fallback_restart, - ) + from agent.conversation_loop import _arm_fallback_restart def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> NousRateGuardVerdict: return NousRateGuardVerdict( @@ -251,9 +218,6 @@ def nous_rate_limit_guard( compression_attempts=compression_attempts, result=result, ) - # ── Nous Portal rate limit guard ────────────────────── - # Skip the call if another session recorded a rate limit: every attempt - # (incl. SDK retries) counts against RPH. if agent.provider == "nous": try: from agent.nous_rate_guard import ( @@ -265,9 +229,7 @@ def nous_rate_limit_guard( f"Nous Portal rate limit active — " f"resets in {_fmt_nous_remaining(_nous_remaining)}." ) - agent._buffer_vprint( - f"⏳ {_nous_msg} Trying fallback..." - ) + agent._buffer_vprint(f"⏳ {_nous_msg} Trying fallback...") agent._buffer_status(f"⏳ {_nous_msg}") if agent._try_activate_fallback(): active_system_prompt = _arm_fallback_restart( @@ -275,8 +237,7 @@ def nous_rate_limit_guard( retry_count = 0 compression_attempts = 0 return _verdict("break") - # No fallback available — surface buffered context - # so user sees the rate-limit message that led here. + # No fallback — surface the buffered rate-limit context that led here. agent._flush_status_buffer() agent._persist_session(messages, conversation_history) return _verdict("return", { @@ -292,8 +253,6 @@ def nous_rate_limit_guard( "failed": True, "error": _nous_msg, }) - except ImportError: - pass except Exception: pass # Never let rate guard break the agent loop return _verdict("fallthrough") diff --git a/agent/turn_response_check.py b/agent/turn_response_check.py index 05499b2030..cc04110267 100644 --- a/agent/turn_response_check.py +++ b/agent/turn_response_check.py @@ -2,9 +2,8 @@ spinner, validate the response shape (retry / eager fallback / terminal invalid-response result), derive ``finish_reason`` per api_mode, route content-policy refusals and ``length`` truncation, fold usage into the compressor, and mark the logical relay call -complete. Extracted from ``run_conversation``; nothing here imports -``agent.conversation_loop`` at module level (cycle) — loop-internal helpers resolve lazily so -existing ``patch("agent.conversation_loop.X")`` sites keep intercepting. +complete. Nothing here imports ``agent.conversation_loop`` at module level (cycle) — +loop-internal helpers resolve lazily so ``patch("agent.conversation_loop.X")`` keeps intercepting. """ from __future__ import annotations @@ -14,6 +13,7 @@ import logging import time from typing import Any, Dict, Optional +from agent.turn_api_call import stop_thinking_spinner from agent.turn_truncation import handle_content_policy_refusal, recover_from_truncation from agent.turn_usage import record_response_usage @@ -42,6 +42,45 @@ class ResponseCheckVerdict: result: Optional[Dict[str, Any]] = None +def _codex_finish_reason(response: Any) -> str: + """Responses API max-output exhaustion is a normal Codex incomplete turn: route it to + the Codex continuation path (``"incomplete"``), not the length rollback.""" + status = getattr(response, "status", None) + if isinstance(status, str): + status = status.strip().lower() + incomplete_details = getattr(response, "incomplete_details", None) + if isinstance(incomplete_details, dict): + incomplete_reason = incomplete_details.get("reason") + else: + incomplete_reason = getattr(incomplete_details, "reason", None) + if incomplete_reason is not None: + incomplete_reason = str(incomplete_reason).strip().lower() + if status == "incomplete" and incomplete_reason in {"max_output_tokens", "length"}: + return "incomplete" + if status == "incomplete" and incomplete_reason == "content_filter": + return "content_filter" + return "stop" + + +def _derive_finish_reason(agent: Any, response: Any, messages: Any) -> str: + if agent.api_mode == "codex_responses": + return _codex_finish_reason(response) + transport = agent._get_transport() + if agent.api_mode == "anthropic_messages": + return transport.map_finish_reason(response.stop_reason) + normalized = transport.normalize_response(response) # Bedrock already normalized at dispatch + finish_reason = normalized.finish_reason + if agent.api_mode != "bedrock_converse" and agent._should_treat_stop_as_truncated( + finish_reason, normalized, messages + ): + agent._vprint( + f"{agent.log_prefix}⚠️ Treating suspicious Ollama/GLM stop response as truncated", + force=True, + ) + return "length" + return finish_reason + + def check_api_response( agent: Any, *, response: Any, _retry: Any, thinking_spinner: Any, messages: Any, api_messages: Any, api_kwargs: Any, active_system_prompt: Any, conversation_history: Any, @@ -54,10 +93,7 @@ def check_api_response( """Verify ``response`` in the original order. The retry buffer is NOT cleared on success (bytes back != usable content); ``_preflight_compression_blocked``/``_last_preflight_pressure`` reset only when the usage fold re-arms the compression budget.""" - from agent.conversation_loop import ( - validate_response_shape, - ) - api_duration = None + from agent.conversation_loop import validate_response_shape def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> ResponseCheckVerdict: return ResponseCheckVerdict( @@ -74,25 +110,17 @@ def check_api_response( api_duration = time.time() - api_start_time - # Stop thinking spinner silently -- the response box or tool - # execution messages that follow are more informative. - if thinking_spinner: - thinking_spinner.stop("") - thinking_spinner = None - if agent.thinking_callback: - agent.thinking_callback("") + # Silent stop: the response box / tool messages that follow are more informative. + thinking_spinner = stop_thinking_spinner(agent, thinking_spinner) if not agent.quiet_mode: agent._vprint(f"{agent.log_prefix}⏱️ API call completed in {api_duration:.2f}s") if agent.verbose_logging: - # Log response with provider info if available resp_model = getattr(response, 'model', 'N/A') if response else 'N/A' logging.debug(f"API Response received - Model: {resp_model}, Usage: {response.usage if hasattr(response, 'usage') else 'N/A'}") - # Validate response shape before proceeding response_invalid, error_details = validate_response_shape(agent, response) - if response_invalid: _iv = retry_invalid_response( agent, response=response, error_details=error_details, _retry=_retry, @@ -112,54 +140,9 @@ def check_api_response( return _verdict(_iv.action, _iv.result) agent._turn_received_provider_response = True + finish_reason = _derive_finish_reason(agent, response, messages) - # Check finish_reason before proceeding - if agent.api_mode == "codex_responses": - status = getattr(response, "status", None) - if isinstance(status, str): - status = status.strip().lower() - incomplete_details = getattr(response, "incomplete_details", None) - incomplete_reason = None - if isinstance(incomplete_details, dict): - incomplete_reason = incomplete_details.get("reason") - else: - incomplete_reason = getattr(incomplete_details, "reason", None) - if incomplete_reason is not None: - incomplete_reason = str(incomplete_reason).strip().lower() - if status == "incomplete" and incomplete_reason in {"max_output_tokens", "length"}: - # Responses API max-output exhaustion is a normal Codex - # incomplete turn: use the Codex continuation path, not the - # length rollback. - finish_reason = "incomplete" - elif status == "incomplete" and incomplete_reason == "content_filter": - finish_reason = "content_filter" - else: - finish_reason = "stop" - elif agent.api_mode == "anthropic_messages": - _tfr = agent._get_transport() - finish_reason = _tfr.map_finish_reason(response.stop_reason) - elif agent.api_mode == "bedrock_converse": - # Bedrock response already normalized at dispatch — use transport - _bt_fr = agent._get_transport() - _bedrock_result = _bt_fr.normalize_response(response) - finish_reason = _bedrock_result.finish_reason - else: - _cc_fr = agent._get_transport() - _finish_result = _cc_fr.normalize_response(response) - finish_reason = _finish_result.finish_reason - assistant_message = _finish_result - if agent._should_treat_stop_as_truncated( - finish_reason, assistant_message, messages - ): - agent._vprint( - f"{agent.log_prefix}⚠️ Treating suspicious Ollama/GLM stop response as truncated", - force=True, - ) - finish_reason = "length" - - # ── Content-policy refusal (HTTP 200) ────────────────── - # Refusal finish reasons (``content_filter``, ``guardrail_intervened``) - # are deterministic: one fallback try, else return the refusal. + # HTTP-200 refusals are deterministic: one fallback try, else return the refusal. if finish_reason == "content_filter": _rv = handle_content_policy_refusal( agent, response, _retry, thinking_spinner=thinking_spinner, messages=messages, @@ -194,12 +177,8 @@ def check_api_response( truncated_tool_call_retries = _tv.truncated_tool_call_retries retry_count = _tv.retry_count compression_attempts = _tv.compression_attempts - if _tv.action == "return": - return _verdict("return", _tv.result) - if _tv.action == "break": - return _verdict("break") - if _tv.action == "continue": - return _verdict("continue") + if _tv.action in ("return", "break", "continue"): + return _verdict(_tv.action, _tv.result) # Fold provider usage into compressor / anchors / session counters / state.db # (agent/turn_usage.py). A rearmed budget also clears the preflight-block latch. @@ -213,10 +192,8 @@ def check_api_response( _preflight_compression_blocked = False _last_preflight_pressure = None - _retry.has_retried_429 = False # Reset on success - # Don't clear the retry buffer: bytes back != usable content; it is - # cleared once genuine content lands. Clearing Nous rate-limit state - # proves the limit reset so other sessions may resume. + _retry.has_retried_429 = False + # Clearing Nous rate-limit state proves the limit reset so other sessions may resume. if agent.provider == "nous": try: from agent.nous_rate_guard import clear_nous_rate_limit @@ -225,12 +202,9 @@ def check_api_response( pass from agent import relay_llm - relay_llm.complete_logical_call( - api_request_id, outcome="success" - ) + relay_llm.complete_logical_call(api_request_id, outcome="success") agent._touch_activity(f"API call #{api_call_count} completed") - return _verdict("break") # Success, exit retry loop - return _verdict("fallthrough") + return _verdict("break") @dataclass @@ -278,20 +252,11 @@ def retry_invalid_response( status_code=getattr(getattr(response, "error", None), "code", None), retry_count=retry_count, max_retries=max_retries, retryable=True, reason="invalid_response", ) - # Stop spinner silently — retry status is now buffered - # and only surfaced if every retry+fallback exhausts. - if thinking_spinner: - thinking_spinner.stop("") - thinking_spinner = None - if agent.thinking_callback: - agent.thinking_callback("") - - # Invalid response — could be rate limiting, provider timeout, - # upstream server error, or malformed response. + # Retry status is buffered and only surfaced if every retry+fallback exhausts. + thinking_spinner = stop_thinking_spinner(agent, thinking_spinner) retry_count += 1 - # Eager fallback: empty/malformed responses often mean rate limiting - # — switch now instead of extended backoff. + # Eager fallback: empty/malformed responses often mean rate limiting. if agent._fallback_index < len(agent._fallback_chain): agent._buffer_status("⚠️ Empty/malformed response — switching to fallback...") if agent._try_activate_fallback(): @@ -304,15 +269,12 @@ def retry_invalid_response( error_msg, provider_name, _failure_hint = describe_invalid_response( agent, response, api_duration ) - agent._buffer_vprint(f"⚠️ Invalid API response (attempt {retry_count}/{max_retries}): {', '.join(error_details)}") agent._buffer_vprint(f" 🏢 Provider: {provider_name}") - cleaned_provider_error = agent._clean_error_message(error_msg) - agent._buffer_vprint(f" 📝 Provider message: {cleaned_provider_error}") + agent._buffer_vprint(f" 📝 Provider message: {agent._clean_error_message(error_msg)}") agent._buffer_vprint(f" ⏱️ {_failure_hint}") if retry_count >= max_retries: - # Try fallback before giving up if agent._has_pending_fallback(): agent._buffer_status(f"⚠️ Max retries ({max_retries}) for invalid responses — trying fallback...") if agent._try_activate_fallback(): @@ -333,17 +295,15 @@ def retry_invalid_response( "completed": False, "api_calls": api_call_count, "error": _final_response, - "failed": True # Mark as failure for filtering + "failed": True, }) - # Backoff before retry — jittered exponential: 5s base, 120s cap wait_time = jittered_backoff(retry_count, base_delay=5.0, max_delay=120.0) agent._buffer_vprint(f"⏳ Retrying in {wait_time:.1f}s ({_failure_hint})...") logger.warning("Invalid API response (retry %d/%d): %s | Provider: %s", retry_count, max_retries, ', '.join(error_details), provider_name) - # A redirect cancels only the live request; the helper preserves the - # pending correction (restart_with_redirected_messages) instead of - # destroying it with clear_interrupt(). + # A redirect cancels only the live request; the helper preserves the pending + # correction (restart_with_redirected_messages) instead of clear_interrupt()-ing it. _interrupted = interruptible_backoff_sleep( agent, wait_time, _retry, messages=messages, conversation_history=conversation_history, api_call_count=api_call_count, @@ -355,5 +315,4 @@ def retry_invalid_response( return _verdict("return", _interrupted) if _retry.restart_with_redirected_messages: return _verdict("break") # rebuild this iteration from the correction - return _verdict("continue") # Retry the API call - return _verdict("fallthrough") + return _verdict("continue") diff --git a/agent/turn_truncation.py b/agent/turn_truncation.py index c7d5e063b0..a2ad233a80 100644 --- a/agent/turn_truncation.py +++ b/agent/turn_truncation.py @@ -1,11 +1,10 @@ """Truncation recovery (``finish_reason == "length"``) for the conversation turn loop. -Extracted from ``run_conversation``. Handles thinking-budget exhaustion, repetition- -dominated truncation (#86581), content-filter stream stalls escalated to the fallback -chain (#32421), text continuation nudges (up to 4, with the ceiling exit that drops the -fragment trail), truncated tool-call retries with max_tokens boosts, and the final -roll-back. Nothing here imports ``agent.conversation_loop`` at module level (cycle); -loop-internal helpers are imported lazily so tests patching them on the loop keep working. +Handles thinking-budget exhaustion, repetition-dominated truncation, content-filter stream +stalls escalated to the fallback chain, text continuation nudges (up to 4, with the ceiling +exit that drops the fragment trail), truncated tool-call retries with max_tokens boosts, and +the final roll-back. Nothing here imports ``agent.conversation_loop`` at module level +(cycle); loop-internal helpers are imported lazily so tests patching them keep working. """ from __future__ import annotations @@ -19,11 +18,80 @@ from agent.error_classifier import FailoverReason from agent.message_metadata import append_message from agent.message_sanitization import close_interrupted_tool_sequence from agent.repetition_guard import is_repetition_dominated +from agent.turn_api_call import stop_thinking_spinner from agent.turn_retry_state import TurnRetryState from hermes_constants import PARTIAL_STREAM_STUB_ID logger = logging.getLogger("agent.conversation_loop") +_CONTINUABLE_MODES = {"chat_completions", "bedrock_converse", "anthropic_messages"} +_THINK_TAG_RE = re.compile(r'<(?:think|thinking|reasoning|REASONING_SCRATCHPAD)[^>]*>', re.IGNORECASE) +_TRUNCATED_FINAL = "Response truncated due to output length limit" +_FIRST_TRUNCATED_FINAL = "First response truncated due to output length limit" + +_THINKING_EXHAUSTED = ( + "💭 Reasoning exhausted the output token budget — no visible response was produced.", + "⚠️ **Thinking Budget Exhausted**\n\n" + "The model used all its output tokens on reasoning " + "and had none left for the actual response.\n\n" + "To fix this:\n" + "→ Lower reasoning effort: `/reasoning low` or `/reasoning minimal`\n" + "→ Or switch to a larger/non-reasoning model with `/model`", + "Model used all output tokens on reasoning with none left " + "for the response. Try lowering reasoning effort or " + "increasing max_tokens.", +) +_REPETITION_DOMINATED = ( + "🔁 Response dominated by repeated text — stopping instead of " + "continuing a degenerate response.", + "⚠️ **Response Stopped — Repetition Detected**\n\n" + "The model fell into a repetition loop while " + "writing this response, so continuing would only " + "produce more repeated text. The partial response " + "was discarded.\n\n" + "→ Switch to a different model with `/model`\n" + "→ Or resend your message (your conversation " + "history is preserved)", + "Model output entered a repetition loop and was " + "truncated mid-loop; refusing to continue a " + "degenerate response.", +) +_CEILING_NO_TEXT = ( + "⚠️ **No visible answer was produced.** The " + "model hit its output-token limit on every " + "continuation attempt — its reasoning " + "consumed the entire budget each time.\n\n" + "To fix this:\n" + "→ Lower reasoning effort: `/reasoning low` " + "or `/reasoning none`\n" + "→ Or raise max_tokens for this model" +) + + +def normalize_response_for_agent(agent: Any, response: Any) -> Any: + """One OpenAI-style message from any transport; Anthropic strips the OAuth tool prefix.""" + if agent.api_mode == "anthropic_messages": + return agent._get_transport().normalize_response( + response, strip_tool_prefix=agent._is_anthropic_oauth + ) + return agent._get_transport().normalize_response(response) + + +def partial_result( + messages: List[Dict[str, Any]], api_call_count: int, final_response: str, + error: Optional[str] = None, *, failed: bool = False, +) -> Dict[str, Any]: + """Typed incomplete-turn result (``partial`` unless ``failed``); ``error`` defaults to + ``final_response``.""" + return { + "final_response": final_response, + "messages": messages, + "api_calls": api_call_count, + "completed": False, + ("failed" if failed else "partial"): True, + "error": final_response if error is None else error, + } + @dataclass class TruncationVerdict: @@ -45,6 +113,220 @@ class TruncationVerdict: compression_attempts: int +@dataclass +class _Trunc: + """Mutable working state for the truncation phases; rebound loop locals are handed + back through ``verdict()``.""" + + agent: Any + response: Any + finish_reason: str + conversation_history: Any + api_call_count: int + effective_task_id: Any + current_turn_user_idx: Any + messages: List[Dict[str, Any]] + length_continue_retries: int + truncated_response_parts: List[str] + truncated_tool_call_retries: int + retry_count: int + compression_attempts: int + + def verdict(self, action: str, result: Optional[Dict[str, Any]] = None) -> TruncationVerdict: + return TruncationVerdict( + action=action, result=result, messages=self.messages, + length_continue_retries=self.length_continue_retries, + truncated_response_parts=self.truncated_response_parts, + truncated_tool_call_retries=self.truncated_tool_call_retries, + retry_count=self.retry_count, compression_attempts=self.compression_attempts, + ) + + def end_turn( + self, final_response: str, error: Optional[str] = None, *, + result_messages: Optional[List[Dict[str, Any]]] = None, cleanup: bool = True, + failed: bool = False, + ) -> TruncationVerdict: + """Persist and end the turn as partial (or ``failed``).""" + agent = self.agent + if cleanup: + agent._cleanup_task_resources(self.effective_task_id) + agent._persist_session(self.messages, self.conversation_history) + return self.verdict("return", partial_result( + self.messages if result_messages is None else result_messages, self.api_call_count, + final_response, error, failed=failed, + )) + + @property + def is_stub(self) -> bool: + return getattr(self.response, "id", "") == PARTIAL_STREAM_STUB_ID + + +def _abort_reason(agent: Any, content: Any, has_tool_calls: bool) -> Optional[tuple]: + """``(vprint, user response, error)`` when continuation must NOT be attempted: + thinking exhausted the budget (reasoning blocks with no visible text after them — + ``content=None`` from non- models is normal truncation), or a repetition loop + burned the budget on one fragment (reasoning stripped first).""" + if has_tool_calls: + return None + if content and _THINK_TAG_RE.search(content) and not agent._has_content_after_think_block(content): + return _THINKING_EXHAUSTED + visible = agent._strip_think_blocks(content) if isinstance(content, str) else content + if visible and is_repetition_dominated(visible): + return _REPETITION_DOMINATED + return None + + +def _content_filter_fallback(st: _Trunc, _retry: TurnRetryState) -> Optional[TruncationVerdict]: + """Content-filter stream stall → fallback. ``_content_filter_terminated`` is + content-deterministic, so escalate before retrying the primary; without a fallback + fall through to normal continuation (best-effort, may loop).""" + agent = st.agent + if not ( + getattr(st.response, "_content_filter_terminated", False) + and agent._fallback_index < len(agent._fallback_chain) + ): + return None + agent._vprint( + f"{agent.log_prefix}🛡️ Content filter terminated " + f"stream — activating fallback provider...", + force=True, + ) + agent._emit_status("Content filter terminated stream; switching to fallback...") + if agent._try_activate_fallback(): + # Roll partial content back to the last clean turn so the fallback gets a + # coherent continuation point; unmark survivors (their text left the partial). + if st.truncated_response_parts: + st.messages = agent._get_messages_up_to_last_assistant(st.messages) + for _frag in st.messages: + if isinstance(_frag, dict): + _frag.pop("_length_continuation_fragment", None) + _frag.pop("_length_continuation_nudge", None) + agent._session_messages = st.messages + st.length_continue_retries = 0 + st.truncated_response_parts = [] + st.retry_count = 0 + st.compression_attempts = 0 + _retry.primary_recovery_attempted = False + _retry.restart_with_rebuilt_messages = True + return st.verdict("break") + agent._vprint( + f"{agent.log_prefix}⚠️ No fallback provider " + f"configured — retrying with same provider " + f"(may re-hit filter)...", + force=True, + ) + return None + + +def _continue_text(st: _Trunc, _retry: TurnRetryState, assistant_message: Any) -> TruncationVerdict: + """Text truncation (no tool calls): append the fragment + a continuation nudge (up to + 4), then the ceiling exit that drops the fragment trail and keeps the stitched partial. + Never appends an interim assistant row with NO visible content — strict providers + reject it with 400 — only the nudge.""" + from agent.conversation_loop import _get_continuation_prompt, _join_truncated_parts + + agent = st.agent + messages = st.messages + st.length_continue_retries += 1 + n = st.length_continue_retries + _interim_content = getattr(assistant_message, "content", None) + if not _interim_content and not st.is_stub: + # Thinking-only truncation: continuing with thinking ON re-burns the budget. + agent._ephemeral_reasoning_off = True + if _interim_content: + interim_msg = agent._build_assistant_message(assistant_message, st.finish_reason) + interim_msg["_length_continuation_fragment"] = True # ceiling exit drops these + append_message(messages, interim_msg) + st.truncated_response_parts.append(_interim_content) + + if n < 4: + _dropped_tools = getattr(st.response, "_dropped_tool_names", None) + if st.is_stub and _dropped_tools: + agent._vprint( + f"{agent.log_prefix}↻ Stream interrupted mid " + f"tool-call ({', '.join(_dropped_tools[:3])}) — requesting " + f"chunked retry " + f"({n}/4)..." + ) + elif st.is_stub: + agent._vprint(f"{agent.log_prefix}↻ Stream interrupted — requesting continuation ({n}/4)...") + else: + agent._vprint(f"{agent.log_prefix}↻ Requesting continuation ({n}/4)...") + append_message(messages, { + "role": "user", "content": _get_continuation_prompt(st.is_stub, _dropped_tools), + "_length_continuation_nudge": True, + }) + agent._session_messages = messages + _retry.restart_with_length_continuation = True + return st.verdict("break") + + partial_response = agent._strip_think_blocks(_join_truncated_parts(st.truncated_response_parts)).strip() + # The one-shot reasoning-off override must not leak into the next turn. + agent._ephemeral_reasoning_off = False + agent._vprint( + f"{agent.log_prefix}⚠️ Response still truncated " + f"after {n} continuation attempts — " + + ("keeping the partial response received so far." if partial_response + else "no visible text was produced."), + force=True, + ) + # Unanswered continue nudges made every later turn re-truncate: drop the trail. + idx = st.current_turn_user_idx + _turn_start = idx + 1 if isinstance(idx, int) and idx >= 0 else 0 + messages[_turn_start:] = [ + m for m in messages[_turn_start:] + if not (isinstance(m, dict) and ( + m.get("_length_continuation_fragment") or m.get("_length_continuation_nudge") + )) + ] + if partial_response: + append_message(messages, { + "role": "assistant", "content": partial_response, "finish_reason": "length" + }) + agent._session_messages = messages + return st.end_turn( + partial_response or _CEILING_NO_TEXT, + "Response remained truncated after 4 continuation attempts", + ) + + +def _retry_truncated_tool_call(st: _Trunc, api_kwargs: Any) -> TruncationVerdict: + """Truncated tool call: re-run the same call (up to 4×) with a boosted max_tokens — + a real output-cap truncation needs it, harmless for a network stall — else refuse to + execute incomplete arguments.""" + agent = st.agent + if st.truncated_tool_call_retries < 4: + st.truncated_tool_call_retries += 1 + n = st.truncated_tool_call_retries + if st.is_stub: + agent._buffer_vprint(f"⚠️ Stream interrupted mid tool-call — retrying ({n}/4)...") + else: + agent._buffer_vprint(f"⚠️ Truncated tool call detected — retrying API call ({n}/4)...") + _tc_boost = (agent.max_tokens if agent.max_tokens else 4096) * (2 ** n) + _tc_requested_cap = agent._requested_output_cap_from_api_kwargs(api_kwargs) + if _tc_requested_cap is not None: + _tc_boost = max(_tc_boost, _tc_requested_cap) + agent._ephemeral_max_output_tokens = min(_tc_boost, max(32768, _tc_requested_cap or 0)) + return st.verdict("continue") # don't append the broken response + agent._flush_status_buffer() + if st.is_stub: + agent._vprint( + f"{agent.log_prefix}⚠️ Stream kept dropping mid tool-call after 4 retries — the action was not executed.", + force=True, + ) + _final_response = "Stream repeatedly dropped mid tool-call (network); the tool was not executed" + else: + agent._vprint( + f"{agent.log_prefix}⚠️ Truncated tool call response detected again — refusing to execute incomplete tool arguments.", + force=True, + ) + _final_response = _TRUNCATED_FINAL + agent._cleanup_task_resources(st.effective_task_id) + # Prior tool batches can leave a tool-result tail; this path never reaches finalize_turn. + close_interrupted_tool_sequence(st.messages, _final_response) + return st.end_turn(_final_response, cleanup=False) + + def recover_from_truncation( agent: Any, response: Any, finish_reason: str, _retry: TurnRetryState, *, messages: List[Dict[str, Any]], conversation_history: Any, api_kwargs: Any, api_call_count: int, @@ -55,408 +337,57 @@ def recover_from_truncation( """Recover from a truncated response. Order is load-bearing: thinking exhaustion and repetition abort BEFORE any continuation; a content-filter stall escalates to the fallback chain BEFORE the primary is retried; text continuation (no tool calls) then - truncated tool-call retry; finally roll back to the last complete assistant turn. - Never appends an interim assistant row with NO visible content (strict providers - reject it with 400) — only the continuation nudge.""" - from agent.conversation_loop import _get_continuation_prompt, _join_truncated_parts - - def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> TruncationVerdict: - return TruncationVerdict( - action=action, result=result, messages=messages, - length_continue_retries=length_continue_retries, - truncated_response_parts=truncated_response_parts, - truncated_tool_call_retries=truncated_tool_call_retries, retry_count=retry_count, - compression_attempts=compression_attempts, - ) - - if getattr(response, "id", "") == PARTIAL_STREAM_STUB_ID: - agent._vprint( - f"{agent.log_prefix}⚠️ Response truncated — stream " - f"ended before completion", - force=True, - ) - else: - agent._vprint( - f"{agent.log_prefix}⚠️ Response truncated " - f"(finish_reason='length') - model hit max output tokens", - force=True, - ) - - # Normalize to one OpenAI-style message so continuation and tool- - # call retry work across transports (Anthropic reuses the loop's - # adapter). - _trunc_msg = None - _trunc_transport = agent._get_transport() - if agent.api_mode == "anthropic_messages": - _trunc_result = _trunc_transport.normalize_response( - response, strip_tool_prefix=agent._is_anthropic_oauth - ) - else: - _trunc_result = _trunc_transport.normalize_response(response) - _trunc_msg = _trunc_result + truncated tool-call retry; finally roll back to the last complete assistant turn.""" + st = _Trunc( + agent=agent, response=response, finish_reason=finish_reason, + conversation_history=conversation_history, api_call_count=api_call_count, + effective_task_id=effective_task_id, current_turn_user_idx=current_turn_user_idx, + messages=messages, length_continue_retries=length_continue_retries, + truncated_response_parts=truncated_response_parts, + truncated_tool_call_retries=truncated_tool_call_retries, retry_count=retry_count, + compression_attempts=compression_attempts, + ) + agent._vprint( + f"{agent.log_prefix}⚠️ Response truncated — stream ended before completion" + if st.is_stub else + f"{agent.log_prefix}⚠️ Response truncated (finish_reason='length') - model hit max output tokens", + force=True, + ) + _trunc_msg = normalize_response_for_agent(agent, response) _trunc_content = getattr(_trunc_msg, "content", None) if _trunc_msg else None _trunc_has_tool_calls = bool(getattr(_trunc_msg, "tool_calls", None)) if _trunc_msg else False - # ── Detect thinking-budget exhaustion ────────────── - # Only when reasoning blocks exist with no visible text after them; - # content=None from non- models is normal truncation. - _has_think_tags = bool( - _trunc_content and re.search( - r'<(?:think|thinking|reasoning|REASONING_SCRATCHPAD)[^>]*>', - _trunc_content, - re.IGNORECASE, - ) - ) - _thinking_exhausted = ( - not _trunc_has_tool_calls - and _has_think_tags - and ( - (_trunc_content is not None and not agent._has_content_after_think_block(_trunc_content)) - or _trunc_content is None - ) - ) + abort = _abort_reason(agent, _trunc_content, _trunc_has_tool_calls) + if abort is not None: + line, user_response, error = abort + agent._vprint(f"{agent.log_prefix}{line}", force=True) + return st.end_turn(user_response, error) - if _thinking_exhausted: - _exhaust_error = ( - "Model used all output tokens on reasoning with none left " - "for the response. Try lowering reasoning effort or " - "increasing max_tokens." - ) - agent._vprint( - f"{agent.log_prefix}💭 Reasoning exhausted the output token budget — " - f"no visible response was produced.", - force=True, - ) - # Return a user-friendly message as the response so CLI and - # gateway display it. - _exhaust_response = ( - "⚠️ **Thinking Budget Exhausted**\n\n" - "The model used all its output tokens on reasoning " - "and had none left for the actual response.\n\n" - "To fix this:\n" - "→ Lower reasoning effort: `/reasoning low` or `/reasoning minimal`\n" - "→ Or switch to a larger/non-reasoning model with `/model`" - ) - agent._cleanup_task_resources(effective_task_id) - agent._persist_session(messages, conversation_history) - return _verdict("return", { - "final_response": _exhaust_response, - "messages": messages, - "api_calls": api_call_count, - "completed": False, - "partial": True, - "error": _exhaust_error, - }) + if agent.api_mode in _CONTINUABLE_MODES: + cf = _content_filter_fallback(st, _retry) + if cf is not None: + return cf + if _trunc_msg is not None: + if not _trunc_has_tool_calls: + return _continue_text(st, _retry, _trunc_msg) + return _retry_truncated_tool_call(st, api_kwargs) - # ── Detect repetition-dominated truncation (#86581) ── - # A repetition loop can burn the whole budget on one fragment; abort - # like _thinking_exhausted (reasoning stripped first). - _visible_trunc = ( - agent._strip_think_blocks(_trunc_content) - if isinstance(_trunc_content, str) - else _trunc_content - ) - _repetition_dominated = ( - not _trunc_has_tool_calls - and bool(_visible_trunc) - and is_repetition_dominated(_visible_trunc) - ) - if _repetition_dominated: - _rep_error = ( - "Model output entered a repetition loop and was " - "truncated mid-loop; refusing to continue a " - "degenerate response." - ) - agent._vprint( - f"{agent.log_prefix}🔁 Response dominated by " - f"repeated text — stopping instead of " - f"continuing a degenerate response.", - force=True, - ) - _rep_response = ( - "⚠️ **Response Stopped — Repetition Detected**\n\n" - "The model fell into a repetition loop while " - "writing this response, so continuing would only " - "produce more repeated text. The partial response " - "was discarded.\n\n" - "→ Switch to a different model with `/model`\n" - "→ Or resend your message (your conversation " - "history is preserved)" - ) - agent._cleanup_task_resources(effective_task_id) - agent._persist_session(messages, conversation_history) - return _verdict("return", { - "final_response": _rep_response, - "messages": messages, - "api_calls": api_call_count, - "completed": False, - "partial": True, - "error": _rep_error, - }) - - if agent.api_mode in {"chat_completions", "bedrock_converse", "anthropic_messages"}: - assistant_message = _trunc_msg - # ── Content-filter stream stall → fallback (#32421) ── - # ``_content_filter_terminated`` is content-deterministic; - # escalate to the fallback before retrying the primary. - _cf_terminated = getattr( - response, "_content_filter_terminated", False - ) - if ( - _cf_terminated and agent._fallback_index < len(agent._fallback_chain) - ): - agent._vprint( - f"{agent.log_prefix}🛡️ Content filter terminated " - f"stream — activating fallback provider...", - force=True, - ) - agent._emit_status( - "Content filter terminated stream; switching to fallback..." - ) - if agent._try_activate_fallback(): - # Roll partial content back to the last clean turn so - # the fallback gets a coherent continuation point. - if truncated_response_parts: - messages = agent._get_messages_up_to_last_assistant(messages) - # Unmark survivors: their text left the stitched partial. - for _frag in messages: - if isinstance(_frag, dict): - _frag.pop("_length_continuation_fragment", None) - _frag.pop("_length_continuation_nudge", None) - agent._session_messages = messages - length_continue_retries = 0 - truncated_response_parts = [] - retry_count = 0 - compression_attempts = 0 - _retry.primary_recovery_attempted = False - _retry.restart_with_rebuilt_messages = True - return _verdict("break") - # No fallback available — fall through to normal - # continuation (best-effort, may loop). - agent._vprint( - f"{agent.log_prefix}⚠️ No fallback provider " - f"configured — retrying with same provider " - f"(may re-hit filter)...", - force=True, - ) - if assistant_message is not None and not _trunc_has_tool_calls: - length_continue_retries += 1 - # Never append an interim assistant message with NO visible - # content: strict providers reject it (HTTP 400), poisoning - # history. Append only the nudge. - _interim_content = getattr(assistant_message, "content", None) - _is_empty_partial_stub = ( - getattr(response, "id", "") == PARTIAL_STREAM_STUB_ID and not _interim_content - ) - if not _interim_content and not _is_empty_partial_stub: - # Thinking-only truncation: continuing with thinking ON - # re-burns the budget, so drop thinking for one request. - agent._ephemeral_reasoning_off = True - if _interim_content: - interim_msg = agent._build_assistant_message(assistant_message, finish_reason) - # Marked so the ceiling exit can drop the fragment trail. - interim_msg["_length_continuation_fragment"] = True - append_message(messages, interim_msg) - truncated_response_parts.append(_interim_content) - - if length_continue_retries < 4: - _is_partial_stream_stub = ( - getattr(response, "id", "") == PARTIAL_STREAM_STUB_ID - ) - _dropped_tools = getattr( - response, "_dropped_tool_names", None - ) - - if _is_partial_stream_stub and _dropped_tools: - _tool_list = ", ".join(_dropped_tools[:3]) - agent._vprint( - f"{agent.log_prefix}↻ Stream interrupted mid " - f"tool-call ({_tool_list}) — requesting " - f"chunked retry " - f"({length_continue_retries}/4)..." - ) - elif _is_partial_stream_stub: - agent._vprint( - f"{agent.log_prefix}↻ Stream interrupted — " - f"requesting continuation " - f"({length_continue_retries}/4)..." - ) - else: - agent._vprint( - f"{agent.log_prefix}↻ Requesting continuation " - f"({length_continue_retries}/4)..." - ) - - _continue_content = _get_continuation_prompt( - _is_partial_stream_stub, _dropped_tools - ) - continue_msg = { - "role": "user", "content": _continue_content, "_length_continuation_nudge": True - } - append_message(messages, continue_msg) - agent._session_messages = messages - _retry.restart_with_length_continuation = True - return _verdict("break") - - partial_response = agent._strip_think_blocks(_join_truncated_parts(truncated_response_parts)).strip() - # The one-shot reasoning-off override must not leak into the - # next turn when the ceiling exit skips the consuming call. - agent._ephemeral_reasoning_off = False - if partial_response: - agent._vprint( - f"{agent.log_prefix}⚠️ Response still truncated " - f"after {length_continue_retries} continuation attempts — keeping the " - f"partial response received so far.", - force=True, - ) - _ceiling_final = partial_response - else: - # Every fragment was empty (e.g. reasoning-only model): - # return an actionable message, not a bare None. - agent._vprint( - f"{agent.log_prefix}⚠️ Response still truncated " - f"after {length_continue_retries} continuation attempts — no visible " - f"text was produced.", - force=True, - ) - _ceiling_final = ( - "⚠️ **No visible answer was produced.** The " - "model hit its output-token limit on every " - "continuation attempt — its reasoning " - "consumed the entire budget each time.\n\n" - "To fix this:\n" - "→ Lower reasoning effort: `/reasoning low` " - "or `/reasoning none`\n" - "→ Or raise max_tokens for this model" - ) - # Unanswered continue nudges made every later turn re-truncate. - _turn_start = ( - current_turn_user_idx + 1 - if isinstance(current_turn_user_idx, int) - and current_turn_user_idx >= 0 - else 0 - ) - messages[_turn_start:] = [ - m for m in messages[_turn_start:] - if not ( - isinstance(m, dict) - and ( - m.get("_length_continuation_fragment") - or m.get("_length_continuation_nudge") - ) - ) - ] - if partial_response: - append_message(messages, { - "role": "assistant", "content": partial_response, "finish_reason": "length" - }) - agent._session_messages = messages - agent._cleanup_task_resources(effective_task_id) - agent._persist_session(messages, conversation_history) - return _verdict("return", { - "final_response": _ceiling_final, - "messages": messages, - "api_calls": api_call_count, - "completed": False, - "partial": True, - "error": "Response remained truncated after 4 continuation attempts", - }) - - if agent.api_mode in {"chat_completions", "bedrock_converse", "anthropic_messages"}: - assistant_message = _trunc_msg - if assistant_message is not None and _trunc_has_tool_calls: - _is_stub_stall = ( - getattr(response, "id", "") == PARTIAL_STREAM_STUB_ID - ) - if truncated_tool_call_retries < 4: - truncated_tool_call_retries += 1 - if _is_stub_stall: - # Stream broke mid tool-call (network), not a real - # output cap — say so. - agent._buffer_vprint( - f"⚠️ Stream interrupted mid tool-call — " - f"retrying ({truncated_tool_call_retries}/4)..." - ) - else: - agent._buffer_vprint( - f"⚠️ Truncated tool call detected — " - f"retrying API call " - f"({truncated_tool_call_retries}/4)..." - ) - # Boost max_tokens per retry: a real output-cap - # truncation needs it; harmless for a stall. - _tc_boost_base = agent.max_tokens if agent.max_tokens else 4096 - _tc_boost = _tc_boost_base * (2 ** truncated_tool_call_retries) - _tc_requested_cap = agent._requested_output_cap_from_api_kwargs(api_kwargs) - if _tc_requested_cap is not None: - _tc_boost = max(_tc_boost, _tc_requested_cap) - _tc_boost_cap = max(32768, _tc_requested_cap or 0) - agent._ephemeral_max_output_tokens = min(_tc_boost, _tc_boost_cap) - # Don't append the broken response; re-run the same call - # from current state. - return _verdict("continue") - agent._flush_status_buffer() - if _is_stub_stall: - agent._vprint( - f"{agent.log_prefix}⚠️ Stream kept dropping mid tool-call after 4 retries — the action was not executed.", - force=True, - ) - else: - agent._vprint( - f"{agent.log_prefix}⚠️ Truncated tool call response detected again — refusing to execute incomplete tool arguments.", - force=True, - ) - agent._cleanup_task_resources(effective_task_id) - _final_response = ( - "Stream repeatedly dropped mid tool-call (network); " - "the tool was not executed" - if _is_stub_stall - else "Response truncated due to output length limit" - ) - # Prior tool batches can leave a tool-result tail; this path - # never reaches finalize_turn (#48879). - close_interrupted_tool_sequence(messages, _final_response) - agent._persist_session(messages, conversation_history) - return _verdict("return", { - "final_response": _final_response, - "messages": messages, - "api_calls": api_call_count, - "completed": False, - "partial": True, - "error": _final_response, - }) - - # If we have prior messages, roll back to last complete state if len(messages) > 1: agent._vprint(f"{agent.log_prefix} ⏪ Rolling back to last complete assistant turn") - rolled_back_messages = agent._get_messages_up_to_last_assistant(messages) + return st.end_turn( + _TRUNCATED_FINAL, result_messages=agent._get_messages_up_to_last_assistant(messages) + ) + # First message was truncated - mark as failed + agent._flush_status_buffer() + agent._vprint(f"{agent.log_prefix}❌ First response truncated - cannot recover", force=True) + return st.end_turn(_FIRST_TRUNCATED_FINAL, cleanup=False, failed=True) - agent._cleanup_task_resources(effective_task_id) - agent._persist_session(messages, conversation_history) - return _verdict("return", { - "final_response": "Response truncated due to output length limit", - "messages": rolled_back_messages, - "api_calls": api_call_count, - "completed": False, - "partial": True, - "error": "Response truncated due to output length limit" - }) - else: - # First message was truncated - mark as failed - agent._flush_status_buffer() - agent._vprint(f"{agent.log_prefix}❌ First response truncated - cannot recover", force=True) - agent._persist_session(messages, conversation_history) - return _verdict("return", { - "final_response": "First response truncated due to output length limit", - "messages": messages, - "api_calls": api_call_count, - "completed": False, - "failed": True, - "error": "First response truncated due to output length limit" - }) - return _verdict("fallthrough") +_CODEX_REPLAY_KEYS = ( + "content", "reasoning", "reasoning_content", "reasoning_details", + "codex_reasoning_items", "codex_message_items", +) def continue_codex_incomplete( @@ -466,128 +397,84 @@ def continue_codex_incomplete( """Codex Responses ``status=incomplete`` continuation (max 3 per turn). Appends the interim assistant message (deduped on visible content only — opaque - provider state drifts per continuation, #52711; ``codex_reasoning_items`` are merged, - not overwritten, because the earlier response holds the only native-compaction + provider state drifts per continuation; ``codex_reasoning_items`` are merged, not + overwritten, because the earlier response holds the only native-compaction checkpoint) and, when a bare retry would be byte-identical, a user-role nudge — only after an assistant row, to preserve role alternation. Returns ``None`` to continue the turn loop, or the terminal ``partial`` result once retries are exhausted.""" from agent.conversation_loop import _CODEX_INCOMPLETE_NUDGE agent._codex_incomplete_retries += 1 + n = agent._codex_incomplete_retries interim_msg = agent._build_assistant_message(assistant_message, finish_reason) interim_has_content = bool((interim_msg.get("content") or "").strip()) - interim_has_reasoning = bool(interim_msg.get("reasoning", "").strip()) if isinstance(interim_msg.get("reasoning"), str) else False + _reasoning = interim_msg.get("reasoning") + interim_has_reasoning = isinstance(_reasoning, str) and bool(_reasoning.strip()) interim_has_codex_reasoning = bool(interim_msg.get("codex_reasoning_items")) interim_has_codex_message_items = bool(interim_msg.get("codex_message_items")) - if ( - interim_has_content - or interim_has_reasoning - or interim_has_codex_reasoning - or interim_has_codex_message_items - ): + if interim_has_content or interim_has_reasoning or interim_has_codex_reasoning or interim_has_codex_message_items: last_msg = messages[-1] if messages else None - # Dedup on visible content only (content + reasoning): opaque - # provider state drifts per continuation and would defeat dedup - # (#52711). - last_interim_visible = ( - agent._interim_assistant_visible_text(last_msg) if isinstance(last_msg, dict) else "" - ) + last_is_dict = isinstance(last_msg, dict) + last_interim_visible = agent._interim_assistant_visible_text(last_msg) if last_is_dict else "" current_interim_visible = agent._interim_assistant_visible_text(interim_msg) if last_interim_visible or current_interim_visible: same_visible_output = last_interim_visible == current_interim_visible else: - # Preserve the existing reasoning-only behavior when - # neither response has text eligible for interim delivery. - same_visible_output = ( + # Neither has text eligible for interim delivery: compare raw content+reasoning. + same_visible_output = last_is_dict and ( (last_msg.get("content") or "") == (interim_msg.get("content") or "") and (last_msg.get("reasoning") or "") == (interim_msg.get("reasoning") or "") - ) if isinstance(last_msg, dict) else False - visible_duplicate = ( - isinstance(last_msg, dict) + ) + if ( + last_is_dict and last_msg.get("role") == "assistant" and last_msg.get("finish_reason") == "incomplete" and same_visible_output - ) - if visible_duplicate: - # Update replay state in-place: keep the latest provider payload - # without re-emitting identical user-visible commentary. - for _key in ( - "content", - "reasoning", - "reasoning_content", - "reasoning_details", - "codex_reasoning_items", - "codex_message_items", - ): - if _key in interim_msg: - if _key == "codex_reasoning_items": - # Merge, don't overwrite: the earlier response's - # native compaction checkpoint is the only copy. See - # merge_interim_reasoning_items. - from agent.native_compaction import ( - merge_interim_reasoning_items, - ) - last_msg[_key] = merge_interim_reasoning_items( - last_msg.get(_key), interim_msg[_key] - ) - else: - last_msg[_key] = interim_msg[_key] + ): + # Duplicate: refresh replay state in place, no re-emitted commentary. + for _key in _CODEX_REPLAY_KEYS: + if _key not in interim_msg: + continue + if _key == "codex_reasoning_items": + from agent.native_compaction import merge_interim_reasoning_items + last_msg[_key] = merge_interim_reasoning_items(last_msg.get(_key), interim_msg[_key]) + else: + last_msg[_key] = interim_msg[_key] else: append_message(messages, interim_msg) agent._emit_interim_assistant_message(interim_msg) - if agent._codex_incomplete_retries < 3: - # If the interim has nothing the Responses converter will replay, a - # bare retry is byte-identical and fails identically; append a - # user-role nudge so the retry differs and asks for the answer. - interim_replayable = ( - interim_has_content or interim_has_codex_reasoning or interim_has_codex_message_items - ) - # Replayable ≠ different: an interim holding only a ``compaction`` - # checkpoint in ``codex_reasoning_items`` is replayable yet re-sends - # identically. One bare retry, then always nudge. - if not interim_replayable or agent._codex_incomplete_retries >= 2: + if n < 3: + # If the interim has nothing the Responses converter will replay, a bare retry is + # byte-identical; a replayable interim holding only a ``compaction`` checkpoint + # ALSO re-sends identically. One bare retry, then always nudge. + interim_replayable = interim_has_content or interim_has_codex_reasoning or interim_has_codex_message_items + if not interim_replayable or n >= 2: _last_msg = messages[-1] if messages else None - _already_nudged = ( - isinstance(_last_msg, dict) - and _last_msg.get("role") == "user" - and _last_msg.get("content") == _CODEX_INCOMPLETE_NUDGE - ) - # Alternation guard: the user-role nudge may only follow an - # assistant message; after a too-empty interim it would create - # user→user / tool→user. - _last_is_assistant = ( - isinstance(_last_msg, dict) and _last_msg.get("role") == "assistant" - ) - if not _already_nudged and _last_is_assistant: - append_message(messages, { - "role": "user", "content": _CODEX_INCOMPLETE_NUDGE - }) + if isinstance(_last_msg, dict): + _already_nudged = ( + _last_msg.get("role") == "user" and _last_msg.get("content") == _CODEX_INCOMPLETE_NUDGE + ) + # Alternation guard: the nudge may only follow an assistant row. + if not _already_nudged and _last_msg.get("role") == "assistant": + append_message(messages, {"role": "user", "content": _CODEX_INCOMPLETE_NUDGE}) if not agent.quiet_mode: - agent._vprint(f"{agent.log_prefix}↻ Codex response incomplete; continuing turn ({agent._codex_incomplete_retries}/3)") - # Show the continuation on the spinner/status line and gateway - # heartbeat; these retries can take minutes and otherwise look like - # infinite thinking (#64434). + agent._vprint(f"{agent.log_prefix}↻ Codex response incomplete; continuing turn ({n}/3)") + # Spinner/heartbeat notice: these retries can take minutes and otherwise look + # like infinite thinking. agent._emit_wait_notice( - f"↻ model returned reasoning with no final answer — " - f"asking it to continue " - f"({agent._codex_incomplete_retries}/3)" + f"↻ model returned reasoning with no final answer — asking it to continue ({n}/3)" ) agent._session_messages = messages return None agent._codex_incomplete_retries = 0 agent._persist_session(messages, conversation_history) - return { - "final_response": "Codex response remained incomplete after 3 continuation attempts", - "messages": messages, - "api_calls": api_call_count, - "completed": False, - "partial": True, - "error": "Codex response remained incomplete after 3 continuation attempts", - } + return partial_result( + messages, api_call_count, "Codex response remained incomplete after 3 continuation attempts" + ) @dataclass @@ -610,25 +497,13 @@ def handle_content_policy_refusal( ) -> RefusalVerdict: """HTTP-200 refusal (``finish_reason`` ``content_filter`` / ``guardrail_intervened``). Deterministic for the unchanged prompt — never retried: one configured-fallback try, - else surface the refusal (explanation may live only in the reasoning channel). The - caller stops its spinner reference; this stops the spinner object.""" + else surface the refusal (explanation may live only in the reasoning channel).""" from agent.conversation_loop import ( _CONTENT_POLICY_RECOVERY_HINT, _arm_fallback_restart, _content_policy_blocked_result ) - def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> RefusalVerdict: - return RefusalVerdict(action=action, result=result, active_system_prompt=active_system_prompt) - - _refusal_transport = agent._get_transport() - if agent.api_mode == "anthropic_messages": - _refusal_result = _refusal_transport.normalize_response( - response, strip_tool_prefix=agent._is_anthropic_oauth - ) - else: - _refusal_result = _refusal_transport.normalize_response(response) + _refusal_result = normalize_response_for_agent(agent, response) _refusal_text = (getattr(_refusal_result, "content", None) or "").strip() - # Some refusals carry the explanation only in the reasoning - # channel; fall back to it so the user sees *something*. if not _refusal_text: _refusal_text = (agent._extract_reasoning(_refusal_result) or "").strip() @@ -640,41 +515,25 @@ def handle_content_policy_refusal( status_code=None, retry_count=retry_count, max_retries=max_retries, retryable=False, reason=FailoverReason.content_policy_blocked.value, ) + stop_thinking_spinner(agent, thinking_spinner) - if thinking_spinner: - thinking_spinner.stop("") - if agent.thinking_callback: - agent.thinking_callback("") - - # Deterministic for the unchanged prompt — never retry. Try a - # configured fallback once; otherwise surface the refusal. if agent._has_pending_fallback(): - agent._buffer_status( - "⚠️ Model declined to respond (safety refusal) — trying fallback..." - ) + agent._buffer_status("⚠️ Model declined to respond (safety refusal) — trying fallback...") if agent._try_activate_fallback(): - active_system_prompt = _arm_fallback_restart( - agent, api_messages, active_system_prompt, _retry) - return _verdict("break") + active_system_prompt = _arm_fallback_restart(agent, api_messages, active_system_prompt, _retry) + return RefusalVerdict("break", None, active_system_prompt) agent._flush_status_buffer() - _refusal_log = ( - _refusal_text[:500] + "..." if len(_refusal_text) > 500 else _refusal_text - ) + _refusal_log = _refusal_text[:500] + "..." if len(_refusal_text) > 500 else _refusal_text logger.warning( "%sModel declined to respond (finish_reason=content_filter). " "model=%s provider=%s refusal=%s", agent.log_prefix, agent.model, agent.provider, _refusal_log or "(no text)", ) - agent._emit_status( - "⚠️ The model declined to respond to this request (safety refusal)." - ) - + agent._emit_status("⚠️ The model declined to respond to this request (safety refusal).") _refusal_detail = ( - f"Model's explanation: {_refusal_text}" - if _refusal_text - else "The model returned no explanation." + f"Model's explanation: {_refusal_text}" if _refusal_text else "The model returned no explanation." ) _refusal_response = ( "⚠️ The model declined to respond to this request " @@ -682,10 +541,9 @@ def handle_content_policy_refusal( f"{_refusal_detail}\n\n" f"{_CONTENT_POLICY_RECOVERY_HINT}" ) - agent._cleanup_task_resources(effective_task_id) agent._persist_session(messages, conversation_history) - return _verdict("return", _content_policy_blocked_result( + return RefusalVerdict("return", _content_policy_blocked_result( messages, api_call_count, final_response=_refusal_response, error_detail=_refusal_text or "model declined (content_filter)", - )) + ), active_system_prompt)