Files
hermes-agent/agent/conversation_loop.py
T

2315 lines
100 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""The agent conversation loop — extracted from ``run_agent.AIAgent``.
``run_conversation(agent, ...)`` drives one user turn (model call, tool dispatch,
retries, fallbacks, compression, post-turn hooks). Symbols that callers patch on
``run_agent`` (``handle_function_call``, ``_set_interrupt``, ``OpenAI``) resolve via
``_ra`` so those patches keep working."""
from __future__ import annotations
import json
import logging
import random
import re
import time
from typing import Any, Dict, List, Optional
from agent.codex_responses_adapter import _summarize_user_message_for_log
from agent.conversation_compression import (
conversation_history_after_compression, # noqa: F401 — resolved lazily by turn_overflow/turn_preflight/turn_recovery (tests patch it here)
)
from agent.display import KawaiiSpinner
from agent.fast_mode import begin_turn as begin_fast_mode_turn
from agent.message_metadata import append_message
from agent.turn_context import (
PreflightCompressionTimedOut,
build_turn_context,
reanchor_current_turn_user_idx,
)
from agent.turn_retry_state import TurnRetryState
from agent.runtime_cwd import resolve_agent_cwd
from agent.message_sanitization import (
_repair_tool_call_arguments,
_sanitize_surrogates,
)
# Must mirror _STALE_TOOL_CALL_MARKER_RE in hermes_state.py; kept local so importing
# hermes_state (module-level DEFAULT_DB_PATH) is not forced at load time.
_STALE_MARKER_RE = re.compile(r"^\[[A-Za-z_][A-Za-z0-9_.-]*\]$")
from agent.model_metadata import (
MINIMUM_CONTEXT_LENGTH,
_estimate_tools_tokens_rough,
estimate_messages_tokens_rough, # noqa: F401 — resolved lazily by agent.turn_request_assembly (tests patch it here)
estimate_request_tokens_rough, # noqa: F401 — resolved lazily by turn_overflow/turn_preflight/turn_recovery (tests patch it here)
save_context_length, # noqa: F401 — resolved lazily by agent.turn_overflow (tests patch it here)
)
from agent.process_bootstrap import _install_safe_stdio
from agent.prompt_caching import (
build_prompt_cache_plan,
effective_cache_ttl,
strip_anthropic_cache_control,
strip_anthropic_tool_cache_control,
)
from agent.retry_utils import ( # noqa: F401 — resolved lazily by agent.turn_* (tests patch them here)
adaptive_rate_limit_backoff,
jittered_backoff,
)
from agent.turn_recovery import ( # noqa: F401 — resolved lazily by agent.turn_response_check
describe_invalid_response,
interruptible_backoff_sleep,
validate_response_shape,
)
# Bind before the turn starts so a source-tree swap cannot load a skewed
# finalizer at turn end.
from agent.turn_finalizer import finalize_turn
from agent.turn_iteration_prep import prepare_iteration
from agent.turn_preflight_gate import run_preflight_gate
from agent.turn_request_assembly import assemble_api_request
from agent.turn_api_request import build_api_request
from agent.turn_api_call import perform_api_call
from agent.turn_response_check import check_api_response
from agent.turn_api_error import handle_api_error
from agent.turn_final_response import finish_text_response
from agent.turn_tool_round import run_tool_round
from agent.turn_response_intake import normalize_model_response
from agent.turn_loop_errors import handle_outer_loop_error
from hermes_logging import set_session_context
from tools.skill_provenance import set_current_write_origin
from utils import base_url_host_matches
logger = logging.getLogger(__name__)
# Scaffold marker used by _apply_active_turn_redirect and the ghost-row filter
# in the api_messages loop. Module-level so both sites can never drift.
_INTERRUPT_SCAFFOLD_MARKER = "[This response was interrupted by a user correction.]"
# One-time wrap-up notice appended when a wall-clock run budget crosses 80%
# (agent.run_budget_seconds / --run-budget): stop new work, deliver current state.
RUN_BUDGET_WRAPUP_NOTICE = (
"[SYSTEM NOTICE — run time budget nearly exhausted] "
"Run time budget nearly exhausted. Stop new discovery/verification work "
"now. Produce the required final deliverable (answer/JSON/summary) from "
"the state you already have, completing only mandatory writes."
)
def _midturn_request_pressure_tokens(
agent: Any,
api_messages: List[Dict[str, Any]],
effective_system: str,
approx_tokens: int,
) -> int:
"""Token figure the mid-turn pre-API compression guard compares.
Returns the pruned native-Responses estimate when native compaction eligibility is
proven (the generic estimate overstates the wire on compacted sessions, #96995),
else the generic message+tools figure. System prompt is counted exactly once."""
try:
from agent.codex_responses_adapter import (
estimate_native_responses_preflight_tokens,
)
native = estimate_native_responses_preflight_tokens(
agent,
api_messages,
system_prompt=effective_system or "",
tools=getattr(agent, "tools", None) or None,
)
if isinstance(native, int) and not isinstance(native, bool) and native >= 0:
return native
except Exception:
logger.debug(
"native Responses mid-turn estimate unavailable; "
"using generic transcript estimate",
exc_info=True,
)
return approx_tokens + (
_estimate_tools_tokens_rough(agent.tools) if agent.tools else 0
)
def _review_input_budget_exhausted(agent: Any) -> bool:
"""True when a detached review fork has replayed its aggregate input budget.
Only forks with an explicit ``_review_input_token_budget`` are gated (#93057). Fires
at the top of the NEXT iteration, so the budget-crossing request completes first."""
budget = getattr(agent, "_review_input_token_budget", None)
if not isinstance(budget, int) or isinstance(budget, bool) or budget <= 0:
return False
used = getattr(agent, "session_input_tokens", 0)
return isinstance(used, int) and not isinstance(used, bool) and used >= budget
def _maybe_inject_run_budget_wrapup(agent: Any, messages: List[Dict[str, Any]]) -> bool:
"""Inject the one-time wall-clock wrap-up notice when past 80% of budget.
Appends to the NEWEST ``role:"tool"`` message (cache-safe, like /steer); latches
``_run_budget_wrapup_injected`` only on a successful append. Returns True when
injected. Dormant unless ``run_budget_seconds`` + ``_run_budget_started_at`` set."""
budget = getattr(agent, "run_budget_seconds", None)
if not budget:
return False
if getattr(agent, "_run_budget_wrapup_injected", False):
return False
started = getattr(agent, "_run_budget_started_at", None)
if not started:
return False
if (time.time() - started) < 0.8 * float(budget):
return False
for i in range(len(messages) - 1, -1, -1):
msg = messages[i]
if isinstance(msg, dict) and msg.get("role") == "tool":
existing = msg.get("content", "")
if isinstance(existing, str):
msg["content"] = existing + f"\n\n{RUN_BUDGET_WRAPUP_NOTICE}"
else:
# Multimodal content blocks — append a text block.
try:
blocks = list(existing) if existing else []
blocks.append({"type": "text", "text": RUN_BUDGET_WRAPUP_NOTICE})
msg["content"] = blocks
except Exception:
return False
agent._run_budget_wrapup_injected = True
logger.info(
"Run budget wrap-up notice injected (budget=%.0fs, elapsed=%.0fs)",
float(budget),
time.time() - started,
)
return True
return False
def _restore_user_after_reference_handoff(
messages: List[Dict[str, Any]], user_message: Any
) -> bool:
"""Re-append this turn's real user ask when compaction left only a handoff.
Returns True when a restore append happened; only decides whether a restorable
ask exists (#80622)."""
if user_message is None:
return False
if isinstance(user_message, str):
if not user_message.strip():
return False
content: Any = user_message
elif isinstance(user_message, list):
if not user_message:
return False
content = user_message
else:
return False
if (
messages
and isinstance(messages[-1], dict)
and messages[-1].get("role") == "user"
and messages[-1].get("content") == content
):
return False
append_message(messages, {"role": "user", "content": content})
return True
def _should_skip_model_call_for_reference_handoff(
messages: List[Dict[str, Any]], user_message: Any
) -> bool:
"""Guard post-compaction continues against sole-handoff active turns (#80622)."""
from agent.context_compressor import reference_handoff_would_drive_next_model_call
if not reference_handoff_would_drive_next_model_call(messages):
return False
if _restore_user_after_reference_handoff(messages, user_message):
# The restored ask is an actionable non-synthetic user row appended
# after the handoff — by construction the handoff no longer drives.
return False
return True
# Fallback final_response for the sole-handoff skip (#80622). Not a replay of the
# last assistant text: finalize_turn appends final_response as a fresh assistant row.
_HANDOFF_SKIP_FINAL_RESPONSE = (
"Context was compacted. The previous response is complete — "
"awaiting your next message."
)
# Terminal final_response when compression hit its host timeout while the request
# was still oversized; resending would only bounce off the overflow error (#98722).
_COMPRESSION_TIMEOUT_FINAL_RESPONSE = (
"Context compression timed out without reducing this conversation. "
"No messages were dropped. Start a fresh session with /new, or check "
"auxiliary.compression before retrying /compress."
)
# Stable prefix of the local interrupt status string; surfaces (ACP, TUI) match on
# it to treat the text as cancellation metadata rather than assistant prose.
INTERRUPT_WAITING_FOR_MODEL_PREFIX = "Operation interrupted: waiting for model response ("
def _should_rearm_compression_budget(
compression_attempts: int,
*,
completed_compaction_pending: bool,
prompt_tokens: int,
threshold_tokens: int,
) -> bool:
"""Return True after a provider proves a completed compaction worked.
Rough estimates cannot rearm the anti-thrash budget; require the completed-
compaction latch and a positive normalized prompt count below the threshold."""
return bool(
compression_attempts
and completed_compaction_pending
and threshold_tokens > 0
and 0 < prompt_tokens < threshold_tokens
)
# Modules whose presence in a traceback (without any API-call module) marks a
# deterministic local bug not worth retrying. NEVER add "conversation_loop" or
# "run_agent": every exception passes through them; _hit_local would be True (#66267)
_LOCAL_PROCESSING_MODULES = frozenset({
"agent_runtime_helpers",
"message_content",
"message_sanitization",
"chat_completion_helpers", # only local when NOT also an API-call module
})
_API_CALL_MODULES = frozenset({
"chat_completion_helpers",
})
# Max outer-loop exceptions per user turn before giving up; only exceptions that
# ESCAPE the inner retry/fallback machinery count, so this can be small (#92450).
_MAX_OUTER_LOOP_ERRORS = 8
def _is_interpreter_shutdown_error(exc: Exception) -> bool:
"""Check if *exc* is a fatal interpreter-shutdown failure.
Delegates to ``tools.interpreter_shutdown`` (one text-matching site for the
shutdown-race bug class) but keeps the RuntimeError type gate: a ValueError
carrying similar text must not match (#93269)."""
if isinstance(exc, RuntimeError):
from tools.interpreter_shutdown import interpreter_shutting_down
return interpreter_shutting_down(exc)
return False
def _moa_client_consumes_prepared_request(client: Any) -> bool:
"""True when ``client`` is the in-process MoA facade.
Only ``MoAChatCompletions`` exposes ``prepare()``; other clients raise TypeError on
``_moa_prepared_request`` even while ``agent.provider`` stays ``"moa"``."""
completions = getattr(getattr(client, "chat", None), "completions", None)
return callable(getattr(completions, "prepare", None))
def _join_truncated_parts(parts: List[str]) -> str:
"""Join continuation fragments, adding a newline where two would glue together (#78577)."""
joined = ""
for part in parts:
if joined and not joined[-1].isspace() and part and not part[0].isspace():
joined += "\n"
joined += part
return joined
def _moa_reference_metrics_for_hook(agent: Any) -> Any:
"""Per-advisor metrics for post_api_request, or None off the MoA path.
MoA returns only the aggregator response, so a plugin sees one generation for
the whole fan-out; this carries the per-slot advisor spend across the hook boundary."""
client = getattr(agent, "client", None)
getter = getattr(client, "last_reference_metrics", None)
if not callable(getter):
return None
try:
return getter()
except Exception:
return None
def _apply_active_turn_redirect(agent: Any, messages: List[Dict[str, Any]], text: str) -> None:
"""Append a provider-safe checkpoint and correction to the live turn.
Keeps only the *visible* text (demoted to plain text) then adds the correction as a
real user message, so role alternation holds and cached messages stay byte-identical.
INVARIANT: raw chain-of-thought never enters replayable content — inlined CoT reads
as a prefill jailbreak and bricks the session with empty-response storms.
INVARIANT: the interruption scaffold is replay text, carried only in the user
correction's ``api_content``; an on-screen-empty placeholder is ``display_kind=hidden``."""
visible = agent._strip_think_blocks(
getattr(agent, "_current_streamed_assistant_text", "") or ""
).strip()
checkpoint_parts = [_INTERRUPT_SCAFFOLD_MARKER]
if visible:
checkpoint_parts.extend(
["Visible response before the interruption:", visible]
)
checkpoint = "\n\n".join(checkpoint_parts)
correction = (
"[Context from the interrupted assistant response]\n"
f"{checkpoint}\n\n"
f"{text}"
)
# The live tail is normally user or tool, so an assistant placeholder + correction
# keeps strict alternation; if the tail is already assistant, fold the checkpoint
# into the user correction instead of creating assistant→assistant.
if messages and messages[-1].get("role") == "assistant":
# Transcript shows the user's own words; the provider replays the
# scaffolded form so it still sees the interrupted context.
append_message(
messages,
{"role": "user", "content": text, "api_content": correction},
)
else:
# Placeholder preserves role alternation only. Scaffold bytes must never land
# here: api_content is substituted back into content on replay (#81841).
placeholder: Dict[str, Any] = {
"role": "assistant",
"content": visible or "",
}
if not visible:
placeholder["display_kind"] = "hidden"
# Hidden row, but a non-empty neutral api_content so the pre-call
# sanitizer does not re-heal it every call (#88955). Never
# _INTERRUPT_SCAFFOLD_MARKER: as assistant text the model echoes it (#81841)
from agent.agent_runtime_helpers import _INTERRUPTED_PLACEHOLDER
placeholder["api_content"] = _INTERRUPTED_PLACEHOLDER
append_message(messages, placeholder)
append_message(
messages,
{"role": "user", "content": text, "api_content": correction},
)
agent._current_streamed_assistant_text = ""
agent._stream_needs_break = True
def _is_copilot_provider(agent: Any) -> bool:
"""Delegate to ``AIAgent._is_copilot_provider`` (single owner of the check).
``agent.provider`` may hold the aliases ``github-copilot`` / ``github``; a bare
``provider == "copilot"`` gate would skip credential recovery for them."""
try:
return bool(agent._is_copilot_provider())
except Exception:
return (getattr(agent, "provider", "") or "").strip().lower() in {
"copilot",
"github-copilot",
"github",
}
def _is_stale_copilot_credential_error(status_code: Optional[int], error_message: str) -> bool:
"""Detect a Copilot 400 that is really a STALE / DEGRADED credential.
Matches status 400 AND ``model_not_available_for_integrator`` or
``model_not_supported`` / "the requested model is not supported", so a wrong model
name never triggers the single-shot re-exchange. Caller enforces scoping/guard."""
lowered = (error_message or "").lower()
is_400 = status_code == 400 or "error code: 400" in lowered
if not is_400:
return False
return (
"model_not_available_for_integrator" in lowered
or "not available for integrator" in lowered
or "model_not_supported" in lowered
or "the requested model is not supported" in lowered
)
def _ollama_context_limit_error(agent: Any, request_tokens: int) -> Optional[str]:
"""Return a user-facing error when Ollama is loaded with too little context."""
if not getattr(agent, "tools", None):
return None
runtime_ctx = getattr(agent, "_ollama_num_ctx", None)
if not isinstance(runtime_ctx, int) or runtime_ctx <= 0:
return None
if runtime_ctx >= MINIMUM_CONTEXT_LENGTH:
return None
model = getattr(agent, "model", "") or "the selected model"
base_url = getattr(agent, "base_url", "") or "unknown base URL"
provider = getattr(agent, "provider", "") or "unknown"
tool_count = len(getattr(agent, "tools", None) or [])
logger.warning(
"Ollama runtime context too small for Hermes tool use: "
"model=%s provider=%s base_url=%s runtime_context=%d "
"minimum_context=%d estimated_request_tokens=%d tool_count=%d "
"session=%s",
model,
provider,
base_url,
runtime_ctx,
MINIMUM_CONTEXT_LENGTH,
request_tokens,
tool_count,
getattr(agent, "session_id", None) or "none",
)
return (
f"Ollama loaded `{model}` with only {runtime_ctx:,} tokens of runtime "
f"context, but Hermes needs at least {MINIMUM_CONTEXT_LENGTH:,} tokens "
"for reliable tool use.\n\n"
"Increase the Ollama context for this model and restart/reload the "
"model before trying again. A known-good starting point is 65,536 "
"tokens. In Hermes config, set `model.ollama_num_ctx: 65536` "
"(and `model.context_length: 65536` if you also override the displayed "
"model context). If you manage the model through an Ollama Modelfile, "
"set `PARAMETER num_ctx 65536` there instead."
)
def _maybe_grow_local_window(agent: Any, compressor: Any,
request_tokens: int) -> Optional[int]:
"""Try growing the managed local model's context window before compressing.
Returns the new window when the ladder granted one, else None (hold / at native /
not a managed local session). Cheap for non-local providers: one compare."""
provider = (getattr(agent, "provider", "") or "").strip().lower()
if provider not in ("llamacpp", "llama.cpp", "llama-cpp", "custom"):
return None
base_url = getattr(agent, "base_url", "") or ""
if "127.0.0.1" not in base_url and "localhost" not in base_url:
return None
try:
from hermes_cli.local_runtime.growth import maybe_grow_window
current_window = int(getattr(compressor, "context_length", 0) or 0)
if current_window <= 0:
return None
return maybe_grow_window(
getattr(agent, "model", "") or "",
base_url=base_url,
session_tokens=int(request_tokens),
current_window=current_window,
)
except Exception as exc: # noqa: BLE001 — growth must never break a turn
logger.debug("local window growth check failed: %s", exc)
return None
def _ra():
"""Lazy ``run_agent`` reference so patches on ``run_agent.handle_function_call`` /
``run_agent._set_interrupt`` / ``run_agent.OpenAI`` reach this code path."""
import run_agent
return run_agent
def _nous_entitlement_message(capability: str) -> str:
try:
from hermes_cli.nous_account import (
format_nous_portal_entitlement_message,
get_nous_portal_account_info,
)
account_info = get_nous_portal_account_info(force_fresh=True)
message = format_nous_portal_entitlement_message(
account_info,
capability=capability,
)
return message or ""
except Exception:
return ""
def _print_nous_entitlement_guidance(agent, capability: str) -> bool:
message = _nous_entitlement_message(capability)
if not message:
return False
for line in message.splitlines():
agent._vprint(f"{agent.log_prefix} 💡 {line}", force=True)
return True
def _system_prompt_for_hooks(api_kwargs: Any, request_messages: Any) -> Any:
"""System prompt as actually sent to the provider, for observability hooks.
Checks ``system`` (Anthropic), ``instructions`` (Responses/Codex), then
``messages[0]``. Returns None when the request carries no system prompt."""
system_prompt = api_kwargs.get("system")
if system_prompt is None:
system_prompt = api_kwargs.get("instructions")
if system_prompt is None and isinstance(request_messages, list) and request_messages:
first = request_messages[0]
if isinstance(first, dict) and first.get("role") == "system":
system_prompt = first.get("content")
return system_prompt
def _is_nous_inference_route(provider: str, base_url: str) -> bool:
provider = (provider or "").strip().lower()
if provider == "nous":
return True
base = str(base_url or "")
return (
base_url_host_matches(base, "inference-api.nousresearch.com")
)
def _billing_or_entitlement_message(
*,
capability: str,
provider: str,
base_url: str,
model: str,
unverified: bool = False,
) -> str:
if _is_nous_inference_route(provider, base_url):
return _nous_entitlement_message(capability)
provider_label = (provider or "").strip() or "the selected provider"
model_label = (model or "").strip() or "the selected model"
# Anthropic Pro/Max OAuth surfaces exhaustion of the "extra usage" bucket as a hard
# 400; point at the settings page and cycle reset — "add credits" does not apply.
if (provider or "").strip().lower() == "anthropic":
# ``unverified`` (#82154): the "out of extra usage" 400 is also returned for a
# server-side content-filter rejection, so hedge and name the other cause.
if unverified:
lines = [
(
f"{provider_label} reported that your Claude subscription usage may be "
f"exhausted for {model_label} (included quota + extra-usage credits) — "
"but this specific error is not proof of a billing problem."
),
"If https://claude.ai/settings/usage still shows quota remaining, this is "
"probably NOT a billing problem: on a Claude subscription (OAuth) token "
"Anthropic returns this same message when its content filter rejects part "
"of the request — typically a phrase in the system prompt.",
"If usage really is exhausted: wait for the billing cycle to reset, or add "
"extra usage at https://claude.ai/settings/usage",
"You can also switch to an Anthropic API key or another provider with "
"/model <model> --provider <provider>.",
# The exhaustion latch replays the stored error without issuing
# a request, so a real fix looks like it didn't work.
"Retry with a fresh credential state: `hermes auth reset anthropic`. Until "
"that cooldown clears, this error can be replayed from cache without "
"contacting the API.",
]
else:
lines = [
(
f"{provider_label} reported that your Claude subscription usage is "
f"exhausted for {model_label} (included quota + extra-usage credits)."
),
"Options: wait for the billing cycle to reset, or add extra usage at "
"https://claude.ai/settings/usage",
"You can also switch to an Anthropic API key or another provider with "
"/model <model> --provider <provider>.",
]
return "\n".join(lines)
# Provider-agnostic billing URL so every text surface (CLI, gateway, TUI) shows the
# same actionable link, not just OpenRouter.
try:
from agent.billing_links import build_billing_block
_link = build_billing_block(provider=provider, base_url=base_url, model=model)
if _link.provider_label:
provider_label = _link.provider_label
billing_url = _link.billing_url
except Exception:
billing_url = None
lines = [
(
f"{provider_label} reported that billing, credits, or account "
f"entitlement is exhausted for {model_label}."
),
"Add credits or update billing with that provider, then retry.",
]
if billing_url:
lines.append(f"{provider_label} billing: {billing_url}")
lines.append("You can switch providers temporarily with /model <model> --provider <provider>.")
return "\n".join(lines)
def _billing_block_dict(
provider, base_url, model, message="", *, unverified: bool = False
) -> Optional[dict]:
"""Best-effort structured billing descriptor (None if billing_links is unavailable)."""
try:
from agent.billing_links import build_billing_block
block = build_billing_block(
provider=provider, base_url=str(base_url), model=model, message=message
).to_dict()
except Exception:
return None
if block is not None and unverified:
# Carry the classifier's ambiguity into the structured descriptor so
# every surface rendering the block can hedge too (#82154).
block["unverified"] = True
return block
def _billing_terminal_label(summary: str, unverified: bool) -> str:
"""Terminal-failure prefix for a billing-classified error.
``unverified`` (#82154): the Anthropic "out of extra usage" 400 can be a
content-filter rejection, so the line must not assert exhaustion as fact."""
if unverified:
return (
"Provider reported usage/credit exhaustion (unverified — the same "
f"error can be a content-filter rejection, not billing): {summary}"
)
return f"Billing or credits exhausted: {summary}"
def _billing_failure_result(
*,
classified,
summary: str,
messages,
api_call_count: int,
provider: str,
base_url,
model: str,
guidance: Optional[str] = None,
) -> dict:
"""Structured terminal result for a billing-classified failure.
Single construction point so label, guidance, structured block and ambiguity flag
stay consistent across the non-retryable abort and max-retries paths (#82154)."""
unverified = bool(getattr(classified, "billing_unverified", False))
if guidance is None:
guidance = _billing_or_entitlement_message(
capability="model access",
provider=provider,
base_url=str(base_url),
model=model,
unverified=unverified,
)
final = _billing_terminal_label(summary, unverified)
if guidance:
final += f"\n\n{guidance}"
return {
"final_response": final,
"messages": messages,
"api_calls": api_call_count,
"completed": False,
"failed": True,
"error": summary,
"failure_reason": classified.reason.value,
# Classifier's own retry verdict so UI (agent/error_surface.py) shows Retry
# only when a re-run can differ, not re-derived from a second taxonomy.
"failure_retryable": bool(classified.retryable),
# The billing verdict may rest on an ambiguous body (#82154) — carry
# that through the structured result, not just the prose.
"billing_unverified": unverified,
"billing_block": _billing_block_dict(
provider, base_url, model, guidance, unverified=unverified
),
}
def _print_billing_or_entitlement_guidance(
agent,
*,
capability: str,
provider: str,
base_url: str,
model: str,
unverified: bool = False,
) -> bool:
message = _billing_or_entitlement_message(
capability=capability,
provider=provider,
base_url=base_url,
model=model,
unverified=unverified,
)
if not message:
return False
for line in message.splitlines():
agent._vprint(f"{agent.log_prefix} 💡 {line}", force=True)
return True
def _restore_or_build_system_prompt(agent, system_message, conversation_history):
"""Restore the cached system prompt from the session DB or build it fresh.
Mutates ``agent._cached_system_prompt`` and persists a freshly-built prompt on first
build. Row states ``missing``/``null``/``empty``/``present`` are logged and DB
failures log at WARNING so silent prefix-cache misses show in ``agent.log``."""
stored_prompt = None
stored_state = "missing"
session_row = None
if conversation_history and agent._session_db:
try:
session_row = agent._session_db.get_session(agent.session_id)
if session_row is not None:
raw_prompt = session_row.get("system_prompt")
if raw_prompt is None:
stored_state = "null"
elif raw_prompt == "":
stored_state = "empty"
else:
stored_prompt = raw_prompt
stored_state = "present"
except Exception as exc:
logger.warning(
"Session DB get_session failed for system-prompt restore "
"(session=%s): %s. Falling back to fresh build — prefix "
"cache will miss for this turn.",
agent.session_id, exc,
)
if stored_prompt and _stored_prompt_matches_runtime(agent, stored_prompt):
# Bot Chat capability epoch: the stored prompt embeds a capability fingerprint;
# a mismatch is a deliberate once-per-change rebuild. Unstamped prompts never
# take this branch; probe failures fail closed to "reuse" so cache is kept.
_bot_stale = False
try:
from tools.bot_mode_probe import (
BOT_CHAT_TITLE,
stored_bot_chat_prompt_needs_upgrade,
stored_prompt_capability_stale,
)
_home_for_epoch = None
try:
from agent.system_prompt import _agent_home
_home_for_epoch = _agent_home(agent)
except Exception:
pass
_bot_stale = stored_prompt_capability_stale(stored_prompt, _home_for_epoch)
if not _bot_stale and getattr(agent, "_bot_mode_protocol", True):
# Legacy upgrade: a Bot Chat prompt predating the epoch mechanism gets
# ONE title-gated migration rebuild; the stamped result cannot re-fire.
_t = str(getattr(agent, "_session_title_hint", "") or "").strip()
if not _t and agent._session_db and agent.session_id:
try:
_t = str(agent._session_db.get_session_title(agent.session_id) or "").strip()
except Exception:
_t = ""
if _t == BOT_CHAT_TITLE:
_bot_stale = stored_bot_chat_prompt_needs_upgrade(stored_prompt, _home_for_epoch)
except Exception:
_bot_stale = False
if _bot_stale:
logger.info(
"Bot Chat capability epoch changed for session %s; rebuilding "
"system prompt to adopt the new capability surface (one-time "
"prefix-cache break).",
agent.session_id,
)
agent._session_title_hint = "Bot Chat"
# The skills index cache (LRU + disk snapshot) does not watch the skills
# dir; a capability refresh must rebuild THROUGH it or new skills are lost.
try:
from agent.prompt_builder import clear_skills_system_prompt_cache
clear_skills_system_prompt_cache(clear_snapshot=True)
except Exception:
pass
agent._cached_system_prompt = agent._build_system_prompt(system_message)
# Persist so the NEXT turn restores the new bytes verbatim (cache break is
# once per capability change). on_session_start not re-fired: continuation.
if agent._session_db:
try:
agent._session_db.update_system_prompt(
agent.session_id, agent._cached_system_prompt
)
except Exception as exc:
logger.warning(
"Session DB update_system_prompt failed after Bot Chat "
"capability refresh (session=%s): %s. The refresh will "
"re-fire next turn.",
agent.session_id, exc,
)
return
# Continuing session — reuse the exact system prompt from the
# previous turn so the Anthropic cache prefix matches.
agent._cached_system_prompt = stored_prompt
# Same contract for tools[]: pin the array to the order this session already
# sent (tools freeze) instead of re-probing every check_fn on a fresh AIAgent.
try:
saved_tools = session_row.get("tool_names") if session_row else None
if saved_tools:
from tools.mcp_tool import restore_agent_tool_prefix
restore_agent_tool_prefix(agent, json.loads(saved_tools))
except Exception:
logger.debug("tool prefix restore skipped", exc_info=True)
# Prompt-section callbacks are new-session-only; recover their frozen bytes
# from the persisted prompt so a compression rebuild keeps them.
from agent.system_prompt import restore_plugin_prompt_sections
restore_plugin_prompt_sections(agent, stored_prompt)
# The static prefix is not persisted; rebuild it for the early cache breakpoint
# or fresh-per-turn gateway agents fall back to the single-breakpoint layout.
# reconstruct_static_prefix gates on _use_prompt_caching, fails open to legacy.
from agent.system_prompt import reconstruct_static_prefix
reconstruct_static_prefix(agent, system_message=system_message)
return
if stored_prompt:
stored_state = "stale_runtime"
logger.info(
"Stored system prompt for session %s has stale runtime identity; "
"rebuilding for model=%s provider=%s.",
agent.session_id,
getattr(agent, "model", "") or "",
getattr(agent, "provider", "") or "",
)
if conversation_history and stored_state in ("null", "empty"):
# Continuing session with an unusable stored prompt: every turn now rebuilds
# and the prefix cache misses every time.
logger.warning(
"Stored system prompt for session %s is %s; rebuilding "
"from scratch this turn. Prefix cache will miss until "
"the rebuild persists. Investigate the previous turn's "
"update_system_prompt write path.",
agent.session_id, stored_state,
)
# First turn of a new session (or recovering from a broken stored
# prompt) — build from scratch.
agent._cached_system_prompt = agent._build_system_prompt(system_message)
# Plugin hook: on_session_start — fired once for a brand-new session, not on
# continuation.
try:
from hermes_cli.lifecycle import invoke_hook as _invoke_hook
_invoke_hook(
"on_session_start",
session_id=agent.session_id,
model=agent.model,
platform=getattr(agent, "platform", None) or "",
)
except Exception as exc:
logger.warning("on_session_start hook failed: %s", exc)
# Cold-start credits seed (L3) fallback for the first-turn path; TUI/desktop seed at
# session open, so this is idempotent (skips when _credits_state exists). Fail-open.
try:
from agent.credits_tracker import seed_credits_at_session_start
seed_credits_at_session_start(agent)
except Exception:
logger.debug("cold-start credits seed failed (fail-open)", exc_info=True)
# Persist the system prompt snapshot; the gateway path (fresh AIAgent per turn)
# reads this row every turn, so a failure here breaks prefix-cache reuse.
if agent._session_db:
try:
agent._session_db.update_system_prompt(agent.session_id, agent._cached_system_prompt)
from tools.mcp_tool import persist_agent_tool_names
persist_agent_tool_names(agent)
except Exception as exc:
logger.warning(
"Session DB update_system_prompt failed for session %s: "
"%s. Subsequent turns will rebuild the system prompt and "
"miss the prefix cache.",
agent.session_id, exc,
)
def _stored_prompt_matches_runtime(agent, prompt: str) -> bool:
"""Return False when the persisted runtime-identity lines are stale."""
def line_value(label: str) -> str:
"""Last matching line wins.
Safe ONLY for fields in the volatile tier at the END of the prompt; embedded
project context could shadow earlier fields — see ``host_info_value``."""
prefix = f"{label}:"
value = ""
for line in prompt.splitlines():
if line.startswith(prefix):
value = line[len(prefix):].strip()
return value
def host_info_value(label: str) -> str:
"""Read a field from the prompt's own host-info block.
Anchors on the FIRST ``User home directory:`` line so a user's ``AGENTS.md`` row
cannot match; a false mismatch would rebuild the prompt every turn."""
prefix = f"{label}:"
lines = prompt.splitlines()
for idx, line in enumerate(lines):
if not line.startswith("User home directory:"):
continue
for candidate in lines[idx + 1: idx + 4]:
if candidate.startswith(prefix):
return candidate[len(prefix):].strip()
return ""
stored_model = line_value("Model")
current_model = str(getattr(agent, "model", "") or "").strip()
if stored_model and current_model and stored_model != current_model:
return False
stored_provider = line_value("Provider")
current_provider = str(getattr(agent, "provider", "") or "").strip()
if stored_provider and current_provider and stored_provider != current_provider:
return False
# cwd drift check. Compare against resolve_agent_cwd() — the SAME resolver used to
# build the prompt — so TERMINAL_CWD sessions are not falsely rejected.
stored_cwd = host_info_value("Current working directory")
if stored_cwd:
if stored_cwd != str(resolve_agent_cwd()):
return False
# Runtime-surface drift: reusing a desktop-built prompt on a terminal session (or
# vice versa) would inject the wrong runtime hints.
stored_platform = line_value("Platform")
current_platform = str(getattr(agent, "platform", "") or "").strip()
if stored_platform and current_platform and stored_platform != current_platform:
return False
return True
# Named constants for the _get_continuation_prompt variants so
# _is_synthetic_compression_user_turn can recognize them by content after a crash
# persists one; SessionDB projection strips the _length_continuation_nudge tag.
_LENGTH_CONTINUATION_NETWORK_STUB = (
"[System: The previous response was cut off by a "
"network error mid-stream. Continue exactly where "
"you left off. Do not restart or repeat prior text. "
"Finish the answer directly.]"
)
_LENGTH_CONTINUATION_OUTPUT_LIMIT = (
"[System: Your previous response was truncated by the output "
"length limit. Continue exactly where you left off. Do not "
"restart or repeat prior text. Finish the answer directly.]"
)
# The dropped-tools variant interpolates tool names, so
# _is_synthetic_compression_user_turn matches this prefix with str.startswith.
_LENGTH_CONTINUATION_DROPPED_TOOLS_PREFIX = "[System: Your previous tool call "
def _get_continuation_prompt(is_partial_stub: bool, dropped_tools: Optional[List[str]] = None) -> str:
if is_partial_stub and dropped_tools:
tool_list = ", ".join(dropped_tools[:3])
return (
f"{_LENGTH_CONTINUATION_DROPPED_TOOLS_PREFIX}"
f"({tool_list}) was too large and "
"the stream timed out before it "
"could be delivered. Do NOT retry "
"the same tool call with the same "
"large content. Instead, break the "
"content into multiple smaller tool "
"calls (e.g. use multiple patch calls "
"or write smaller files). Each tool "
"call's arguments must be under ~8K "
"tokens to avoid stream timeouts.]"
)
elif is_partial_stub:
return _LENGTH_CONTINUATION_NETWORK_STUB
else:
return _LENGTH_CONTINUATION_OUTPUT_LIMIT
# Nudge for Codex/Responses turns that returned only internal reasoning: a bare retry
# would be byte-identical (nothing replayable emitted), so the model repeats it.
_CODEX_INCOMPLETE_NUDGE = (
"[System: Your previous response contained only internal reasoning and "
"never produced a visible answer or tool call. Do not keep thinking. "
"Produce your final answer as plain text now (or make the tool call "
"you were planning).]"
)
# Re-prompt after an acknowledgment-only Codex/Responses reply; named so
# _is_synthetic_compression_user_turn can recognize it like _CODEX_INCOMPLETE_NUDGE.
_CODEX_ACK_CONTINUATION_NUDGE = (
"[System: Continue now. Execute the required tool calls and only "
"send your final answer after completing the task.]"
)
# Re-prompt for finish_reason="tool_calls" with empty tool_calls. Named like
# _CODEX_ACK_CONTINUATION_NUDGE: an interrupt mid-retry can persist it.
_DROPPED_TOOLCALL_NUDGE_CONTENT = (
"Your previous turn indicated a tool call but none was "
"included. Do not narrate a plan or restate intent — issue "
"the actual tool call now to continue the task."
)
# Re-prompt for an empty response after tool calls (#9400). Named because its
# _empty_recovery_synthetic metadata flag does not survive SessionDB projection.
_EMPTY_TOOL_RESPONSE_NUDGE = (
"You just executed tool calls but returned an "
"empty response. Please process the tool "
"results above and continue with the task."
)
# Shared recovery trailer for both content-policy refusal paths (HTTP-200
# content_filter and the content_policy_blocked exception) so guidance cannot drift.
_CONTENT_POLICY_RECOVERY_HINT = (
"Try rephrasing the request, narrowing the context, or "
"adding a fallback provider with `hermes fallback add`."
)
# Memo for send-path tool-call argument canonicalization, which re-runs on every
# historical call each iteration. Sound: canonicalization is pure and deterministic;
# malformed strings raise before being stored, so the repair fallback is never memoized.
_CANON_ARGS_CACHE: Dict[str, str] = {}
_CANON_ARGS_CACHE_MAX = 4096
# Count bound alone does not bound MEMORY: argument strings can run 100KB+, so a byte
# budget bounds the worst case while keeping the memo effective for ~0.5-2KB args.
_CANON_ARGS_CACHE_MAX_BYTES = 32 * 1024 * 1024
_canon_args_cache_bytes = 0
def _canonicalize_tool_call_arguments(arg_str: str) -> str:
"""Return the canonical wire form of a tool-call arguments JSON string.
Raises whatever ``json.loads`` raises on malformed input; the caller falls back to
``_repair_tool_call_arguments``."""
global _canon_args_cache_bytes
cached = _CANON_ARGS_CACHE.get(arg_str)
if cached is not None:
return cached
canonical = json.dumps(
json.loads(arg_str), separators=(",", ":"), sort_keys=True,
)
_CANON_ARGS_CACHE[arg_str] = canonical
_canon_args_cache_bytes += len(arg_str) + len(canonical)
while len(_CANON_ARGS_CACHE) > _CANON_ARGS_CACHE_MAX or (
_canon_args_cache_bytes > _CANON_ARGS_CACHE_MAX_BYTES
and len(_CANON_ARGS_CACHE) > 1
):
try:
evicted_key = next(iter(_CANON_ARGS_CACHE))
evicted_val = _CANON_ARGS_CACHE.pop(evicted_key)
_canon_args_cache_bytes -= len(evicted_key) + len(evicted_val)
except (StopIteration, KeyError, RuntimeError):
break
return canonical
def _clone_message_for_send(msg):
"""Structural clone of a history message for the per-call API copy.
Clones every dict/list recursively while sharing immutable leaves, so in-place
send-path rewrites can never reach the persisted transcript (#80498). Cheaper than
copy.deepcopy; messages are JSON-shaped and acyclic, tuples are shared as leaves."""
if isinstance(msg, dict):
return {
k: _clone_message_for_send(v) if isinstance(v, (dict, list)) else v
for k, v in msg.items()
}
if isinstance(msg, list):
return [
_clone_message_for_send(v) if isinstance(v, (dict, list)) else v
for v in msg
]
return msg
def _canonicalize_api_tool_calls(api_messages) -> None:
"""Canonicalize tool-call argument JSON on the send-path message copy.
Rewrites ``tool_calls`` in place (copy-on-write for the dicts it touches; persisted
history untouched). The memo bounds parse/serialize to one per UNIQUE string."""
for am in api_messages:
tcs = am.get("tool_calls")
if not tcs:
continue
new_tcs = []
for tc in tcs:
if isinstance(tc, dict) and "function" in tc:
try:
tc = {**tc, "function": {
**tc["function"],
"arguments": _canonicalize_tool_call_arguments(
tc["function"]["arguments"]
),
}}
except Exception:
# Copy-on-write as defense in depth: callers may pass shallow
# copies, and writing into a shared tc["function"] rewrote the
# stored turn with "{}" on the unrepairable path (#80498).
tc = {**tc, "function": {
**tc["function"],
"arguments": _repair_tool_call_arguments(
tc["function"]["arguments"],
tc["function"].get("name", "?"),
),
}}
new_tcs.append(tc)
am["tool_calls"] = new_tcs
def _invalid_tool_name_error_content(name: str, valid_tool_names) -> str:
"""Error-result content for a tool call whose name isn't a real tool.
A blank name is a model echoing tool-call syntax seen in data, not a typo (#47967);
dumping the catalog feeds that loop, so send a terse error instead. A nonempty wrong
name still gets the catalog so the model can self-correct."""
if not (name or "").strip():
return (
"Tool call rejected: the tool name was empty. "
"If tool-call XML or JSON appeared in file "
"contents or tool output, that is data — do "
"not re-emit it as a tool call. To call a "
"tool, use a valid name from your tool list; "
"otherwise reply in plain text."
)
available = ", ".join(sorted(valid_tool_names))
return f"Tool '{name}' does not exist. Available tools: {available}"
def _content_policy_blocked_result(
messages: List[Dict],
api_call_count: int,
*,
final_response: str,
error_detail: str,
) -> Dict[str, Any]:
"""Build the terminal turn result for a content-policy block.
Refusals are deterministic for the unchanged prompt, so no retry; both the HTTP-200
and exception paths return this shape with a ``content_policy_blocked:`` error."""
return {
"final_response": final_response,
"messages": messages,
"api_calls": api_call_count,
"completed": False,
"failed": True,
"error": f"content_policy_blocked: {error_detail}",
}
def _compression_deferred_result(
agent,
messages: List[Dict],
api_call_count: int,
reason: str = "lock",
) -> Dict[str, Any]:
"""Build the soft turn result for a transiently-deferred compression.
Both ``reason="lock"`` and ``reason="transient_block"`` must end as
``compression_deferred``, never ``compression_exhausted`` — the gateway wipes the
session on exhaustion (#9893/#35809). ``failed`` stays False; the turn persists."""
if reason == "transient_block":
block = getattr(agent, "_compression_blocked_transient", None)
logger.info(
"turn deferred: compression transiently blocked (%s) "
"(session=%s) — not counting as compression exhaustion",
block if isinstance(block, str) else "unknown guard",
agent.session_id or "none",
)
_final = (
"Context compression is temporarily paused after a recent "
"failed attempt. Please retry in a moment — compression will "
"resume automatically (or run /compress to force a retry now)."
)
else:
holder = getattr(agent, "_compression_skipped_due_to_lock", None)
logger.info(
"turn deferred: compression lock held by another path "
"(session=%s holder=%s) — not counting as compression exhaustion",
agent.session_id or "none",
holder if isinstance(holder, str) else "unconfirmed",
)
_final = (
"Context compression is already running for this session. "
"Please retry in a moment — your next message will be processed "
"once the concurrent compression finishes."
)
try:
agent._flush_status_buffer()
except Exception:
pass
return {
"final_response": _final,
"messages": messages,
"completed": False,
"api_calls": api_call_count,
"error": _final,
"partial": True,
"failed": False,
"compression_deferred": True,
"session_id": agent.session_id,
}
def _provider_overflow_exhausted_result(
agent,
messages: List[Dict],
conversation_history,
api_call_count: int,
request_pressure_tokens: int,
max_compression_attempts: int,
) -> Dict[str, Any]:
"""Fail closed when a rebuilt request is still too large after recovery."""
agent._flush_status_buffer()
logger.error(
"%sContext compression failed after %d attempts; rebuilt request "
"remains over threshold at ~%s tokens.",
agent.log_prefix,
max_compression_attempts,
f"{request_pressure_tokens:,}",
)
agent._persist_session(messages, conversation_history)
final_response = (
"Context length exceeded: compression could not reduce the rebuilt "
"request below the safe threshold."
)
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_exhausted",
}
def _rewrite_system_content_blocks(system_message: dict, effective: str) -> bool:
"""Rewrite a cache-decorated system message in place, keeping its blocks.
Assigning a bare string over the ``[static prefix, volatile tail]`` block list drops
both cache_control breakpoints. Only the LAST ``Model:``/``Provider:`` lines change.
Returns False when the shape cannot be safely patched."""
content = system_message.get("content")
if not isinstance(content, list) or not content:
return False
if not all(
isinstance(part, dict) and part.get("type") == "text" for part in content
):
return False
if len(content) == 1:
content[0]["text"] = effective
return True
if len(content) == 2:
head = content[0].get("text") or ""
if head and effective.startswith(head):
tail = effective[len(head):]
if tail:
content[1]["text"] = tail
return True
return False
def _sync_failover_system_message(agent, api_messages, active_system_prompt):
"""Refresh the in-flight system message after a provider failover.
``try_activate_fallback`` rewrites the identity lines on ``_cached_system_prompt``,
but this call block's ``api_messages`` were built pre-failover and are reused each
retry. Mutates ``api_messages[0]`` in place; returns the new ``active_system_prompt``."""
sp = getattr(agent, "_cached_system_prompt", None)
if not isinstance(sp, str) or not sp:
return active_system_prompt
if api_messages and api_messages[0].get("role") == "system":
effective = sp
if agent.ephemeral_system_prompt:
effective = (effective + "\n\n" + agent.ephemeral_system_prompt).strip()
if not _rewrite_system_content_blocks(api_messages[0], effective):
api_messages[0]["content"] = effective
return sp
def _arm_fallback_restart(agent, api_messages, active_system_prompt, _retry):
"""After ``_try_activate_fallback`` succeeded: sync the system message to the new
provider and arm ``restart_with_rebuilt_messages`` (re-issue against the fallback,
refunding the stalled attempt). Callers also reset ``retry_count`` /
``compression_attempts`` to 0 and ``break`` the retry loop."""
active_system_prompt = _sync_failover_system_message(
agent, api_messages, active_system_prompt)
_retry.primary_recovery_attempted = False
_retry.restart_with_rebuilt_messages = True
return active_system_prompt
def _ensure_cached_system_prompt_static(agent, system_message=None) -> None:
"""Rebuild ``_cached_system_prompt_static`` when caching becomes active (#72626).
Sessions restored under a cache-off primary skip the static-prefix rebuild; a later
failover to a cache-on provider would otherwise silently fall back to the legacy
system-plus-3 layout. Wraps ``reconstruct_static_prefix`` (memoizes failures)."""
from agent.system_prompt import reconstruct_static_prefix
reconstruct_static_prefix(
agent, system_message=system_message, log_label="failover redecoration"
)
def _peel_moa_guidance(
messages: List[Dict[str, Any]],
guidance: Any,
) -> List[Dict[str, Any]]:
"""Remove MoA reference guidance attached by ``_attach_reference_guidance``.
Kept adjacent to the attach so the forward/inverse shapes evolve together."""
from agent.moa_loop import peel_reference_guidance
return peel_reference_guidance(messages, guidance)
def _redecorate_prompt_cache_for_provider(
agent,
api_messages: List[Dict[str, Any]],
*,
system_message=None,
moa_prepared: Optional[Dict[str, Any]] = None,
tools_for_api: Optional[List[Dict[str, Any]]] = None,
) -> tuple[List[Dict[str, Any]], Optional[Dict[str, Any]]] | tuple[List[Dict[str, Any]], Optional[Dict[str, Any]], List[Dict[str, Any]]]:
"""Strip and re-apply cache_control for the *current* provider policy.
Decoration runs once per call block for the primary provider, but failover
``continue`` paths reuse ``api_messages`` (#72626), so reshape at the top of each
retry from the mutated in-flight request. MoA guidance is peeled and rebased."""
messages: List[Dict[str, Any]] = [
dict(m) if isinstance(m, dict) else m for m in (api_messages or [])
]
prepared = moa_prepared
guidance = prepared.get("guidance") if isinstance(prepared, dict) else None
if guidance:
messages = _peel_moa_guidance(messages, guidance)
strip_anthropic_cache_control(messages)
planned_tools = strip_anthropic_tool_cache_control(
tools_for_api if tools_for_api is not None else getattr(agent, "tools", [])
)
if prepared is not None and getattr(agent, "provider", None) == "moa":
# Prepared MoA state is canonical: the synchronous acting-aggregator
# sender owns its destination-local cache plan after it resolves the slot.
completions = getattr(getattr(agent.client, "chat", None), "completions", None)
rebase = getattr(completions, "rebase_prepared_request", None)
if callable(rebase):
prepared = rebase(prepared, messages)
messages = prepared["messages"]
if tools_for_api is None:
return messages, prepared
return messages, prepared, planned_tools
# Direct attribute access, not getattr: the flags are always initialized on
# AIAgent, and a default would mask a real init bug as silent cache-off.
if agent._use_prompt_caching:
_ensure_cached_system_prompt_static(agent, system_message=system_message)
static = getattr(agent, "_cached_system_prompt_static", None)
direct_tool_cache = getattr(
agent,
"_direct_native_anthropic_tool_cache_capability",
lambda: False,
)()
from agent.prompt_caching import envelope_tool_part_cache_markers_supported
plan = build_prompt_cache_plan(
messages,
planned_tools,
# Clamp per-destination: a configured 1h regresses to 5m on
# Qwen/Alibaba routes, whose context cache is 5m-only (#84733).
cache_ttl=effective_cache_ttl(
agent._cache_ttl,
provider=agent.provider,
model=agent.model,
),
native_anthropic=agent._use_native_cache_layout,
static_system_prefix=static if isinstance(static, str) else None,
direct_native_tool_cache=direct_tool_cache,
# LiteLLM-style envelope routes forward part-level markers into
# tool_result.content[] → non-retryable 400 (#89886).
tool_part_markers=envelope_tool_part_cache_markers_supported(
getattr(agent, "provider", ""), getattr(agent, "base_url", "")
),
)
messages = plan.messages
planned_tools = plan.tools
if tools_for_api is None:
return messages, prepared
return messages, prepared, planned_tools
def _apply_context_engine_selection(
agent: Any,
api_messages: List[Dict[str, Any]],
conversation_messages: List[Dict[str, Any]],
incoming_message: Optional[Dict[str, Any]],
*,
logger: Any,
) -> List[Dict[str, Any]]:
"""Run the optional per-turn ``ContextEngine.select_context()`` hook.
Returns the (possibly replaced) request list. Fail-open: a missing hook, exception,
or invalid return yields ``api_messages`` unchanged; history is never mutated."""
engine = getattr(agent, "context_compressor", None)
if engine is None or not hasattr(engine, "select_context"):
return api_messages
# Skip the no-op base ``select_context`` so non-implementing engines pay nothing;
# ``hasattr`` is not enough: the ABC defines a default. Lazy import avoids a cycle.
try:
from agent.context_engine import ContextEngine as _CE
if getattr(engine.select_context, "__func__", None) is _CE.select_context:
return api_messages
except Exception:
pass
session_label = getattr(agent, "session_id", None) or "-"
# Structural clones: the engine must not be able to write through nested
# containers into persisted history; only the request list is acted on (#80498).
_conv_copy = [_clone_message_for_send(m) for m in conversation_messages] \
if conversation_messages is not None else None
_incoming_copy = _clone_message_for_send(incoming_message) if isinstance(incoming_message, dict) else incoming_message
try:
selected = engine.select_context(
api_messages,
conversation_messages=_conv_copy,
incoming_message=_incoming_copy,
budget_tokens=getattr(engine, "context_length", 0) or 0,
)
except Exception:
logger.warning(
"Context engine select_context hook failed; using unmodified "
"request messages (session=%s)",
session_label,
exc_info=True,
)
return api_messages
if selected is None:
return api_messages
# Require a NON-EMPTY list of dicts: ``all([])`` is ``True``, so a ``[]`` from a
# buggy engine would otherwise replace the request instead of failing open.
if isinstance(selected, list) and selected and all(isinstance(m, dict) for m in selected):
return selected
logger.warning(
"Context engine select_context returned an invalid value "
"(not a non-empty list of dicts); ignoring (session=%s)",
session_label,
)
return api_messages
def _notify_context_engine_turn_complete(
agent: Any,
messages: List[Dict[str, Any]],
*,
usage: Optional[Dict[str, Any]] = None,
logger: Any,
**meta: Any,
) -> None:
"""Notify the active context engine that a user turn has finished.
Fail-open: a missing/no-op hook or any exception is swallowed. ``messages`` is
passed as a copy so the engine cannot mutate the persisted transcript."""
engine = getattr(agent, "context_compressor", None)
hook = getattr(engine, "on_turn_complete", None)
if engine is None or not callable(hook):
return
# Skip the no-op base ``on_turn_complete`` so non-implementing engines pay nothing
# per turn. Lazy import avoids an import cycle with agent.context_engine.
try:
from agent.context_engine import ContextEngine as _CE
if getattr(hook, "__func__", None) is _CE.on_turn_complete:
return
except Exception:
pass
try:
hook(
# Structural clones: dict(m) would let a hook write into nested containers
# of the persisted transcript (#80498).
[_clone_message_for_send(m) for m in messages],
usage=usage,
**meta,
)
except Exception:
logger.warning(
"Context engine on_turn_complete hook failed (session=%s)",
getattr(agent, "session_id", None) or "-",
exc_info=True,
)
def run_conversation(
agent,
user_message: Any,
system_message: str = None,
conversation_history: List[Dict[str, Any]] = None,
task_id: str = None,
stream_callback: Optional[callable] = None,
persist_user_message: Optional[Any] = None,
persist_user_timestamp: Optional[float] = None,
persist_user_display_kind: Optional[str] = None,
persist_user_display_metadata: Optional[Dict[str, Any]] = None,
persist_user_platform_id: Optional[str] = None,
moa_config: Optional[dict[str, Any]] = None,
) -> Dict[str, Any]:
"""Run a complete conversation with tool calling until completion.
Args:
stream_callback: per-text-delta callback (TTS); None uses the non-streaming path.
persist_user_message: clean text to store when ``user_message`` carries API-only
synthetic prefixes; ``persist_user_timestamp`` / ``persist_user_platform_id``
are stored as metadata (platform id lets restart drain recovery dedup).
persist_user_display_kind/metadata: display-only event rendering (``auto_continue``,
``model_switch``); the model still receives the message unchanged.
Returns: dict with the final response and message history."""
if moa_config is None:
try:
from hermes_cli.moa_config import decode_moa_turn
_decoded_message, _decoded_moa_config = decode_moa_turn(user_message)
if _decoded_moa_config is not None:
user_message = _decoded_message
moa_config = _decoded_moa_config
if persist_user_message is None:
persist_user_message = _decoded_message
except Exception:
pass
# The gateway caches agents across turns; compression state is per-turn, or a stale
# in-place boundary would make a later uncompressed result look compacted.
agent._last_compaction_in_place = False
agent._last_compression_attempt_recorded = False
agent._last_compression_attempt_in_place = None
begin_fast_mode_turn(agent, conversation_history)
# Adopt ~/.hermes/.env credential/base-url edits made since the last turn — a
# Settings save updates .env, not this worker's client (#67821). No-op if unchanged.
try:
agent._try_refresh_env_client_credentials()
except Exception:
logger.debug("per-turn env credential refresh failed", exc_info=True)
# ── Per-turn setup (the prologue) ──
# All once-per-turn setup lives in ``build_turn_context`` (agent/turn_context.py);
# it mutates ``agent`` as the inline code did and returns the locals the loop reads.
try:
_ctx = build_turn_context(
agent,
user_message,
system_message,
conversation_history,
task_id,
stream_callback,
persist_user_message,
persist_user_timestamp,
persist_user_display_kind=persist_user_display_kind,
persist_user_display_metadata=persist_user_display_metadata,
persist_user_platform_id=persist_user_platform_id,
restore_or_build_system_prompt=_restore_or_build_system_prompt,
install_safe_stdio=_install_safe_stdio,
sanitize_surrogates=_sanitize_surrogates,
summarize_user_message_for_log=_summarize_user_message_for_log,
set_session_context=set_session_context,
set_current_write_origin=set_current_write_origin,
ra=_ra,
# MoA turns append per-call aggregated context to the API copy of the
# user message, so no byte-stable api_content sidecar can be stamped.
moa_active=bool(moa_config),
)
except PreflightCompressionTimedOut as _preflight_timeout_exc:
# Preflight compression timed out; no provider call sent (#98424). Return the
# typed recovery result: surfaces hide raw exception text, which would bury the
# actionable guidance and skip the compression_exhausted recovery contract.
logger.warning(
"Turn-start preflight compression timed out — ending turn with "
"typed recovery result: %s",
_preflight_timeout_exc,
)
# Clear the tripwire slot note_turn_start registered; the early return skips the
# persist funnel that clears it. The user row is deliberately NOT persisted:
# the gateway skips persistence for compression_exhausted results (#7100).
from agent.agent_runtime_helpers import note_turn_persisted
note_turn_persisted(agent)
# Not _COMPRESSION_TIMEOUT_FINAL_RESPONSE — that describes a different state
# (compression ran, could not reduce); the exception text carries the guidance.
_final_response = str(_preflight_timeout_exc)
return {
"final_response": _final_response,
"messages": list(conversation_history or []),
"completed": False,
"api_calls": 0,
"error": _final_response,
"partial": True,
"failed": True,
"compression_exhausted": True,
"turn_exit_reason": "context_compression_timeout",
}
user_message = _ctx.user_message
original_user_message = _ctx.original_user_message
messages = _ctx.messages
conversation_history = _ctx.conversation_history
active_system_prompt = _ctx.active_system_prompt
effective_task_id = _ctx.effective_task_id
turn_id = _ctx.turn_id
current_turn_user_idx = _ctx.current_turn_user_idx
_should_review_memory = _ctx.should_review_memory
_plugin_user_context = _ctx.plugin_user_context
_ext_prefetch_cache = _ctx.ext_prefetch_cache
# Commentary deduplication spans all provider continuations and tool calls
# within one user turn, but must not suppress the same phrase next turn.
agent._delivered_interim_texts = set()
# A configured SessionDB append failure halts only the affected turn. A
# cached gateway agent must recover on the next message if storage did.
agent._incremental_persistence_failed = False
# Cause of the last persistence failure this turn ('locked'/'disk'/'unknown', see
# hermes_state.classify_persistence_error). Reset so a prior diagnosis cannot leak.
agent._last_persistence_error_cause = None
# Per-turn diagnostic: a failed compression-tip adoption in a previous
# turn's flush must not be reported against this turn.
agent._compression_adoption_failed = False
# Main conversation loop counters (pure locals consumed by the loop below).
api_call_count = 0
final_response = None
interrupted = False
failed = False
codex_ack_continuations = 0
length_continue_retries = 0
# Turn-scoped one-shot: armed by a thinking-only truncation, consumed by
# build_api_kwargs; must not survive an interrupted turn into the next one.
agent._ephemeral_reasoning_off = False
# Total outer-loop exceptions this turn (#92450) — see _MAX_OUTER_LOOP_ERRORS.
_outer_error_count = 0
truncated_tool_call_retries = 0
truncated_response_parts: List[str] = []
compression_attempts = 0
# Per-turn compression attempt cap shared by the pre-API gate, 413 handlers and
# post-tool compaction; a consecutive-ineffective-attempt backstop, rearmed only
# after a provider response reports a prompt below threshold. Default 3 if unset.
max_compression_attempts = getattr(agent, "max_compression_attempts", 3)
_last_preflight_pressure: Optional[int] = None
_preflight_compression_blocked = _ctx.preflight_compression_blocked
# A provider overflow outweighs the rough-estimate calibration that defers preflight
# after compaction: stay armed until the rebuilt request is below the threshold.
_provider_overflow_recovery_pending = False
# Armed when a compression host-timeout ends the turn; finalize reuses the gateway
# context-recovery contract (error/partial/compression_exhausted) (#98722).
_compression_timeout_exhausted = False
_turn_exit_reason = "unknown" # Diagnostic: why the loop ended
# Last answer held back by a verification gate: if the continuation exhausts the
# budget this is the best user-facing result, distinct from error/recovery text.
_pending_verification_response = None
# Whether the pending verification candidate was already streamed as interim.
# ``_response_was_previewed`` is set ONLY if it becomes the final response (#65919).
_pending_verification_response_previewed = False
# If pre-API compression fires after MoA advisors ran, retain their guidance and
# rebase it onto the compacted transcript next iteration — no second fan-out.
pending_moa_prepared_request = None
# Per-turn tally of credential-pool refreshes by (provider, pool-entry-id): caps
# same-entry refreshes on a persistent 401 so fallback takes over (#26080).
agent._auth_pool_refresh_counts = {}
# Per-turn usage forwarded to the context engine's on_turn_complete() hook; left
# None on turns that never reach a response so the hook never sees stale usage.
agent._last_turn_usage = None
# Opt-in runtime: api_mode == codex_app_server hands the whole turn to the codex
# app-server subprocess (see agent/transports/codex_app_server_session.py).
if agent.api_mode == "codex_app_server":
return agent._run_codex_app_server_turn(
user_message=user_message,
original_user_message=original_user_message,
messages=messages,
effective_task_id=effective_task_id,
should_review_memory=_should_review_memory,
)
while (api_call_count < agent.max_iterations and agent.iteration_budget.remaining > 0) or agent._budget_grace_call:
_redirect_text = agent._drain_pending_redirect()
if _redirect_text:
_apply_active_turn_redirect(agent, messages, _redirect_text)
if isinstance(original_user_message, str):
original_user_message = (
f"{original_user_message}\n\n"
f"User correction during the turn: {_redirect_text}"
)
agent._persist_session(messages, conversation_history)
# Reset per-turn checkpoint dedup so each iteration can take one snapshot
agent._checkpoint_mgr.new_turn()
# Check for interrupt request (e.g., user sent new message)
if agent._interrupt_requested:
interrupted = True
_turn_exit_reason = "interrupted_by_user"
if not agent.quiet_mode:
agent._safe_print("\n⚡ Breaking out of tool loop due to interrupt...")
break
# Aggregate input budget for detached auxiliary forks: bounds the whole review,
# not each request. Checked between iterations so the crossing request's writes
# have landed, mirroring the iteration-budget exit (#93057).
if _review_input_budget_exhausted(agent):
_turn_exit_reason = "review_input_budget_exhausted"
if not agent.quiet_mode:
agent._safe_print(
f"\n⏹️ Review input budget exhausted "
f"({int(agent.session_input_tokens):,} tokens) — stopping "
f"the review tool loop before the next provider call."
)
break
api_call_count += 1
agent._api_call_count = api_call_count
agent._touch_activity(f"starting API call #{api_call_count}")
# Grace call: budget exhausted but the model gets one more call. Consume the
# flag so the loop exits after this iteration regardless of outcome.
if agent._budget_grace_call:
agent._budget_grace_call = False
elif not agent.iteration_budget.consume():
_turn_exit_reason = "budget_exhausted"
if not agent.quiet_mode:
agent._safe_print(f"\n⚠️ Iteration budget exhausted ({agent.iteration_budget.used}/{agent.iteration_budget.max_total} iterations used)")
break
_ip = prepare_iteration(
agent,
messages=messages,
api_call_count=api_call_count,
)
messages = _ip.messages
request_logger = _ip.request_logger
_rr = assemble_api_request(
agent,
messages=messages,
current_turn_user_idx=current_turn_user_idx,
_ext_prefetch_cache=_ext_prefetch_cache,
_plugin_user_context=_plugin_user_context,
moa_config=moa_config,
active_system_prompt=active_system_prompt,
original_user_message=original_user_message,
pending_moa_prepared_request=pending_moa_prepared_request,
request_logger=request_logger,
)
api_messages = _rr.api_messages
tools_for_api = _rr.tools_for_api
_moa_prepared_request = _rr._moa_prepared_request
pending_moa_prepared_request = _rr.pending_moa_prepared_request
approx_tokens = _rr.approx_tokens
request_pressure_tokens = _rr.request_pressure_tokens
total_chars = _rr.total_chars
_pg = run_preflight_gate(
agent,
request_pressure_tokens=request_pressure_tokens,
_moa_prepared_request=_moa_prepared_request,
pending_moa_prepared_request=pending_moa_prepared_request,
messages=messages,
system_message=system_message,
user_message=user_message,
active_system_prompt=active_system_prompt,
conversation_history=conversation_history,
api_call_count=api_call_count,
compression_attempts=compression_attempts,
max_compression_attempts=max_compression_attempts,
effective_task_id=effective_task_id,
final_response=final_response,
failed=failed,
_turn_exit_reason=_turn_exit_reason,
_compression_timeout_exhausted=_compression_timeout_exhausted,
_preflight_compression_blocked=_preflight_compression_blocked,
_provider_overflow_recovery_pending=_provider_overflow_recovery_pending,
_last_preflight_pressure=_last_preflight_pressure,
)
pending_moa_prepared_request = _pg.pending_moa_prepared_request
messages = _pg.messages
active_system_prompt = _pg.active_system_prompt
conversation_history = _pg.conversation_history
api_call_count = _pg.api_call_count
compression_attempts = _pg.compression_attempts
final_response = _pg.final_response
failed = _pg.failed
_turn_exit_reason = _pg._turn_exit_reason
_compression_timeout_exhausted = _pg._compression_timeout_exhausted
_preflight_compression_blocked = _pg._preflight_compression_blocked
_provider_overflow_recovery_pending = _pg._provider_overflow_recovery_pending
_last_preflight_pressure = _pg._last_preflight_pressure
if _pg.action == "return":
return _pg.result
if _pg.action == "break":
break
if _pg.action == "continue":
continue
# Thinking spinner for quiet mode (animated during API call)
thinking_spinner = None
if not agent.quiet_mode:
agent._vprint(f"\n{agent.log_prefix}🔄 Making API call #{api_call_count}/{agent.max_iterations}...")
agent._vprint(f"{agent.log_prefix} 📊 Request size: {len(api_messages)} messages, ~{approx_tokens:,} tokens (~{total_chars:,} chars)")
agent._vprint(f"{agent.log_prefix} 🔧 Available tools: {len(agent.tools) if agent.tools else 0}")
else:
# Animated thinking spinner in quiet mode
face = random.choice(KawaiiSpinner.get_thinking_faces())
verb = random.choice(KawaiiSpinner.get_thinking_verbs())
if agent.thinking_callback:
# CLI TUI mode: use prompt_toolkit widget instead of raw spinner
# (works in both streaming and non-streaming modes)
agent.thinking_callback(f"{face} {verb}...")
elif not agent._has_stream_consumers() and agent._should_start_quiet_spinner():
# Raw KawaiiSpinner only when no streaming consumers and the
# spinner output has a safe sink.
spinner_type = random.choice(['brain', 'sparkle', 'pulse', 'moon', 'star'])
thinking_spinner = KawaiiSpinner(f"{face} {verb}...", spinner_type=spinner_type, print_fn=agent._print_fn)
thinking_spinner.start()
# Log request details if verbose
if agent.verbose_logging:
logging.debug(f"API Request - Model: {agent.model}, Messages: {len(messages)}, Tools: {len(agent.tools) if agent.tools else 0}")
logging.debug(f"Last message role: {messages[-1]['role'] if messages else 'none'}")
logging.debug(f"Total message size: ~{approx_tokens:,} tokens")
api_start_time = time.time()
retry_count = 0
max_retries = agent._api_max_retries
_retry = TurnRetryState()
finish_reason = "stop"
response = None # Guard against UnboundLocalError if all retries fail
api_kwargs = None # Guard against UnboundLocalError in except handler
api_request_id = f"{turn_id}:api:{api_call_count}"
agent._current_api_request_id = api_request_id
while retry_count < max_retries:
# ── Nous Portal rate limit guard ──────────────────────
# Skip the call if another session recorded a rate limit: every attempt
# (incl. SDK retries) counts against RPH.
if agent.provider == "nous":
try:
from agent.nous_rate_guard import (
nous_rate_limit_remaining,
format_remaining as _fmt_nous_remaining,
)
_nous_remaining = nous_rate_limit_remaining()
if _nous_remaining is not None and _nous_remaining > 0:
_nous_msg = (
f"Nous Portal rate limit active — "
f"resets in {_fmt_nous_remaining(_nous_remaining)}."
)
agent._buffer_vprint(
f"⏳ {_nous_msg} Trying fallback..."
)
agent._buffer_status(f"⏳ {_nous_msg}")
if agent._try_activate_fallback():
active_system_prompt = _arm_fallback_restart(
agent, api_messages, active_system_prompt, _retry)
retry_count = 0
compression_attempts = 0
break
# No fallback available — surface buffered context
# so user sees the rate-limit message that led here.
agent._flush_status_buffer()
agent._persist_session(messages, conversation_history)
return {
"final_response": (
f"⏳ {_nous_msg}\n\n"
"No fallback provider available. "
"Try again after the reset, or add a "
"fallback provider in config.yaml."
),
"messages": messages,
"api_calls": api_call_count,
"completed": False,
"failed": True,
"error": _nous_msg,
}
except ImportError:
pass
except Exception:
pass # Never let rate guard break the agent loop
try:
_rq = build_api_request(
agent,
api_messages=api_messages,
_moa_prepared_request=_moa_prepared_request,
tools_for_api=tools_for_api,
system_message=system_message,
messages=messages,
original_user_message=original_user_message,
approx_tokens=approx_tokens,
total_chars=total_chars,
retry_count=retry_count,
api_call_count=api_call_count,
api_request_id=api_request_id,
api_start_time=api_start_time,
effective_task_id=effective_task_id,
turn_id=turn_id,
)
api_messages = _rq.api_messages
_moa_prepared_request = _rq._moa_prepared_request
tools_for_api = _rq.tools_for_api
api_kwargs = _rq.api_kwargs
_original_api_kwargs = _rq._original_api_kwargs
_llm_middleware_trace = _rq._llm_middleware_trace
_ac = perform_api_call(
agent,
api_kwargs=api_kwargs,
_original_api_kwargs=_original_api_kwargs,
_llm_middleware_trace=_llm_middleware_trace,
_moa_prepared_request=_moa_prepared_request,
_retry=_retry,
thinking_spinner=thinking_spinner,
retry_count=retry_count,
api_call_count=api_call_count,
api_request_id=api_request_id,
effective_task_id=effective_task_id,
turn_id=turn_id,
interrupted=interrupted,
)
response = _ac.response
thinking_spinner = _ac.thinking_spinner
interrupted = _ac.interrupted
if _ac.action == "break":
break
_rc = check_api_response(
agent,
response=response,
_retry=_retry,
thinking_spinner=thinking_spinner,
messages=messages,
api_messages=api_messages,
api_kwargs=api_kwargs,
active_system_prompt=active_system_prompt,
conversation_history=conversation_history,
finish_reason=finish_reason,
retry_count=retry_count,
max_retries=max_retries,
compression_attempts=compression_attempts,
max_compression_attempts=max_compression_attempts,
length_continue_retries=length_continue_retries,
truncated_response_parts=truncated_response_parts,
truncated_tool_call_retries=truncated_tool_call_retries,
current_turn_user_idx=current_turn_user_idx,
api_call_count=api_call_count,
api_request_id=api_request_id,
api_start_time=api_start_time,
effective_task_id=effective_task_id,
turn_id=turn_id,
_preflight_compression_blocked=_preflight_compression_blocked,
_last_preflight_pressure=_last_preflight_pressure,
)
thinking_spinner = _rc.thinking_spinner
messages = _rc.messages
active_system_prompt = _rc.active_system_prompt
finish_reason = _rc.finish_reason
retry_count = _rc.retry_count
compression_attempts = _rc.compression_attempts
length_continue_retries = _rc.length_continue_retries
truncated_response_parts = _rc.truncated_response_parts
truncated_tool_call_retries = _rc.truncated_tool_call_retries
_preflight_compression_blocked = _rc._preflight_compression_blocked
_last_preflight_pressure = _rc._last_preflight_pressure
api_duration = _rc.api_duration
if _rc.action == "return":
return _rc.result
if _rc.action == "break":
break
if _rc.action == "continue":
continue
except InterruptedError:
if thinking_spinner:
thinking_spinner.stop("")
thinking_spinner = None
if agent.thinking_callback:
agent.thinking_callback("")
if agent._has_pending_redirect():
# redirect() cancelled only this request: keep the correction
# queued, clear the cancellation bit, let the outer loop rebuild.
# Never materialize incomplete signed/encrypted reasoning items.
if agent.clear_interrupt(preserve_redirect=True):
_retry.restart_with_redirected_messages = True
break
api_elapsed = time.time() - api_start_time
agent._vprint(f"{agent.log_prefix}⚡ Interrupted during API call.", force=True)
interrupted = True
# Keep assistant text already streamed before the stop, else the next
# turn has no record of the half-finished reply.
_partial = agent._strip_think_blocks(
getattr(agent, "_current_streamed_assistant_text", "") or ""
).strip()
if _partial:
append_message(messages, {"role": "assistant", "content": _partial})
final_response = _partial
else:
final_response = f"{INTERRUPT_WAITING_FOR_MODEL_PREFIX}{api_elapsed:.1f}s elapsed)."
agent._persist_session(messages, conversation_history)
break
except Exception as api_error:
_ae = handle_api_error(
agent,
api_error=api_error,
_retry=_retry,
thinking_spinner=thinking_spinner,
messages=messages,
api_messages=api_messages,
api_kwargs=api_kwargs,
system_message=system_message,
active_system_prompt=active_system_prompt,
conversation_history=conversation_history,
approx_tokens=approx_tokens,
retry_count=retry_count,
max_retries=max_retries,
compression_attempts=compression_attempts,
max_compression_attempts=max_compression_attempts,
api_call_count=api_call_count,
api_request_id=api_request_id,
api_start_time=api_start_time,
effective_task_id=effective_task_id,
turn_id=turn_id,
)
thinking_spinner = _ae.thinking_spinner
messages = _ae.messages
active_system_prompt = _ae.active_system_prompt
conversation_history = _ae.conversation_history
approx_tokens = _ae.approx_tokens
retry_count = _ae.retry_count
max_retries = _ae.max_retries
compression_attempts = _ae.compression_attempts
if _ae._provider_overflow_recovery_pending:
_provider_overflow_recovery_pending = True
if _ae.action == "return":
return _ae.result
if _ae.action == "break":
break
if _ae.action == "continue":
continue
if _retry.restart_with_redirected_messages:
# Cancelled request produced no valid assistant item: reuse the same logical
# iteration after the outer loop appends partial context + correction.
api_call_count -= 1
agent.iteration_budget.refund()
_retry.restart_with_redirected_messages = False
continue
# If the API call was interrupted, skip response processing
if interrupted:
_turn_exit_reason = "interrupted_during_api_call"
break
if _retry.restart_with_compressed_messages:
api_call_count -= 1
agent.iteration_budget.refund()
# Compression restarts count toward the retry limit so a compression that
# shrinks messages but not enough can't loop forever.
retry_count += 1
_retry.restart_with_compressed_messages = False
if _should_skip_model_call_for_reference_handoff(
messages, user_message
):
logger.info(
"Skipping compressed-restart 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"
break
# In-loop compression rebuilt `messages`; re-anchor the current-turn index
# like the prologue, AFTER the handoff guard (it may re-append this turn's
# ask). A stale anchor injects prefetch into a historical row.
current_turn_user_idx = reanchor_current_turn_user_idx(
messages, user_message
)
agent._persist_user_message_idx = current_turn_user_idx
continue
if _retry.restart_with_rebuilt_messages:
# A stall/failure escalated to the fallback chain: re-issue against the
# active fallback provider, refunding budget/count for the stalled attempt.
api_call_count -= 1
agent.iteration_budget.refund()
_retry.restart_with_rebuilt_messages = False
# Failover shrank the compressor window: clear the preflight block so
# preflight re-runs before the first fallback call. Hoisted to the single
# consumer. (#84733)
_preflight_compression_blocked = False
continue
if _retry.restart_with_length_continuation:
# Boost output budget per retry: 2×, 4×, 8×, 16× base, capped at 32 768, via
# _ephemeral_max_output_tokens. Keep a larger original provider/model
# default as the floor so retries never downshift.
_boost_base = agent.max_tokens if agent.max_tokens else 4096
_boost = _boost_base * (2 ** length_continue_retries)
_requested_cap = agent._requested_output_cap_from_api_kwargs(api_kwargs)
if _requested_cap is not None:
_boost = max(_boost, _requested_cap)
_boost_cap = max(32768, _requested_cap or 0)
agent._ephemeral_max_output_tokens = min(_boost, _boost_cap)
continue
# All retries may exhaust with `response` still None; break out cleanly.
if response is None:
_turn_exit_reason = "all_retries_exhausted_no_response"
print(f"{agent.log_prefix}❌ All API retries exhausted with no successful response.")
agent._persist_session(messages, conversation_history)
break
try:
_ri = normalize_model_response(
agent,
response=response,
messages=messages,
api_messages=api_messages,
conversation_history=conversation_history,
api_call_count=api_call_count,
api_duration=api_duration,
api_start_time=api_start_time,
api_request_id=api_request_id,
effective_task_id=effective_task_id,
turn_id=turn_id,
)
assistant_message = _ri.assistant_message
finish_reason = _ri.finish_reason
if _ri.action == "return":
return _ri.result
if _ri.action == "continue":
continue
# Check for tool calls
if assistant_message.tool_calls:
_tr = run_tool_round(
agent,
assistant_message=assistant_message,
finish_reason=finish_reason,
messages=messages,
conversation_history=conversation_history,
api_call_count=api_call_count,
effective_task_id=effective_task_id,
user_message=user_message,
system_message=system_message,
active_system_prompt=active_system_prompt,
compression_attempts=compression_attempts,
max_compression_attempts=max_compression_attempts,
final_response=final_response,
failed=failed,
_turn_exit_reason=_turn_exit_reason,
truncated_tool_call_retries=truncated_tool_call_retries,
)
messages = _tr.messages
conversation_history = _tr.conversation_history
active_system_prompt = _tr.active_system_prompt
compression_attempts = _tr.compression_attempts
final_response = _tr.final_response
failed = _tr.failed
_turn_exit_reason = _tr._turn_exit_reason
truncated_tool_call_retries = _tr.truncated_tool_call_retries
if _tr.action == "return":
return _tr.result
if _tr.action == "break":
break
if _tr.action == "continue":
continue
else:
_fr = finish_text_response(
agent,
assistant_message=assistant_message,
response=response,
finish_reason=finish_reason,
messages=messages,
api_messages=api_messages,
conversation_history=conversation_history,
api_call_count=api_call_count,
user_message=user_message,
active_system_prompt=active_system_prompt,
final_response=final_response,
_turn_exit_reason=_turn_exit_reason,
_preflight_compression_blocked=_preflight_compression_blocked,
codex_ack_continuations=codex_ack_continuations,
truncated_response_parts=truncated_response_parts,
length_continue_retries=length_continue_retries,
_pending_verification_response=_pending_verification_response,
_pending_verification_response_previewed=_pending_verification_response_previewed,
)
active_system_prompt = _fr.active_system_prompt
final_response = _fr.final_response
_turn_exit_reason = _fr._turn_exit_reason
_preflight_compression_blocked = _fr._preflight_compression_blocked
codex_ack_continuations = _fr.codex_ack_continuations
truncated_response_parts = _fr.truncated_response_parts
length_continue_retries = _fr.length_continue_retries
_pending_verification_response = _fr._pending_verification_response
_pending_verification_response_previewed = _fr._pending_verification_response_previewed
if _fr.action == "return":
return _fr.result
if _fr.action == "break":
break
if _fr.action == "continue":
continue
except Exception as e:
_oe = handle_outer_loop_error(
agent,
e=e,
_outer_error_count=_outer_error_count,
api_call_count=api_call_count,
messages=messages,
conversation_history=conversation_history,
_turn_exit_reason=_turn_exit_reason,
failed=failed,
final_response=final_response,
)
_outer_error_count = _oe._outer_error_count
_turn_exit_reason = _oe._turn_exit_reason
failed = _oe.failed
final_response = _oe.final_response
if _oe.action == "break":
break
# Post-loop finalization lives in agent/turn_finalizer.finalize_turn.
result = finalize_turn(
agent,
final_response=final_response,
api_call_count=api_call_count,
interrupted=interrupted,
failed=failed,
messages=messages,
conversation_history=conversation_history,
effective_task_id=effective_task_id,
turn_id=turn_id,
user_message=user_message,
original_user_message=original_user_message,
_should_review_memory=_should_review_memory,
_turn_exit_reason=_turn_exit_reason,
_pending_verification_response=_pending_verification_response,
_pending_verification_response_previewed=_pending_verification_response_previewed,
)
if _compression_timeout_exhausted:
# Reuse the gateway's context-recovery contract: transcript stays intact while
# future input can move to a clean session (#98722).
result["error"] = _COMPRESSION_TIMEOUT_FINAL_RESPONSE
result["partial"] = True
result["compression_exhausted"] = True
return result
__all__ = ["run_conversation"]