diff --git a/agent/turn_api_request.py b/agent/turn_api_request.py index 6aca189616..fbbb321a85 100644 --- a/agent/turn_api_request.py +++ b/agent/turn_api_request.py @@ -2,15 +2,15 @@ reasoning echo pad and prompt-cache decoration for the CURRENT provider (a fallback may differ from the primary), build ``api_kwargs``, run the surrogate/ASCII chokepoints, Codex preflight, OpenRouter cache bypass, Copilot ``x-initiator``, the LLM request middleware, -the ``pre_api_request`` hook and the debug dump. Extracted from ``run_conversation``; -nothing here imports ``agent.conversation_loop`` at module level. +the ``pre_api_request`` hook and the debug dump. Nothing here imports +``agent.conversation_loop`` at module level. """ from __future__ import annotations from dataclasses import dataclass import logging -from typing import Any, Dict, Optional +from typing import Any from agent.message_sanitization import ( _sanitize_structure_non_ascii, _sanitize_structure_surrogates @@ -34,6 +34,63 @@ class ApiRequestBuild: _llm_middleware_trace: Any +def _set_extra_header(api_kwargs: Any, key: str, value: str) -> None: + """Copy-on-write header set (the dict may be shared with the transport's defaults).""" + _xh = dict(api_kwargs.get("extra_headers") or {}) + _xh[key] = value + api_kwargs["extra_headers"] = _xh + + +def _fire_pre_api_request_hook( + agent: Any, api_kwargs: Any, api_messages: Any, _llm_middleware_trace: Any, *, messages: Any, + original_user_message: Any, approx_tokens: Any, total_chars: Any, retry_count: Any, + api_call_count: Any, api_request_id: Any, api_start_time: Any, effective_task_id: Any, + turn_id: Any, +) -> None: + from agent.conversation_loop import _system_prompt_for_hooks + + try: + from hermes_cli.lifecycle import has_hook, invoke_hook as _invoke_hook + if has_hook("pre_api_request"): + request_messages = api_kwargs.get("messages") + if not isinstance(request_messages, list): + request_messages = api_kwargs.get("input") + if not isinstance(request_messages, list): + request_messages = api_messages + # Shallow copies: plugins may retain the lists; deepcopy is costly. + # ``request_messages``/``conversation_history`` are raw langfuse passthroughs. + # Anthropic (``system``) and Responses/Codex (``instructions``) move the system + # prompt out of messages; pass it for observability. + _invoke_hook( + "pre_api_request", + task_id=effective_task_id, + turn_id=turn_id, + api_request_id=api_request_id, + session_id=agent.session_id or "", + user_message=original_user_message, + conversation_history=list(messages), + platform=agent.platform or "", + model=agent.model, + provider=agent.provider, + base_url=agent.base_url, + api_mode=agent.api_mode, + api_call_count=api_call_count, + retry_count=retry_count, + request_messages=list(request_messages) if isinstance(request_messages, list) else [], + system_prompt=_system_prompt_for_hooks(api_kwargs, request_messages), + message_count=len(api_messages), + tool_count=len(agent.tools or []), + approx_input_tokens=approx_tokens, + request_char_count=total_chars, + max_tokens=agent.max_tokens, + started_at=api_start_time, + middleware_trace=list(_llm_middleware_trace), + request=agent._api_request_payload_for_hook(api_kwargs), + ) + except Exception: + pass + + def build_api_request( agent: Any, *, api_messages: Any, _moa_prepared_request: Any, tools_for_api: Any, system_message: Any, messages: Any, original_user_message: Any, approx_tokens: Any, @@ -44,29 +101,15 @@ def build_api_request( middleware/hooks/debug dumps observe the payload).""" from agent.conversation_loop import ( _moa_client_consumes_prepared_request, _redecorate_prompt_cache_for_provider, - _system_prompt_for_hooks, ) - api_kwargs = None - _original_api_kwargs = None - _llm_middleware_trace = [] - - def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> ApiRequestBuild: - return ApiRequestBuild( - action=action, api_messages=api_messages, _moa_prepared_request=_moa_prepared_request, - tools_for_api=tools_for_api, api_kwargs=api_kwargs, - _original_api_kwargs=_original_api_kwargs, _llm_middleware_trace=_llm_middleware_trace, - ) agent._reset_stream_delivery_tracking() - # Per-attempt first-chunk timestamp so a stale value never leaks into - # post_api_request. + # Per-attempt first-chunk timestamp so a stale value never leaks into post_api_request. agent._last_api_first_chunk_at = None - # api_messages was built for the primary; a fallback (DeepSeek / Kimi / - # MiMo) may require reasoning_content. Re-apply the echo-back pad - # (idempotent). + # api_messages was built for the primary; a fallback (DeepSeek / Kimi / MiMo) may + # require reasoning_content — re-apply the echo-back pad (idempotent) and re-render + # the prompt-cache decoration for the current provider. agent._reapply_reasoning_echo_for_provider(api_messages) - # Same for prompt-cache decoration (#72626): strip the primary's - # breakpoints and re-render for the current provider. api_messages, _moa_prepared_request, tools_for_api = ( _redecorate_prompt_cache_for_provider( agent, api_messages, system_message=system_message, moa_prepared=_moa_prepared_request, @@ -76,12 +119,9 @@ def build_api_request( if tools_for_api == agent.tools: api_kwargs = agent._build_api_kwargs(api_messages) else: - api_kwargs = agent._build_api_kwargs( - api_messages, tools_for_api=tools_for_api - ) - # Surrogate chokepoint (#50959): tool descriptions, extra_body and - # kwargs strings can carry invalid code points (HTTP 400). One walk - # makes the payload json.dumps()-safe. + api_kwargs = agent._build_api_kwargs(api_messages, tools_for_api=tools_for_api) + # Surrogate chokepoint: tool descriptions, extra_body and kwargs strings can carry + # invalid code points (HTTP 400). One walk makes the payload json.dumps()-safe. _sanitize_structure_surrogates(api_kwargs) if agent._force_ascii_payload: _sanitize_structure_non_ascii(api_kwargs) @@ -90,18 +130,14 @@ def build_api_request( api_kwargs, allow_stream=False, is_github_responses=agent._is_copilot_url(), sanitize_harmony_tokens=agent._is_codex_backend(), ) - # OpenRouter caching replays identical responses, even empty ones; an - # empty-response retry must bypass the cache. + # OpenRouter caching replays identical responses, even empty ones; an empty-response + # retry must bypass the cache. if agent._empty_content_retries > 0 and agent._is_openrouter_url(): - _xh = dict(api_kwargs.get("extra_headers") or {}) - _xh["X-OpenRouter-Cache"] = "false" - api_kwargs["extra_headers"] = _xh - # Copilot x-initiator: first call of a user turn is "user" (billed - # premium); tool-loop follow-ups keep the default "agent" (#3040). + _set_extra_header(api_kwargs, "X-OpenRouter-Cache", "false") + # Copilot x-initiator: first call of a user turn is "user" (billed premium); + # tool-loop follow-ups keep the default "agent". if getattr(agent, "_is_user_initiated_turn", False) and agent._is_copilot_url(): - _xh = dict(api_kwargs.get("extra_headers") or {}) - _xh["x-initiator"] = "user" - api_kwargs["extra_headers"] = _xh + _set_extra_header(api_kwargs, "x-initiator", "user") agent._is_user_initiated_turn = False try: from hermes_cli.middleware import apply_llm_request_middleware @@ -119,66 +155,22 @@ def build_api_request( _original_api_kwargs = dict(api_kwargs) _llm_middleware_trace = [] - try: - from hermes_cli.lifecycle import ( - has_hook, invoke_hook as _invoke_hook - ) - if has_hook("pre_api_request"): - request_messages = api_kwargs.get("messages") - if not isinstance(request_messages, list): - request_messages = api_kwargs.get("input") - if not isinstance(request_messages, list): - request_messages = api_messages - # Shallow copy: plugins may retain the list; deepcopy is costly. - # ``request_messages``/``conversation_history`` are raw langfuse - # passthroughs. - _request_payload = agent._api_request_payload_for_hook(api_kwargs) - # Anthropic (``system``) and Responses/Codex (``instructions``) - # move the system prompt out of messages; pass it for - # observability. - system_prompt_for_hooks = _system_prompt_for_hooks( - api_kwargs, request_messages - ) - _invoke_hook( - "pre_api_request", - task_id=effective_task_id, - turn_id=turn_id, - api_request_id=api_request_id, - session_id=agent.session_id or "", - user_message=original_user_message, - conversation_history=list(messages), - platform=agent.platform or "", - model=agent.model, - provider=agent.provider, - base_url=agent.base_url, - api_mode=agent.api_mode, - api_call_count=api_call_count, - retry_count=retry_count, - request_messages=list(request_messages) - if isinstance(request_messages, list) - else [], - system_prompt=system_prompt_for_hooks, - message_count=len(api_messages), - tool_count=len(agent.tools or []), - approx_input_tokens=approx_tokens, - request_char_count=total_chars, - max_tokens=agent.max_tokens, - started_at=api_start_time, - middleware_trace=list(_llm_middleware_trace), - request=_request_payload, - ) - except Exception: - pass + _fire_pre_api_request_hook( + agent, api_kwargs, api_messages, _llm_middleware_trace, messages=messages, + original_user_message=original_user_message, approx_tokens=approx_tokens, + total_chars=total_chars, retry_count=retry_count, api_call_count=api_call_count, + api_request_id=api_request_id, api_start_time=api_start_time, + effective_task_id=effective_task_id, turn_id=turn_id, + ) if env_var_enabled("HERMES_DUMP_REQUESTS"): agent._dump_api_request_debug(api_kwargs, reason="preflight") - # Private to the in-process MoA facade; add after middleware/hooks/debug - # dumps so none serializes it into the provider payload. + # Private to the in-process MoA facade; added after middleware/hooks/debug dumps so + # none serializes it into the provider payload. Re-read the live client: + # rotation/fallback/cleanup rebuild agent.client between attempts; a native OpenAI + # client rejects this key (TypeError). if _moa_prepared_request is not None and agent.provider == "moa": - # Re-read the live client: rotation/fallback/cleanup rebuild - # agent.client between attempts; a native OpenAI client rejects this - # key (TypeError). if _moa_client_consumes_prepared_request(agent.client): api_kwargs["_moa_prepared_request"] = _moa_prepared_request else: @@ -187,4 +179,7 @@ def build_api_request( "prepared prompt without the MoA handshake", type(agent.client).__name__, ) - return _verdict("fallthrough") + return ApiRequestBuild( + "fallthrough", api_messages, _moa_prepared_request, tools_for_api, api_kwargs, + _original_api_kwargs, _llm_middleware_trace, + ) diff --git a/agent/turn_empty_response.py b/agent/turn_empty_response.py index 75a8ab063f..7f4543f739 100644 --- a/agent/turn_empty_response.py +++ b/agent/turn_empty_response.py @@ -1,11 +1,10 @@ """Empty / thinking-only final-response recovery ladder for the conversation turn loop. -Extracted from ``run_conversation``. Runs when the model returned no visible text after -```` blocks. Ladder order is load-bearing: partial-stream recovery → reuse prior -turn content (housekeeping tools only) → one post-tool-call nudge (#9400) → thinking-only -prefill continuation (×2) → empty-response retries (budgeted, deterministic-empty -short-circuit) → fallback provider → terminal ``(empty)`` sentinel. Nothing here imports -``agent.conversation_loop`` at module level (cycle); loop-internal helpers resolve lazily. +Runs when the model returned no visible text after ```` blocks. Ladder order is +load-bearing: partial-stream recovery → reuse prior turn content (housekeeping tools only) +→ one post-tool-call nudge → thinking-only prefill continuation (×2) → empty-response +retries (budgeted, deterministic-empty short-circuit) → fallback provider → terminal +``(empty)`` sentinel. Nothing here imports ``agent.conversation_loop`` at module level. """ from __future__ import annotations @@ -21,6 +20,8 @@ from agent.turn_recovery import interruptible_backoff_sleep logger = logging.getLogger("agent.conversation_loop") +_INLINE_THINK_RE = re.compile(r'||', re.IGNORECASE) + @dataclass class EmptyResponseVerdict: @@ -40,6 +41,112 @@ class EmptyResponseVerdict: preflight_compression_blocked: bool +def _retry_empty( + agent: Any, response: Any, finish_reason: str, empty_candidate: bool, *, messages: Any, + conversation_history: Any, api_call_count: int, +) -> tuple: + """Budgeted empty-response retry. Each empty attempt re-bills the full input, so the + signature is recorded and deterministic empties stop burning paid retries (fails + open: missing usage or any output keeps the budget). Returns + ``(action_or_None, interrupt_result, deterministic_empty)``.""" + from agent.conversation_loop import jittered_backoff + + if empty_candidate: + _empty_guard.record_empty_attempt(agent, finish_reason=finish_reason, response=response) + budget = ( + _empty_guard.empty_retry_budget(agent, response) + if empty_candidate else _empty_guard.DEFAULT_EMPTY_RETRY_BUDGET + ) + deterministic = empty_candidate and _empty_guard.deterministic_empty(agent) + if not (empty_candidate and agent._empty_content_retries < budget and not deterministic): + return None, None, deterministic + agent._empty_content_retries += 1 + n = agent._empty_content_retries + wait_time = jittered_backoff(n, base_delay=5.0, max_delay=60.0) + logger.warning( + "Empty response (no content or reasoning) — " + "retry %d/%d in %.1fs (model=%s)", + n, budget, wait_time, agent.model, + ) + _budget_note = ( + " — high-cost request, reduced retry budget" + if budget < _empty_guard.DEFAULT_EMPTY_RETRY_BUDGET else "" + ) + agent._buffer_status( + f"⚠️ Empty response from model — retrying " + f"({n}/{budget}) " + f"in {wait_time:.0f}s{_budget_note}" + ) + _interrupted = interruptible_backoff_sleep( + agent, wait_time, None, + messages=messages, + conversation_history=conversation_history, + api_call_count=api_call_count, + abort_message="Interrupt detected during empty-response retry wait, aborting.", + interrupt_text=( + f"Operation interrupted: retrying empty response from model " + f"(retry {n}/{budget})." + ), + activity_label=f"empty response retry backoff ({n}/{budget})", + ) + if _interrupted is not None: + return "return", _interrupted, deterministic + return "continue", None, deterministic + + +def _terminal_empty(agent: Any, assistant_message: Any, finish_reason: str, messages: Any) -> str: + """Retries and fallback exhausted: persist the ``(empty)`` sentinel row and return the + delivery text. Reasoning is surfaced ONLY here, for delivery — the persisted row keeps + the sentinel so later "continue" turns don't replay it and loop on empties.""" + _streak_cost = _empty_guard.streak_cost_usd(agent) + if _streak_cost is not None: + agent._buffer_status( + f"ℹ️ Estimated cost of these empty attempts: " + f"~${_streak_cost:.2f} (input tokens are billed " + f"per attempt even when no answer is produced)" + ) + agent._flush_status_buffer() + reasoning_text = agent._extract_reasoning(assistant_message) + agent._drop_trailing_empty_response_scaffolding(messages) + assistant_msg = agent._build_assistant_message(assistant_message, finish_reason) + assistant_msg["content"] = "(empty)" + assistant_msg["_empty_terminal_sentinel"] = True + append_message(messages, assistant_msg) + + if not reasoning_text: + logger.warning( + "Empty response (no content or reasoning) " + "after %d retries. No fallback available. " + "model=%s provider=%s", + agent._empty_content_retries, agent.model, + agent.provider, + ) + agent._emit_status( + "❌ Model returned no content after all retries" + + (" and fallback attempts." if agent._fallback_chain else + ". No fallback providers configured.") + ) + return "(empty)" + + reasoning_preview = reasoning_text[:500] + "..." if len(reasoning_text) > 500 else reasoning_text + logger.warning( + "Reasoning-only response (no visible content) " + "after exhausting retries and fallback. " + "Reasoning: %s", reasoning_preview, + ) + agent._emit_status( + "⚠️ Model produced reasoning but no visible " + "response after all retries. Returning empty." + ) + return ( + "⚠️ The model produced only internal reasoning and " + "no final answer, despite retries" + + (" and fallback" if agent._fallback_chain else "") + + ". Its last reasoning, which may contain the " + "answer:\n\n" + reasoning_preview + ) + + def recover_empty_response( agent: Any, assistant_message: Any, response: Any, finish_reason: str, *, final_response: Any, messages: List[Dict[str, Any]], api_messages: Any, conversation_history: Any, @@ -48,12 +155,8 @@ def recover_empty_response( ) -> EmptyResponseVerdict: """Recover from a final response with no visible content (see module docstring for the ladder). Role alternation is preserved: the post-tool nudge appends the empty - assistant row BEFORE the user-level hint (APIs reject tool→user). Reasoning is - surfaced only at the terminal step, for delivery — the persisted row keeps the - ``(empty)`` sentinel.""" - from agent.conversation_loop import ( - _EMPTY_TOOL_RESPONSE_NUDGE, _sync_failover_system_message, jittered_backoff - ) + assistant row BEFORE the user-level hint (APIs reject tool→user).""" + from agent.conversation_loop import _EMPTY_TOOL_RESPONSE_NUDGE, _sync_failover_system_message _turn_exit_reason = turn_exit_reason _preflight_compression_blocked = preflight_compression_blocked @@ -65,11 +168,9 @@ def recover_empty_response( preflight_compression_blocked=_preflight_compression_blocked, ) - # Partial stream recovery: content streamed before the connection - # died becomes the final response instead of fallback or retries. - _partial_streamed = ( - getattr(agent, "_current_streamed_assistant_text", "") or "" - ) + # Partial stream recovery: content streamed before the connection died becomes the + # final response instead of fallback or retries. + _partial_streamed = getattr(agent, "_current_streamed_assistant_text", "") or "" if agent._has_content_after_think_block(_partial_streamed): _turn_exit_reason = "partial_stream_recovery" _recovered = agent._strip_think_blocks(_partial_streamed).strip() @@ -78,19 +179,16 @@ def recover_empty_response( "— using as final response", len(_recovered), ) - agent._emit_status( - "↻ Stream interrupted — using delivered content " "as final response" - ) + agent._emit_status("↻ Stream interrupted — using delivered content " "as final response") final_response = _recovered - # A streamed fragment isn't a confirmed preview: keep - # response_previewed false so gateway fallback delivery can - # send the text plus the abnormal-turn explanation. + # A streamed fragment isn't a confirmed preview: gateway fallback delivery + # sends the text plus the abnormal-turn explanation. agent._response_was_previewed = False return _verdict("break") - # Prior turn had real content + ONLY housekeeping tools: model is - # done, reuse it. With substantive tools it was mid-task narration - # and the empty reply is a choke; let the post-tool nudge handle it. + # Prior turn had real content + ONLY housekeeping tools: model is done, reuse it. + # With substantive tools it was mid-task narration and the empty reply is a choke; + # let the post-tool nudge handle it. fallback = getattr(agent, '_last_content_with_tools', None) if fallback and getattr(agent, '_last_content_tools_all_housekeeping', False): _turn_exit_reason = "fallback_prior_turn_content" @@ -99,42 +197,28 @@ def recover_empty_response( agent._last_content_with_tools = None agent._last_content_tools_all_housekeeping = False agent._empty_content_retries = 0 - # Do NOT modify the assistant message content (injected text - # poisoned history); use the fallback as the response and break. + # Do NOT modify the assistant message content (injected text poisoned history). final_response = agent._strip_think_blocks(fallback).strip() agent._response_was_previewed = True return _verdict("break") - # ── Post-tool-call empty response nudge ─────────── - # Empty after tool results (no prior content, or only mid-task - # narration): nudge once via a user-level hint. (#9400) - _prior_was_tool = any( - m.get("role") == "tool" - for m in messages[-5:] # check recent messages - ) - # Ollama puts in content, not reasoning_content, so - # _has_structured misses it; detect here to route to prefill. - _has_inline_thinking = bool( - re.search( r'||', final_response or "", re.IGNORECASE ) - ) + # Post-tool-call empty (no prior content, or only mid-task narration): nudge once. + _prior_was_tool = any(m.get("role") == "tool" for m in messages[-5:]) + # Ollama puts in content, not reasoning_content, so _has_structured misses + # it; detect here to route to prefill. + _has_inline_thinking = bool(_INLINE_THINK_RE.search(final_response or "")) if ( _prior_was_tool and not getattr(agent, "_post_tool_empty_retried", False) and not _has_inline_thinking # thinking model still working — let prefill handle ): agent._post_tool_empty_retried = True - # Clear stale narration so it doesn't resurface - # on a later empty response after the nudge. + # Clear stale narration so it doesn't resurface on a later empty response. agent._last_content_with_tools = None agent._last_content_tools_all_housekeeping = False - logger.info( - "Empty response after tool calls — nudging model " "to continue processing" - ) - agent._buffer_status( - "⚠️ Model returned empty after tool calls — " "nudging to continue" - ) - # Append the empty assistant first so the sequence stays valid: - # tool → assistant("(empty)") → user (APIs reject tool→user). + logger.info("Empty response after tool calls — nudging model " "to continue processing") + agent._buffer_status("⚠️ Model returned empty after tool calls — " "nudging to continue") + # tool → assistant("(empty)") → user keeps the sequence valid. _nudge_msg = agent._build_assistant_message(assistant_message, finish_reason) _nudge_msg["content"] = "(empty)" _nudge_msg["_empty_recovery_synthetic"] = True @@ -144,9 +228,8 @@ def recover_empty_response( }) return _verdict("continue") - # ── Thinking-only prefill continuation ────────── - # Reasoning but no text: append as-is and continue so the model sees - # its own reasoning and writes text. Covers _has_inline_thinking. + # Thinking-only prefill: append the reasoning as-is and continue so the model sees + # its own reasoning and writes text. _has_structured = bool( getattr(assistant_message, "reasoning", None) or getattr(assistant_message, "reasoning_content", None) @@ -164,81 +247,22 @@ def recover_empty_response( f"↻ Thinking-only response — prefilling to continue " f"({agent._thinking_prefill_retries}/2)" ) - interim_msg = agent._build_assistant_message( - assistant_message, "incomplete" - ) + interim_msg = agent._build_assistant_message(assistant_message, "incomplete") interim_msg["_thinking_prefill"] = True append_message(messages, interim_msg) agent._session_messages = messages return _verdict("continue") - # ── Empty response retry ────────────────────── - # Retry up to 3 times before fallback; covers truly empty replies - # AND reasoning-only replies after prefill exhaustion. - _truly_empty = not agent._strip_think_blocks( - final_response - ).strip() - _prefill_exhausted = ( - _has_structured and agent._thinking_prefill_retries >= 2 + # Empty-response retries: truly empty replies AND reasoning-only replies after + # prefill exhaustion. + _truly_empty = not agent._strip_think_blocks(final_response).strip() + _empty_candidate = _truly_empty and (not _has_structured or agent._thinking_prefill_retries >= 2) + action, interrupt_result, _deterministic_empty = _retry_empty( + agent, response, finish_reason, _empty_candidate, messages=messages, + conversation_history=conversation_history, api_call_count=api_call_count, ) - _empty_candidate = _truly_empty and ( - not _has_structured or _prefill_exhausted - ) - if _empty_candidate: - # Each empty attempt re-bills the full input; record its - # signature so deterministic empties stop burning paid retries. - # Fails open: missing usage or any output keeps the budget. - _empty_guard.record_empty_attempt( - agent, finish_reason=finish_reason, response=response - ) - _empty_retry_budget = ( - _empty_guard.empty_retry_budget(agent, response) - if _empty_candidate - else _empty_guard.DEFAULT_EMPTY_RETRY_BUDGET - ) - _deterministic_empty = _empty_candidate and ( - _empty_guard.deterministic_empty(agent) - ) - if ( - _empty_candidate - and agent._empty_content_retries < _empty_retry_budget - and not _deterministic_empty - ): - agent._empty_content_retries += 1 - wait_time = jittered_backoff( - agent._empty_content_retries, base_delay=5.0, max_delay=60.0 - ) - logger.warning( - "Empty response (no content or reasoning) — " - "retry %d/%d in %.1fs (model=%s)", - agent._empty_content_retries, - _empty_retry_budget, wait_time, agent.model, - ) - _budget_note = ( - " — high-cost request, reduced retry budget" - if _empty_retry_budget < _empty_guard.DEFAULT_EMPTY_RETRY_BUDGET - else "" - ) - agent._buffer_status( - f"⚠️ Empty response from model — retrying " - f"({agent._empty_content_retries}/{_empty_retry_budget}) " - f"in {wait_time:.0f}s{_budget_note}" - ) - _interrupted = interruptible_backoff_sleep( - agent, wait_time, None, - messages=messages, - conversation_history=conversation_history, - api_call_count=api_call_count, - abort_message="Interrupt detected during empty-response retry wait, aborting.", - interrupt_text=( - f"Operation interrupted: retrying empty response from model " - f"(retry {agent._empty_content_retries}/{_empty_retry_budget})." - ), - activity_label=f"empty response retry backoff ({agent._empty_content_retries}/{_empty_retry_budget})", - ) - if _interrupted is not None: - return _verdict("return", _interrupted) - return _verdict("continue") + if action is not None: + return _verdict(action, interrupt_result) if _truly_empty and _deterministic_empty: logger.warning( @@ -254,8 +278,7 @@ def recover_empty_response( "to avoid repeat charges" ) - # ── Exhausted retries — try fallback provider ── - # Before "(empty)", switch to the next provider in the chain. + # Exhausted retries — try the next provider in the chain before "(empty)". if _truly_empty and agent._fallback_chain: logger.warning( "Empty response after %d retries — " @@ -263,85 +286,21 @@ def recover_empty_response( agent._empty_content_retries, agent.model, agent.provider, ) - agent._buffer_status( - "⚠️ Model returning empty responses — " "switching to fallback provider..." - ) + agent._buffer_status("⚠️ Model returning empty responses — " "switching to fallback provider...") if agent._try_activate_fallback(): - active_system_prompt = _sync_failover_system_message( - agent, api_messages, active_system_prompt) + active_system_prompt = _sync_failover_system_message(agent, api_messages, active_system_prompt) agent._empty_content_retries = 0 - agent._buffer_status( - f"↻ Switched to fallback: {agent.model} " f"({agent.provider})" - ) + agent._buffer_status(f"↻ Switched to fallback: {agent.model} " f"({agent.provider})") logger.info( "Fallback activated after empty responses: " "now using %s on %s", agent.model, agent.provider, ) - # OUTER loop: `continue` re-runs preflight against the - # fallback's window; `break` would end the turn without - # calling the fallback. Clear the preflight block. (#84733) + # OUTER loop: `continue` re-runs preflight against the fallback's window; + # `break` would end the turn without calling the fallback. _preflight_compression_blocked = False return _verdict("continue") - # Retries and fallback exhausted — fall through to "(empty)". - # Surface the buffered retry trace and, if known, what the empty - # streak cost (each attempt re-billed the full input). - _streak_cost = _empty_guard.streak_cost_usd(agent) - if _streak_cost is not None: - agent._buffer_status( - f"ℹ️ Estimated cost of these empty attempts: " - f"~${_streak_cost:.2f} (input tokens are billed " - f"per attempt even when no answer is produced)" - ) - agent._flush_status_buffer() _turn_exit_reason = "empty_response_exhausted" - reasoning_text = agent._extract_reasoning(assistant_message) - agent._drop_trailing_empty_response_scaffolding(messages) - assistant_msg = agent._build_assistant_message(assistant_message, finish_reason) - assistant_msg["content"] = "(empty)" - # Gateway failure sentinel, not content: persisting it lets later - # "continue" turns replay assistant("(empty)") and loop on empties. - assistant_msg["_empty_terminal_sentinel"] = True - append_message(messages, assistant_msg) - - if reasoning_text: - reasoning_preview = reasoning_text[:500] + "..." if len(reasoning_text) > 500 else reasoning_text - logger.warning( - "Reasoning-only response (no visible content) " - "after exhausting retries and fallback. " - "Reasoning: %s", reasoning_preview, - ) - agent._emit_status( - "⚠️ Model produced reasoning but no visible " - "response after all retries. Returning empty." - ) - else: - logger.warning( - "Empty response (no content or reasoning) " - "after %d retries. No fallback available. " - "model=%s provider=%s", - agent._empty_content_retries, agent.model, - agent.provider, - ) - agent._emit_status( - "❌ Model returned no content after all retries" - + (" and fallback attempts." if agent._fallback_chain else - ". No fallback providers configured.") - ) - - # Delivery-only: show labeled reasoning instead of bare "(empty)" - # when the model thought but wrote no text. The persisted row keeps - # the sentinel; reasoning is never promoted earlier in the ladder. - if reasoning_text: - final_response = ( - "⚠️ The model produced only internal reasoning and " - "no final answer, despite retries" - + (" and fallback" if agent._fallback_chain else "") - + ". Its last reasoning, which may contain the " - "answer:\n\n" + reasoning_preview - ) - else: - final_response = "(empty)" + final_response = _terminal_empty(agent, assistant_message, finish_reason, messages) return _verdict("break") - return _verdict("fallthrough") diff --git a/agent/turn_request_assembly.py b/agent/turn_request_assembly.py index d5c18d06e3..cb5490e368 100644 --- a/agent/turn_request_assembly.py +++ b/agent/turn_request_assembly.py @@ -2,17 +2,16 @@ from the transcript, append MoA context, inject prefills, run the context-engine selection hook and the send-time sanitizers, canonicalize for bit-perfect cache prefixes, build the request-local prompt-cache plan LAST (after every transcript mutation), prepare the -persistent-MoA request, then measure request pressure. Extracted from -``run_conversation``; nothing here imports ``agent.conversation_loop`` at module level -(cycle) — loop-internal helpers resolve lazily so ``patch("agent.conversation_loop.X")`` -sites keep intercepting. +persistent-MoA request, then measure request pressure. Nothing here imports +``agent.conversation_loop`` at module level (cycle) — loop-internal helpers resolve lazily +so ``patch("agent.conversation_loop.X")`` sites keep intercepting. """ from __future__ import annotations from dataclasses import dataclass import logging -from typing import Any, Dict, Optional +from typing import Any from agent.message_sanitization import _sanitize_messages_surrogates from agent.model_metadata import anchored_context_tokens @@ -38,6 +37,72 @@ class AssembledRequest: total_chars: Any +def _append_moa_context(agent: Any, api_messages: Any, moa_config: Any, original_user_message: Any) -> None: + """Run the MoA reference models and append their aggregated context to the last user + message (as a trailing text part on multimodal turns). Fail-open.""" + try: + from agent.message_content import flatten_message_text as _flatten_mt + from agent.moa_loop import _preset_temperature, aggregate_moa_context + + _moa_context = aggregate_moa_context( + user_prompt=( + original_user_message + if isinstance(original_user_message, str) + # Multimodal content list: extract visible text rather than + # str()-ing parts, which would leak base64 image payloads. + else _flatten_mt(original_user_message) + ), + api_messages=api_messages, + reference_models=moa_config.get("reference_models") or [], + aggregator=moa_config.get("aggregator") or {}, + temperature=_preset_temperature(moa_config, "reference_temperature"), + aggregator_temperature=_preset_temperature(moa_config, "aggregator_temperature"), + reference_max_tokens=moa_config.get("reference_max_tokens"), + # None = no per-preset override; inherit auxiliary.moa_reference.timeout. + reference_timeout=( + float(moa_config["reference_timeout"]) + if moa_config.get("reference_timeout") + else None + ), + degraded_reference_policy=str( + moa_config.get("degraded_reference_policy") or "loud" + ), + agent=agent, + ) + if not _moa_context: + return + for _msg in reversed(api_messages): + if _msg.get("role") == "user": + _base = _msg.get("content", "") + if isinstance(_base, str): + _msg["content"] = _base + "\n\n" + _moa_context + elif isinstance(_base, list): + _msg["content"] = [*_base, {"type": "text", "text": "\n\n" + _moa_context}] + break + except Exception as _moa_exc: + logger.warning("MoA context aggregation failed: %s", _moa_exc) + + +def _prepare_moa_request(agent: Any, api_messages: Any, pending_moa_prepared_request: Any) -> tuple: + """Persistent-MoA request: rebase the pending prepared request onto the new messages + when the client supports it, else prepare a fresh one. Returns + ``(prepared_request, api_messages, pending_moa_prepared_request)``.""" + _moa_completions = getattr(getattr(agent.client, "chat", None), "completions", None) + prepared: Any = None + if pending_moa_prepared_request is not None: + _rebase = getattr(_moa_completions, "rebase_prepared_request", None) + if callable(_rebase): + prepared = _rebase(pending_moa_prepared_request, api_messages) + pending_moa_prepared_request = None + if prepared is None: + _prepare = getattr(_moa_completions, "prepare", None) + if callable(_prepare): + prepared = _prepare(api_messages) + if prepared is not None: + api_messages = prepared["messages"] + return prepared, api_messages, pending_moa_prepared_request + + def assemble_api_request( agent: Any, *, messages: Any, current_turn_user_idx: Any, _ext_prefetch_cache: Any, _plugin_user_context: Any, moa_config: Any, active_system_prompt: Any, @@ -51,19 +116,6 @@ def assemble_api_request( _midturn_request_pressure_tokens, estimate_messages_tokens_rough, ) - def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> AssembledRequest: - return AssembledRequest( - action=action, - api_messages=api_messages, - tools_for_api=tools_for_api, - _moa_prepared_request=_moa_prepared_request, - pending_moa_prepared_request=pending_moa_prepared_request, - approx_tokens=approx_tokens, - request_pressure_tokens=request_pressure_tokens, - total_chars=total_chars, - - ) - api_messages, effective_system = build_api_messages( agent, messages, current_turn_user_idx=current_turn_user_idx, ext_prefetch_cache=_ext_prefetch_cache, plugin_user_context=_plugin_user_context, @@ -71,54 +123,9 @@ def assemble_api_request( ) if moa_config: - try: - from agent.message_content import flatten_message_text as _flatten_mt - from agent.moa_loop import _preset_temperature, aggregate_moa_context + _append_moa_context(agent, api_messages, moa_config, original_user_message) - _moa_context = aggregate_moa_context( - user_prompt=( - original_user_message - if isinstance(original_user_message, str) - # Multimodal content list: extract visible text rather than - # str()-ing parts, which would leak base64 image payloads. - else _flatten_mt(original_user_message) - ), - api_messages=api_messages, - reference_models=moa_config.get("reference_models") or [], - aggregator=moa_config.get("aggregator") or {}, - temperature=_preset_temperature(moa_config, "reference_temperature"), - aggregator_temperature=_preset_temperature(moa_config, "aggregator_temperature"), - reference_max_tokens=moa_config.get("reference_max_tokens"), - # None = no per-preset override; inherit - # auxiliary.moa_reference.timeout via call_llm. - reference_timeout=( - float(moa_config["reference_timeout"]) - if moa_config.get("reference_timeout") - else None - ), - degraded_reference_policy=str( - moa_config.get("degraded_reference_policy") or "loud" - ), - agent=agent, - ) - if _moa_context: - for _msg in reversed(api_messages): - if _msg.get("role") == "user": - _base = _msg.get("content", "") - if isinstance(_base, str): - _msg["content"] = _base + "\n\n" + _moa_context - elif isinstance(_base, list): - # Multimodal turn: append MoA context as a trailing text - # part instead of silently dropping it. - _msg["content"] = [ - *_base, {"type": "text", "text": "\n\n" + _moa_context} - ] - break - except Exception as _moa_exc: - logger.warning("MoA context aggregation failed: %s", _moa_exc) - - # Inject ephemeral prefill messages right after the system prompt - # but before conversation history. Same API-call-time-only pattern. + # Ephemeral prefill messages go right after the system prompt, API-call-time only. if agent.prefill_messages: sys_offset = 1 if (api_messages and api_messages[0].get("role") == "system") else 0 for idx, pfm in enumerate(agent.prefill_messages): @@ -140,17 +147,14 @@ def assemble_api_request( api_messages = agent._sanitize_api_messages(api_messages) # One-time repeated-heal notice goes out via the status/warning callback, NEVER - # appended to messages: the cached prompt prefix stays byte-identical (#96870). + # appended to messages: the cached prompt prefix stays byte-identical. try: - from agent.agent_runtime_helpers import ( - consume_pending_sanitizer_heal_notice, - ) + from agent.agent_runtime_helpers import consume_pending_sanitizer_heal_notice _heal_notice = consume_pending_sanitizer_heal_notice() if _heal_notice: agent._emit_warning(_heal_notice) except Exception: - # A notice hiccup must never break the send path. logger.debug("sanitizer heal notice delivery failed", exc_info=True) # Drop thinking-only assistant turns + merge adjacent users, API copy only: @@ -171,24 +175,21 @@ def assemble_api_request( _sanitize_messages_surrogates(api_messages) # No send-time pad loop here: ``repair_empty_non_final_messages`` (inside - # ``_sanitize_api_messages``) is the single owner of empty-turn repair, and its - # non-whitespace placeholder survives normalization regardless of ordering. + # ``_sanitize_api_messages``) is the single owner of empty-turn repair. # Build the request-local cache sections LAST, after every transcript mutation; # the canonical tool registry stays undecorated. Marked ``content`` becomes text # blocks the whitespace pass skips, so the same row's bytes vary across turns. tools_for_api = agent.tools if agent._use_prompt_caching and agent.provider != "moa": - from agent.prompt_caching import ( - envelope_tool_part_cache_markers_supported, - ) + from agent.prompt_caching import envelope_tool_part_cache_markers_supported _static_system_prefix = getattr(agent, "_cached_system_prompt_static", None) _initial_cache_plan = build_prompt_cache_plan( api_messages, tools_for_api, # Clamp per-destination: a configured 1h regresses to 5m on - # Qwen/Alibaba routes, whose context cache is 5m-only (#84733). + # Qwen/Alibaba routes, whose context cache is 5m-only. cache_ttl=effective_cache_ttl( agent._cache_ttl, provider=agent.provider, model=agent.model ), @@ -198,7 +199,7 @@ def assemble_api_request( ), direct_native_tool_cache=agent._direct_native_anthropic_tool_cache_capability(), # LiteLLM-style envelope routes forward part-level markers into - # tool_result.content[] → non-retryable 400 (#89886). + # tool_result.content[] → non-retryable 400. tool_part_markers=envelope_tool_part_cache_markers_supported( getattr(agent, "provider", ""), getattr(agent, "base_url", "") ), @@ -211,50 +212,35 @@ def assemble_api_request( # prepared request instead of running the advisors again. _moa_prepared_request = None if agent.provider == "moa": - _moa_completions = getattr(getattr(agent.client, "chat", None), "completions", None) - if pending_moa_prepared_request is not None: - _rebase_moa_request = getattr(_moa_completions, "rebase_prepared_request", None) - if callable(_rebase_moa_request): - _moa_prepared_request = _rebase_moa_request( - pending_moa_prepared_request, api_messages - ) - pending_moa_prepared_request = None - if _moa_prepared_request is None: - _prepare_moa_request = getattr(_moa_completions, "prepare", None) - if callable(_prepare_moa_request): - _moa_prepared_request = _prepare_moa_request(api_messages) - if _moa_prepared_request is not None: - api_messages = _moa_prepared_request["messages"] + _moa_prepared_request, api_messages, pending_moa_prepared_request = _prepare_moa_request( + agent, api_messages, pending_moa_prepared_request + ) # One image-stripped estimate feeds both figures; tools counted separately (50+ # tools ≈ 20-30K tokens); total_chars is a rough proxy for logs/hooks only. - # Charge stale thinking only when the active route replays it (#84371). + # Charge stale thinking only when the active route replays it. from agent.turn_context import _agent_stale_thinking_on_wire if _agent_stale_thinking_on_wire(agent): approx_tokens = estimate_messages_tokens_rough(api_messages) else: - approx_tokens = estimate_messages_tokens_rough( - api_messages, charge_stale_thinking=False - ) + approx_tokens = estimate_messages_tokens_rough(api_messages, charge_stale_thinking=False) # Route-aware: native Responses compaction prunes the wire payload, so the raw - # history figure overstates it and fires needless local compression (#96995). + # history figure overstates it and fires needless local compression. request_pressure_tokens = _midturn_request_pressure_tokens( agent, api_messages, effective_system or "", approx_tokens ) # Usage-anchored override: real prompt_tokens (incl. system + tool schemas) + # delta estimate replaces the whole-history heuristic when the anchor is fresh. - _anchored_pressure = anchored_context_tokens( - messages, getattr(agent, "_usage_anchor", None) - ) + _anchored_pressure = anchored_context_tokens(messages, getattr(agent, "_usage_anchor", None)) if _anchored_pressure is not None: request_pressure_tokens = _anchored_pressure - total_chars = approx_tokens * 4 # Stash the rough estimate so update_from_response() can pair it with the real # count (should_defer_preflight_to_real_usage). getattr: test doubles lack it. - _note_rough = getattr( - agent.context_compressor, "note_request_rough_estimate", None - ) + _note_rough = getattr(agent.context_compressor, "note_request_rough_estimate", None) if callable(_note_rough): _note_rough(request_pressure_tokens) - return _verdict("fallthrough") + return AssembledRequest( + "fallthrough", api_messages, tools_for_api, _moa_prepared_request, + pending_moa_prepared_request, approx_tokens, request_pressure_tokens, approx_tokens * 4, + ) diff --git a/agent/turn_response_intake.py b/agent/turn_response_intake.py index 8ea50d6fce..8c4b5949ea 100644 --- a/agent/turn_response_intake.py +++ b/agent/turn_response_intake.py @@ -1,8 +1,7 @@ """Response intake for the conversation turn loop: normalize the raw provider response into the assistant message, splice agent-as-provider projections, fire ``post_api_request``, relay reasoning to the progress callback, and apply the incomplete-scratchpad / Codex-incomplete -continuation guards. Extracted from ``run_conversation``; nothing here imports -``agent.conversation_loop`` at module level (cycle). +continuation guards. Nothing here imports ``agent.conversation_loop`` at module level (cycle). """ from __future__ import annotations @@ -15,10 +14,12 @@ from typing import Any, Dict, Optional from agent.provider_projection import splice_provider_projection from agent.trajectory import has_incomplete_scratchpad -from agent.turn_truncation import continue_codex_incomplete +from agent.turn_truncation import continue_codex_incomplete, normalize_response_for_agent, partial_result logger = logging.getLogger("agent.conversation_loop") +_REASONING_TAG_RE = re.compile(r'') + @dataclass class ResponseIntakeVerdict: @@ -32,68 +33,34 @@ class ResponseIntakeVerdict: result: Optional[Dict[str, Any]] = None -def normalize_model_response( - agent: Any, *, response: Any, messages: Any, api_messages: Any, conversation_history: Any, +def _coerce_content_text(raw: Any) -> str: + """Some OpenAI-compatible servers (llama-server) return content as dict/list, which + crashes downstream ``.strip()``; normalize to str (multimodal lists → text parts).""" + if isinstance(raw, dict): + return raw.get("text", "") or raw.get("content", "") or json.dumps(raw) + if isinstance(raw, list): + parts = [] + for part in raw: + if isinstance(part, str): + parts.append(part) + elif isinstance(part, dict) and part.get("type") == "text": + parts.append(part.get("text", "")) + elif isinstance(part, dict) and "text" in part: + parts.append(str(part["text"])) + return "\n".join(parts) + return str(raw) + + +def _fire_post_api_request_hook( + agent: Any, response: Any, assistant_message: Any, finish_reason: Any, *, api_messages: Any, api_call_count: Any, api_duration: Any, api_start_time: Any, api_request_id: Any, effective_task_id: Any, turn_id: Any, -) -> ResponseIntakeVerdict: - """Normalize ``response`` into ``assistant_message`` (str content, never dict/list) and run - the post-response hooks and continuation guards, in the original order.""" - from agent.conversation_loop import ( - _moa_reference_metrics_for_hook, - ) - assistant_message = None - finish_reason = None - - def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> ResponseIntakeVerdict: - return ResponseIntakeVerdict( - action=action, assistant_message=assistant_message, finish_reason=finish_reason, - result=result, - ) - - _transport = agent._get_transport() - _normalize_kwargs = {} - if agent.api_mode == "anthropic_messages": - _normalize_kwargs["strip_tool_prefix"] = agent._is_anthropic_oauth - normalized = _transport.normalize_response(response, **_normalize_kwargs) - assistant_message = normalized - finish_reason = normalized.finish_reason - - # Some OpenAI-compatible servers (llama-server) return content as dict/list, - # which crashes downstream .strip(); normalize to str. - if assistant_message.content is not None and not isinstance(assistant_message.content, str): - raw = assistant_message.content - if isinstance(raw, dict): - assistant_message.content = raw.get("text", "") or raw.get("content", "") or json.dumps(raw) - elif isinstance(raw, list): - # Multimodal content list — extract text parts - parts = [] - for part in raw: - if isinstance(part, str): - parts.append(part) - elif isinstance(part, dict) and part.get("type") == "text": - parts.append(part.get("text", "")) - elif isinstance(part, dict) and "text" in part: - parts.append(str(part["text"])) - assistant_message.content = "\n".join(parts) - else: - assistant_message.content = str(raw) - - # ── Agent-as-provider projection ────────────────────────────── - # Splice the provider-agent's own tool work in as call/result rows before - # this turn's assistant message; no-op for ordinary providers. - splice_provider_projection(agent, response, messages) +) -> None: + from agent.conversation_loop import _moa_reference_metrics_for_hook try: - from hermes_cli.lifecycle import ( - has_hook, invoke_hook as _invoke_hook - ) + from hermes_cli.lifecycle import has_hook, invoke_hook as _invoke_hook if has_hook("post_api_request"): - _assistant_tool_calls = ( - getattr(assistant_message, "tool_calls", None) or [] - ) - _assistant_text = assistant_message.content or "" - _api_ended_at = api_start_time + api_duration _invoke_hook( "post_api_request", task_id=effective_task_id, @@ -108,13 +75,10 @@ def normalize_model_response( api_call_count=api_call_count, api_duration=api_duration, started_at=api_start_time, - ended_at=_api_ended_at, - # First stream chunk time (epoch s) from - # interruptible_streaming_api_call; None if not streamed / no - # chunk. TTFB = first_chunk_at - started_at. - first_chunk_at=getattr( - agent, "_last_api_first_chunk_at", None - ), + ended_at=api_start_time + api_duration, + # First stream chunk time (epoch s); None if not streamed / no chunk. + # TTFB = first_chunk_at - started_at. + first_chunk_at=getattr(agent, "_last_api_first_chunk_at", None), finish_reason=finish_reason, message_count=len(api_messages), response_model=getattr(response, "model", None), @@ -123,73 +87,86 @@ def normalize_model_response( ), usage=agent._usage_summary_for_api_request_hook(response), assistant_message=assistant_message, - assistant_content_chars=len(_assistant_text), - assistant_tool_call_count=len(_assistant_tool_calls), + assistant_content_chars=len(assistant_message.content or ""), + assistant_tool_call_count=len(getattr(assistant_message, "tool_calls", None) or []), moa_references=_moa_reference_metrics_for_hook(agent), ) except Exception: pass - # Handle assistant response - if assistant_message.content and not agent.quiet_mode: + +def _relay_thinking(agent: Any, content: str) -> None: + """Relay the model's text to the progress callback: subagents send the first line to + the parent display; any agent with a structured callback gets ``reasoning.available``.""" + _think_text = _REASONING_TAG_RE.sub('', content.strip()).strip() + first_line = _think_text.split('\n')[0][:80] if _think_text else "" + if first_line and getattr(agent, '_delegate_depth', 0) > 0: + try: + agent.tool_progress_callback("_thinking", first_line) + except Exception: + pass + elif _think_text: + try: + agent.tool_progress_callback("reasoning.available", "_thinking", _think_text[:500], None) + except Exception: + pass + + +def normalize_model_response( + agent: Any, *, response: Any, messages: Any, api_messages: Any, conversation_history: Any, + api_call_count: Any, api_duration: Any, api_start_time: Any, api_request_id: Any, + effective_task_id: Any, turn_id: Any, +) -> ResponseIntakeVerdict: + """Normalize ``response`` into ``assistant_message`` (str content, never dict/list) and run + the post-response hooks and continuation guards, in the original order.""" + assistant_message = normalize_response_for_agent(agent, response) + finish_reason = assistant_message.finish_reason + + def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> ResponseIntakeVerdict: + return ResponseIntakeVerdict( + action=action, assistant_message=assistant_message, finish_reason=finish_reason, + result=result, + ) + + if assistant_message.content is not None and not isinstance(assistant_message.content, str): + assistant_message.content = _coerce_content_text(assistant_message.content) + + # Agent-as-provider projection: splice the provider-agent's own tool work in as + # call/result rows before this turn's assistant message; no-op for ordinary providers. + splice_provider_projection(agent, response, messages) + + _fire_post_api_request_hook( + agent, response, assistant_message, finish_reason, api_messages=api_messages, + api_call_count=api_call_count, api_duration=api_duration, api_start_time=api_start_time, + api_request_id=api_request_id, effective_task_id=effective_task_id, turn_id=turn_id, + ) + + content = assistant_message.content + if content and not agent.quiet_mode: if agent.verbose_logging: - agent._vprint(f"{agent.log_prefix}🤖 Assistant: {assistant_message.content}") + agent._vprint(f"{agent.log_prefix}🤖 Assistant: {content}") else: - agent._vprint(f"{agent.log_prefix}🤖 Assistant: {assistant_message.content[:100]}{'...' if len(assistant_message.content) > 100 else ''}") + agent._vprint(f"{agent.log_prefix}🤖 Assistant: {content[:100]}{'...' if len(content) > 100 else ''}") + if content and agent.tool_progress_callback: + _relay_thinking(agent, content) - # Notify progress callback of model's thinking (used by subagent - # delegation to relay the child's reasoning to the parent display). - if (assistant_message.content and agent.tool_progress_callback): - _think_text = assistant_message.content.strip() - # Strip reasoning XML tags that shouldn't leak to parent display - _think_text = re.sub( - r'', '', _think_text - ).strip() - # For subagents: relay first line to parent display (existing behaviour). - # For all agents with a structured callback: emit reasoning.available event. - first_line = _think_text.split('\n')[0][:80] if _think_text else "" - if first_line and getattr(agent, '_delegate_depth', 0) > 0: - try: - agent.tool_progress_callback("_thinking", first_line) - except Exception: - pass - elif _think_text: - try: - agent.tool_progress_callback("reasoning.available", "_thinking", _think_text[:500], None) - except Exception: - pass - - # Check for incomplete (opened but never closed) - # This means the model ran out of output tokens mid-reasoning — retry up to 2 times - if has_incomplete_scratchpad(assistant_message.content or ""): + # Incomplete (opened, never closed): the model ran out of + # output tokens mid-reasoning — retry up to 2 times, then save as partial. + if has_incomplete_scratchpad(content or ""): agent._incomplete_scratchpad_retries += 1 - agent._buffer_vprint("⚠️ Incomplete detected (opened but never closed)") - if agent._incomplete_scratchpad_retries <= 2: agent._buffer_vprint(f"🔄 Retrying API call ({agent._incomplete_scratchpad_retries}/2)...") - # Don't add the broken message, just retry - return _verdict("continue") - else: - # Max retries - discard this turn and save as partial - agent._flush_status_buffer() - agent._vprint(f"{agent.log_prefix}❌ Max retries (2) for incomplete scratchpad. Saving as partial.", force=True) - agent._incomplete_scratchpad_retries = 0 - - rolled_back_messages = agent._get_messages_up_to_last_assistant(messages) - agent._cleanup_task_resources(effective_task_id) - agent._persist_session(messages, conversation_history) - - return _verdict("return", { - "final_response": "Incomplete REASONING_SCRATCHPAD after 2 retries", - "messages": rolled_back_messages, - "api_calls": api_call_count, - "completed": False, - "partial": True, - "error": "Incomplete REASONING_SCRATCHPAD after 2 retries" - }) - - # Reset incomplete scratchpad counter on clean response + return _verdict("continue") # don't add the broken message + agent._flush_status_buffer() + agent._vprint(f"{agent.log_prefix}❌ Max retries (2) for incomplete scratchpad. Saving as partial.", force=True) + agent._incomplete_scratchpad_retries = 0 + rolled_back_messages = agent._get_messages_up_to_last_assistant(messages) + agent._cleanup_task_resources(effective_task_id) + agent._persist_session(messages, conversation_history) + return _verdict("return", partial_result( + rolled_back_messages, api_call_count, "Incomplete REASONING_SCRATCHPAD after 2 retries" + )) agent._incomplete_scratchpad_retries = 0 if agent.api_mode == "codex_responses" and finish_reason == "incomplete": @@ -200,6 +177,6 @@ def normalize_model_response( if _codex_result is not None: return _verdict("return", _codex_result) return _verdict("continue") - elif hasattr(agent, "_codex_incomplete_retries"): + if hasattr(agent, "_codex_incomplete_retries"): agent._codex_incomplete_retries = 0 return _verdict("fallthrough")