1950 lines
80 KiB
Python
1950 lines
80 KiB
Python
"""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 inspect
|
|
import json
|
|
import logging
|
|
import re
|
|
import time
|
|
from dataclasses import dataclass, field, fields
|
|
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.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,
|
|
)
|
|
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 (
|
|
announce_api_call,
|
|
apply_retry_restarts,
|
|
begin_iteration,
|
|
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 handle_api_interrupt, nous_rate_limit_guard, 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 _decode_inline_moa_turn(user_message, persist_user_message):
|
|
"""Decode a MoA preset encoded into ``user_message`` (``hermes_cli.moa_config``).
|
|
|
|
Returns ``(user_message, moa_config, persist_user_message)``; unchanged with
|
|
``moa_config=None`` when nothing is encoded or decoding fails."""
|
|
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:
|
|
if persist_user_message is None:
|
|
persist_user_message = _decoded_message
|
|
return _decoded_message, _decoded_moa_config, persist_user_message
|
|
except Exception:
|
|
pass
|
|
return user_message, None, persist_user_message
|
|
|
|
|
|
def _preflight_timeout_result(agent, exc, conversation_history) -> Dict[str, Any]:
|
|
"""Typed recovery result when turn-start preflight compression timed out (#98424):
|
|
no provider call was sent. 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",
|
|
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(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",
|
|
}
|
|
|
|
|
|
@dataclass
|
|
class _LoopState:
|
|
"""Every local the turn loop threads through the phase helpers in ``agent/turn_*.py``.
|
|
|
|
Each helper takes the loop locals it needs as keyword arguments named exactly like
|
|
these fields and returns a verdict dataclass whose non-``action``/``result`` fields
|
|
carry the same names; :func:`_run_phase` passes and copies them back by name, so a
|
|
field added to a helper's signature or verdict needs a field here and nothing else.
|
|
Per-iteration slots (``response`` … ``assistant_message``) are rebound by the phases
|
|
before any later phase reads them, exactly as the former inline locals were."""
|
|
|
|
# Fixed for the turn.
|
|
user_message: Any
|
|
system_message: Any
|
|
moa_config: Any
|
|
original_user_message: Any
|
|
conversation_history: Any
|
|
effective_task_id: Any
|
|
turn_id: Any
|
|
_should_review_memory: Any
|
|
_plugin_user_context: Any
|
|
_ext_prefetch_cache: Any
|
|
# Turn-scoped state (rebound by the phases).
|
|
messages: Any
|
|
active_system_prompt: Any
|
|
current_turn_user_idx: Any
|
|
_preflight_compression_blocked: Any
|
|
# 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: Any
|
|
api_call_count: int = 0
|
|
final_response: Any = None
|
|
interrupted: bool = False
|
|
failed: bool = False
|
|
codex_ack_continuations: int = 0
|
|
length_continue_retries: int = 0
|
|
# Total outer-loop exceptions this turn (#92450) — see _MAX_OUTER_LOOP_ERRORS.
|
|
_outer_error_count: int = 0
|
|
truncated_tool_call_retries: int = 0
|
|
truncated_response_parts: List[str] = field(default_factory=list)
|
|
compression_attempts: int = 0
|
|
_last_preflight_pressure: Optional[int] = None
|
|
# 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: bool = 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: bool = False
|
|
_turn_exit_reason: str = "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: Any = None
|
|
# Whether that candidate was already streamed as interim; ``_response_was_previewed``
|
|
# is set ONLY if it becomes the final response (#65919).
|
|
_pending_verification_response_previewed: bool = 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: Any = None
|
|
# Per-iteration slots.
|
|
request_logger: Any = None
|
|
api_messages: Any = None
|
|
tools_for_api: Any = None
|
|
_moa_prepared_request: Any = None
|
|
approx_tokens: Any = None
|
|
request_pressure_tokens: Any = None
|
|
total_chars: Any = None
|
|
thinking_spinner: Any = None
|
|
api_start_time: Any = None
|
|
retry_count: int = 0
|
|
max_retries: Any = None
|
|
_retry: Any = None
|
|
finish_reason: str = "stop"
|
|
response: Any = None # None when every retry failed
|
|
api_kwargs: Any = None # None until built; read by the except handlers
|
|
api_request_id: Any = None
|
|
_original_api_kwargs: Any = None
|
|
_llm_middleware_trace: Any = None
|
|
api_duration: Any = None
|
|
assistant_message: Any = None
|
|
|
|
|
|
# Keyword names each phase helper takes (minus ``agent``), cached per function object.
|
|
_PHASE_PARAMS: Dict[Any, tuple] = {}
|
|
# Verdict fields the loop latches (only ever sets True) instead of copying back:
|
|
# ``handle_api_error`` reports overflow recovery per call and must not clear an earlier arm.
|
|
_LATCHED_VERDICT_FIELDS = {"handle_api_error": frozenset({"_provider_overflow_recovery_pending"})}
|
|
|
|
|
|
def _run_phase(fn, agent, state: _LoopState, **extra):
|
|
"""Call phase helper ``fn`` with the loop locals it names, copy its verdict fields back.
|
|
|
|
``extra`` supplies non-state arguments (the caught exception). Returns the verdict so
|
|
the caller can act on ``.action`` / ``.result``."""
|
|
params = _PHASE_PARAMS.get(fn)
|
|
if params is None:
|
|
params = _PHASE_PARAMS[fn] = tuple(
|
|
p for p in inspect.signature(fn).parameters if p != "agent"
|
|
)
|
|
verdict = fn(agent, **{
|
|
name: extra[name] if name in extra else getattr(state, name) for name in params
|
|
})
|
|
latched = _LATCHED_VERDICT_FIELDS.get(getattr(fn, "__name__", ""), ())
|
|
for f in fields(verdict):
|
|
if f.name in ("action", "result"):
|
|
continue
|
|
value = getattr(verdict, f.name)
|
|
if f.name not in latched:
|
|
setattr(state, f.name, value)
|
|
elif value:
|
|
setattr(state, f.name, True)
|
|
return verdict
|
|
|
|
|
|
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:
|
|
user_message, moa_config, persist_user_message = _decode_inline_moa_turn(
|
|
user_message, persist_user_message
|
|
)
|
|
|
|
# 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): ``build_turn_context`` (agent/turn_context.py)
|
|
# 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:
|
|
return _preflight_timeout_result(agent, _preflight_timeout_exc, conversation_history)
|
|
|
|
# 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
|
|
# 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
|
|
# 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
|
|
|
|
s = _LoopState(
|
|
user_message=_ctx.user_message,
|
|
system_message=system_message,
|
|
moa_config=moa_config,
|
|
original_user_message=_ctx.original_user_message,
|
|
conversation_history=_ctx.conversation_history,
|
|
effective_task_id=_ctx.effective_task_id,
|
|
turn_id=_ctx.turn_id,
|
|
_should_review_memory=_ctx.should_review_memory,
|
|
_plugin_user_context=_ctx.plugin_user_context,
|
|
_ext_prefetch_cache=_ctx.ext_prefetch_cache,
|
|
messages=_ctx.messages,
|
|
active_system_prompt=_ctx.active_system_prompt,
|
|
current_turn_user_idx=_ctx.current_turn_user_idx,
|
|
_preflight_compression_blocked=_ctx.preflight_compression_blocked,
|
|
max_compression_attempts=getattr(agent, "max_compression_attempts", 3),
|
|
)
|
|
|
|
# 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=s.user_message,
|
|
original_user_message=s.original_user_message,
|
|
messages=s.messages,
|
|
effective_task_id=s.effective_task_id,
|
|
should_review_memory=s._should_review_memory,
|
|
)
|
|
|
|
while (s.api_call_count < agent.max_iterations and agent.iteration_budget.remaining > 0) or agent._budget_grace_call:
|
|
if _run_phase(begin_iteration, agent, s).action == "break":
|
|
break
|
|
_run_phase(prepare_iteration, agent, s)
|
|
_run_phase(assemble_api_request, agent, s)
|
|
_pg = _run_phase(run_preflight_gate, agent, s)
|
|
if _pg.action == "return":
|
|
return _pg.result
|
|
if _pg.action == "break":
|
|
break
|
|
if _pg.action == "continue":
|
|
continue
|
|
_run_phase(announce_api_call, agent, s)
|
|
|
|
s.api_start_time = time.time()
|
|
s.retry_count = 0
|
|
s.max_retries = agent._api_max_retries
|
|
s._retry = TurnRetryState()
|
|
s.finish_reason = "stop"
|
|
s.response = None
|
|
s.api_kwargs = None
|
|
s.api_request_id = f"{s.turn_id}:api:{s.api_call_count}"
|
|
agent._current_api_request_id = s.api_request_id
|
|
|
|
while s.retry_count < s.max_retries:
|
|
_ng = _run_phase(nous_rate_limit_guard, agent, s)
|
|
if _ng.action == "return":
|
|
return _ng.result
|
|
if _ng.action == "break":
|
|
break
|
|
|
|
try:
|
|
_run_phase(build_api_request, agent, s)
|
|
if _run_phase(perform_api_call, agent, s).action == "break":
|
|
break
|
|
_rc = _run_phase(check_api_response, agent, s)
|
|
if _rc.action == "return":
|
|
return _rc.result
|
|
if _rc.action == "break":
|
|
break
|
|
if _rc.action == "continue":
|
|
continue
|
|
except InterruptedError:
|
|
if _run_phase(handle_api_interrupt, agent, s).action == "break":
|
|
break
|
|
except Exception as api_error:
|
|
_ae = _run_phase(handle_api_error, agent, s, api_error=api_error)
|
|
if _ae.action == "return":
|
|
return _ae.result
|
|
if _ae.action == "break":
|
|
break
|
|
if _ae.action == "continue":
|
|
continue
|
|
|
|
_rs = _run_phase(apply_retry_restarts, agent, s)
|
|
if _rs.action == "break":
|
|
break
|
|
if _rs.action == "continue":
|
|
continue
|
|
|
|
try:
|
|
_ri = _run_phase(normalize_model_response, agent, s)
|
|
if _ri.action == "return":
|
|
return _ri.result
|
|
if _ri.action == "continue":
|
|
continue
|
|
if s.assistant_message.tool_calls:
|
|
_tr = _run_phase(run_tool_round, agent, s)
|
|
if _tr.action == "return":
|
|
return _tr.result
|
|
if _tr.action == "break":
|
|
break
|
|
if _tr.action == "continue":
|
|
continue
|
|
else:
|
|
_fr = _run_phase(finish_text_response, agent, s)
|
|
if _fr.action == "return":
|
|
return _fr.result
|
|
if _fr.action == "break":
|
|
break
|
|
if _fr.action == "continue":
|
|
continue
|
|
except Exception as e:
|
|
if _run_phase(handle_outer_loop_error, agent, s, e=e).action == "break":
|
|
break
|
|
|
|
# Post-loop finalization lives in agent/turn_finalizer.finalize_turn.
|
|
result = finalize_turn(
|
|
agent,
|
|
final_response=s.final_response,
|
|
api_call_count=s.api_call_count,
|
|
interrupted=s.interrupted,
|
|
failed=s.failed,
|
|
messages=s.messages,
|
|
conversation_history=s.conversation_history,
|
|
effective_task_id=s.effective_task_id,
|
|
turn_id=s.turn_id,
|
|
user_message=s.user_message,
|
|
original_user_message=s.original_user_message,
|
|
_should_review_memory=s._should_review_memory,
|
|
_turn_exit_reason=s._turn_exit_reason,
|
|
_pending_verification_response=s._pending_verification_response,
|
|
_pending_verification_response_previewed=s._pending_verification_response_previewed,
|
|
)
|
|
if s._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"]
|