refactor(agent/turn): lift hook-firing, MoA context/prepare, empty-retry and terminal-empty phases into helpers
This commit is contained in:
+88
-93
@@ -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,
|
||||
)
|
||||
|
||||
+153
-194
@@ -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
|
||||
``<think>`` 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 ``<think>`` 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'<think>|<thinking>|<reasoning>', 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 <think> in content, not reasoning_content, so
|
||||
# _has_structured misses it; detect here to route to prefill.
|
||||
_has_inline_thinking = bool(
|
||||
re.search( r'<think>|<thinking>|<reasoning>', 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 <think> 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")
|
||||
|
||||
+90
-104
@@ -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,
|
||||
)
|
||||
|
||||
+100
-123
@@ -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'</?(?:REASONING_SCRATCHPAD|think|reasoning)>')
|
||||
|
||||
|
||||
@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'</?(?:REASONING_SCRATCHPAD|think|reasoning)>', '', _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 <REASONING_SCRATCHPAD> (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 <REASONING_SCRATCHPAD> (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 <REASONING_SCRATCHPAD> 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")
|
||||
|
||||
Reference in New Issue
Block a user