Merge branch 'simp/agent-loop' into simp/integration

This commit is contained in:
Teknium
2026-09-02 14:19:36 -07:00
17 changed files with 6028 additions and 6838 deletions
+896 -6019
View File
File diff suppressed because it is too large Load Diff
+276 -487
View File
File diff suppressed because it is too large Load Diff
+377
View File
@@ -0,0 +1,377 @@
"""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.
"""
from __future__ import annotations
import logging
import re
from dataclasses import dataclass
from typing import Any, Dict, List, Optional
from agent import empty_response_guard as _empty_guard
from agent.message_metadata import append_message
from agent.turn_recovery import interruptible_backoff_sleep
logger = logging.getLogger("agent.conversation_loop")
@dataclass
class EmptyResponseVerdict:
"""Outcome of ``recover_empty_response``.
``action``: ``"break"`` (turn is done — ``final_response`` is set), ``"continue"``
(re-enter the OUTER turn loop: a nudge/prefill row was appended, a retry wait
elapsed, or a fallback was activated and preflight must re-run), ``"return"``
(interrupted during a retry wait — return ``result``) or ``"fallthrough"``
(unreachable: every path exits; kept for the contract)."""
action: str
result: Optional[Dict[str, Any]]
final_response: Any
turn_exit_reason: Any
active_system_prompt: Any
preflight_compression_blocked: bool
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,
active_system_prompt: Any,
api_call_count: int,
turn_exit_reason: Any,
preflight_compression_blocked: bool,
) -> 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,
)
_turn_exit_reason = turn_exit_reason
_preflight_compression_blocked = preflight_compression_blocked
def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> EmptyResponseVerdict:
return EmptyResponseVerdict(
action=action,
result=result,
final_response=final_response,
turn_exit_reason=_turn_exit_reason,
active_system_prompt=active_system_prompt,
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 ""
)
if agent._has_content_after_think_block(_partial_streamed):
_turn_exit_reason = "partial_stream_recovery"
_recovered = agent._strip_think_blocks(_partial_streamed).strip()
logger.info(
"Partial stream content delivered (%d chars) "
"— using as final response",
len(_recovered),
)
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.
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.
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"
logger.info("Empty follow-up after tool calls — using prior turn content as final response")
agent._emit_status("↻ Empty response after tool calls — using earlier content as final answer")
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.
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,
)
)
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.
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).
_nudge_msg = agent._build_assistant_message(assistant_message, finish_reason)
_nudge_msg["content"] = "(empty)"
_nudge_msg["_empty_recovery_synthetic"] = True
append_message(messages, _nudge_msg)
append_message(messages, {
"role": "user",
"content": _EMPTY_TOOL_RESPONSE_NUDGE,
"_empty_recovery_synthetic": True,
})
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.
_has_structured = bool(
getattr(assistant_message, "reasoning", None)
or getattr(assistant_message, "reasoning_content", None)
or getattr(assistant_message, "reasoning_details", None)
or _has_inline_thinking
)
if _has_structured and agent._thinking_prefill_retries < 2:
agent._thinking_prefill_retries += 1
logger.info(
"Thinking-only response (no visible content) — "
"prefilling to continue (%d/2)",
agent._thinking_prefill_retries,
)
agent._buffer_status(
f"↻ Thinking-only response — prefilling to continue "
f"({agent._thinking_prefill_retries}/2)"
)
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_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 _truly_empty and _deterministic_empty:
logger.warning(
"Deterministic empty response detected "
"(consecutive zero-output completions, "
"model=%s provider=%s finish_reason=%s) — "
"skipping remaining retries",
agent.model, agent.provider, finish_reason,
)
agent._buffer_status(
"⚠️ Model is deterministically returning empty "
"(zero output tokens) — skipping further retries "
"to avoid repeat charges"
)
# ── Exhausted retries — try fallback provider ──
# Before "(empty)", switch to the next provider in the chain.
if _truly_empty and agent._fallback_chain:
logger.warning(
"Empty response after %d retries — "
"attempting fallback (model=%s, provider=%s)",
agent._empty_content_retries, agent.model,
agent.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)
agent._empty_content_retries = 0
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)
_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)"
return _verdict("break")
return _verdict("fallthrough")
+101 -273
View File
@@ -1,24 +1,9 @@
"""Post-loop turn finalization for ``run_conversation``.
Extracted from ``agent/conversation_loop.py`` as part of the god-file
decomposition campaign (``~/.hermes/plans/god-file-decomposition.md``, Phase 1
step 4 — the post-loop ``TurnFinalizer`` seam). ``run_conversation``'s tail
(everything after the main tool-calling ``while`` loop) is lifted here verbatim:
budget-exhaustion summary, trajectory save, session persist, turn diagnostics,
response transforms, result-dict assembly, steer drain, and the memory/skill
review trigger.
Behavior-neutral: the body is moved unchanged. All ``agent.*`` side effects fire
exactly as before; only the post-loop *locals* are passed in as keyword args, and
the assembled ``result`` dict is returned to ``run_conversation`` which returns it
to the caller. The function is synchronous with a single return — mirroring the
region it replaces (no awaits, no early returns).
Module ``logger`` is imported lazily inside the body (``from
agent.conversation_loop import logger``) so this module never imports
``agent.conversation_loop`` at import time -> no import cycle, and the log records
keep the exact logger name (``"agent.conversation_loop"``).
"""
Lifted verbatim: budget summary, trajectory save, persist, diagnostics, response
transforms, result assembly, steer drain, memory/skill review. Synchronous, single
return. ``logger`` is imported lazily from ``agent.conversation_loop`` (no cycle,
same logger name)."""
from __future__ import annotations
@@ -54,11 +39,9 @@ def _fill_assistant_tail_content(agent, tail: dict, final_response) -> None:
agent._db_flush_scan_prefix = None
# Verification continuation scaffolding flags: verify-on-stop / pre_verify
# inject a synthetic user nudge to keep the agent going one more turn.
# These nudges must be stripped from returned/live history to avoid
# role-alternation breaks and poisoning the resumed transcript. The
# assistant response is real content and is not flagged. (#65919 §7)
# Verification-continuation nudges (verify-on-stop / pre_verify) must be stripped from
# returned/live history to avoid role-alternation breaks; the assistant response is
# real content and is not flagged. (#65919)
_VERIFICATION_CONTINUATION_FLAGS = (
"_verification_stop_synthetic",
"_pre_verify_synthetic",
@@ -71,14 +54,10 @@ def _record_kanban_budget_exhausted(
max_iterations: int,
logger: logging.Logger,
) -> None:
"""Record a terminal ``timed_out`` outcome for a kanban worker that
exhausted its iteration budget.
"""Record a terminal ``timed_out`` outcome for a kanban worker out of budget.
This is a bounded fallback (#87096): the CAS invariant in ``_end_run``
(``WHERE ended_at IS NULL``) guarantees idempotence — if another path
already closed the run this is a no-op — so it is safe to call from
multiple exit paths.
"""
Idempotent via the ``_end_run`` CAS (``WHERE ended_at IS NULL``): a no-op if
another path already closed the run, so safe from multiple exit paths (#87096)."""
try:
from hermes_cli import kanban_db as _kb
_conn = _kb.connect()
@@ -116,10 +95,8 @@ def _record_kanban_budget_exhausted(
def _drop_verification_continuation_scaffolding(messages) -> None:
"""Remove verification-continuation nudge messages from *messages* in place.
Only the synthetic nudges carry these flags, so this strips just the
nudges while preserving the real attempted-final-answer that was
persisted to state.db.
"""
Only the synthetic nudges carry these flags, so the real attempted-final-answer
persisted to state.db survives."""
messages[:] = [
m for m in messages
if not (isinstance(m, dict) and any(m.get(f) for f in _VERIFICATION_CONTINUATION_FLAGS))
@@ -153,11 +130,7 @@ def finalize_turn(
_pending_verification_response=None,
_pending_verification_response_previewed=False,
):
"""Run the post-loop finalization and return the turn ``result`` dict.
Lifted verbatim from ``run_conversation`` (the region after the main agent
loop). See module docstring.
"""
"""Run the post-loop finalization and return the turn ``result`` dict."""
from agent.conversation_loop import logger
budget_exhausted = (
@@ -179,24 +152,19 @@ def finalize_turn(
iteration_limit_fallback = False
preserved_verification_fallback = False
if continuation_budget_exhausted:
# A verification/continuation gate deliberately withheld a composed
# answer, then consumed the remaining budget before producing a newer
# one. Preserve that exact answer instead of replacing it with another
# fallible model call. The explicit pending value is the provenance
# guard: unrelated error/recovery exits can never enter this branch.
# A verification gate withheld a composed answer, then the budget ran out:
# preserve it rather than make another fallible call. The explicit pending
# value is the provenance guard; unrelated error exits never enter here.
final_response = _pending_verification_response
# Mark the turn as previewed only when the reused candidate was
# actually streamed to the user as interim content. (#65919 review:
# response-loss blocker)
# Previewed only if the reused candidate was actually streamed as interim.
if _pending_verification_response_previewed:
agent._response_was_previewed = True
_turn_exit_reason = f"max_iterations_reached({api_call_count}/{agent.max_iterations})"
iteration_limit_fallback = True
preserved_verification_fallback = True
elif final_response is None and budget_fallback_eligible:
# Budget exhausted — ask the model for a summary via one extra
# API call with tools stripped. _handle_max_iterations injects a
# user message and makes a single toolless request.
# Budget exhausted: _handle_max_iterations makes one extra toolless request
# for a summary.
_turn_exit_reason = f"max_iterations_reached({api_call_count}/{agent.max_iterations})"
agent._emit_status(
f"⚠️ Iteration budget exhausted ({api_call_count}/{agent.max_iterations}) "
@@ -211,29 +179,18 @@ def finalize_turn(
iteration_limit_fallback = True
if iteration_limit_fallback:
# If running as a kanban worker, signal the dispatcher that the
# worker could not complete (rather than treating it as a
# protocol violation). This applies whether the user-facing fallback
# came from the summary call or an explicitly pending continuation;
# both exhausted the task budget and must advance the failure circuit.
#
# We route through ``_record_task_failure(outcome="timed_out")``
# rather than ``kanban_block`` so this counts toward the dispatcher's
# consecutive-failure circuit breaker (#29747 gap 2).
# Kanban worker: signal the dispatcher the worker could not complete. Route
# via ``_record_task_failure(outcome="timed_out")`` (not ``kanban_block``) so
# it counts toward the consecutive-failure circuit breaker (#29747).
_kanban_task = os.environ.get("HERMES_KANBAN_TASK")
if _kanban_task:
_record_kanban_budget_exhausted(
_kanban_task, api_call_count, agent.max_iterations, logger,
)
elif budget_exhausted:
# Bounded fallback (#87096): budget was exhausted but none of the
# normal fallback paths were eligible (interrupted / failed /
# anomalous exit_reason). If running as a kanban worker we must
# still record a terminal outcome so the task does not remain in
# an ambiguous lifecycle state. The worker's run is closed via
# ``_record_task_failure`` (compare-and-swap receipt path) which
# is a no-op if another path closed it — the CAS invariant in
# ``_end_run`` (``WHERE ended_at IS NULL``) guarantees idempotence.
# Bounded fallback: budget exhausted with no eligible fallback path. A kanban
# worker must still record a terminal outcome; the ``_end_run`` CAS makes it
# idempotent if another path already closed the run (#87096).
_kanban_task = os.environ.get("HERMES_KANBAN_TASK")
if _kanban_task:
_record_kanban_budget_exhausted(
@@ -251,18 +208,9 @@ def finalize_turn(
)
)
# Preflight can seed the display count before the provider receives the
# request. Roll that estimate back only when an interrupt wins the race
# before any successful provider response. Compaction state remains owned
# by the real-usage/post-compaction path, including its ``-1`` sentinel.
# Guard rules (test-double density on this path is high):
# - snapshot is type-pinned to a real int — MagicMock agents auto-create
# truthy Mock attributes that must never arm the rollback;
# - the received-response flag is pinned to ``is not True`` — its real
# domain is True/False, and only a literal True means a provider
# response completed;
# - the compressor method gets a getattr+callable guard — SimpleNamespace
# compressor doubles and plugin context engines lack it.
# Roll back the preflight-seeded display count only when an interrupt wins
# before any provider response; compaction state (incl. ``-1``) stays with the
# real-usage path. Type-pinned guards keep MagicMock/SimpleNamespace doubles inert.
_preflight_snapshot = getattr(
agent, "_turn_preflight_display_snapshot", None
)
@@ -281,17 +229,9 @@ def finalize_turn(
if callable(_rollback_fn):
_rollback_fn(_preflight_snapshot)
# Post-loop cleanup must never lose the response. Trajectory save,
# resource teardown, and session persistence all touch fallible
# surfaces — file I/O / JSON serialization (_save_trajectory), remote
# VM/browser teardown over the network (_cleanup_task_resources), and
# SQLite writes (_persist_session). A raise from any of them used to
# propagate straight out of run_conversation, discarding the partial
# final_response the caller is waiting for (subprocess wrappers saw an
# empty stdout with no traceback — #8049). Each step is now guarded
# independently so one failure can't skip the others, and any errors
# are surfaced on the result dict via ``cleanup_errors`` rather than
# killing the turn.
# Post-loop cleanup must never lose the response: trajectory save, teardown,
# and session persist are guarded independently and errors surface via
# ``cleanup_errors`` rather than killing the turn (#8049).
_cleanup_errors = []
# Save trajectory if enabled. ``user_message`` may be a multimodal
@@ -309,22 +249,17 @@ def finalize_turn(
_cleanup_errors.append(f"cleanup_task_resources: {_cleanup_err}")
logger.error("finalize_turn: _cleanup_task_resources failed: %s", _cleanup_err, exc_info=True)
# Persist session to both JSON log and SQLite only after private retry
# scaffolding has been removed. Otherwise a later user "continue" turn
# can replay assistant("(empty)") / recovery nudges and fall into the
# same empty-response loop again.
# Persist only after private retry scaffolding is removed, or a later "continue"
# replays assistant("(empty)") / recovery nudges into the same empty-response loop.
try:
agent._drop_trailing_empty_response_scaffolding(messages)
# Drop verification-continuation nudges (synthetic user messages)
# from the live history before the tail-assistant check — only the
# nudges need stripping; the assistant candidate persists in
# state.db. (#65919 §7)
# Strip only the synthetic verification nudges before the tail-assistant
# check; the assistant candidate persists in state.db. (#65919)
_drop_verification_continuation_scaffolding(messages)
# #95514: an empty terminal completion is not authoritative when the
# stream already delivered text. Recover before persist so a blank
# assistant tail is filled instead of frozen as content=''.
# An empty terminal completion is not authoritative when the stream already
# delivered text; recover before persist so a blank tail isn't frozen (#95514).
_recovered_from_stream = False
if not interrupted and not failed:
_streamed = getattr(agent, "_current_streamed_assistant_text", "") or ""
@@ -337,39 +272,16 @@ def finalize_turn(
final_response = _streamed
_recovered_from_stream = True
# When the turn was interrupted and the last message is a tool
# result, append a synthetic assistant message to close the
# tool-call sequence. Without this, the session persists a
# ``tool → user`` alternation that strict providers (Gemini,
# Claude) reject, causing them to hallucinate a continuation of
# the user's message on the next turn (#48879).
#
# ``_drop_trailing_empty_response_scaffolding`` only rewinds the
# tool tail when an empty-response scaffolding flag is present; a
# clean ``/stop`` interrupt after a successful tool sets no such
# flag, so the tool result survives as the tail and we close it
# here instead. On an interrupt ``final_response`` is typically
# empty, so fall back to an explicit placeholder rather than
# persisting an empty-content assistant turn.
# An interrupt can leave a tool result as the tail (no scaffolding flag rewinds
# it); close the sequence so strict providers don't see ``tool → user``. An
# explicit placeholder is used since final_response is usually empty (#48879).
if interrupted:
from agent.message_sanitization import close_interrupted_tool_sequence
close_interrupted_tool_sequence(messages, final_response)
# Some recovery/fallback paths return a real final_response without
# adding a closing assistant message to the transcript (e.g. the
# partial-stream and prior-turn-content recovery ``break`` sites in
# ``conversation_loop``). If persisted as-is, the durable session can
# end at a tool/user message even though the caller — and the gateway
# platform — already saw a completed assistant response. The next turn
# then replays a user-only backlog and the model re-answers every
# "unanswered" message. Close the durable turn at the source, at the
# single chokepoint every recovery ``break`` flows through, so the
# invariant "delivered final_response ⇒ assistant row in transcript"
# holds regardless of which path produced it. (#43849 / #44100)
#
# Compare content (not just role) so a verification candidate that
# matches the final response is not duplicated at budget
# exhaustion. (#65919 §7)
# Recovery ``break`` sites can return a final_response with no closing
# assistant row; enforce "delivered final_response ⇒ assistant row" here.
# Compare content, not role, so a matching verification candidate isn't dup'd.
if final_response and not interrupted:
try:
_tail = messages[-1] if messages else None
@@ -394,66 +306,45 @@ def finalize_turn(
)
)
):
# The tail IS an assistant row, but a *pure tool-call turn* or
# a blank assistant tail whose content was recovered from the
# stream buffer (#95514). Fill that row's content instead of
# appending, so the durable turn ends with the answer without
# creating an assistant→assistant pair.
# Tail is an assistant row (pure tool-call turn or stream-recovered
# blank, #95514): fill its content rather than append a second row.
_fill_assistant_tail_content(agent, _tail, final_response)
# The model has completed its request, so replace API-local
# voice/model/skill guidance with the clean user input before writing the
# final durable snapshot and returning the continuation history. Earlier
# turn-start flushes use the DB-only override because their messages are
# still needed for the API request; this finalizer runs after that request
# is complete (#48677 / #63766).
# Request is complete, so replace API-local voice/model/skill guidance with
# the clean user input before the durable snapshot; earlier flushes used the
# DB-only override as their messages were still needed (#48677 / #63766).
_apply_override = getattr(agent, "_apply_persist_user_message_override", None)
if callable(_apply_override):
_apply_override(messages)
# ── Post-turn micro-compaction ────────────────────────────
# After the assistant response is finalized but before the session is
# persisted, run micro-compaction to absorb the oldest uncompacted
# exchange into the rolling summary. This amortizes compression
# across turns rather than batching it into one big pause.
# Absorb the oldest uncompacted exchange into the rolling summary before
# persist, amortizing compression across turns instead of one big pause.
if not interrupted and not failed:
try:
_compressor = getattr(agent, "context_compressor", None)
# Strict `is True` + isinstance gates: plugin context engines
# (and MagicMock compressors in tests) satisfy getattr/duck
# checks with truthy auto-attributes — a bare truthiness check
# here called _micro_compact on a mock and spliced its (empty-
# iterating) return value over the transcript, wiping it.
# Strict `is True` + isinstance gates: plugin context engines and
# MagicMock compressors pass duck checks and would wipe the transcript.
if (
_compressor
and getattr(_compressor, '_micro_compact_enabled', False) is True
and callable(getattr(_compressor, '_micro_compact', None))
and final_response
# compression.checkpoint_required: agent init already
# forces _micro_compact_enabled off, but the compressor
# attribute is plain state a future path could flip on a
# live agent. Micro-compaction has no checkpoint hook in
# its path, so it must never run while the gate is armed.
# Micro-compaction has no checkpoint hook, so it must never run
# while compression.checkpoint_required is armed.
and getattr(
agent, "compression_checkpoint_required", False
) is not True
# Persistence-isolated agents (background review fork)
# must not micro-compact: the pass burns a real aux-LLM
# call on a throwaway replay transcript, and if the
# compressor ever holds a session_db binding it would
# archive_and_compact the CANONICAL session rows — the
# exact write class _persist_disabled exists to stop.
# Persistence-isolated agents (background review fork) must not
# micro-compact: it burns an aux-LLM call on a throwaway transcript
# and could archive_and_compact the CANONICAL session rows.
and not getattr(agent, "_persist_disabled", False)
):
_before = len(messages)
_compacted = _compressor._micro_compact(messages)
# Micro-compaction defrag rewrites the newest MICRO
# marker's content and pops _db_persisted from the live
# dict in place — the sibling of the pop site above. The
# compressor has no agent reference, so it raises a flag
# for us to invalidate the bounded flush-scan cursor;
# otherwise the rewritten marker row is identity-skipped
# and the stale summary persists to state.db.
# Defrag rewrites the newest MICRO marker in place and pops
# _db_persisted; the compressor flags us to invalidate the flush-
# scan cursor, else the rewritten row is identity-skipped (stale).
if getattr(
_compressor, "_flush_scan_cursor_invalidated", False
):
@@ -475,18 +366,16 @@ def finalize_turn(
_cleanup_errors.append(f"persist_session: {_persist_err}")
logger.error("finalize_turn: _persist_session failed: %s", _persist_err, exc_info=True)
# The gateway owns a separate in-memory history snapshot. Keep it current
# even when finalization reports a cleanup error: a later prompt must not be
# sent with the pre-turn snapshot while the durable DB already has this turn.
# Keep the gateway's separate in-memory history snapshot current even on
# cleanup error, so a later prompt isn't sent with a pre-turn snapshot.
try:
agent._session_messages = messages
except Exception:
pass
# ── Turn-exit diagnostic log ─────────────────────────────────────
# Always logged at INFO so agent.log captures WHY every turn ended.
# When the last message is a tool result (agent was mid-work), log
# at WARNING — this is the "just stops" scenario users report.
# Always INFO so agent.log captures WHY every turn ended; WARNING when the last
# message is a tool result (the "just stops" scenario).
_last_msg_role = messages[-1].get("role") if messages else None
_last_tool_name = None
if _last_msg_role == "tool":
@@ -527,21 +416,9 @@ def finalize_turn(
else:
logger.info(_diag_msg, *_diag_args)
# File-mutation verifier footer.
# If one or more ``write_file`` / ``patch`` calls failed during this
# turn and were never superseded by a successful write to the same
# path, append an advisory footer to the assistant response. This
# catches the specific case — reported by Ben Eng (#15524-adjacent)
# — where a model issues a batch of parallel patches, half of them
# fail with "Could not find old_string", and the model summarises
# the turn claiming every file was edited. The user then has to
# manually run ``git status`` to catch the lie. With this footer
# the truth is surfaced on every turn, so over-claiming is
# structurally impossible past the model.
#
# Gate: only applied when a real text response exists for this
# turn and the user didn't interrupt. Empty/interrupted turns
# already have other surface text that shouldn't be augmented.
# File-mutation verifier footer: if ``write_file`` / ``patch`` calls failed and
# were never superseded by a successful write to the same path, append an
# advisory so over-claiming is surfaced. Only on real, uninterrupted responses.
if final_response and not interrupted:
try:
_failed = getattr(agent, "_turn_failed_file_mutations", None) or {}
@@ -552,30 +429,16 @@ def finalize_turn(
except Exception as _ver_err:
logger.debug("file-mutation verifier footer failed: %s", _ver_err)
# Turn-completion explainer.
# When a turn ends abnormally after substantive work — empty content
# after retries, a partial/truncated stream, a still-pending tool
# result, or an iteration/budget limit — the user otherwise gets a
# blank or fragmentary response box with no consolidated reason why
# the agent stopped (#34452). Surface a single user-visible
# explanation derived from ``_turn_exit_reason``, mirroring the
# file-mutation verifier footer pattern above.
#
# Gate carefully so healthy turns stay quiet:
# - ``text_response(...)`` exits never produce an explanation
# (handled inside the formatter), so a terse ``Done.`` is silent.
# - We only ACT when there is no genuinely usable reply this turn:
# an empty response, the "(empty)" terminal sentinel, or a
# suspiciously short partial fragment with no terminating
# punctuation (e.g. "The"). A real short answer keeps its text.
# Turn-completion explainer: on abnormal exits, surface one explanation from
# ``_turn_exit_reason``. Only acts when no usable reply exists (empty, "(empty)",
# or a short unpunctuated fragment); ``text_response(...)`` exits stay silent.
if not interrupted:
try:
if agent._turn_completion_explainer_enabled():
_stripped = (final_response or "").strip()
_is_empty_terminal = _stripped == "" or _stripped == "(empty)"
# A short fragment that is not a normal text_response exit
# and lacks sentence-ending punctuation is treated as a
# truncated partial (the "The" case from #34452).
# A short fragment not from a text_response exit and lacking sentence-
# ending punctuation is treated as a truncated partial (#34452).
_is_partial_fragment = (
not _is_empty_terminal
and not preserved_verification_fallback
@@ -601,9 +464,7 @@ def finalize_turn(
# the actionable explanation.
final_response = _explanation
else:
# Keep the partial fragment, append the reason so
# the user sees both what arrived and why it
# stopped.
# Keep the partial fragment and append why it stopped.
final_response = (
_stripped + "\n\n" + _explanation
)
@@ -613,10 +474,8 @@ def finalize_turn(
_response_transformed = False
_pre_transform_response = None
# Plugin hook: transform_llm_output
# Fired once per turn after the tool-calling loop completes.
# Plugins can transform the LLM's output text before it's returned.
# First hook to return a string wins; None/empty return leaves text unchanged.
# Plugin hook: transform_llm_output — fired once per turn after the tool loop.
# First hook to return a string wins; None/empty leaves the text unchanged.
if final_response and not interrupted:
try:
from hermes_cli.lifecycle import invoke_hook as _invoke_hook
@@ -636,10 +495,8 @@ def finalize_turn(
except Exception as exc:
logger.warning("transform_llm_output hook failed: %s", exc)
# Plugin hook: post_llm_call
# Fired once per turn after the tool-calling loop completes.
# Plugins can use this to persist conversation data (e.g. sync
# to an external memory system).
# Plugin hook: post_llm_call — fired once per turn after the tool loop (e.g. sync
# conversation data to an external memory system).
if final_response and not interrupted:
try:
from hermes_cli.lifecycle import invoke_hook as _invoke_hook
@@ -657,18 +514,12 @@ def finalize_turn(
except Exception as exc:
logger.warning("post_llm_call hook failed: %s", exc)
# Context engine observation hook: notify the active engine that this
# turn has finished, with the finalized transcript. Complements the
# per-request select_context() hook (selection before the request;
# observation after the turn). No-op default, fail-open.
# Context engine observation hook (complements per-request select_context()):
# notify the engine the turn finished with the finalized transcript. Fail-open.
try:
from agent.conversation_loop import _notify_context_engine_turn_complete
# Forward the turn's canonical usage when the host has it. The loop
# stashes the most recent API response's usage dict (the same
# canonical buckets fed to ``update_from_response``) on the agent as
# ``_last_turn_usage``. It is ``None`` on turns that never reached a
# provider response (early failure / interrupt), which is exactly the
# contract: real usage when available, ``None`` otherwise.
# ``_last_turn_usage`` holds the last API response's canonical usage dict, or
# ``None`` on turns that never reached a provider response — by contract.
_turn_usage = getattr(agent, "_last_turn_usage", None)
_notify_context_engine_turn_complete(
agent,
@@ -685,15 +536,9 @@ def finalize_turn(
except Exception as exc:
logger.warning("on_turn_complete notification failed: %s", exc)
# Extract reasoning from the CURRENT turn only. Walk backwards
# but stop at the user message that started this turn — anything
# earlier is from a prior turn and must not leak into the reasoning
# box (confusing stale display; #17055). Within the current turn
# we still want the *most recent* non-empty reasoning: many
# providers (Claude thinking, DeepSeek v4, Codex Responses) emit
# reasoning on the tool-call step and leave the final-answer step
# with reasoning=None, so picking only the last assistant would
# silently drop legitimate same-turn reasoning.
# Reasoning from the CURRENT turn only: stop at this turn's user message
# (#17055), but take the most recent non-empty reasoning since many providers
# emit it on the tool-call step and leave the final step with reasoning=None.
last_reasoning = None
for msg in reversed(messages):
if msg.get("role") == "user":
@@ -702,15 +547,9 @@ def finalize_turn(
last_reasoning = msg["reasoning"]
break
# Class-level surrogate chokepoint (#80366, #55143, #55309, #19819):
# ``final_response`` is often the RAW SDK content
# (``assistant_message.content``), not the sanitized copy stored in
# history by ``build_assistant_message``. Any lone UTF-16 surrogate
# (U+D800–U+DFFF) in it crashes downstream consumers — oneshot stdout
# writes, Telegram's ``utf16_len`` length check, Signal formatting,
# JSON envelope encodes — on every provider (Ollama, NVIDIA NIM, …).
# Scrub once here, where model text leaves the conversation loop, so
# every delivery surface receives valid Unicode.
# Surrogate chokepoint: ``final_response`` may be RAW SDK content, and a lone UTF-16
# surrogate crashes downstream consumers (stdout, Telegram ``utf16_len``, JSON).
# Scrub once where model text leaves the loop (#80366).
if isinstance(final_response, str):
final_response = _sanitize_surrogates(final_response)
@@ -752,30 +591,27 @@ def finalize_turn(
}
if agent._tool_guardrail_halt_decision is not None:
result["guardrail"] = agent._tool_guardrail_halt_decision.to_metadata()
# Persistence failures already set failed=True + an explanation in
# final_response; also stamp `error` so gateway surfaces status="error"
# (and desktop can toast the cause) instead of a quiet complete frame.
# Persistence failures already set failed=True; also stamp `error` so the gateway
# surfaces status="error" (and desktop can toast) instead of a quiet complete frame.
if failed and str(_turn_exit_reason) == "session_persistence_failed":
result["error"] = final_response or (
"session storage could not be written — check the state database "
"health (`hermes doctor`), then send your message again"
)
# Machine-readable cause for the gateway/desktop: exactly
# 'session_persistence_failed:<locked|compression|turn_lease|corrupt|replaced|disk|unknown>'.
# Machine-readable cause for the gateway/desktop, exactly
# 'session_persistence_failed:<locked|compression|turn_lease|corrupt|...>'.
# Never clobber a failure_reason another path already stamped.
if "failure_reason" not in result:
_cause = getattr(agent, "_last_persistence_error_cause", None)
result["failure_reason"] = (
"session_persistence_failed:" + (_cause or "unknown")
)
# Surface any post-loop cleanup failures so the caller can distinguish a
# clean turn from one whose trajectory/session/resource teardown raised
# (the response is still returned either way — #8049).
# Surface post-loop cleanup failures so the caller can tell a clean turn from one
# whose teardown raised; the response is returned either way (#8049).
if _cleanup_errors:
result["cleanup_errors"] = _cleanup_errors
# If a /steer landed after the final assistant turn (no more tool
# batches to drain into), hand it back to the caller so it can be
# delivered as the next user turn instead of being silently lost.
# A /steer landing after the final assistant turn has no tool batch to drain into;
# hand it back so it becomes the next user turn instead of being lost.
_leftover_steer = agent._drain_pending_steer()
if _leftover_steer:
result["pending_steer"] = _leftover_steer
@@ -807,11 +643,9 @@ def finalize_turn(
messages=messages,
)
# Background memory/skill review — runs AFTER the response is delivered
# so it never competes with the user's task for model attention.
# Suppressed when skip_background_review=True (e.g. cron) — review forks
# spawn another AIAgent (~30K tokens / event) and cron sessions have no
# human-in-the-loop benefit from the review.
# Background memory/skill review runs AFTER delivery so it never competes with the
# user's task. Suppressed by skip_background_review (e.g. cron): the fork costs
# ~30K tokens / event with no human-in-the-loop benefit.
if (
final_response
and not interrupted
@@ -829,16 +663,10 @@ def finalize_turn(
except Exception:
pass # Background review is best-effort
# Note: Memory provider on_session_end() + shutdown_all() are NOT
# called here — run_conversation() is called once per user message in
# multi-turn sessions. Shutting down after every turn would kill the
# provider before the second message. Actual session-end cleanup is
# handled by the CLI (atexit / /reset) and gateway (session expiry /
# _reset_session).
# Memory provider on_session_end()/shutdown_all() are NOT called here:
# run_conversation() runs once per message; CLI/gateway own session-end cleanup.
# Plugin hook: on_session_end
# Fired at the very end of every run_conversation call.
# Plugins can use this for cleanup, flushing buffers, etc.
# Plugin hook: on_session_end — fired at the end of every run_conversation call.
try:
from hermes_cli.lifecycle import invoke_hook as _invoke_hook
_invoke_hook(
+560
View File
@@ -0,0 +1,560 @@
"""Overflow recovery for the conversation turn loop: 413 payload-too-large and
context-length errors after ``classify_api_error``.
Extracted from ``run_conversation``'s ``except`` branch. Each path either compresses and
signals a restart, defers softly (compression lock / transient block), or ends the turn
with a typed result. Nothing here imports ``agent.conversation_loop`` at module level
(cycle); loop-internal helpers and the token estimators that tests patch on the loop
module are imported lazily inside the handler so they keep resolving through the loop.
"""
from __future__ import annotations
import logging
import time
from dataclasses import dataclass
from typing import Any, Dict, List, Optional
from agent.conversation_compression import (
COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE,
COMPRESSION_RETRY_TOKENS_STATUS_TEMPLATE,
COMPRESSION_RETRY_TOO_LARGE_STATUS_TEMPLATE,
compression_blocked_transiently,
compression_skipped_due_to_lock,
context_compression_timed_out,
)
from agent.error_classifier import FailoverReason
from agent.message_sanitization import serialized_messages_bytes
from agent.model_metadata import (
get_context_length_from_provider_error,
is_output_cap_error,
parse_available_output_tokens_from_error,
)
from agent.turn_retry_state import TurnRetryState
from utils import base_url_host_matches
logger = logging.getLogger("agent.conversation_loop")
@dataclass
class OverflowVerdict:
"""Outcome of ``recover_from_overflow``.
``action`` is one of ``"return"`` (end the turn with ``result``), ``"break"``
(restart the API call — a ``_retry.restart_with_*`` flag is set), ``"continue"``
(retry the call immediately) or ``"fallthrough"`` (not an overflow error, or
overflow recovery declined — continue generic error handling). The remaining
fields are the loop locals the handler may have rebound."""
action: str
result: Optional[Dict[str, Any]]
messages: List[Dict[str, Any]]
active_system_prompt: Any
conversation_history: Any
approx_tokens: int
compression_attempts: int
provider_overflow_recovery_pending: bool
is_context_length_error: bool
def recover_from_overflow(
agent: Any,
api_error: Exception,
classified: Any,
_retry: TurnRetryState,
*,
status_code: Optional[int],
error_msg: str,
wrapped_output_cap_budget: Optional[int],
messages: List[Dict[str, Any]],
api_messages: Any,
system_message: Any,
active_system_prompt: Any,
conversation_history: Any,
approx_tokens: int,
compression_attempts: int,
max_compression_attempts: int,
api_call_count: int,
effective_task_id: Any,
) -> OverflowVerdict:
"""413 payload-too-large and context-length recovery (compress + retry, output-cap
clamp, provider-reported context limit, GitHub Models free-tier hint). Order is
load-bearing: 413 is checked BEFORE the generic 4xx handler, and context-length
errors (incl. relay-wrapped output-cap 429s) BEFORE non-retryable client errors.
Compression progress is scored in payload BYTES for 413 (never the byte-blind token
estimate) and in tokens/message count for context overflow."""
# Token estimators + loop-internal helpers resolve through the loop module so
# existing ``patch("agent.conversation_loop.X")`` mocks keep intercepting.
from agent.conversation_loop import (
_COMPRESSION_TIMEOUT_FINAL_RESPONSE,
_compression_deferred_result,
conversation_history_after_compression,
estimate_messages_tokens_rough,
estimate_request_tokens_rough,
save_context_length,
)
_provider_overflow_recovery_pending = False
is_context_length_error = False
_wrapped_output_cap_budget = wrapped_output_cap_budget
def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> OverflowVerdict:
return OverflowVerdict(
action=action,
result=result,
messages=messages,
active_system_prompt=active_system_prompt,
conversation_history=conversation_history,
approx_tokens=approx_tokens,
compression_attempts=compression_attempts,
provider_overflow_recovery_pending=_provider_overflow_recovery_pending,
is_context_length_error=is_context_length_error,
)
is_payload_too_large = (
classified.reason == FailoverReason.payload_too_large
)
# GitHub Models free tier caps requests at 8K tokens, under the system
# prompt + tool schema floor; compression can't help, so say so.
if (
status_code == 413
and isinstance(agent.base_url, str)
and base_url_host_matches(agent.base_url, "models.inference.ai.azure.com")
):
agent._vprint(
f"{agent.log_prefix} 💡 GitHub Models free tier (models.inference.ai.azure.com) caps every",
force=True,
)
agent._vprint(
f"{agent.log_prefix} request at ~8K tokens. Hermes' system prompt + tool schemas baseline",
force=True,
)
agent._vprint(
f"{agent.log_prefix} exceeds that floor, so this endpoint cannot run an agentic loop.",
force=True,
)
agent._vprint(
f"{agent.log_prefix} Use the `copilot` provider with a Copilot subscription token (`hermes",
force=True,
)
agent._vprint(
f"{agent.log_prefix} setup` → GitHub Copilot), or pick any other provider.",
force=True,
)
if is_payload_too_large:
compression_attempts += 1
if compression_attempts > max_compression_attempts:
# Terminal — surface the buffered retry trace.
agent._flush_status_buffer()
agent._vprint(f"{agent.log_prefix}❌ Max compression attempts ({max_compression_attempts}) reached for payload-too-large error.", force=True)
agent._vprint(f"{agent.log_prefix} 💡 Try /new to start a fresh conversation, or /compress to retry compression.", force=True)
logger.error("%s413 compression failed after %d attempts.", agent.log_prefix, max_compression_attempts)
agent._persist_session(messages, conversation_history)
_final_response = f"Request payload too large: max compression attempts ({max_compression_attempts}) reached."
return _verdict("return", {
"final_response": _final_response,
"messages": messages,
"completed": False,
"api_calls": api_call_count,
"error": _final_response,
"partial": True,
"failed": True,
"compression_exhausted": True,
})
agent._buffer_status(f"⚠️ Request payload too large (413) — compression attempt {compression_attempts}/{max_compression_attempts}...")
original_len = len(messages)
# A 413 is a BYTE-size error: score progress in payload bytes,
# never the token estimate, which is deliberately byte-blind to
# images and wedged sessions on "no progress" (#88960 / #47339).
original_bytes = serialized_messages_bytes(messages)
_overflow_input = messages
# Option A (LCM issue 441): overhead-aware request size so recovery arms on the
# true request (msgs + tools + system), not the tool-blind message count.
messages, active_system_prompt = agent._compress_context(
messages, system_message,
approx_tokens=estimate_request_tokens_rough(api_messages, tools=agent.tools or None),
task_id=effective_task_id,
# Provider proved the request doesn't fit: ignore the
# summary-failure cooldown for this ONE attempt (#100661).
bypass_cooldown=True,
)
if messages is _overflow_input and compression_skipped_due_to_lock(agent):
# Lock-skip: another path holds the compression lock. A
# temporary defer, not exhaustion — refund the attempt and
# end softly so the gateway does NOT auto-reset (#69870).
compression_attempts -= 1
agent._persist_session(messages, conversation_history)
return _verdict("return", _compression_deferred_result(
agent, messages, api_call_count
))
if messages is _overflow_input and compression_blocked_transiently(agent):
# Transient-block: a timed guard no-oped compression. A
# defer, never compression_exhausted (auto-reset) (#97488).
compression_attempts -= 1
agent._persist_session(messages, conversation_history)
return _verdict("return", _compression_deferred_result(
agent, messages, api_call_count,
reason="transient_block",
))
conversation_history = conversation_history_after_compression(
agent, messages, conversation_history
)
# Re-measure: same-count compression and media aging can shrink
# the request without shrinking the array. Bytes are the yardstick
# for a 413; tokens only for status display.
new_tokens = estimate_messages_tokens_rough(messages)
approx_tokens = new_tokens # update for downstream logging
new_bytes = serialized_messages_bytes(messages)
made_progress = (
len(messages) < original_len
or (new_bytes > 0 and new_bytes < original_bytes * 0.95)
)
if made_progress:
if len(messages) < original_len:
agent._buffer_status(COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE.format(before=original_len, after=len(messages)))
else:
agent._buffer_status(
f"🗜️ Compressed {original_bytes:,} → {new_bytes:,} "
f"payload bytes, retrying..."
)
time.sleep(2) # Brief pause between compression retries
_retry.restart_with_compressed_messages = True
return _verdict("break")
else:
if agent._try_strip_image_parts_from_tool_messages(
api_messages,
remember_model=False,
):
agent._buffer_status(
"📐 Compression could not reduce the request further — "
"removed retained vision payloads and retrying..."
)
return _verdict("continue")
# Terminal — surface buffered context so the user
# sees what compression attempts were made.
agent._flush_status_buffer()
agent._vprint(f"{agent.log_prefix}❌ Payload too large and cannot compress further.", force=True)
agent._vprint(f"{agent.log_prefix} 💡 Try /new to start a fresh conversation, or /compress to retry compression.", force=True)
logger.error("%s413 payload too large. Cannot compress further.", agent.log_prefix)
agent._persist_session(messages, conversation_history)
_final_response = "Request payload too large (413). Cannot compress further."
return _verdict("return", {
"final_response": _final_response,
"messages": messages,
"completed": False,
"api_calls": api_call_count,
"error": _final_response,
"partial": True,
"failed": True,
"compression_exhausted": True,
})
# Check context-length errors BEFORE the generic 4xx handler; the
# classifier also covers 400/disconnect + large-session heuristics.
is_context_length_error = (
classified.reason == FailoverReason.context_overflow
# Relay-wrapped output-cap 429s (parsed above) go to the clamp
# below, not failover or generic retries (#72281).
or _wrapped_output_cap_budget is not None
)
if is_context_length_error:
compressor = agent.context_compressor
old_ctx = compressor.context_length
# Two errors: "prompt too long" = INPUT overflows the window (shrink
# context_length + compress); "max_tokens too large" = input fits
# but input + max_tokens > window (shrink OUTPUT cap only).
available_out = parse_available_output_tokens_from_error(error_msg)
if available_out is not None:
# Output-cap error: provider available_tokens is the
# authoritative bound; also estimate the real request shape
# (API-only content), use the smaller minus a margin.
request_input_estimate = estimate_request_tokens_rough(
api_messages, tools=agent.tools or None,
)
local_available_out = old_ctx - request_input_estimate
if local_available_out > 0:
safe_out = max(1, min(available_out, local_available_out) - 64)
else:
# Local estimate can overshoot; fall back to the
# authoritative provider-reported budget.
safe_out = max(1, available_out - 64)
agent._ephemeral_max_output_tokens = safe_out
agent._buffer_vprint(
f"⚠️ Output cap too large for current prompt — "
f"retrying with max_tokens={safe_out:,} "
f"(provider_available={available_out:,}, "
f"estimated_request_tokens={request_input_estimate:,}; "
f"context_length unchanged at {old_ctx:,})"
)
# Still count against compression_attempts so we don't
# loop forever if the error keeps recurring.
compression_attempts += 1
if compression_attempts > max_compression_attempts:
agent._flush_status_buffer()
agent._vprint(f"{agent.log_prefix}❌ Max compression attempts ({max_compression_attempts}) reached.", force=True)
agent._vprint(f"{agent.log_prefix} 💡 Try /new to start a fresh conversation, or /compress to retry compression.", force=True)
logger.error("%sContext compression failed after %d attempts.", agent.log_prefix, max_compression_attempts)
agent._persist_session(messages, conversation_history)
_final_response = f"Context length exceeded: max compression attempts ({max_compression_attempts}) reached."
return _verdict("return", {
"final_response": _final_response,
"messages": messages,
"completed": False,
"api_calls": api_call_count,
"error": _final_response,
"partial": True,
"failed": True,
"compression_exhausted": True,
})
# Also compress history so the output-cap retry doesn't spin on
# max_tokens alone; dropping the middle window makes the total
# fit. (#55546)
try:
original_len = len(messages)
original_tokens = estimate_messages_tokens_rough(messages)
_overflow_input = messages
messages, active_system_prompt = agent._compress_context(
messages, system_message,
approx_tokens=request_input_estimate,
task_id=effective_task_id,
bypass_cooldown=True, # #100661 provider-proven overflow
)
if messages is _overflow_input and compression_skipped_due_to_lock(agent):
compression_attempts -= 1
agent._persist_session(messages, conversation_history)
return _verdict("return", _compression_deferred_result(
agent, messages, api_call_count
))
if messages is _overflow_input and compression_blocked_transiently(agent):
# #97488: timed transient guard — defer, never
# exhaustion (gateway auto-reset).
compression_attempts -= 1
agent._persist_session(messages, conversation_history)
return _verdict("return", _compression_deferred_result(
agent, messages, api_call_count,
reason="transient_block",
))
conversation_history = conversation_history_after_compression(
agent, messages, conversation_history
)
new_tokens = estimate_messages_tokens_rough(messages)
if len(messages) < original_len:
agent._buffer_status(COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE.format(before=original_len, after=len(messages)))
elif new_tokens > 0 and new_tokens < original_tokens * 0.95:
agent._buffer_status(COMPRESSION_RETRY_TOKENS_STATUS_TEMPLATE.format(before=original_tokens, after=new_tokens))
except Exception:
# Compression must never turn an output-cap error
# fatal — fall through and retry on max_tokens alone.
logger.warning(
"%sOutput-cap compression hit an error; retrying on max_tokens only.",
agent.log_prefix,
)
_retry.restart_with_compressed_messages = True
return _verdict("break")
# Output-cap error with unparseable budget: compression can't help
# (input already fits) and would death-loop on the same 400. Fail
# fast. (#55546)
if is_output_cap_error(error_msg):
agent._flush_status_buffer()
agent._vprint(
f"{agent.log_prefix}❌ The provider rejected the request because "
f"max_tokens exceeds its output cap for this model.",
force=True,
)
agent._vprint(
f"{agent.log_prefix} 💡 Lower model.max_tokens in your config.yaml to "
f"at or below the model's max-output limit. "
f"(This is an output-cap error, not a context overflow — "
f"compression cannot fix it.)",
force=True,
)
logger.error(
f"{agent.log_prefix}Output-cap error not routed into compression "
f"(max_tokens over provider cap): {error_msg[:200]}"
)
agent._persist_session(messages, conversation_history)
_final_response = (
"max_tokens exceeds the provider's output cap for this model. "
"Lower model.max_tokens in config.yaml."
)
return _verdict("return", {
"final_response": _final_response,
"messages": messages,
"completed": False,
"api_calls": api_call_count,
"error": _final_response,
"partial": True,
"failed": True,
})
# Input too large: shrink context_length only when the provider
# reports the real limit; else keep the window and compress. Guessed
# probe tiers can turn a configured 1M window into 256K/128K/64K.
new_ctx = get_context_length_from_provider_error(error_msg, old_ctx)
_provider_lower = (getattr(agent, "provider", "") or "").lower()
_base_lower = (getattr(agent, "base_url", "") or "").rstrip("/").lower()
is_minimax_provider = (
_provider_lower in {"minimax", "minimax-cn"}
or _base_lower.startswith((
"https://api.minimax.io/anthropic",
"https://api.minimaxi.com/anthropic",
))
)
minimax_delta_only_overflow = (
is_minimax_provider
and new_ctx is None
and "context window exceeds limit (" in error_msg
)
if new_ctx is not None:
agent._buffer_vprint(f"Context limit detected from API: {new_ctx:,} tokens (was {old_ctx:,})")
compressor.update_model(
model=agent.model,
context_length=new_ctx,
base_url=agent.base_url,
api_key=getattr(agent, "api_key", ""),
provider=agent.provider,
api_mode=agent.api_mode,
)
# Persist the provider-reported limit before compression/retry:
# rate limit, missing usage, or restart must not lose confirmed
# metadata. Probe flags remain a fallback if this write fails.
save_context_length(agent.model, agent.base_url, new_ctx)
# Probe flags only on the built-in compressor (plugin engines
# manage their own); provider-sourced value, so safe to cache.
if hasattr(compressor, "_context_probed"):
compressor._context_probed = True
compressor._context_probe_persistable = True
agent._buffer_vprint(f"⚠️ Context length exceeded — using provider limit: {old_ctx:,} → {new_ctx:,} tokens")
elif minimax_delta_only_overflow:
agent._buffer_vprint(
f"Provider reported overflow amount only; "
f"keeping context_length at {old_ctx:,} tokens and compressing."
)
else:
agent._buffer_vprint(
f"⚠️ Context length exceeded, but provider did not report a max context length; "
f"keeping context_length at {old_ctx:,} tokens and compressing."
)
compression_attempts += 1
if compression_attempts > max_compression_attempts:
agent._flush_status_buffer()
agent._vprint(f"{agent.log_prefix}❌ Max compression attempts ({max_compression_attempts}) reached.", force=True)
agent._vprint(f"{agent.log_prefix} 💡 Try /new to start a fresh conversation, or /compress to retry compression.", force=True)
logger.error("%sContext compression failed after %d attempts.", agent.log_prefix, max_compression_attempts)
agent._persist_session(messages, conversation_history)
_final_response = f"Context length exceeded: max compression attempts ({max_compression_attempts}) reached."
return _verdict("return", {
"final_response": _final_response,
"messages": messages,
"completed": False,
"api_calls": api_call_count,
"error": _final_response,
"partial": True,
"failed": True,
"compression_exhausted": True,
})
agent._buffer_status(COMPRESSION_RETRY_TOO_LARGE_STATUS_TEMPLATE.format(tokens=approx_tokens, attempt=compression_attempts, cap=max_compression_attempts))
original_len = len(messages)
original_tokens = estimate_messages_tokens_rough(messages)
_overflow_input = messages
# Pass the OVERHEAD-AWARE size (msgs + tool schemas + system) so LCM
# forced-overflow recovery arms on the TRUE request; approx_tokens
# stays for status. See hermes-lcm _should_force_overflow_recovery.
messages, active_system_prompt = agent._compress_context(
messages, system_message,
approx_tokens=estimate_request_tokens_rough(api_messages, tools=agent.tools or None),
task_id=effective_task_id,
# Provider proved the request doesn't fit: ignore the
# summary-failure cooldown for this ONE attempt (bounded by
# max_compression_attempts). (#100661)
bypass_cooldown=True,
)
if messages is _overflow_input and compression_skipped_due_to_lock(agent):
# Lock-skip: another path holds the compression lock, so this
# pass no-oped. Temporary defer, not exhaustion — refund the
# attempt, end the turn softly, no auto-reset. (#69870)
compression_attempts -= 1
agent._persist_session(messages, conversation_history)
return _verdict("return", _compression_deferred_result(
agent, messages, api_call_count
))
if messages is _overflow_input and compression_blocked_transiently(agent):
# Transient block: a timed guard (host-timeout cooldown /
# structural backoff) no-oped this pass — defer softly, never
# compression_exhausted (auto-reset). (#97488)
compression_attempts -= 1
agent._persist_session(messages, conversation_history)
return _verdict("return", _compression_deferred_result(
agent, messages, api_call_count,
reason="transient_block",
))
if context_compression_timed_out(agent):
# Host timeout: recovery spent its wait budget with no committed
# summary. Re-sending would hit the same overflow; end the turn
# via the typed recovery contract. (#98722)
agent._persist_session(messages, conversation_history)
_final_response = _COMPRESSION_TIMEOUT_FINAL_RESPONSE
return _verdict("return", {
"final_response": _final_response,
"messages": messages,
"completed": False,
"api_calls": api_call_count,
"error": _final_response,
"partial": True,
"failed": True,
"compression_exhausted": True,
"turn_exit_reason": "context_compression_timeout",
})
conversation_history = conversation_history_after_compression(
agent, messages, conversation_history
)
# Re-estimate after compression: same-message-count compression
# (tool-result pruning, in-place summarization) can shrink the
# request. (#39550)
new_tokens = estimate_messages_tokens_rough(messages)
approx_tokens = new_tokens # update for downstream logging
if len(messages) < original_len or (new_tokens > 0 and new_tokens < original_tokens * 0.95) or (new_ctx and new_ctx < old_ctx):
if len(messages) < original_len:
agent._buffer_status(COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE.format(before=original_len, after=len(messages)))
elif new_tokens > 0 and new_tokens < original_tokens * 0.95:
agent._buffer_status(COMPRESSION_RETRY_TOKENS_STATUS_TEMPLATE.format(before=original_tokens, after=new_tokens))
time.sleep(2) # Brief pause between compression retries
# Rebuild the full request and force normal preflight to honor
# it; message count alone doesn't prove system/tool-inclusive
# pressure fell.
_provider_overflow_recovery_pending = True
_retry.restart_with_compressed_messages = True
return _verdict("break")
else:
# Can't compress further and already at minimum tier
agent._flush_status_buffer()
agent._vprint(f"{agent.log_prefix}❌ Context length exceeded and cannot compress further.", force=True)
agent._vprint(f"{agent.log_prefix} 💡 The conversation has accumulated too much content. Try /new to start fresh, or /compress to manually trigger compression.", force=True)
logger.error("%sContext length exceeded: %s tokens. Cannot compress further.", agent.log_prefix, f"{new_tokens:,}")
agent._persist_session(messages, conversation_history)
_final_response = f"Context length exceeded ({new_tokens:,} tokens). Cannot compress further."
return _verdict("return", {
"final_response": _final_response,
"messages": messages,
"completed": False,
"api_calls": api_call_count,
"error": _final_response,
"partial": True,
"failed": True,
"compression_exhausted": True,
})
return _verdict("fallthrough")
+506
View File
@@ -0,0 +1,506 @@
"""Pre-API preflight compression gate for the conversation turn loop.
Extracted from ``run_conversation``. Runs once per API call after the request pressure
is measured: grow a managed llama.cpp window (last resort), or compress when over
threshold (deferring on noisy estimates, in failure cooldown, or when the review fork's
first request is pending), handle the provider-proven overflow re-check fail-closed,
and emit the deduped blocked/uncompressed overflow warnings. Nothing here imports
``agent.conversation_loop`` at module level (cycle); loop-internal helpers resolve lazily.
"""
from __future__ import annotations
import logging
from dataclasses import dataclass
from typing import Any, Dict, List, Optional
from agent.context_engine import automatic_compaction_status_message
from agent.conversation_compression import (
PRE_API_COMPRESSION_STATUS_TEMPLATE,
compression_blocked_transiently,
compression_skipped_due_to_lock,
context_compression_timed_out,
conversation_history_after_compression,
)
from agent.turn_context import _review_fork_first_request_pending
logger = logging.getLogger("agent.conversation_loop")
@dataclass
class PreflightVerdict:
"""Outcome of ``run_preflight_compression``.
``action``: ``"proceed"`` (make the API call), ``"continue"`` (window grown or
history compacted — the call/budget was refunded, re-enter the turn loop and
re-measure), ``"break"`` (turn ends: compression timeout or non-actionable
compaction handoff — ``final_response``/``failed``/``turn_exit_reason`` set) or
``"return"`` (typed deferred/exhausted result in ``result``). The remaining fields
are the loop locals the gate may have rebound."""
action: str
result: Optional[Dict[str, Any]]
messages: List[Dict[str, Any]]
active_system_prompt: Any
conversation_history: Any
api_call_count: int
compression_attempts: int
pending_moa_prepared_request: Any
last_preflight_pressure: Optional[int]
final_response: Any
failed: bool
compression_timeout_exhausted: bool
turn_exit_reason: Any
def run_preflight_compression(
agent: Any,
*,
compressor: Any,
request_pressure_tokens: int,
provider_overflow_preflight: bool,
preflight_compression_blocked: bool,
defer_preflight: Any,
moa_prepared_request: Any,
pending_moa_prepared_request: Any,
messages: List[Dict[str, Any]],
system_message: Any,
user_message: Any,
active_system_prompt: Any,
conversation_history: Any,
api_call_count: int,
compression_attempts: int,
max_compression_attempts: int,
effective_task_id: Any,
final_response: Any,
failed: bool,
compression_timeout_exhausted: bool,
turn_exit_reason: Any,
) -> PreflightVerdict:
"""Mirror of the turn-prologue guard chain (defer on noisy estimate → skip in failure
cooldown → ``should_compress``), #11529. A compression pass that never reaches the
provider refunds the call/budget in every branch (skip, re-run, timeout) so
``api_call_count`` never over-reports; a lock/transient skip refunds the attempt and
leaves the progress blocker unarmed (#69870, #97488). A forced provider-overflow
preflight that any gate blocks fails closed (llama.cpp may silently truncate)."""
from agent.conversation_loop import (
_COMPRESSION_TIMEOUT_FINAL_RESPONSE,
_HANDOFF_SKIP_FINAL_RESPONSE,
_compression_deferred_result,
_maybe_grow_local_window,
_provider_overflow_exhausted_result,
_should_skip_model_call_for_reference_handoff,
)
_compressor = compressor
_provider_overflow_preflight = provider_overflow_preflight
_preflight_compression_blocked = preflight_compression_blocked
_defer_preflight = defer_preflight
_moa_prepared_request = moa_prepared_request
_last_preflight_pressure = None
_compression_timeout_exhausted = compression_timeout_exhausted
_turn_exit_reason = turn_exit_reason
def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> PreflightVerdict:
return PreflightVerdict(
action=action,
result=result,
messages=messages,
active_system_prompt=active_system_prompt,
conversation_history=conversation_history,
api_call_count=api_call_count,
compression_attempts=compression_attempts,
pending_moa_prepared_request=pending_moa_prepared_request,
last_preflight_pressure=_last_preflight_pressure,
final_response=final_response,
failed=failed,
compression_timeout_exhausted=_compression_timeout_exhausted,
turn_exit_reason=_turn_exit_reason,
)
_compression_cooldown = getattr(
_compressor, "get_active_compression_failure_cooldown", lambda: None
)()
if (
agent.compression_enabled
and not _review_fork_first_request_pending(agent)
and len(messages) > 1
and compression_attempts < max_compression_attempts
and (
not _preflight_compression_blocked
or _provider_overflow_preflight
)
and (
not _defer_preflight(request_pressure_tokens)
or _provider_overflow_preflight
)
and not _compression_cooldown
and _compressor.should_compress(request_pressure_tokens)
):
# Managed local runtime: grow the context window before compressing (last
# resort). Only for a llamacpp provider at the supervised base_url.
_grown_window = _maybe_grow_local_window(
agent, _compressor, request_pressure_tokens
)
if _grown_window:
# Bigger window granted: recalibrate the compressor and skip compression
# this pass.
_compressor.update_model(
agent.model,
_grown_window,
base_url=getattr(agent, "base_url", "") or "",
api_key=getattr(agent, "api_key", "") or "",
provider=getattr(agent, "provider", "") or "",
api_mode=getattr(agent, "api_mode", "") or "",
)
agent._buffer_status(
f"📈 Context window grown to {_grown_window // 1024}K "
f"(local model; conversation continues uncompressed)"
)
# Never reached the provider — refund the call/budget like the
# compression path does before its continue.
api_call_count -= 1
agent._api_call_count = api_call_count
agent.iteration_budget.refund()
return _verdict("continue")
if _moa_prepared_request is not None:
pending_moa_prepared_request = _moa_prepared_request
compression_attempts += 1
# Compression is running: reset the blocked-overflow warning dedup so a
# later blocked turn warns again (#62625). getattr: test doubles lack it.
_clear_warn = getattr(agent, "_clear_context_overflow_warn", None)
if callable(_clear_warn):
_clear_warn()
logger.info(
"Pre-API compression: ~%s request tokens >= %s threshold "
"(context=%s, attempt=%s/%s)",
f"{request_pressure_tokens:,}",
f"{int(getattr(_compressor, 'threshold_tokens', 0) or 0):,}",
f"{int(getattr(_compressor, 'context_length', 0) or 0):,}"
if getattr(_compressor, "context_length", 0) else "unknown",
compression_attempts,
max_compression_attempts,
)
_pre_api_status = automatic_compaction_status_message(
_compressor,
phase="pre_api",
default_message=PRE_API_COMPRESSION_STATUS_TEMPLATE.format(
tokens=request_pressure_tokens
),
approx_tokens=request_pressure_tokens,
threshold_tokens=int(
getattr(_compressor, "threshold_tokens", 0) or 0
),
context_length=int(
getattr(_compressor, "context_length", 0) or 0
),
model=agent.model,
attempt=compression_attempts,
max_attempts=max_compression_attempts,
)
if _pre_api_status:
agent._emit_status(_pre_api_status)
_last_preflight_pressure = request_pressure_tokens
_pre_api_input = messages
messages, active_system_prompt = agent._compress_context(
messages,
system_message,
approx_tokens=request_pressure_tokens,
task_id=effective_task_id,
)
if context_compression_timed_out(agent):
# Progress-aware timeout (#98722): never reached the provider — refund
# the call/budget and stop; an overflow retry would only re-compress.
api_call_count -= 1
agent._api_call_count = api_call_count
agent.iteration_budget.refund()
final_response = _COMPRESSION_TIMEOUT_FINAL_RESPONSE
failed = True
_compression_timeout_exhausted = True
_turn_exit_reason = "context_compression_timeout"
return _verdict("break")
if messages is _pre_api_input and (
compression_skipped_due_to_lock(agent)
or compression_blocked_transiently(agent)
):
# Temporary DEFER (lock held / cooldown), not evidence about
# compressibility: refund the attempt, leave the progress blocker
# unarmed and proceed (#69870, #97488).
compression_attempts -= 1
_last_preflight_pressure = None
if pending_moa_prepared_request is _moa_prepared_request:
pending_moa_prepared_request = None
else:
# Reset retry/empty-response state so the compacted request gets a fresh
# chance.
agent._empty_content_retries = 0
agent._thinking_prefill_retries = 0
agent._last_content_with_tools = None
agent._last_content_tools_all_housekeeping = False
agent._mute_post_response = False
# Re-baseline the flush cursor: rotation returns None (child flushes
# whole); in-place returns list(messages) — None would re-append
# persisted rows. See conversation_history_after_compression().
conversation_history = conversation_history_after_compression(
agent, messages, conversation_history
)
# Never reaches the provider on skip or re-run — refund the call/budget
# in BOTH cases, else budget leaks and api_call_count over-reports.
api_call_count -= 1
agent._api_call_count = api_call_count
agent.iteration_budget.refund()
if _should_skip_model_call_for_reference_handoff(
messages, user_message
):
# Reference-only handoff must not become the active turn
# after a completed assistant response (#80622).
logger.info(
"Skipping post-compaction model call: reference-only "
"handoff would be the sole active user turn (#80622)"
)
if not final_response:
final_response = _HANDOFF_SKIP_FINAL_RESPONSE
_turn_exit_reason = "compaction_handoff_not_actionable"
return _verdict("break")
return _verdict("continue")
elif _provider_overflow_preflight and _compression_cooldown:
# Provider proved the request cannot fit and the compressor is unavailable:
# don't resend; let the next user turn retry after cooldown.
agent._persist_session(messages, conversation_history)
return _verdict("return", _compression_deferred_result(
agent,
messages,
api_call_count,
reason="transient_block",
))
elif (
_provider_overflow_preflight
and compression_attempts >= max_compression_attempts
):
# All recovery passes consumed and still over threshold: fail closed —
# llama.cpp may silently truncate an oversized retry.
return _verdict("return", _provider_overflow_exhausted_result(
agent,
messages,
conversation_history,
api_call_count,
request_pressure_tokens,
max_compression_attempts,
))
elif (
agent.compression_enabled
and len(messages) > 1
and compression_attempts < max_compression_attempts
and not _defer_preflight(request_pressure_tokens)
and _compression_cooldown
):
# Summary-LLM cooldown blocks compression: deduped warning only when over
# threshold (should_compress_info reason is None below it) (#62625).
_block_reason = None
try:
_block_reason = _compressor.should_compress_info(
request_pressure_tokens
)[1]
except Exception:
_block_reason = None
if _block_reason:
agent._warn_context_overflow_blocked(
_block_reason,
request_pressure_tokens,
int(getattr(_compressor, "threshold_tokens", 0) or 0),
)
elif not agent.compression_enabled and len(messages) > 1:
# Uncompressed session guard (#89297): compression is disabled, so warn
# (deduped) when the request exceeds the context window; the turn-context
# preflight re-arms the dedup.
_ctx_len = getattr(
getattr(agent, "context_compressor", None), "context_length", None
)
if (
isinstance(_ctx_len, int)
and _ctx_len > 0
and request_pressure_tokens > _ctx_len
):
_warn_fn = getattr(
agent, "_warn_uncompressed_context_overflow", None
)
if callable(_warn_fn):
_warn_fn(request_pressure_tokens, _ctx_len)
if _provider_overflow_preflight:
# Any other gate blocking the forced preflight (e.g. uncompressible one-
# message request) must fail closed: the request is proven not to fit.
return _verdict("return", _provider_overflow_exhausted_result(
agent,
messages,
conversation_history,
api_call_count,
request_pressure_tokens,
max_compression_attempts,
))
return _verdict("proceed")
@dataclass
class PostToolCompressionVerdict:
"""``end_turn`` True → a reference-only compaction handoff would be the sole active
user turn (#80622): stop without another model call (``final_response`` /
``turn_exit_reason`` set)."""
end_turn: bool
messages: List[Dict[str, Any]]
active_system_prompt: Any
conversation_history: Any
compression_attempts: int
final_response: Any
turn_exit_reason: Any
def compress_after_tool_results(
agent: Any,
*,
messages: List[Dict[str, Any]],
system_message: Any,
user_message: Any,
active_system_prompt: Any,
conversation_history: Any,
compression_attempts: int,
max_compression_attempts: int,
effective_task_id: Any,
final_response: Any,
turn_exit_reason: Any,
) -> PostToolCompressionVerdict:
"""Post-tool-call compression decision. Pressure comes from API-reported
``prompt_tokens`` (a tight lower bound; thinking models inflate completion tokens,
#12026), ``0`` right after compression (no real count yet), else the route-aware
overhead-inclusive estimate (#14695). Over threshold but blocked → deduped warning
(#62625) plus the deterministic tool-result-only prune, committed only when the
engine returns a NEW list (never rebuild ``conversation_history`` for it)."""
from agent.conversation_loop import (
_HANDOFF_SKIP_FINAL_RESPONSE,
_midturn_request_pressure_tokens,
_should_skip_model_call_for_reference_handoff,
estimate_request_tokens_rough,
)
_turn_exit_reason = turn_exit_reason
def _verdict(end_turn: bool) -> PostToolCompressionVerdict:
return PostToolCompressionVerdict(
end_turn=end_turn,
messages=messages,
active_system_prompt=active_system_prompt,
conversation_history=conversation_history,
compression_attempts=compression_attempts,
final_response=final_response,
turn_exit_reason=_turn_exit_reason,
)
# Decide compression from API-reported prompt tokens (tight lower bound;
# tool results get counted on the next call). If last_prompt_tokens is 0
# (disconnect / no usage data) fall back to a rough estimate. (#2153)
_compressor = agent.context_compressor
if _compressor.last_prompt_tokens > 0:
# Only prompt_tokens: thinking models inflate completion_tokens with
# reasoning that uses no context → premature compression. (#12026)
_real_tokens = _compressor.last_prompt_tokens
elif _compressor.last_prompt_tokens == -1:
# Compression just ran, no API prompt count yet: don't treat a rough
# schema-heavy post-compression estimate as real context pressure.
_real_tokens = 0
else:
# Include tool schemas (20-30K tokens the messages-only estimate
# misses) and stay route-aware: on a compacted native-Codex session
# the generic durable-history figure would false-trigger. (#14695)
_real_tokens = _midturn_request_pressure_tokens(
agent,
messages,
active_system_prompt or "",
estimate_request_tokens_rough(
messages, tools=agent.tools or None
),
)
if (
agent.compression_enabled
and compression_attempts < max_compression_attempts
and _compressor.should_compress(_real_tokens)
):
compression_attempts += 1
# Compression is running: reset blocked-overflow warning dedup so a
# future blocked turn can warn again. getattr: test doubles lack it.
_clear_warn = getattr(agent, "_clear_context_overflow_warn", None)
if callable(_clear_warn):
_clear_warn()
agent._safe_print(" ⟳ compacting context…")
_post_tool_input = messages
# Pass overhead-aware _real_tokens, not last_prompt_tokens (0 in
# the no-usage fallback), so the overflow guard sees the true size.
messages, active_system_prompt = agent._compress_context(
messages, system_message,
approx_tokens=_real_tokens,
task_id=effective_task_id,
)
if (
messages is _post_tool_input
and compression_skipped_due_to_lock(agent)
):
# Lock-skip no-op is a temporary defer, not evidence about
# compressibility: refund so a lock-loser loop doesn't burn the
# budget toward compression_exhausted. (#69870)
compression_attempts -= 1
else:
conversation_history = conversation_history_after_compression(
agent, messages, conversation_history
)
if _should_skip_model_call_for_reference_handoff(
messages, user_message
):
logger.info(
"Skipping post-tool compaction model call: "
"reference-only handoff would be the sole "
"active user turn (#80622)"
)
if not final_response:
final_response = _HANDOFF_SKIP_FINAL_RESPONSE
_turn_exit_reason = "compaction_handoff_not_actionable"
return _verdict(True)
elif agent.compression_enabled:
# Over threshold but compression blocked (cooldown/anti-thrash):
# deduped warning so context can't silently overflow. (#62625)
_block_reason = None
_info = getattr(_compressor, "should_compress_info", None)
if _info is not None:
try:
_block_reason = _info(_real_tokens)[1]
except Exception:
_block_reason = None
if _block_reason:
agent._warn_context_overflow_blocked(
_block_reason,
_real_tokens,
int(getattr(_compressor, "threshold_tokens", 0) or 0),
)
# Proactive tool-result prune (deterministic, no LLM, keeps tail):
# no-op unless proactive_prune_tokens is exceeded; commits only past
# proactive_prune_min_reclaim_tokens so cache breaks stay episodic.
_prune = getattr(_compressor, "prune_tool_results_only", None)
if callable(_prune):
try:
_pruned_msgs, _pruned_n = _prune(
messages, current_tokens=_real_tokens
)
except Exception:
logger.debug(
"proactive tool-result prune failed; skipping",
exc_info=True,
)
_pruned_msgs, _pruned_n = messages, 0
# Standard no-op caller contract: only commit when the
# engine returned a NEW list object with a non-zero count.
if _pruned_n and _pruned_msgs is not messages:
# Do NOT rebuild conversation_history: rows already carry
# _DB_PERSISTED_MARKER, and on a stale in-place flag the
# helper could seed unpersisted rows into history_ids.
messages = _pruned_msgs
return _verdict(False)
File diff suppressed because it is too large Load Diff
+15 -45
View File
@@ -1,28 +1,8 @@
"""Per-attempt recovery bookkeeping for the conversation turn loop.
"""Per-attempt recovery bookkeeping (``TurnRetryState``) for the conversation turn loop.
The inner retry loop in ``run_conversation`` (``while retry_count <
max_retries``) makes several distinct recovery attempts on a single model API
call: a credential-pool 429 retry, a per-provider OAuth refresh (codex,
anthropic, nous, copilot), a long-context compression restart, a length-
continuation restart, and a handful of format-recovery branches (thinking-
signature stripping, multimodal-tool-content stripping, llama.cpp grammar
fallback, image shrink, invalid-encrypted-content, 1M-beta header).
Each of those branches is guarded by a one-shot boolean so it fires at most
once per attempt. They used to be ~16 bare ``*_attempted`` / ``has_retried_*``
/ ``restart_with_*`` locals declared inline before the loop and threaded
through its 2,400-line body. ``TurnRetryState`` collapses them into one object
the loop mutates in place (``state.codex_auth_retry_attempted = True``), giving
the recovery bookkeeping a single named, testable home.
Loop-control variables (``retry_count``, ``max_retries``,
``max_compression_attempts``) intentionally stay as plain locals — they are the
``while`` mechanics, not recovery bookkeeping, and putting them on the object
would add indirection without clarifying anything.
This module is dependency-free so it can be unit-tested in isolation and
imported by the turn loop without an import cycle.
"""
Each one-shot recovery branch of the inner retry loop is guarded by a flag here so it
fires at most once per attempt. Loop-control (``retry_count``, ``max_retries``) stays
as plain locals. Dependency-free so it imports without a cycle."""
from __future__ import annotations
@@ -33,11 +13,8 @@ from dataclasses import dataclass, fields
class TurnRetryState:
"""One-shot recovery guards + restart signals for a single API-call attempt.
A fresh instance is created for each iteration of the outer turn loop
(once per ``api_call_count``). Each guard fires its recovery branch at most
once; the ``restart_with_*`` signals are read by the loop after the attempt
to decide whether to rebuild the request and retry.
"""
A fresh instance is created per ``api_call_count`` iteration; each guard fires at
most once, and ``restart_with_*`` signals tell the loop to rebuild and retry."""
# ── Per-provider OAuth / credential refresh guards ───────────────────
codex_auth_retry_attempted: bool = False
@@ -46,12 +23,8 @@ class TurnRetryState:
nous_paid_entitlement_refresh_attempted: bool = False
copilot_auth_retry_attempted: bool = False
# Copilot surfaces a stale/degraded credential as a 400
# ``model_not_available_for_integrator`` / ``model_not_supported`` instead
# of a clean 401 (e.g. a raw OAuth token seeded when the token exchange
# degraded at startup, routing the request to the restricted
# ``copilot-language-server`` integrator). Guard a single-shot forced
# re-exchange + client rebuild for that case, separate from the 401 guard
# so both can fire within one attempt if needed.
# ``model_not_available_for_integrator`` / ``model_not_supported``, not a 401.
# Single-shot forced re-exchange + rebuild, separate from the 401 guard.
copilot_stale_cred_retry_attempted: bool = False
vertex_auth_retry_attempted: bool = False
@@ -69,22 +42,19 @@ class TurnRetryState:
has_retried_429: bool = False
# ── Auth-failure provider failover ───────────────────────────────────
# Set once we've escalated a persistent 401/403 (after the per-provider
# credential-refresh attempt above failed) to the fallback chain, so we
# don't loop on the same auth failover within one attempt.
# Set once a persistent 401/403 has been escalated to the fallback chain, so
# we don't loop on the same auth failover within one attempt.
auth_failover_attempted: bool = False
# ── Restart signals (read by the outer loop after the attempt) ───────
restart_with_compressed_messages: bool = False
restart_with_length_continuation: bool = False
# Set when a content-filter stream stall (e.g. MiniMax "new_sensitive")
# has been escalated to the fallback chain: the partial-stream content
# was rolled back off ``messages`` and the loop should re-issue the API
# call against the newly-activated provider (#32421).
# Set when a content-filter stream stall (e.g. MiniMax "new_sensitive") was
# escalated to the fallback chain: partial content was rolled back off
# ``messages``; re-issue the call against the new provider (#32421).
restart_with_rebuilt_messages: bool = False
# A user correction cancelled the in-flight provider request. The outer
# loop must append a role-safe checkpoint + user message, rebuild the API
# payload, and retry the same logical iteration.
# A user correction cancelled the in-flight request: append a role-safe checkpoint +
# user message, rebuild the payload, and retry the same logical iteration.
restart_with_redirected_messages: bool = False
def __iter__(self):
+207
View File
@@ -0,0 +1,207 @@
"""Text-response stop gates for the conversation turn loop.
Extracted from ``run_conversation``. When the model stops with a text answer, three
gates may instead append the answer as an interim row plus a synthetic user-role nudge
and continue the turn: verify-on-stop (#65919), the ``pre_verify`` plugin hook after code
edits, and the kanban worker terminal-tool guard. Each keeps the candidate answer as a
budget-exhaustion fallback (``pending_verification_response``) and clears
``final_response`` so the finalizer can tell this gate from error exits (#61631).
Nothing here imports ``agent.conversation_loop`` at module level (cycle).
"""
from __future__ import annotations
import logging
import os
from dataclasses import dataclass
from typing import Any, Dict, List
from agent.message_metadata import append_message
logger = logging.getLogger("agent.conversation_loop")
@dataclass
class StopGateVerdict:
"""``continue_turn`` True → a nudge was appended; re-enter the turn loop with
``final_response=None`` and the pending-verification fields updated."""
continue_turn: bool
final_response: Any
pending_verification_response: Any
pending_verification_response_previewed: Any
def apply_stop_gates(
agent: Any,
final_msg: Dict[str, Any],
*,
final_response: Any,
messages: List[Dict[str, Any]],
conversation_history: Any,
pending_verification_response: Any,
pending_verification_response_previewed: Any,
) -> StopGateVerdict:
"""Run verify-on-stop → pre_verify hook → kanban stop guard, in that order. Nudges
are user-role rows appended only after the assistant answer row, so role alternation
holds. Hook lookups are imported lazily from their origin modules (tests patch them
there)."""
_pending_verification_response = pending_verification_response
_pending_verification_response_previewed = pending_verification_response_previewed
def _verdict(continue_turn: bool) -> StopGateVerdict:
return StopGateVerdict(
continue_turn=continue_turn,
final_response=None if continue_turn else final_response,
pending_verification_response=_pending_verification_response,
pending_verification_response_previewed=_pending_verification_response_previewed,
)
try:
from agent.verification_stop import (
build_verify_on_stop_nudge,
verify_on_stop_enabled,
)
if verify_on_stop_enabled():
_verify_nudge = build_verify_on_stop_nudge(
session_id=getattr(agent, "session_id", None),
changed_paths=getattr(agent, "_turn_file_mutation_paths", set()),
attempts=getattr(agent, "_verification_stop_nudges", 0),
)
else:
_verify_nudge = None
except Exception:
logger.debug("verification stop-loop check failed", exc_info=True)
_verify_nudge = None
if _verify_nudge:
agent._verification_stop_nudges = (
getattr(agent, "_verification_stop_nudges", 0) + 1
)
final_msg["finish_reason"] = "verification_required"
# Real content: persist and emit as interim so the user sees the
# attempted answer; only the nudge is flagged synthetic. (#65919)
agent._emit_interim_assistant_message(final_msg)
append_message(messages, final_msg)
try:
agent._flush_messages_to_session_db(messages, conversation_history)
except Exception:
logger.debug("verify-on-stop interim flush failed", exc_info=True)
append_message(messages, {
"role": "user",
"content": _verify_nudge,
"_verification_stop_synthetic": True,
})
agent._session_messages = messages
# Internal nudge: stay silent on the terminal, debug-log only.
logger.debug("verification stop-loop nudge issued (attempt %d)",
agent._verification_stop_nudges)
# Keep the answer only as a budget-exhaustion fallback; clear
# ``final_response`` so the finalizer can tell this gate from error
# exits. Mark previewed only if the candidate is reused. (#61631)
_pending_verification_response = final_response
_pending_verification_response_previewed = (
agent._interim_content_was_streamed(final_response or "")
)
return _verdict(True)
# pre_verify hook gate: after code edits a registered hook may keep the
# agent going one more turn; no default continuation cost.
_verify_nudge2 = None
_edited = sorted(getattr(agent, "_turn_file_mutation_paths", set()) or [])
_attempt = getattr(agent, "_pre_verify_nudges", 0)
try:
from agent.verify_hooks import max_verify_nudges
from hermes_cli.lifecycle import has_hook
from hermes_cli.plugins import get_pre_verify_continue_message
if _edited and has_hook("pre_verify") and _attempt < max_verify_nudges():
# Posture is fixed for the session — resolve once + cache.
coding = getattr(agent, "_resolved_is_coding", None)
if coding is None:
from agent.coding_context import is_coding_context
coding = bool(is_coding_context(platform=getattr(agent, "platform", "") or ""))
agent._resolved_is_coding = coding
_verify_nudge2 = get_pre_verify_continue_message(
session_id=getattr(agent, "session_id", None) or "",
platform=getattr(agent, "platform", "") or "",
model=getattr(agent, "model", "") or "",
coding=coding,
attempt=_attempt,
final_response=final_response,
changed_paths=_edited,
)
except Exception:
logger.debug("pre_verify hook check failed", exc_info=True)
_verify_nudge2 = None
if _verify_nudge2:
agent._pre_verify_nudges = _attempt + 1
final_msg["finish_reason"] = "verify_hook_continue"
# Real content: persist and emit as interim so the user sees the
# attempted answer; only the nudge is flagged synthetic. (#65919)
agent._emit_interim_assistant_message(final_msg)
append_message(messages, final_msg)
try:
agent._flush_messages_to_session_db(messages, conversation_history)
except Exception:
logger.debug("pre_verify interim flush failed", exc_info=True)
append_message(messages, {
"role": "user",
"content": _verify_nudge2,
"_pre_verify_synthetic": True,
})
agent._session_messages = messages
logger.debug("pre_verify nudge issued (attempt %d)",
agent._pre_verify_nudges)
_pending_verification_response = final_response
_pending_verification_response_previewed = (
agent._interim_content_was_streamed(final_response or "")
)
return _verdict(True)
# ── Kanban worker terminal-tool stop guard ─────────────
# Workers must end with kanban_complete / kanban_block; a narrated stop
# is recorded as protocol_violation, so nudge once or twice first.
try:
from agent.kanban_stop import build_kanban_stop_nudge
_kanban_nudge = build_kanban_stop_nudge(
messages=messages,
attempts=getattr(agent, "_kanban_stop_nudges", 0),
)
except Exception:
logger.debug("kanban stop-loop check failed", exc_info=True)
_kanban_nudge = None
if _kanban_nudge:
agent._kanban_stop_nudges = (
getattr(agent, "_kanban_stop_nudges", 0) + 1
)
final_msg["finish_reason"] = "kanban_terminal_required"
final_msg["_kanban_stop_synthetic"] = True
append_message(messages, final_msg)
append_message(messages, {
"role": "user",
"content": _kanban_nudge,
"_kanban_stop_synthetic": True,
})
agent._session_messages = messages
logger.info(
"kanban stop-loop nudge issued (attempt %d) task=%s",
agent._kanban_stop_nudges,
os.environ.get("HERMES_KANBAN_TASK", ""),
)
agent._emit_status(
"⚠️ Kanban worker tried to exit without "
"kanban_complete/kanban_block — nudging to finish"
)
# Same finalizer contract as verify-on-stop: clear final_response so
# budget exhaustion doesn't treat the narrated stop as an answer.
_pending_verification_response = final_response
_pending_verification_response_previewed = (
agent._interim_content_was_streamed(final_response or "")
)
return _verdict(True)
return _verdict(False)
+248
View File
@@ -0,0 +1,248 @@
"""Tool-call validation for the conversation turn loop: unknown tool names (with
auto-repair and the 3-strike partial exit) and malformed JSON arguments (retry, then
recovery tool results).
Extracted from ``run_conversation``. Role alternation is preserved on every path: an
invalid batch is answered with tool-role error results (never a user message), and
the exits close any open tool-result tail (#48879). Nothing here imports
``agent.conversation_loop`` at module level (cycle); loop-internal helpers resolve lazily.
"""
from __future__ import annotations
import json
import logging
from dataclasses import dataclass
from typing import Any, Dict, List, Optional
from agent.message_metadata import append_message
from agent.message_sanitization import close_interrupted_tool_sequence, coalesce_tool_call_id
logger = logging.getLogger("agent.conversation_loop")
@dataclass
class ToolValidationVerdict:
"""Outcome of ``validate_tool_calls``.
``action``: ``"ok"`` (dispatch the calls), ``"continue"`` (re-issue the API call —
error results / retry state were recorded) or ``"return"`` (terminal partial
result in ``result``). ``mixed_invalid_batch`` is True when the batch contains BOTH
valid and unknown tool names: only the invalid calls get error results, the valid
ones run."""
action: str
result: Optional[Dict[str, Any]]
mixed_invalid_batch: bool
def validate_tool_calls(
agent: Any,
assistant_message: Any,
finish_reason: str,
*,
messages: List[Dict[str, Any]],
conversation_history: Any,
api_call_count: int,
effective_task_id: Any,
) -> ToolValidationVerdict:
"""Validate ``assistant_message.tool_calls`` in place (ids uniquified, names
repaired, dict/empty args normalized to JSON strings). Strikes for invalid names
advance only when a turn has NO valid call, so a degenerate model still halts at
3; args cut off mid-stream (routers rewrite ``length`` → ``tool_calls``) are refused
outright rather than retried."""
from agent.conversation_loop import _invalid_tool_name_error_content
_mixed_invalid_batch = False
def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> ToolValidationVerdict:
return ToolValidationVerdict(action=action, result=result, mixed_invalid_batch=_mixed_invalid_batch)
# Uniquify duplicate tool-call ids BEFORE any downstream consumer: the
# pre-API sanitizer keeps only the first call/result per id. See
# _uniquify_tool_call_ids.
agent._uniquify_tool_call_ids(assistant_message.tool_calls)
# Validate tool call names - detect model hallucinations
# Repair mismatched tool names before validating
for tc in assistant_message.tool_calls:
if tc.function.name not in agent.valid_tool_names:
repaired = agent._repair_tool_call(tc.function.name)
if repaired:
print(f"{agent.log_prefix}🔧 Auto-repaired tool name: '{tc.function.name}' -> '{repaired}'")
tc.function.name = repaired
invalid_tool_calls = [
tc.function.name for tc in assistant_message.tool_calls
if tc.function.name not in agent.valid_tool_names
]
# Mixed batch: error-result ONLY the invalid calls and run the valid
# ones; voiding the turn discards real work. Strikes advance only when a
# turn has NO valid call, so a degenerate model still halts at 3.
_mixed_invalid_batch = bool(invalid_tool_calls) and any(
tc.function.name in agent.valid_tool_names
for tc in assistant_message.tool_calls
)
if _mixed_invalid_batch:
agent._invalid_tool_retries = 0
invalid_name = invalid_tool_calls[0]
invalid_preview = invalid_name[:80] + "..." if len(invalid_name) > 80 else invalid_name
_n_valid = sum(
1 for tc in assistant_message.tool_calls
if tc.function.name in agent.valid_tool_names
)
agent._buffer_vprint(
f"⚠️ Unknown tool '{invalid_preview}' in batch — erroring that call, "
f"executing {_n_valid} valid call(s)"
)
elif invalid_tool_calls:
# Track retries for invalid tool calls
agent._invalid_tool_retries += 1
# Return helpful error to model — model can agent-correct next turn
invalid_name = invalid_tool_calls[0]
invalid_preview = invalid_name[:80] + "..." if len(invalid_name) > 80 else invalid_name
agent._buffer_vprint(f"⚠️ Unknown tool '{invalid_preview}' — sending error to model for agent-correction ({agent._invalid_tool_retries}/3)")
if agent._invalid_tool_retries >= 3:
agent._flush_status_buffer()
agent._vprint(f"{agent.log_prefix}❌ Max retries (3) for invalid tool calls exceeded. Stopping as partial.", force=True)
agent._invalid_tool_retries = 0
_final_response = f"Model generated invalid tool call: {invalid_preview}"
# Prior retries or an earlier tool batch leave a tool-result
# tail; close it as interrupt aborts do so the next turn is not
# tool→user. (#48879)
close_interrupted_tool_sequence(messages, _final_response)
agent._persist_session(messages, conversation_history)
return _verdict("return", {
"final_response": _final_response,
"messages": messages,
"api_calls": api_call_count,
"completed": False,
"partial": True,
"error": _final_response
})
assistant_msg = agent._build_assistant_message(assistant_message, finish_reason)
append_message(messages, assistant_msg)
for tc in assistant_message.tool_calls:
_tc_name = tc.function.name
if _tc_name not in agent.valid_tool_names:
# See _invalid_tool_name_error_content for the
# blank-name anti-priming rationale (#47967).
content = _invalid_tool_name_error_content(
_tc_name, agent.valid_tool_names
)
else:
content = "Skipped: another tool call in this turn used an invalid name. Please retry this tool call."
append_message(messages, {
"role": "tool",
"name": tc.function.name,
"tool_call_id": coalesce_tool_call_id(tc),
"content": content,
})
return _verdict("continue")
# Reset retry counter on successful tool call validation
agent._invalid_tool_retries = 0
# Validate tool call arguments are valid JSON
# Handle empty strings as empty objects (common model quirk)
invalid_json_args = []
for tc in assistant_message.tool_calls:
args = tc.function.arguments
if isinstance(args, (dict, list)):
tc.function.arguments = json.dumps(args)
continue
if args is not None and not isinstance(args, str):
tc.function.arguments = str(args)
args = tc.function.arguments
# Treat empty/whitespace strings as empty object
if not args or not args.strip():
tc.function.arguments = "{}"
continue
try:
json.loads(args)
except json.JSONDecodeError as e:
if (
_mixed_invalid_batch
and tc.function.name not in agent.valid_tool_names
):
# This call never executes (invalid-name error result
# below); don't let its broken args trigger the whole-turn
# JSON retry.
continue
invalid_json_args.append((tc.function.name, str(e)))
if invalid_json_args:
# Routers may rewrite finish_reason "length" → "tool_calls", hiding
# truncation; args not ending in } or ] (stripped) were cut off
# mid-stream.
_truncated = any(
not (tc.function.arguments or "").rstrip().endswith(("}", "]"))
for tc in assistant_message.tool_calls
if tc.function.name in {n for n, _ in invalid_json_args}
)
if _truncated:
agent._vprint(
f"{agent.log_prefix}⚠️ Truncated tool call arguments detected "
f"(finish_reason={finish_reason!r}) — refusing to execute.",
force=True,
)
agent._invalid_json_retries = 0
agent._cleanup_task_resources(effective_task_id)
_final_response = "Response truncated due to output length limit"
# Same tool-tail close as interrupt / invalid-tool
# exhaustion — this path never reaches finalize_turn.
close_interrupted_tool_sequence(messages, _final_response)
agent._persist_session(messages, conversation_history)
return _verdict("return", {
"final_response": _final_response,
"messages": messages,
"api_calls": api_call_count,
"completed": False,
"partial": True,
"error": _final_response,
})
# Track retries for invalid JSON arguments
agent._invalid_json_retries += 1
tool_name, error_msg = invalid_json_args[0]
agent._buffer_vprint(f"⚠️ Invalid JSON in tool call arguments for '{tool_name}': {error_msg}")
if agent._invalid_json_retries < 3:
agent._buffer_vprint(f"🔄 Retrying API call ({agent._invalid_json_retries}/3)...")
# Don't add anything to messages, just retry the API call
return _verdict("continue")
else:
# Instead of returning partial, inject tool error results so the model can recover.
# Using tool results (not user messages) preserves role alternation.
agent._buffer_vprint("⚠️ Injecting recovery tool results for invalid JSON...")
agent._invalid_json_retries = 0 # Reset for next attempt
# Append the assistant message with its (broken) tool_calls
recovery_assistant = agent._build_assistant_message(assistant_message, finish_reason)
append_message(messages, recovery_assistant)
# Respond with tool error results for each tool call
invalid_names = {name for name, _ in invalid_json_args}
for tc in assistant_message.tool_calls:
if tc.function.name in invalid_names:
err = next(e for n, e in invalid_json_args if n == tc.function.name)
tool_result = (
f"Error: Invalid JSON arguments. {err}. "
f"For tools with no required parameters, use an empty object: {{}}. "
f"Please retry with valid JSON."
)
else:
tool_result = "Skipped: other tool call in this response had invalid JSON."
append_message(messages, {
"role": "tool",
"name": tc.function.name,
"tool_call_id": coalesce_tool_call_id(tc),
"content": tool_result,
})
return _verdict("continue")
# Reset retry counter on successful JSON validation
agent._invalid_json_retries = 0
return _verdict("ok")
+748
View File
@@ -0,0 +1,748 @@
"""Truncation recovery (``finish_reason == "length"``) for the conversation turn loop.
Extracted from ``run_conversation``. Handles thinking-budget exhaustion, repetition-
dominated truncation (#86581), content-filter stream stalls escalated to the fallback
chain (#32421), text continuation nudges (up to 4, with the ceiling exit that drops the
fragment trail), truncated tool-call retries with max_tokens boosts, and the final
roll-back. Nothing here imports ``agent.conversation_loop`` at module level (cycle);
loop-internal helpers are imported lazily so tests patching them on the loop keep working.
"""
from __future__ import annotations
import logging
import re
from dataclasses import dataclass
from typing import Any, Dict, List, Optional
from agent.error_classifier import FailoverReason
from agent.message_metadata import append_message
from agent.message_sanitization import close_interrupted_tool_sequence
from agent.repetition_guard import is_repetition_dominated
from agent.turn_retry_state import TurnRetryState
from hermes_constants import PARTIAL_STREAM_STUB_ID
logger = logging.getLogger("agent.conversation_loop")
@dataclass
class TruncationVerdict:
"""Outcome of ``recover_from_truncation``.
``action``: ``"return"`` (end the turn with ``result``), ``"break"`` (a
``_retry.restart_with_*`` flag is set — restart the API call), ``"continue"``
(re-issue the same call immediately) or ``"fallthrough"`` (unreachable in practice:
every path exits, kept for the contract). The remaining fields are the loop locals
the handler may have rebound."""
action: str
result: Optional[Dict[str, Any]]
messages: List[Dict[str, Any]]
length_continue_retries: int
truncated_response_parts: List[str]
truncated_tool_call_retries: int
retry_count: int
compression_attempts: int
def recover_from_truncation(
agent: Any,
response: Any,
finish_reason: str,
_retry: TurnRetryState,
*,
messages: List[Dict[str, Any]],
conversation_history: Any,
api_kwargs: Any,
api_call_count: int,
effective_task_id: Any,
current_turn_user_idx: Any,
length_continue_retries: int,
truncated_response_parts: List[str],
truncated_tool_call_retries: int,
retry_count: int,
compression_attempts: int,
) -> TruncationVerdict:
"""Recover from a truncated response. Order is load-bearing: thinking exhaustion and
repetition abort BEFORE any continuation; a content-filter stall escalates to the
fallback chain BEFORE the primary is retried; text continuation (no tool calls) then
truncated tool-call retry; finally roll back to the last complete assistant turn.
Never appends an interim assistant row with NO visible content (strict providers
reject it with 400) — only the continuation nudge."""
from agent.conversation_loop import _get_continuation_prompt, _join_truncated_parts
def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> TruncationVerdict:
return TruncationVerdict(
action=action,
result=result,
messages=messages,
length_continue_retries=length_continue_retries,
truncated_response_parts=truncated_response_parts,
truncated_tool_call_retries=truncated_tool_call_retries,
retry_count=retry_count,
compression_attempts=compression_attempts,
)
if getattr(response, "id", "") == PARTIAL_STREAM_STUB_ID:
agent._vprint(
f"{agent.log_prefix}⚠️ Response truncated — stream "
f"ended before completion",
force=True,
)
else:
agent._vprint(
f"{agent.log_prefix}⚠️ Response truncated "
f"(finish_reason='length') - model hit max output tokens",
force=True,
)
# Normalize to one OpenAI-style message so continuation and tool-
# call retry work across transports (Anthropic reuses the loop's
# adapter).
_trunc_msg = None
_trunc_transport = agent._get_transport()
if agent.api_mode == "anthropic_messages":
_trunc_result = _trunc_transport.normalize_response(
response, strip_tool_prefix=agent._is_anthropic_oauth
)
else:
_trunc_result = _trunc_transport.normalize_response(response)
_trunc_msg = _trunc_result
_trunc_content = getattr(_trunc_msg, "content", None) if _trunc_msg else None
_trunc_has_tool_calls = bool(getattr(_trunc_msg, "tool_calls", None)) if _trunc_msg else False
# ── Detect thinking-budget exhaustion ──────────────
# Only when reasoning blocks exist with no visible text after them;
# content=None from non-<think> models is normal truncation.
_has_think_tags = bool(
_trunc_content and re.search(
r'<(?:think|thinking|reasoning|REASONING_SCRATCHPAD)[^>]*>',
_trunc_content,
re.IGNORECASE,
)
)
_thinking_exhausted = (
not _trunc_has_tool_calls
and _has_think_tags
and (
(_trunc_content is not None and not agent._has_content_after_think_block(_trunc_content))
or _trunc_content is None
)
)
if _thinking_exhausted:
_exhaust_error = (
"Model used all output tokens on reasoning with none left "
"for the response. Try lowering reasoning effort or "
"increasing max_tokens."
)
agent._vprint(
f"{agent.log_prefix}💭 Reasoning exhausted the output token budget — "
f"no visible response was produced.",
force=True,
)
# Return a user-friendly message as the response so CLI and
# gateway display it.
_exhaust_response = (
"⚠️ **Thinking Budget Exhausted**\n\n"
"The model used all its output tokens on reasoning "
"and had none left for the actual response.\n\n"
"To fix this:\n"
"→ Lower reasoning effort: `/reasoning low` or `/reasoning minimal`\n"
"→ Or switch to a larger/non-reasoning model with `/model`"
)
agent._cleanup_task_resources(effective_task_id)
agent._persist_session(messages, conversation_history)
return _verdict("return", {
"final_response": _exhaust_response,
"messages": messages,
"api_calls": api_call_count,
"completed": False,
"partial": True,
"error": _exhaust_error,
})
# ── Detect repetition-dominated truncation (#86581) ──
# A repetition loop can burn the whole budget on one fragment; abort
# like _thinking_exhausted (reasoning stripped first).
_visible_trunc = (
agent._strip_think_blocks(_trunc_content)
if isinstance(_trunc_content, str)
else _trunc_content
)
_repetition_dominated = (
not _trunc_has_tool_calls
and bool(_visible_trunc)
and is_repetition_dominated(_visible_trunc)
)
if _repetition_dominated:
_rep_error = (
"Model output entered a repetition loop and was "
"truncated mid-loop; refusing to continue a "
"degenerate response."
)
agent._vprint(
f"{agent.log_prefix}🔁 Response dominated by "
f"repeated text — stopping instead of "
f"continuing a degenerate response.",
force=True,
)
_rep_response = (
"⚠️ **Response Stopped — Repetition Detected**\n\n"
"The model fell into a repetition loop while "
"writing this response, so continuing would only "
"produce more repeated text. The partial response "
"was discarded.\n\n"
"→ Switch to a different model with `/model`\n"
"→ Or resend your message (your conversation "
"history is preserved)"
)
agent._cleanup_task_resources(effective_task_id)
agent._persist_session(messages, conversation_history)
return _verdict("return", {
"final_response": _rep_response,
"messages": messages,
"api_calls": api_call_count,
"completed": False,
"partial": True,
"error": _rep_error,
})
if agent.api_mode in {"chat_completions", "bedrock_converse", "anthropic_messages"}:
assistant_message = _trunc_msg
# ── Content-filter stream stall → fallback (#32421) ──
# ``_content_filter_terminated`` is content-deterministic;
# escalate to the fallback before retrying the primary.
_cf_terminated = getattr(
response, "_content_filter_terminated", False
)
if (
_cf_terminated
and agent._fallback_index < len(agent._fallback_chain)
):
agent._vprint(
f"{agent.log_prefix}🛡️ Content filter terminated "
f"stream — activating fallback provider...",
force=True,
)
agent._emit_status(
"Content filter terminated stream; switching to fallback..."
)
if agent._try_activate_fallback():
# Roll partial content back to the last clean turn so
# the fallback gets a coherent continuation point.
if truncated_response_parts:
messages = agent._get_messages_up_to_last_assistant(messages)
# Unmark survivors: their text left the stitched partial.
for _frag in messages:
if isinstance(_frag, dict):
_frag.pop("_length_continuation_fragment", None)
_frag.pop("_length_continuation_nudge", None)
agent._session_messages = messages
length_continue_retries = 0
truncated_response_parts = []
retry_count = 0
compression_attempts = 0
_retry.primary_recovery_attempted = False
_retry.restart_with_rebuilt_messages = True
return _verdict("break")
# No fallback available — fall through to normal
# continuation (best-effort, may loop).
agent._vprint(
f"{agent.log_prefix}⚠️ No fallback provider "
f"configured — retrying with same provider "
f"(may re-hit filter)...",
force=True,
)
if assistant_message is not None and not _trunc_has_tool_calls:
length_continue_retries += 1
# Never append an interim assistant message with NO visible
# content: strict providers reject it (HTTP 400), poisoning
# history. Append only the nudge.
_interim_content = getattr(assistant_message, "content", None)
_is_empty_partial_stub = (
getattr(response, "id", "") == PARTIAL_STREAM_STUB_ID
and not _interim_content
)
if not _interim_content and not _is_empty_partial_stub:
# Thinking-only truncation: continuing with thinking ON
# re-burns the budget, so drop thinking for one request.
agent._ephemeral_reasoning_off = True
if _interim_content:
interim_msg = agent._build_assistant_message(assistant_message, finish_reason)
# Marked so the ceiling exit can drop the fragment trail.
interim_msg["_length_continuation_fragment"] = True
append_message(messages, interim_msg)
truncated_response_parts.append(_interim_content)
if length_continue_retries < 4:
_is_partial_stream_stub = (
getattr(response, "id", "") == PARTIAL_STREAM_STUB_ID
)
_dropped_tools = getattr(
response, "_dropped_tool_names", None
)
if _is_partial_stream_stub and _dropped_tools:
_tool_list = ", ".join(_dropped_tools[:3])
agent._vprint(
f"{agent.log_prefix}↻ Stream interrupted mid "
f"tool-call ({_tool_list}) — requesting "
f"chunked retry "
f"({length_continue_retries}/4)..."
)
elif _is_partial_stream_stub:
agent._vprint(
f"{agent.log_prefix}↻ Stream interrupted — "
f"requesting continuation "
f"({length_continue_retries}/4)..."
)
else:
agent._vprint(
f"{agent.log_prefix}↻ Requesting continuation "
f"({length_continue_retries}/4)..."
)
_continue_content = _get_continuation_prompt(
_is_partial_stream_stub, _dropped_tools
)
continue_msg = {
"role": "user",
"content": _continue_content,
"_length_continuation_nudge": True,
}
append_message(messages, continue_msg)
agent._session_messages = messages
_retry.restart_with_length_continuation = True
return _verdict("break")
partial_response = agent._strip_think_blocks(_join_truncated_parts(truncated_response_parts)).strip()
# The one-shot reasoning-off override must not leak into the
# next turn when the ceiling exit skips the consuming call.
agent._ephemeral_reasoning_off = False
if partial_response:
agent._vprint(
f"{agent.log_prefix}⚠️ Response still truncated "
f"after {length_continue_retries} continuation attempts — keeping the "
f"partial response received so far.",
force=True,
)
_ceiling_final = partial_response
else:
# Every fragment was empty (e.g. reasoning-only model):
# return an actionable message, not a bare None.
agent._vprint(
f"{agent.log_prefix}⚠️ Response still truncated "
f"after {length_continue_retries} continuation attempts — no visible "
f"text was produced.",
force=True,
)
_ceiling_final = (
"⚠️ **No visible answer was produced.** The "
"model hit its output-token limit on every "
"continuation attempt — its reasoning "
"consumed the entire budget each time.\n\n"
"To fix this:\n"
"→ Lower reasoning effort: `/reasoning low` "
"or `/reasoning none`\n"
"→ Or raise max_tokens for this model"
)
# Unanswered continue nudges made every later turn re-truncate.
_turn_start = (
current_turn_user_idx + 1
if isinstance(current_turn_user_idx, int)
and current_turn_user_idx >= 0
else 0
)
messages[_turn_start:] = [
m for m in messages[_turn_start:]
if not (
isinstance(m, dict)
and (
m.get("_length_continuation_fragment")
or m.get("_length_continuation_nudge")
)
)
]
if partial_response:
append_message(messages, {
"role": "assistant",
"content": partial_response,
"finish_reason": "length",
})
agent._session_messages = messages
agent._cleanup_task_resources(effective_task_id)
agent._persist_session(messages, conversation_history)
return _verdict("return", {
"final_response": _ceiling_final,
"messages": messages,
"api_calls": api_call_count,
"completed": False,
"partial": True,
"error": "Response remained truncated after 4 continuation attempts",
})
if agent.api_mode in {"chat_completions", "bedrock_converse", "anthropic_messages"}:
assistant_message = _trunc_msg
if assistant_message is not None and _trunc_has_tool_calls:
_is_stub_stall = (
getattr(response, "id", "") == PARTIAL_STREAM_STUB_ID
)
if truncated_tool_call_retries < 4:
truncated_tool_call_retries += 1
if _is_stub_stall:
# Stream broke mid tool-call (network), not a real
# output cap — say so.
agent._buffer_vprint(
f"⚠️ Stream interrupted mid tool-call — "
f"retrying ({truncated_tool_call_retries}/4)..."
)
else:
agent._buffer_vprint(
f"⚠️ Truncated tool call detected — "
f"retrying API call "
f"({truncated_tool_call_retries}/4)..."
)
# Boost max_tokens per retry: a real output-cap
# truncation needs it; harmless for a stall.
_tc_boost_base = agent.max_tokens if agent.max_tokens else 4096
_tc_boost = _tc_boost_base * (2 ** truncated_tool_call_retries)
_tc_requested_cap = agent._requested_output_cap_from_api_kwargs(api_kwargs)
if _tc_requested_cap is not None:
_tc_boost = max(_tc_boost, _tc_requested_cap)
_tc_boost_cap = max(32768, _tc_requested_cap or 0)
agent._ephemeral_max_output_tokens = min(_tc_boost, _tc_boost_cap)
# Don't append the broken response; re-run the same call
# from current state.
return _verdict("continue")
agent._flush_status_buffer()
if _is_stub_stall:
agent._vprint(
f"{agent.log_prefix}⚠️ Stream kept dropping mid tool-call after 4 retries — the action was not executed.",
force=True,
)
else:
agent._vprint(
f"{agent.log_prefix}⚠️ Truncated tool call response detected again — refusing to execute incomplete tool arguments.",
force=True,
)
agent._cleanup_task_resources(effective_task_id)
_final_response = (
"Stream repeatedly dropped mid tool-call (network); "
"the tool was not executed"
if _is_stub_stall
else "Response truncated due to output length limit"
)
# Prior tool batches can leave a tool-result tail; this path
# never reaches finalize_turn (#48879).
close_interrupted_tool_sequence(messages, _final_response)
agent._persist_session(messages, conversation_history)
return _verdict("return", {
"final_response": _final_response,
"messages": messages,
"api_calls": api_call_count,
"completed": False,
"partial": True,
"error": _final_response,
})
# If we have prior messages, roll back to last complete state
if len(messages) > 1:
agent._vprint(f"{agent.log_prefix} ⏪ Rolling back to last complete assistant turn")
rolled_back_messages = agent._get_messages_up_to_last_assistant(messages)
agent._cleanup_task_resources(effective_task_id)
agent._persist_session(messages, conversation_history)
return _verdict("return", {
"final_response": "Response truncated due to output length limit",
"messages": rolled_back_messages,
"api_calls": api_call_count,
"completed": False,
"partial": True,
"error": "Response truncated due to output length limit"
})
else:
# First message was truncated - mark as failed
agent._flush_status_buffer()
agent._vprint(f"{agent.log_prefix}❌ First response truncated - cannot recover", force=True)
agent._persist_session(messages, conversation_history)
return _verdict("return", {
"final_response": "First response truncated due to output length limit",
"messages": messages,
"api_calls": api_call_count,
"completed": False,
"failed": True,
"error": "First response truncated due to output length limit"
})
return _verdict("fallthrough")
def continue_codex_incomplete(
agent: Any,
assistant_message: Any,
finish_reason: str,
*,
messages: List[Dict[str, Any]],
conversation_history: Any,
api_call_count: int,
) -> Optional[Dict[str, Any]]:
"""Codex Responses ``status=incomplete`` continuation (max 3 per turn).
Appends the interim assistant message (deduped on visible content only — opaque
provider state drifts per continuation, #52711; ``codex_reasoning_items`` are merged,
not overwritten, because the earlier response holds the only native-compaction
checkpoint) and, when a bare retry would be byte-identical, a user-role nudge — only
after an assistant row, to preserve role alternation. Returns ``None`` to continue
the turn loop, or the terminal ``partial`` result once retries are exhausted."""
from agent.conversation_loop import _CODEX_INCOMPLETE_NUDGE
agent._codex_incomplete_retries += 1
interim_msg = agent._build_assistant_message(assistant_message, finish_reason)
interim_has_content = bool((interim_msg.get("content") or "").strip())
interim_has_reasoning = bool(interim_msg.get("reasoning", "").strip()) if isinstance(interim_msg.get("reasoning"), str) else False
interim_has_codex_reasoning = bool(interim_msg.get("codex_reasoning_items"))
interim_has_codex_message_items = bool(interim_msg.get("codex_message_items"))
if (
interim_has_content
or interim_has_reasoning
or interim_has_codex_reasoning
or interim_has_codex_message_items
):
last_msg = messages[-1] if messages else None
# Dedup on visible content only (content + reasoning): opaque
# provider state drifts per continuation and would defeat dedup
# (#52711).
last_interim_visible = (
agent._interim_assistant_visible_text(last_msg)
if isinstance(last_msg, dict)
else ""
)
current_interim_visible = agent._interim_assistant_visible_text(interim_msg)
if last_interim_visible or current_interim_visible:
same_visible_output = last_interim_visible == current_interim_visible
else:
# Preserve the existing reasoning-only behavior when
# neither response has text eligible for interim delivery.
same_visible_output = (
(last_msg.get("content") or "") == (interim_msg.get("content") or "")
and (last_msg.get("reasoning") or "") == (interim_msg.get("reasoning") or "")
) if isinstance(last_msg, dict) else False
visible_duplicate = (
isinstance(last_msg, dict)
and last_msg.get("role") == "assistant"
and last_msg.get("finish_reason") == "incomplete"
and same_visible_output
)
if visible_duplicate:
# Update replay state in-place: keep the latest provider payload
# without re-emitting identical user-visible commentary.
for _key in (
"content",
"reasoning",
"reasoning_content",
"reasoning_details",
"codex_reasoning_items",
"codex_message_items",
):
if _key in interim_msg:
if _key == "codex_reasoning_items":
# Merge, don't overwrite: the earlier response's
# native compaction checkpoint is the only copy. See
# merge_interim_reasoning_items.
from agent.native_compaction import (
merge_interim_reasoning_items,
)
last_msg[_key] = merge_interim_reasoning_items(
last_msg.get(_key), interim_msg[_key]
)
else:
last_msg[_key] = interim_msg[_key]
else:
append_message(messages, interim_msg)
agent._emit_interim_assistant_message(interim_msg)
if agent._codex_incomplete_retries < 3:
# If the interim has nothing the Responses converter will replay, a
# bare retry is byte-identical and fails identically; append a
# user-role nudge so the retry differs and asks for the answer.
interim_replayable = (
interim_has_content
or interim_has_codex_reasoning
or interim_has_codex_message_items
)
# Replayable ≠ different: an interim holding only a ``compaction``
# checkpoint in ``codex_reasoning_items`` is replayable yet re-sends
# identically. One bare retry, then always nudge.
if not interim_replayable or agent._codex_incomplete_retries >= 2:
_last_msg = messages[-1] if messages else None
_already_nudged = (
isinstance(_last_msg, dict)
and _last_msg.get("role") == "user"
and _last_msg.get("content") == _CODEX_INCOMPLETE_NUDGE
)
# Alternation guard: the user-role nudge may only follow an
# assistant message; after a too-empty interim it would create
# user→user / tool→user.
_last_is_assistant = (
isinstance(_last_msg, dict)
and _last_msg.get("role") == "assistant"
)
if not _already_nudged and _last_is_assistant:
append_message(messages, {
"role": "user",
"content": _CODEX_INCOMPLETE_NUDGE,
})
if not agent.quiet_mode:
agent._vprint(f"{agent.log_prefix}↻ Codex response incomplete; continuing turn ({agent._codex_incomplete_retries}/3)")
# Show the continuation on the spinner/status line and gateway
# heartbeat; these retries can take minutes and otherwise look like
# infinite thinking (#64434).
agent._emit_wait_notice(
f"↻ model returned reasoning with no final answer — "
f"asking it to continue "
f"({agent._codex_incomplete_retries}/3)"
)
agent._session_messages = messages
return None
agent._codex_incomplete_retries = 0
agent._persist_session(messages, conversation_history)
return {
"final_response": "Codex response remained incomplete after 3 continuation attempts",
"messages": messages,
"api_calls": api_call_count,
"completed": False,
"partial": True,
"error": "Codex response remained incomplete after 3 continuation attempts",
}
@dataclass
class RefusalVerdict:
"""Outcome of ``handle_content_policy_refusal``: ``"break"`` (fallback activated —
restart armed on ``_retry``; caller resets retry/compression counters) or
``"return"`` (the typed content-policy result in ``result``). ``active_system_prompt``
is the possibly re-synced system prompt."""
action: str
result: Optional[Dict[str, Any]]
active_system_prompt: Any
def handle_content_policy_refusal(
agent: Any,
response: Any,
_retry: TurnRetryState,
*,
thinking_spinner: Any,
messages: List[Dict[str, Any]],
api_messages: Any,
api_kwargs: Any,
active_system_prompt: Any,
conversation_history: Any,
api_call_count: int,
effective_task_id: Any,
turn_id: Any,
api_request_id: Any,
api_start_time: float,
retry_count: int,
max_retries: int,
) -> RefusalVerdict:
"""HTTP-200 refusal (``finish_reason`` ``content_filter`` / ``guardrail_intervened``).
Deterministic for the unchanged prompt — never retried: one configured-fallback try,
else surface the refusal (explanation may live only in the reasoning channel). The
caller stops its spinner reference; this stops the spinner object."""
from agent.conversation_loop import (
_CONTENT_POLICY_RECOVERY_HINT,
_arm_fallback_restart,
_content_policy_blocked_result,
)
def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> RefusalVerdict:
return RefusalVerdict(action=action, result=result, active_system_prompt=active_system_prompt)
_refusal_transport = agent._get_transport()
if agent.api_mode == "anthropic_messages":
_refusal_result = _refusal_transport.normalize_response(
response, strip_tool_prefix=agent._is_anthropic_oauth
)
else:
_refusal_result = _refusal_transport.normalize_response(response)
_refusal_text = (getattr(_refusal_result, "content", None) or "").strip()
# Some refusals carry the explanation only in the reasoning
# channel; fall back to it so the user sees *something*.
if not _refusal_text:
_refusal_text = (agent._extract_reasoning(_refusal_result) or "").strip()
agent._invoke_api_request_error_hook(
task_id=effective_task_id,
turn_id=turn_id,
api_request_id=api_request_id,
api_call_count=api_call_count,
api_start_time=api_start_time,
api_kwargs=api_kwargs,
error_type="ContentPolicyBlocked",
error_message=_refusal_text or "model declined to respond (content_filter)",
status_code=None,
retry_count=retry_count,
max_retries=max_retries,
retryable=False,
reason=FailoverReason.content_policy_blocked.value,
)
if thinking_spinner:
thinking_spinner.stop("")
if agent.thinking_callback:
agent.thinking_callback("")
# Deterministic for the unchanged prompt — never retry. Try a
# configured fallback once; otherwise surface the refusal.
if agent._has_pending_fallback():
agent._buffer_status(
"⚠️ Model declined to respond (safety refusal) — trying fallback..."
)
if agent._try_activate_fallback():
active_system_prompt = _arm_fallback_restart(
agent, api_messages, active_system_prompt, _retry)
return _verdict("break")
agent._flush_status_buffer()
_refusal_log = (
_refusal_text[:500] + "..."
if len(_refusal_text) > 500
else _refusal_text
)
logger.warning(
"%sModel declined to respond (finish_reason=content_filter). "
"model=%s provider=%s refusal=%s",
agent.log_prefix, agent.model, agent.provider,
_refusal_log or "(no text)",
)
agent._emit_status(
"⚠️ The model declined to respond to this request (safety refusal)."
)
_refusal_detail = (
f"Model's explanation: {_refusal_text}"
if _refusal_text
else "The model returned no explanation."
)
_refusal_response = (
"⚠️ The model declined to respond to this request "
"(safety refusal — not a Hermes/gateway failure).\n\n"
f"{_refusal_detail}\n\n"
f"{_CONTENT_POLICY_RECOVERY_HINT}"
)
agent._cleanup_task_resources(effective_task_id)
agent._persist_session(messages, conversation_history)
return _verdict("return", _content_policy_blocked_result(
messages,
api_call_count,
final_response=_refusal_response,
error_detail=_refusal_text or "model declined (content_filter)",
))
+312
View File
@@ -0,0 +1,312 @@
"""Per-response usage accounting for the conversation turn loop.
After every successful model API call, ``run_conversation`` folds the provider's
``response.usage`` into: the context compressor (``update_from_response`` + the
compression-budget rearm latch), the usage anchor for display/compression math,
per-session token/cost counters, the state.db token-delta queue, and the
observability log line. MoA sessions additionally fold advisor fan-out usage into
the reported counts and price the aggregator at its REAL model/provider.
``record_response_usage`` owns that block. It mutates ``agent`` exactly as the inline
code did and returns the loop-visible verdict (compression budget counter, and
whether a provider-confirmed recovery rearmed it) as ``ResponseUsageOutcome``.
Logger name stays ``agent.conversation_loop`` for caplog/log-routing parity.
"""
from __future__ import annotations
import logging
from dataclasses import dataclass
from typing import Any, Dict, List
from agent.model_metadata import capture_usage_anchor
from agent.usage_pricing import estimate_usage_cost, normalize_usage
logger = logging.getLogger("agent.conversation_loop")
@dataclass
class ResponseUsageOutcome:
"""What the loop reads back after usage accounting.
``compression_attempts`` is the (possibly rearmed-to-zero) budget counter;
``rearmed`` tells the loop to also clear its preflight-block latch."""
compression_attempts: int
rearmed: bool = False
def _loop_mod():
"""Lazy ``agent.conversation_loop`` so tests patching
``agent.conversation_loop.save_context_length`` still intercept, and so this
module never imports the loop at load time (cycle)."""
import agent.conversation_loop as _cl
return _cl
def record_response_usage(
agent: Any,
response: Any,
*,
messages: List[Dict[str, Any]],
api_call_count: int,
api_duration: float,
compression_attempts: int,
max_compression_attempts: int,
) -> ResponseUsageOutcome:
"""Fold ``response.usage`` into compressor, anchors, session counters, state.db
and the API-call log line (see module docstring). No-usage responses only
consume a pending compaction verdict. Returns the loop-visible outcome."""
rearmed = False
# Track actual token usage from response for context management
if hasattr(response, 'usage') and response.usage:
canonical_usage = normalize_usage(
response.usage,
provider=agent.provider,
api_mode=agent.api_mode,
)
# Aggregator-only usage kept for pricing: advisor tokens are priced
# at each advisor's OWN model rate and added as dollars below.
aggregator_usage = canonical_usage
# MoA: fold advisor fan-out usage into REPORTED token counts — only
# aggregator usage is returned, so advisor spend would be invisible.
_moa_ref_cost = None
_moa_client = getattr(agent, "client", None)
if _moa_client is not None and hasattr(_moa_client, "consume_reference_usage"):
try:
_ref_usage, _moa_ref_cost = _moa_client.consume_reference_usage()
if _ref_usage is not None:
canonical_usage = canonical_usage + _ref_usage
except Exception as _moa_acct_exc: # pragma: no cover - defensive
logger.debug("MoA reference usage accounting failed: %s", _moa_acct_exc)
# Flush the full-turn MoA trace when moa.save_traces is on; on the
# streaming path pass the streamed acting text so the trace is self-
# contained.
if _moa_client is not None and hasattr(_moa_client, "consume_and_save_trace"):
try:
_agg_streamed_text = (
getattr(agent, "_current_streamed_assistant_text", "") or ""
)
_moa_client.consume_and_save_trace(
agent.session_id,
aggregator_output_fallback=_agg_streamed_text or None,
)
except Exception as _moa_trace_exc: # pragma: no cover - defensive
logger.debug("MoA trace flush failed: %s", _moa_trace_exc)
prompt_tokens = canonical_usage.prompt_tokens
completion_tokens = canonical_usage.output_tokens
total_tokens = canonical_usage.total_tokens
# Forward canonical token + cache buckets for context engines;
# legacy keys stay for back-compat.
usage_dict = {
"prompt_tokens": prompt_tokens,
"completion_tokens": completion_tokens,
"total_tokens": total_tokens,
"input_tokens": canonical_usage.input_tokens,
"output_tokens": canonical_usage.output_tokens,
"cache_read_tokens": canonical_usage.cache_read_tokens,
"cache_write_tokens": canonical_usage.cache_write_tokens,
"reasoning_tokens": canonical_usage.reasoning_tokens,
}
# Capture the boundary latch before update_from_response() consumes
# it: only the real prompt count right after a compaction rearms the
# budget.
_completed_compaction_pending = bool(
getattr(
agent.context_compressor,
"_verify_compaction_cleared_threshold",
False,
)
)
agent.context_compressor.update_from_response(usage_dict)
# Usage-anchored accounting: snapshot exact provider usage against
# the durable transcript; main-loop ONLY. MoA uses pre-fold
# aggregator usage.
_new_anchor = capture_usage_anchor(
aggregator_usage.prompt_tokens,
aggregator_usage.output_tokens,
messages,
)
if _new_anchor is not None:
agent._usage_anchor = _new_anchor
# Anchor the display meter on the turn's FIRST response:
# later same-turn responses inflate prompt_tokens with replayed
# thinking. Display-only; compression math uses real usage.
if api_call_count == 1:
agent._turn_base_usage_anchor = _new_anchor
_compression_threshold = int(
getattr(agent.context_compressor, "threshold_tokens", 0)
or 0
)
if _loop_mod()._should_rearm_compression_budget(
compression_attempts,
completed_compaction_pending=_completed_compaction_pending,
prompt_tokens=prompt_tokens,
threshold_tokens=_compression_threshold,
):
logger.info(
"Compression budget rearmed after provider-confirmed "
"recovery: prompt=%s < threshold=%s (attempts were %s/%s)",
f"{prompt_tokens:,}",
f"{_compression_threshold:,}",
compression_attempts,
max_compression_attempts,
)
compression_attempts = 0
# Confirmed recovery also clears the loop's stale insufficient-progress
# verdict (``_preflight_compression_blocked``), else it stays armed all
# turn and a later pressure spike grows unchecked.
rearmed = True
# Stash canonical usage for on_turn_complete() (same shape as
# update_from_response); keep the latest call's — last request.
agent._last_turn_usage = dict(usage_dict)
elif getattr(
agent.context_compressor,
"awaiting_real_usage_after_compression",
False,
):
# No usage -> cannot adjudicate the prior compaction; consume the
# pending verdict so later readings aren't charged to it and
# preflight deferral isn't latched indefinitely.
agent.context_compressor.update_from_response({})
if hasattr(response, 'usage') and response.usage:
# Persist only provider-confirmed context lengths, not probe tiers.
if getattr(agent.context_compressor, "_context_probed", False):
ctx = agent.context_compressor.context_length
if getattr(agent.context_compressor, "_context_probe_persistable", False):
_loop_mod().save_context_length(agent.model, agent.base_url, ctx)
agent._safe_print(f"{agent.log_prefix}💾 Cached context length: {ctx:,} tokens for {agent.model}")
agent.context_compressor._context_probed = False
agent.context_compressor._context_probe_persistable = False
agent.session_prompt_tokens += prompt_tokens
agent.session_completion_tokens += completion_tokens
agent.session_total_tokens += total_tokens
agent.session_api_calls += 1
agent.session_input_tokens += canonical_usage.input_tokens
agent.session_output_tokens += canonical_usage.output_tokens
agent.session_cache_read_tokens += canonical_usage.cache_read_tokens
agent.session_cache_write_tokens += canonical_usage.cache_write_tokens
agent.session_reasoning_tokens += canonical_usage.reasoning_tokens
# Rolling history for status-bar averages (last 10).
try:
hist = getattr(agent, "_api_latency_history", None)
if hist is not None:
hist.append(float(api_duration))
ohist = getattr(agent, "_api_output_history", None)
if ohist is not None:
ohist.append(int(canonical_usage.output_tokens or 0))
except Exception:
pass
# Log API call details for debugging/observability
_cache_pct = ""
if canonical_usage.cache_read_tokens and prompt_tokens:
_cache_pct = f" cache={canonical_usage.cache_read_tokens}/{prompt_tokens} ({100*canonical_usage.cache_read_tokens/prompt_tokens:.0f}%)"
logger.info(
"API call #%d: model=%s provider=%s in=%d out=%d total=%d latency=%.1fs%s",
agent.session_api_calls, agent.model, agent.provider or "unknown",
prompt_tokens, completion_tokens, total_tokens,
api_duration, _cache_pct,
)
# MoA: agent.model/provider are the virtual preset/"moa" with no
# pricing entry, silently dropping aggregator spend. Price at the
# REAL model/provider from the MoA client's aggregator slot.
_agg_cost_model = agent.model
_agg_cost_provider = agent.provider
_agg_cost_base_url = agent.base_url
_agg_slot = getattr(_moa_client, "last_aggregator_slot", None) if _moa_client is not None else None
if _agg_slot and _agg_slot.get("model"):
_agg_cost_model = _agg_slot["model"]
_agg_cost_provider = _agg_slot.get("provider") or agent.provider
_agg_cost_base_url = _agg_slot.get("base_url") or agent.base_url
cost_result = estimate_usage_cost(
_agg_cost_model,
aggregator_usage,
provider=_agg_cost_provider,
base_url=_agg_cost_base_url,
api_key=getattr(agent, "api_key", ""),
)
if cost_result.amount_usd is not None:
agent.session_estimated_cost_usd += float(cost_result.amount_usd)
# Add MoA advisor cost (already priced per-advisor at each
# advisor's own model rate) on top of the aggregator cost.
if _moa_ref_cost is not None:
try:
agent.session_estimated_cost_usd += float(_moa_ref_cost)
except (TypeError, ValueError): # pragma: no cover - defensive
pass
agent.session_cost_status = cost_result.status
agent.session_cost_source = cost_result.source
# Persist per-call token deltas for any session_id so non-CLI runs
# can't lose accounting; gateway/session-store writes use absolute
# totals and safely overwrite these deltas.
if agent._session_db and agent.session_id:
try:
# Ensure the row exists: under concurrent SQLite load the
# initial _ensure_db_session() may fail, and UPDATE on a
# missing row silently affects 0 rows.
if not agent._session_db_created:
agent._ensure_db_session()
# Cost delta = aggregator + MoA advisor cost so state.db's
# estimated_cost_usd matches the folded token counts.
_cost_delta = None
if cost_result.amount_usd is not None:
_cost_delta = float(cost_result.amount_usd)
if _moa_ref_cost is not None:
try:
_cost_delta = (_cost_delta or 0.0) + float(_moa_ref_cost)
except (TypeError, ValueError): # pragma: no cover
pass
# Enqueued, not written: a cold state.db UPDATE here stalled
# the tool loop. Drained at finalize via _persist_session.
agent._session_db.queue_token_counts(
agent.session_id,
input_tokens=canonical_usage.input_tokens,
output_tokens=canonical_usage.output_tokens,
cache_read_tokens=canonical_usage.cache_read_tokens,
cache_write_tokens=canonical_usage.cache_write_tokens,
reasoning_tokens=canonical_usage.reasoning_tokens,
estimated_cost_usd=_cost_delta,
cost_status=cost_result.status,
cost_source=cost_result.source,
billing_provider=agent.provider,
billing_base_url=agent.base_url,
billing_mode="subscription_included"
if cost_result.status == "included" else None,
model=agent.model,
api_call_count=1,
)
except Exception as e:
# Log failures — silent loss here undercounts analytics.
logger.debug(
"Token persistence failed (session=%s, tokens=%d): %s",
agent.session_id, total_tokens, e,
)
if agent.verbose_logging:
logging.debug(f"Token usage: prompt={usage_dict['prompt_tokens']:,}, completion={usage_dict['completion_tokens']:,}, total={usage_dict['total_tokens']:,}")
# Report cache stats for any provider that returns
# ``prompt_tokens_details.cached_tokens``, not only when we inject
# cache_control markers. ``canonical_usage`` is already normalised.
cached = canonical_usage.cache_read_tokens
written = canonical_usage.cache_write_tokens
prompt = usage_dict["prompt_tokens"]
if (cached or written) and not agent.quiet_mode:
hit_pct = (cached / prompt * 100) if prompt > 0 else 0
agent._vprint(
f"{agent.log_prefix} 💾 Cache: "
f"{cached:,}/{prompt:,} tokens "
f"({hit_pct:.0f}% hit, {written:,} written)"
)
return ResponseUsageOutcome(
compression_attempts=compression_attempts,
rearmed=rearmed,
)
@@ -197,7 +197,7 @@ class TestConversationLoopWiring:
def test_413_branch_scores_bytes_not_tokens(self):
import inspect
import agent.conversation_loop as loop
import agent.turn_overflow as loop # 413 handler lives in recover_from_overflow
src = inspect.getsource(loop)
# The byte measurement is taken before and after the 413 compression
+6 -6
View File
@@ -1,5 +1,5 @@
"""Tests for the Nous OAuth 401 actionable-guidance branch in
``agent.conversation_loop.run_conversation``.
``agent.turn_recovery.nonretryable_client_error_result`` (the run_conversation terminal branch).
Source-inspection style (matches ``test_gemini_fast_fallback.py``): we assert
that the guidance strings exist in the function body so that the user-facing
@@ -16,14 +16,14 @@ from __future__ import annotations
import inspect
from agent import conversation_loop
from agent import turn_recovery
def test_nous_provider_is_in_oauth_401_set():
"""The provider-set gate that selects OAuth-specific guidance must
include ``nous`` alongside ``openai-codex`` and ``xai-oauth``.
"""
source = inspect.getsource(conversation_loop.run_conversation)
source = inspect.getsource(turn_recovery.nonretryable_client_error_result)
# Be flexible about set element ordering — assert all three are listed
# near each other in the gating expression.
@@ -33,16 +33,16 @@ def test_nous_provider_is_in_oauth_401_set():
# And the gate string itself must mention all three so future refactors
# that split nous off into its own gate still get caught.
needle = "_provider in {\"openai-codex\", \"xai-oauth\", \"nous\"}"
needle = "provider in {\"openai-codex\", \"xai-oauth\", \"nous\"}"
assert needle in source, (
"Expected nous to be co-gated with the other OAuth providers in the "
"actionable-401-guidance branch of run_conversation."
"actionable-401-guidance branch of nonretryable_client_error_result."
)
def test_nous_401_guidance_strings_present():
"""User-facing remediation strings for Nous OAuth 401s must exist."""
source = inspect.getsource(conversation_loop.run_conversation)
source = inspect.getsource(turn_recovery.nonretryable_client_error_result)
# Must tell the user it's an OAuth token problem, NOT an API key problem
# (Nous Portal has no API key path — auth_type=oauth_device_code only).
@@ -172,12 +172,19 @@ class TestSendPathBuildIsWiredToTheClone:
import ast
import inspect
source = inspect.getsource(cl)
tree = ast.parse(source)
import agent.turn_context as tc
# The history build lives in turn_context.build_api_messages; the
# prefill insert stays in conversation_loop. Scan both homes.
nodes = [
node
for mod in (cl, tc)
for node in ast.walk(ast.parse(inspect.getsource(mod)))
]
clone_calls = []
shallow_copies = []
for node in ast.walk(tree):
for node in nodes:
if isinstance(node, ast.Call):
fn = node.func
if isinstance(fn, ast.Name) and fn.id == "_clone_message_for_send":
@@ -447,12 +447,12 @@ class TestCanonicalHistoryIsolation:
import inspect
import re as _re
import agent.conversation_loop as loop_mod
import agent.turn_recovery as loop_mod
src = inspect.getsource(loop_mod)
src = inspect.getsource(loop_mod.recover_after_classification)
# Locate the image_corrupt recovery block and inspect its calls.
block = _re.search(
r"image_corrupt:\n(.*?)\n\s*(?:continue|else)", src, _re.S
r"image_corrupt:\n(.*?)\n\s*(?:continue|return|else)", src, _re.S
)
assert block is not None, "image_corrupt recovery branch not found"
body = block.group(1)
@@ -22,7 +22,7 @@ import sys
from types import SimpleNamespace
from agent.conversation_loop import _image_error_max_dimension
from agent.turn_recovery import _image_error_max_dimension
from agent.error_classifier import FailoverReason, classify_api_error