init_agent (2711 LOC) becomes a ~280-line ordered orchestrator over _resolve_api_mode / _finalize_routing / _init_* / _build_client / _load_tools / _parse_compression_config -> CompressionSettings / _resolve_context_length / _build_context_engine / ... phase helpers. _build_client is further split per wire mode (_init_anthropic_client, _init_moa_client, _init_bedrock_client, _init_openai_client with _explicit_client_kwargs / _routed_client_kwargs). Statement order and every side effect on the agent are preserved (AST body-parity checked against origin/main). Dedupe/dead code: drop _relay_moa_reference_event/_moa_reference_output_allowed (zero callers; only their own test) and their test file; alias _normalize_route_base_url; _parse_config_int replaces three copies of the strict int parser; _cfg_flag replaces four inline truthy-set checks; _client_kwargs_from_routed + _fallback_entries replace duplicated routed-client/fallback-entry blocks; _warn_invalid_config_int unifies the three log+stderr invalid-int warnings (byte-identical text); _bedrock_region_from_url; _memory_provider_init_kwargs; the host->default_headers if/elif chain becomes the _HOST_DEFAULT_HEADERS dispatch table; callback params assigned from _CALLBACK_PARAMS. Comments/docstrings hand-compacted to their rationale (invariants, ordering, failure modes kept; issue numbers and narrative dropped). test_pre_compress_checkpoint_contract source-check repointed at the CompressionSettings field names. Verified: tests/run_agent (2066 passed) + all agent_init-referencing tests (1134 passed), get_tool_definitions() byte-identical vs origin/main, import smokes for cli/run_agent/gateway.run/hermes_cli.main/agent.conversation_loop/ tui_gateway.server.
3074 lines
139 KiB
Python
3074 lines
139 KiB
Python
"""Implementation of :meth:`AIAgent.__init__` as ``init_agent(agent, ...)``.
|
||
|
||
``init_agent`` is a thin, ordered orchestrator over ``_init_*`` / ``_build_*``
|
||
phase helpers (routing → callbacks → client → tools → session → config
|
||
sections → compression → context engine). Phase ORDER is load-bearing: later
|
||
phases read attributes earlier ones set. Symbols that tests patch on
|
||
``run_agent.*`` (``OpenAI``, ``get_tool_definitions``, ``logger``, …) are
|
||
resolved through :func:`_ra` so the patch contract is preserved.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
import os
|
||
import re
|
||
import sys
|
||
import threading
|
||
import time
|
||
import uuid
|
||
from datetime import datetime
|
||
from dataclasses import dataclass
|
||
from typing import Any, Callable, Dict, List, Optional
|
||
from urllib.parse import parse_qs, urlparse, urlunparse
|
||
|
||
from agent.context_compressor import ContextCompressor
|
||
from agent.iteration_budget import IterationBudget
|
||
from agent.memory_manager import StreamingContextScrubber
|
||
from agent.session_activity import ActivityProvenance
|
||
from agent.model_metadata import (
|
||
MINIMUM_CONTEXT_LENGTH,
|
||
fetch_model_metadata,
|
||
is_local_endpoint,
|
||
query_ollama_num_ctx,
|
||
)
|
||
from agent.process_bootstrap import _install_safe_stdio
|
||
from agent.subdirectory_hints import SubdirectoryHintTracker
|
||
from agent.think_scrubber import StreamingThinkScrubber
|
||
from agent.tool_guardrails import (
|
||
ToolCallGuardrailConfig,
|
||
ToolCallGuardrailController,
|
||
ToolGuardrailDecision,
|
||
)
|
||
from hermes_cli.config import cfg_get
|
||
from hermes_cli.route_identity import normalize_route_base_url
|
||
from hermes_cli.timeouts import get_provider_request_timeout
|
||
from hermes_constants import get_hermes_home
|
||
from utils import base_url_host_matches, is_truthy_value
|
||
|
||
# Same logger name as run_agent so caplog/patches on "run_agent" see our records.
|
||
logger = logging.getLogger("run_agent")
|
||
|
||
|
||
# Memory providers we've already warned are unavailable. Deduped because the
|
||
# gateway builds a fresh AIAgent per message, so an un-deduped warning would
|
||
# fire on every turn.
|
||
_warned_unavailable_providers: set[str] = set()
|
||
|
||
|
||
def _warn_memory_provider_unavailable(name: str, reason: str = "") -> None:
|
||
"""Warn (once per provider) when a configured memory provider is unavailable.
|
||
|
||
``is_available()`` is a fast, side-effect-free hot-path check, so it can't
|
||
log for itself. Without this warning a provider whose credentials/config are
|
||
missing is silently dropped — the user has ``memory.provider`` set but gets
|
||
no memory and no diagnostic. A common trigger is systemd/gateway services
|
||
not inheriting ``~/.hermes/.env``. See NousResearch/hermes-agent#2765.
|
||
|
||
``reason`` is the provider's ``unavailable_reason()`` — a provider-specific,
|
||
actionable hint (e.g. which package to install). Because an unavailable
|
||
provider is never initialized, this is the only place such a hint can reach
|
||
the user, so it is appended to the warning when present (#7718).
|
||
"""
|
||
if name in _warned_unavailable_providers:
|
||
return
|
||
_warned_unavailable_providers.add(name)
|
||
logger.warning(
|
||
"Memory provider %r is selected but reports unavailable — external memory "
|
||
"is disabled for this session (built-in memory still works). Check the "
|
||
"provider's credentials/config with 'hermes memory status'. Note: "
|
||
"systemd/gateway services do not inherit ~/.hermes/.env automatically; set "
|
||
"any required variables in the service environment.%s",
|
||
name,
|
||
f" {reason}" if reason else "",
|
||
)
|
||
|
||
|
||
def _ra():
|
||
"""Lazy reference to ``run_agent`` so callers can patch
|
||
``run_agent.OpenAI`` / ``run_agent.cleanup_vm`` / ... and have those
|
||
patches reach this code path.
|
||
"""
|
||
import run_agent
|
||
return run_agent
|
||
|
||
|
||
# Canonicalize an endpoint URL for model-route identity comparisons.
|
||
_normalize_route_base_url = normalize_route_base_url
|
||
|
||
|
||
def _provider_default_routes(provider: str) -> set[str]:
|
||
"""Return known exact default routes for a canonical provider id."""
|
||
routes: set[str] = set()
|
||
try:
|
||
from hermes_cli.providers import HERMES_OVERLAYS, get_provider
|
||
|
||
overlay = HERMES_OVERLAYS.get(provider)
|
||
provider_def = get_provider(provider, allow_network=False)
|
||
for value in (
|
||
getattr(overlay, "base_url_override", ""),
|
||
getattr(provider_def, "base_url", ""),
|
||
):
|
||
route = _normalize_route_base_url(value)
|
||
if route:
|
||
routes.add(route)
|
||
except Exception:
|
||
pass
|
||
|
||
try:
|
||
from providers import get_provider_profile
|
||
|
||
profile = get_provider_profile(provider)
|
||
route = _normalize_route_base_url(
|
||
getattr(profile, "base_url", "")
|
||
)
|
||
if route:
|
||
routes.add(route)
|
||
except Exception:
|
||
pass
|
||
|
||
try:
|
||
from hermes_cli.auth import PROVIDER_REGISTRY
|
||
from hermes_cli.models import normalize_provider as normalize_model_provider
|
||
from hermes_cli.providers import normalize_provider as normalize_registry_provider
|
||
|
||
for provider_id, config in PROVIDER_REGISTRY.items():
|
||
canonical_id = normalize_registry_provider(
|
||
normalize_model_provider(provider_id)
|
||
)
|
||
if canonical_id != provider:
|
||
continue
|
||
route = _normalize_route_base_url(
|
||
getattr(config, "inference_base_url", "")
|
||
)
|
||
if route:
|
||
routes.add(route)
|
||
except Exception:
|
||
pass
|
||
|
||
if provider == "gemini":
|
||
routes.update(
|
||
f"{route.rstrip('/')}/openai"
|
||
for route in list(routes)
|
||
)
|
||
return routes
|
||
|
||
|
||
def _context_route_mismatch(
|
||
configured_base_url: Any,
|
||
active_base_url: Any,
|
||
configured_provider: Any,
|
||
active_provider: Any,
|
||
*,
|
||
already_normalized: bool = False,
|
||
) -> bool:
|
||
"""Return whether a context pin's configured route differs from runtime."""
|
||
if already_normalized:
|
||
configured_route = str(configured_base_url or "")
|
||
active_route = str(active_base_url or "")
|
||
else:
|
||
configured_route = _normalize_route_base_url(configured_base_url)
|
||
active_route = _normalize_route_base_url(active_base_url)
|
||
if configured_route:
|
||
return configured_route != active_route
|
||
|
||
configured_provider = str(configured_provider or "").strip()
|
||
active_provider = str(active_provider or "").strip()
|
||
if not configured_provider:
|
||
return False
|
||
try:
|
||
from hermes_cli.models import normalize_provider as normalize_model_provider
|
||
|
||
configured_provider = normalize_model_provider(configured_provider)
|
||
active_provider = normalize_model_provider(active_provider)
|
||
except Exception:
|
||
configured_provider = configured_provider.lower()
|
||
active_provider = active_provider.lower()
|
||
try:
|
||
from hermes_cli.providers import normalize_provider as normalize_registry_provider
|
||
|
||
configured_provider = normalize_registry_provider(configured_provider)
|
||
active_provider = normalize_registry_provider(active_provider)
|
||
except Exception:
|
||
pass
|
||
|
||
if active_route:
|
||
configured_routes = _provider_default_routes(configured_provider)
|
||
if configured_routes:
|
||
return active_route not in configured_routes
|
||
# Named/custom providers have no catalog default routes: an empty
|
||
# configured URL with a matching provider identity is the same route
|
||
# (gateway display paths compare the raw empty model.base_url and must
|
||
# not drop model.context_length to family defaults).
|
||
return not (active_provider and configured_provider == active_provider)
|
||
return bool(
|
||
configured_provider
|
||
and active_provider
|
||
and configured_provider != active_provider
|
||
)
|
||
|
||
|
||
def _normalize_custom_provider_name(value: Any) -> str:
|
||
"""Mirror runtime normalization for a requested custom-provider identity."""
|
||
return str(value or "").strip().lower().replace(" ", "-")
|
||
|
||
|
||
def _custom_provider_runtime_ids(value: Any) -> set[str]:
|
||
"""Return raw/menu identities that runtime accepts for a configured name."""
|
||
normalized = _normalize_custom_provider_name(value)
|
||
if not normalized:
|
||
return set()
|
||
return {normalized, f"custom:{normalized}"}
|
||
|
||
|
||
def _build_codex_gpt5_autoraise_notice(
|
||
autoraise: Dict[str, Any], context_length: Optional[int] = None
|
||
) -> str:
|
||
"""Build the one-time notice shown when Codex gpt-5.x raises compaction.
|
||
|
||
``autoraise`` is ``{"model": <slug>, "from": <old_ratio>, "to": <new_ratio>}``.
|
||
``context_length`` is the live-resolved window from the context compressor
|
||
(Codex's /models catalog is authoritative and can change server-side, e.g.
|
||
the gpt-5.6 family's 272K → 372K → 272K shifts in July 2026), so the banner
|
||
reports what this session actually got rather than a hardcoded cap. The
|
||
same text is printed inline for CLI users and replayed via
|
||
``status_callback`` for gateway users, so it must be self-contained and
|
||
include the exact opt-back-out command.
|
||
"""
|
||
model = str(autoraise.get("model") or "gpt-5.4/5.5").strip().lower().rsplit("/", 1)[-1]
|
||
if isinstance(context_length, int) and context_length > 0:
|
||
cap = f"{round(context_length / 1000)}K"
|
||
else:
|
||
# Static fallback when the resolved window isn't available:
|
||
# gpt-5.3-codex-spark has a native 128K window; the gpt-5.4/5.5/5.6
|
||
# family is capped at 272K by the Codex OAuth backend.
|
||
cap = "128K" if model.startswith("gpt-5.3-codex-spark") else "272K"
|
||
from_pct = int(round(autoraise["from"] * 100))
|
||
to_pct = int(round(autoraise["to"] * 100))
|
||
return (
|
||
f"ℹ Codex {model} caps context at {cap}, so auto-compaction was raised "
|
||
f"to {to_pct}% (from {from_pct}%) to use more of the window before "
|
||
f"summarizing.\n"
|
||
f" Opt back out: hermes config set compression.codex_gpt55_autoraise false"
|
||
)
|
||
|
||
|
||
def _resolve_compression_threshold(
|
||
global_threshold: float,
|
||
model_cthresh: Optional[float],
|
||
*,
|
||
model: Optional[str] = None,
|
||
is_codex_autoraise: bool,
|
||
) -> tuple[float, Optional[Dict[str, Any]]]:
|
||
"""Combine the user's global compaction threshold with a per-model override.
|
||
|
||
Returns ``(effective_threshold, autoraise_notice)``. ``autoraise_notice`` is
|
||
``{"model": <slug>, "from": <old>, "to": <new>}`` only when a Codex
|
||
autoraise (gpt-5.4/5.5 272K family or gpt-5.3-codex-spark) actually raises
|
||
the threshold, otherwise ``None``.
|
||
|
||
The Codex overrides are *autoraises*: they must never LOWER a higher
|
||
user-configured threshold. A user who already set ``compression.threshold``
|
||
above the raised value deliberately keeps more raw context, and silently
|
||
dropping them would both waste usable window and contradict the feature's
|
||
purpose (use more of the window). Other overrides (e.g. Arcee Trinity)
|
||
keep their existing unconditional behaviour.
|
||
"""
|
||
if model_cthresh is None:
|
||
return global_threshold, None
|
||
if is_codex_autoraise:
|
||
if model_cthresh <= global_threshold + 1e-9:
|
||
# Autoraise never lowers; keep the user's higher/equal threshold.
|
||
return global_threshold, None
|
||
return model_cthresh, {
|
||
"model": model,
|
||
"from": global_threshold,
|
||
"to": model_cthresh,
|
||
}
|
||
return model_cthresh, None
|
||
|
||
|
||
def _codex_gpt55_autoraise_notice_marker():
|
||
"""Path to the per-profile marker recording that the autoraise notice ran.
|
||
|
||
Lives under ``$HERMES_HOME`` (which is profile-scoped) alongside the other
|
||
internal markers like ``.container-mode`` — so it is not a user-facing config
|
||
key, and every profile tracks its own notice state independently.
|
||
"""
|
||
return get_hermes_home() / ".codex_gpt55_autoraise_notice"
|
||
|
||
|
||
def _codex_gpt55_autoraise_notice_state(autoraise: Dict[str, Any]) -> str:
|
||
"""Stable identity for one autoraise notice, keyed on what it displays.
|
||
|
||
Uses the model slug plus the same from→to percentages the notice text
|
||
shows, so an unchanged threshold stays silent across restarts while a
|
||
later change (the user edits their global ``threshold``, or switches to a
|
||
different autoraised Codex model) re-notifies once.
|
||
"""
|
||
model = str(autoraise.get("model") or "").strip().lower().rsplit("/", 1)[-1]
|
||
from_pct = int(round(float(autoraise["from"]) * 100))
|
||
to_pct = int(round(float(autoraise["to"]) * 100))
|
||
return f"{model}:{from_pct}:{to_pct}"
|
||
|
||
|
||
def _codex_gpt55_autoraise_notice_seen(autoraise: Dict[str, Any]) -> bool:
|
||
"""True if this exact autoraise notice was already shown for this profile.
|
||
|
||
A missing/unreadable marker (or one recording a different threshold) reads
|
||
as unseen, so the notice shows.
|
||
"""
|
||
try:
|
||
current = _codex_gpt55_autoraise_notice_state(autoraise)
|
||
return _codex_gpt55_autoraise_notice_marker().read_text(
|
||
encoding="utf-8"
|
||
).strip() == current
|
||
except (OSError, KeyError, TypeError, ValueError):
|
||
return False
|
||
|
||
|
||
def _record_codex_gpt55_autoraise_notice(autoraise: Dict[str, Any]) -> None:
|
||
"""Persist that the autoraise notice was shown for this profile/config state.
|
||
|
||
Best-effort: a read-only or missing ``$HERMES_HOME`` just means the notice
|
||
may show again next init, which is preferable to breaking agent init.
|
||
"""
|
||
try:
|
||
marker = _codex_gpt55_autoraise_notice_marker()
|
||
marker.parent.mkdir(parents=True, exist_ok=True)
|
||
marker.write_text(
|
||
_codex_gpt55_autoraise_notice_state(autoraise), encoding="utf-8"
|
||
)
|
||
except (OSError, KeyError, TypeError, ValueError):
|
||
pass
|
||
|
||
|
||
def _normalized_custom_base_url(value: Any) -> str:
|
||
if not isinstance(value, str):
|
||
return ""
|
||
return value.strip().rstrip("/")
|
||
|
||
|
||
def _custom_provider_model_matches(agent_model: str, entry: Dict[str, Any]) -> bool:
|
||
agent_model_norm = str(agent_model or "").strip().lower()
|
||
# Multi-model entries (`providers.<name>.models` mapping / legacy `models:`
|
||
# list): matching ANY catalog entry counts. Otherwise a provider whose
|
||
# `model` differs from the session model drops its extra_body (e.g. OpenAI
|
||
# service_tier) and bills the whole session at the wrong tier.
|
||
models = entry.get("models")
|
||
catalog: List[str] = []
|
||
if isinstance(models, dict):
|
||
catalog = [str(k).strip().lower() for k in models]
|
||
elif isinstance(models, (list, tuple)):
|
||
catalog = [str(m).strip().lower() for m in models]
|
||
if catalog and agent_model_norm in catalog:
|
||
return True
|
||
provider_model = str(entry.get("model", "") or "").strip().lower()
|
||
if not provider_model and not catalog:
|
||
return True
|
||
return provider_model == agent_model_norm
|
||
|
||
|
||
def _custom_provider_extra_body_for_agent(
|
||
*,
|
||
provider: str,
|
||
model: str,
|
||
base_url: str,
|
||
custom_providers: List[Dict[str, Any]],
|
||
) -> Optional[Dict[str, Any]]:
|
||
provider_norm = (provider or "").strip().lower()
|
||
if provider_norm == "custom":
|
||
provider_key_filter = ""
|
||
elif provider_norm.startswith("custom:"):
|
||
provider_key_filter = provider_norm.split(":", 1)[1].strip()
|
||
else:
|
||
return None
|
||
|
||
target_url = _normalized_custom_base_url(base_url)
|
||
if not target_url:
|
||
return None
|
||
|
||
fallback: Optional[Dict[str, Any]] = None
|
||
for entry in custom_providers or []:
|
||
if not isinstance(entry, dict):
|
||
continue
|
||
if provider_key_filter:
|
||
entry_keys = {
|
||
str(entry.get("provider_key", "") or "").strip().lower(),
|
||
str(entry.get("name", "") or "").strip().lower(),
|
||
}
|
||
if provider_key_filter not in entry_keys:
|
||
continue
|
||
if _normalized_custom_base_url(entry.get("base_url")) != target_url:
|
||
continue
|
||
extra_body = entry.get("extra_body")
|
||
if not isinstance(extra_body, dict) or not extra_body:
|
||
continue
|
||
provider_model = str(entry.get("model", "") or "").strip()
|
||
if provider_model:
|
||
if _custom_provider_model_matches(model, entry):
|
||
return dict(extra_body)
|
||
elif fallback is None:
|
||
fallback = dict(extra_body)
|
||
|
||
return fallback
|
||
|
||
|
||
def _merge_custom_provider_extra_body(agent, custom_providers: List[Dict[str, Any]]) -> None:
|
||
extra_body = _custom_provider_extra_body_for_agent(
|
||
provider=agent.provider,
|
||
model=agent.model,
|
||
base_url=agent.base_url,
|
||
custom_providers=custom_providers,
|
||
)
|
||
if not extra_body:
|
||
return
|
||
|
||
overrides = dict(getattr(agent, "request_overrides", {}) or {})
|
||
merged_extra_body = dict(extra_body)
|
||
existing_extra_body = overrides.get("extra_body")
|
||
if isinstance(existing_extra_body, dict):
|
||
merged_extra_body.update(existing_extra_body)
|
||
overrides["extra_body"] = merged_extra_body
|
||
agent.request_overrides = overrides
|
||
|
||
|
||
def _normalize_run_budget_seconds(value) -> Optional[float]:
|
||
"""Normalize a wall-clock run budget value to a positive float or None.
|
||
|
||
None / absent / non-numeric / non-positive all resolve to ``None``
|
||
(feature off) so a malformed config value can never activate the
|
||
deadline machinery, only leave it dormant. ``bool`` is rejected because
|
||
YAML ``true`` would otherwise become a 1-second budget.
|
||
"""
|
||
if value is None or isinstance(value, bool):
|
||
return None
|
||
try:
|
||
seconds = float(value)
|
||
except (TypeError, ValueError):
|
||
return None
|
||
if seconds != seconds or seconds <= 0: # NaN or non-positive
|
||
return None
|
||
return seconds
|
||
|
||
|
||
def _refuse_checkpoint_required_on_codex_app_server(
|
||
checkpoint_required: bool, api_mode: Optional[str]
|
||
) -> None:
|
||
"""Fail closed at init when the checkpoint gate cannot be honored.
|
||
|
||
The codex app-server owns its thread and compacts it without a truthful
|
||
pre-compaction transcript boundary (in "native" auto-compaction mode —
|
||
the default — Hermes never even initiates the compaction), so no
|
||
pre-compress checkpoint can be guaranteed on this API mode. Refusing here
|
||
keeps a turn from ever reaching a codex-owned compaction boundary; the
|
||
compress_context() guard alone cannot cover native turns that bypass
|
||
Hermes compression entirely.
|
||
"""
|
||
if checkpoint_required and api_mode == "codex_app_server":
|
||
raise RuntimeError(
|
||
"BLOCKED_MISSING_PREREQUISITE: compression.checkpoint_required "
|
||
"is incompatible with the codex_app_server API mode: the codex "
|
||
"agent compacts its own thread without a truthful pre-compaction "
|
||
"transcript boundary, so a required pre-compress checkpoint "
|
||
"cannot be guaranteed. Disable compression.checkpoint_required "
|
||
"or use a non-app-server API mode."
|
||
)
|
||
|
||
|
||
def _parse_config_int(raw: Any, default: int) -> int:
|
||
"""Strict int coercion for config knobs (compression.max_attempts & co).
|
||
|
||
Rejects booleans (bool subclasses int — YAML ``true`` would coerce to 1)
|
||
and fractional floats ("4.7 attempts" is a mistake, not a request for 4);
|
||
accepts ints, integral floats and numeric strings; else ``default``.
|
||
"""
|
||
if isinstance(raw, bool):
|
||
return default
|
||
if isinstance(raw, int):
|
||
return raw
|
||
if isinstance(raw, float):
|
||
return int(raw) if raw.is_integer() else default
|
||
try:
|
||
return int(str(raw).strip())
|
||
except (TypeError, ValueError):
|
||
return default
|
||
|
||
|
||
def _cfg_flag(cfg: Dict[str, Any], key: str, default: bool) -> bool:
|
||
"""Legacy string-set truthiness used by the ``compression`` section."""
|
||
return str(cfg.get(key, default)).lower() in {"true", "1", "yes"}
|
||
|
||
|
||
@dataclass
|
||
class CompressionSettings:
|
||
"""Parsed ``compression`` config section (see ``_parse_compression_config``)."""
|
||
threshold: Any
|
||
autoraise_notice_enabled: Any
|
||
enabled: Any
|
||
target_ratio: Any
|
||
protect_last: Any
|
||
tail_mode: Any
|
||
min_tail_users: Any
|
||
max_attempts: Any
|
||
proactive_prune_tokens: Any
|
||
proactive_prune_min_chars: Any
|
||
proactive_prune_min_reclaim: Any
|
||
protect_first: Any
|
||
abort_on_summary_failure: Any
|
||
model_thresholds: Any
|
||
threshold_tokens: Any
|
||
checkpoint_required: Any
|
||
in_place: Any
|
||
micro_compact: Any
|
||
micro_compact_every_n_turns: Any
|
||
micro_compact_defrag_tokens: Any
|
||
codex_app_server_auto: Any
|
||
codex_responses_native: Any
|
||
codex_responses_compact_threshold: Any
|
||
idle_compact_after_seconds: Any
|
||
|
||
|
||
def _resolve_api_mode(agent, api_mode, provider_name, base_url):
|
||
if api_mode in {"chat_completions", "codex_responses", "anthropic_messages", "bedrock_converse", "codex_app_server"}:
|
||
agent.api_mode = api_mode
|
||
elif agent.provider in {"openai-codex", "xai", "xai-oauth"}:
|
||
agent.api_mode = "codex_responses"
|
||
elif (provider_name is None) and (
|
||
agent._base_url_hostname == "chatgpt.com"
|
||
and "/backend-api/codex" in agent._base_url_lower
|
||
):
|
||
agent.api_mode = "codex_responses"
|
||
agent.provider = "openai-codex"
|
||
elif (provider_name is None) and agent._base_url_hostname == "api.x.ai":
|
||
agent.api_mode = "codex_responses"
|
||
agent.provider = "xai"
|
||
elif agent.provider == "anthropic" or (provider_name is None and agent._base_url_hostname == "api.anthropic.com"):
|
||
agent.api_mode = "anthropic_messages"
|
||
agent.provider = "anthropic"
|
||
elif agent._base_url_lower.rstrip("/").endswith("/anthropic"):
|
||
# Third-party Anthropic-compatible endpoints (e.g. MiniMax, DashScope)
|
||
# use a URL convention ending in /anthropic. Auto-detect these so the
|
||
# Anthropic Messages API adapter is used instead of chat completions.
|
||
agent.api_mode = "anthropic_messages"
|
||
elif agent.provider == "bedrock" or (
|
||
agent._base_url_hostname.startswith("bedrock-runtime.")
|
||
and base_url_host_matches(agent._base_url_lower, "amazonaws.com")
|
||
):
|
||
# AWS Bedrock — auto-detect from provider name or base URL
|
||
# (bedrock-runtime.<region>.amazonaws.com).
|
||
agent.api_mode = "bedrock_converse"
|
||
elif agent.provider in {"nous", "nous-portal", "nousresearch"}:
|
||
# Portal is dual-wire: anthropic/* → Messages, everything else →
|
||
# chat_completions. Callers that already pass api_mode win above;
|
||
# this covers direct AIAgent construction without a resolved runtime.
|
||
from hermes_cli.providers import nous_api_mode
|
||
|
||
agent.api_mode = nous_api_mode(agent.model)
|
||
else:
|
||
# Host-mandated wire check — LAST, so the provider-slug rewrites above
|
||
# (e.g. api.anthropic.com → provider="anthropic") always win. Covers
|
||
# api.meta.ai → codex_responses (prompt caching: 0% on chat vs 93-99%).
|
||
# Deliberately URL-driven, not provider-name-driven: `providers.meta`
|
||
# may point at any OpenAI-compatible endpoint without the Responses API.
|
||
try:
|
||
from hermes_cli.providers import host_mandated_api_mode as _host_mandated_api_mode
|
||
|
||
_mandated = _host_mandated_api_mode(base_url or "")
|
||
except Exception:
|
||
_mandated = None
|
||
if _mandated is not None:
|
||
agent.api_mode = _mandated
|
||
else:
|
||
agent.api_mode = "chat_completions"
|
||
|
||
|
||
def _finalize_routing(agent, api_mode, credential_pool):
|
||
# Credential-pool validation runs AFTER provider auto-detection so a pool
|
||
# scoped to "anthropic" isn't rejected for provider=None + anthropic.com URL.
|
||
if credential_pool is not None:
|
||
try:
|
||
from agent.credential_pool import credential_pool_matches_provider
|
||
|
||
if not credential_pool_matches_provider(
|
||
credential_pool,
|
||
agent.provider,
|
||
base_url=agent.base_url,
|
||
):
|
||
agent._credential_pool = None
|
||
except Exception:
|
||
agent._credential_pool = None
|
||
|
||
# Eagerly warm the transport cache so import errors surface at init,
|
||
# not mid-conversation. Also validates the api_mode is registered.
|
||
try:
|
||
agent._get_transport()
|
||
except Exception:
|
||
pass # Non-fatal — transport may not exist for all modes yet
|
||
|
||
try:
|
||
from hermes_cli.model_normalize import (
|
||
_AGGREGATOR_PROVIDERS,
|
||
normalize_model_for_provider,
|
||
)
|
||
|
||
if agent.provider not in _AGGREGATOR_PROVIDERS:
|
||
agent.model = normalize_model_for_provider(agent.model, agent.provider)
|
||
except Exception:
|
||
pass
|
||
|
||
# Auto-upgrade to Responses for GPT-5.x-style models and direct OpenAI
|
||
# URLs, unless: api_mode was explicit (user knows their endpoint), the
|
||
# runtime is ACP (`acp://` scheme — ACP clients route themselves and lack
|
||
# the Responses surface), or the URL is Azure OpenAI (serves gpt-5.x on
|
||
# /chat/completions only). Provider exceptions (e.g. Copilot gpt-5-mini)
|
||
# live in _provider_model_requires_responses_api.
|
||
if (
|
||
api_mode is None
|
||
and agent.api_mode == "chat_completions"
|
||
and agent.provider != "copilot-acp"
|
||
and not str(agent.base_url or "").lower().startswith("acp://")
|
||
and not str(agent.base_url or "").lower().startswith("acp+tcp://")
|
||
and not agent._is_azure_openai_url()
|
||
and (
|
||
agent._is_direct_openai_url()
|
||
or agent._provider_model_requires_responses_api(
|
||
agent.model,
|
||
provider=agent.provider,
|
||
)
|
||
)
|
||
):
|
||
agent.api_mode = "codex_responses"
|
||
# Invalidate the eager-warmed transport cache — api_mode changed
|
||
# from chat_completions to codex_responses after the warm at __init__.
|
||
if hasattr(agent, "_transport_cache"):
|
||
agent._transport_cache.clear()
|
||
|
||
# Pre-warm the OpenRouter model metadata cache (1h TTL) off-thread so the
|
||
# first pricing estimate doesn't block. Process-level Event guard: the
|
||
# gateway builds an AIAgent per message, and an unguarded spawn leaks one
|
||
# OS thread per message until "can't start new thread".
|
||
if (agent.provider == "openrouter" or agent._is_openrouter_url()) and \
|
||
not _ra()._openrouter_prewarm_done.is_set():
|
||
_ra()._openrouter_prewarm_done.set()
|
||
threading.Thread(
|
||
target=fetch_model_metadata,
|
||
daemon=True,
|
||
name="openrouter-prewarm",
|
||
).start()
|
||
|
||
|
||
def _init_control_state(agent):
|
||
|
||
# Tool execution state — allows _vprint during tool execution
|
||
# even when stream consumers are registered (no tokens streaming then)
|
||
agent._executing_tools = False
|
||
agent._tool_guardrails = ToolCallGuardrailController()
|
||
agent._tool_guardrail_halt_decision: ToolGuardrailDecision | None = None
|
||
|
||
# Interrupt mechanism for breaking out of tool loops
|
||
agent._interrupt_requested = False
|
||
agent._interrupt_message = None # Optional message that triggered interrupt
|
||
# Explicit hard cancellation is separate from redirect/message state. A
|
||
# thread-safe Event makes the cause atomic for auxiliary stream pollers.
|
||
agent._hard_interrupt_requested = threading.Event()
|
||
agent._execution_thread_id: int | None = None # Set at run_conversation() start
|
||
agent._interrupt_thread_signal_pending = False
|
||
agent._client_lock = threading.RLock()
|
||
agent._model_request_active = threading.Event()
|
||
agent._supports_active_turn_redirect = True
|
||
|
||
# /steer — inject a user note into the next tool result without
|
||
# interrupting: the drain hook appends it to the last tool result after
|
||
# the current batch finishes, preserving role alternation (no new user turn).
|
||
agent._pending_steer: Optional[str] = None
|
||
agent._pending_steer_lock = threading.Lock()
|
||
|
||
# Active-turn redirect mechanism. A regular follow-up sent while the model
|
||
# is generating is different from a hard /stop: preserve the valid turn
|
||
# prefix, cancel only the in-flight model request, and rebuild its tail with
|
||
# the correction. The loop drains this slot at a role-safe boundary.
|
||
agent._pending_redirect: Optional[str] = None
|
||
agent._pending_redirect_lock = threading.Lock()
|
||
|
||
# Concurrent-tool worker tids: `_set_interrupt` on `_execution_thread_id`
|
||
# alone doesn't reach ThreadPoolExecutor workers, so interrupt() /
|
||
# clear_interrupt() fan out to these explicitly.
|
||
agent._tool_worker_threads: set[int] = set()
|
||
agent._tool_worker_threads_lock = threading.Lock()
|
||
|
||
# Subagent delegation state
|
||
agent._delegate_depth = 0 # 0 = top-level agent, incremented for children
|
||
agent._active_children = [] # Running child AIAgents (for interrupt propagation)
|
||
agent._active_children_lock = threading.Lock()
|
||
|
||
# Background memory/skill review state (agent/background_review.py).
|
||
# ``_background_review_run`` is installed before the worker starts and
|
||
# fences its first provider-capable phase; the direct agent pointer keeps
|
||
# normal interrupt propagation available once the fork is constructed.
|
||
agent._background_review_agent = None
|
||
agent._background_review_run = None
|
||
agent._background_review_lock = threading.Lock()
|
||
|
||
|
||
def _init_prompt_cache_config(agent):
|
||
# Anthropic prompt caching: auto-enabled for Claude on native Anthropic,
|
||
# OpenRouter and anthropic_messages gateways (~75% input savings). Four
|
||
# breakpoints: static system prefix, full system prompt, last two messages.
|
||
# See ``_anthropic_prompt_cache_policy`` for layout-vs-transport.
|
||
agent._use_prompt_caching, agent._use_native_cache_layout = (
|
||
agent._anthropic_prompt_cache_policy()
|
||
)
|
||
agent._cache_disabled = False
|
||
# prompt_caching.cache_ttl: "5m" (default) or "1h" (2x write cost, pays off
|
||
# with >5-minute pauses); unknown values keep "5m". A falsy value (false /
|
||
# null / "off" / "disabled" / "no" / "none") disables caching entirely —
|
||
# for OAuth plans billing cache writes as extra usage, or proxies that add
|
||
# their own cache_control. The disable survives /model switches and
|
||
# fallback re-derivation via anthropic_prompt_cache_policy().
|
||
agent._cache_ttl = "5m"
|
||
try:
|
||
from hermes_cli.config import load_config_readonly as _load_pc_cfg
|
||
|
||
from agent.agent_runtime_helpers import cache_ttl_means_disabled
|
||
|
||
_pc_cfg = _load_pc_cfg().get("prompt_caching", {}) or {}
|
||
_ttl = _pc_cfg.get("cache_ttl", "5m")
|
||
if _ttl in {"5m", "1h"}:
|
||
agent._cache_ttl = _ttl
|
||
elif cache_ttl_means_disabled(_ttl):
|
||
agent._use_prompt_caching = False
|
||
agent._use_native_cache_layout = False
|
||
agent._cache_ttl = None
|
||
agent._cache_disabled = True
|
||
except Exception:
|
||
pass
|
||
|
||
|
||
def _init_turn_state(agent, run_budget_seconds):
|
||
# Iteration budget: notify the LLM only on actual exhaustion (ONE message,
|
||
# one grace call, then a forced summarise request). No intermediate
|
||
# pressure warnings — they made models give up early on complex tasks.
|
||
agent._budget_exhausted_injected = False
|
||
agent._budget_grace_call = False
|
||
|
||
# Optional wall-clock run budget (seconds per run_conversation turn).
|
||
# Explicit constructor arg wins; else resolved from config.yaml
|
||
# (agent.run_budget_seconds) further below. None = feature fully off:
|
||
# no clock reads, no injection, no stale-timeout capping.
|
||
agent.run_budget_seconds = _normalize_run_budget_seconds(run_budget_seconds)
|
||
# Wall-clock start of the CURRENT run_conversation turn. Set by
|
||
# turn_context.prepare_turn when a run budget is active; None otherwise.
|
||
agent._run_budget_started_at = None
|
||
# One-shot latch for the 80% wrap-up notice (reset each turn).
|
||
agent._run_budget_wrapup_injected = False
|
||
|
||
# Activity tracking — updated on each API call, tool execution, and
|
||
# stream chunk. Used by the gateway timeout handler to report what the
|
||
# agent was doing when it was killed, and by the "still working"
|
||
# notifications to show progress.
|
||
agent._last_activity_ts: float = time.time()
|
||
agent._last_activity_desc: str = "initializing"
|
||
# Default / unmigrated paths and _touch_activity stamp unknown; named
|
||
# provenances are stamped by compression writers (heartbeat / timeout / cooldown).
|
||
agent._last_activity_provenance = ActivityProvenance.UNKNOWN
|
||
# Rate-limit durable SessionDB activity stamps from _touch_activity (#72016).
|
||
agent._session_activity_last_persist_mono: float = 0.0
|
||
agent._current_tool: str | None = None
|
||
agent._api_call_count: int = 0
|
||
# Opt-out flag for the between-turns MCP tool refresh (build_turn_context).
|
||
# Set on internal forks (e.g. background_review) that must keep ``tools[]``
|
||
# byte-identical to a parent for provider cache parity.
|
||
agent._skip_mcp_refresh = False
|
||
# Registry generation the current tool snapshot was derived from. Lets a
|
||
# late/concurrent refresh reject a stale (older-generation) rebuild instead
|
||
# of clobbering a newer one. Set adjacent to the tool snapshot below.
|
||
agent._tool_snapshot_generation = 0
|
||
# Rate limit tracking — updated from x-ratelimit-* response headers
|
||
# after each API call. Accessed by /usage slash command.
|
||
agent._rate_limit_state: Optional["RateLimitState"] = None
|
||
|
||
# Credits tracking (dev-only, L0 usage-aware-credits) — updated from
|
||
# x-nous-credits-* response headers after each API call. Session-start
|
||
# remaining is latched the first time a header is ever seen so we can
|
||
# report cumulative micros spent. Surfaced behind HERMES_DEV_CREDITS.
|
||
agent._credits_state = None
|
||
agent._credits_session_start_micros = None
|
||
# Threshold-notice latch (L4): active sticky-notice keys + the crossing gates.
|
||
from agent.credits_tracker import new_credits_latch
|
||
|
||
agent._credits_latch = new_credits_latch()
|
||
|
||
# OpenRouter response cache hit counter — incremented when
|
||
# X-OpenRouter-Cache-Status: HIT is seen in streaming response headers.
|
||
agent._or_cache_hits: int = 0
|
||
|
||
|
||
def _setup_logging(agent):
|
||
# Centralized logging — agent.log (INFO+) and errors.log (WARNING+)
|
||
# both live under ~/.hermes/logs/. Idempotent, so gateway mode
|
||
# (which creates a new AIAgent per message) won't duplicate handlers.
|
||
from hermes_logging import setup_logging, setup_verbose_logging
|
||
setup_logging(hermes_home=_ra()._hermes_home)
|
||
|
||
if agent.verbose_logging:
|
||
setup_verbose_logging()
|
||
_ra().logger.info("Verbose logging enabled (third-party library logs suppressed)")
|
||
# Quiet mode (CLI default) deliberately does NOT raise per-logger levels:
|
||
# that would starve the root file handlers (agent.log, errors.log), since
|
||
# isEnabledFor() is checked before propagation. setup_logging() installs no
|
||
# console handler in quiet mode; any noise reduction belongs in hermes_logging.
|
||
|
||
|
||
def _init_stream_state(agent):
|
||
# Internal stream callback (set during streaming TTS).
|
||
# Initialized here so _vprint can reference it before run_conversation.
|
||
agent._stream_callback = None
|
||
# Deferred paragraph break flag — set after tool iterations so a
|
||
# single "\n\n" is prepended to the next real text delta.
|
||
agent._stream_needs_break = False
|
||
# Stateful scrubber for <memory-context> spans split across stream
|
||
# deltas (#5719). sanitize_context() alone can't survive chunk
|
||
# boundaries because the block regex needs both tags in one string.
|
||
agent._stream_context_scrubber = StreamingContextScrubber()
|
||
# Stateful scrubber for reasoning/thinking tags split across deltas (a
|
||
# per-delta regex erased '<think>' in delta1 so delta2 leaked as content).
|
||
agent._stream_think_scrubber = StreamingThinkScrubber()
|
||
# Visible assistant text already delivered through live token callbacks
|
||
# during the current model response. Used to avoid re-sending the same
|
||
# commentary when the provider later returns it as a completed interim
|
||
# assistant message.
|
||
agent._current_streamed_assistant_text = ""
|
||
# Completed interim messages delivered during the current user turn.
|
||
# Unlike token-stream tracking, this spans Codex continuation/tool calls so
|
||
# repeated commentary is not re-sent before normalization can deduplicate it.
|
||
agent._delivered_interim_texts: set[str] = set()
|
||
|
||
# Single-writer guard for the streaming delta sink: a superseded stream
|
||
# (reconnected past, socket abort raced) must not interleave tokens with
|
||
# the retry's stream. Each attempt claims a monotonic writer token; the
|
||
# sink drops chunks from threads holding a stale one. Threads that never
|
||
# claimed are never fenced, so only a superseded stream can be dropped.
|
||
agent._stream_writer_lock = threading.Lock()
|
||
agent._stream_writer_token = 0
|
||
agent._stream_writer_tls = threading.local()
|
||
agent._stream_writer_dropped = 0
|
||
|
||
# Optional current-turn user-message override used when the API-facing
|
||
# user message intentionally differs from the persisted transcript
|
||
# (e.g. CLI voice mode adds a temporary prefix for the live call only).
|
||
agent._persist_user_message_idx = None
|
||
agent._persist_user_message_override = None
|
||
agent._persist_user_message_timestamp = None
|
||
|
||
# Cache anthropic image-to-text fallbacks per image payload/URL so a
|
||
# single tool loop does not repeatedly re-run auxiliary vision on the
|
||
# same image history.
|
||
agent._anthropic_image_fallback_cache: Dict[str, str] = {}
|
||
|
||
|
||
def _bedrock_region_from_url(base_url) -> str:
|
||
"""AWS region from a bedrock-runtime.<region>.amazonaws.com URL (default us-east-1)."""
|
||
m = re.search(r"bedrock-runtime\.([a-z0-9-]+)\.", base_url or "")
|
||
return m.group(1) if m else "us-east-1"
|
||
|
||
|
||
def _init_anthropic_client(agent, api_key, base_url, _provider_timeout):
|
||
"""anthropic_messages: native Anthropic SDK (or AnthropicBedrock for Bedrock+Claude)."""
|
||
from agent.anthropic_adapter import build_anthropic_client, resolve_anthropic_token
|
||
# Bedrock + Claude → use AnthropicBedrock SDK for full feature parity
|
||
# (prompt caching, thinking budgets, adaptive thinking).
|
||
_is_bedrock_anthropic = agent.provider == "bedrock"
|
||
if _is_bedrock_anthropic:
|
||
from agent.anthropic_adapter import build_anthropic_bedrock_client
|
||
_br_region = agent._bedrock_region = _bedrock_region_from_url(base_url)
|
||
agent._anthropic_client = build_anthropic_bedrock_client(_br_region)
|
||
agent._anthropic_api_key = "aws-sdk"
|
||
agent._anthropic_base_url = base_url
|
||
agent._is_anthropic_oauth = False
|
||
agent.api_key = "aws-sdk"
|
||
agent.client = None
|
||
agent._client_kwargs = {}
|
||
if not agent.quiet_mode:
|
||
print(f"🤖 AI Agent initialized with model: {agent.model} (AWS Bedrock + AnthropicBedrock SDK, {_br_region})")
|
||
else:
|
||
# Only fall back to ANTHROPIC_TOKEN when the provider is actually Anthropic.
|
||
# Other anthropic_messages providers (MiniMax, Alibaba, etc.) must use their own API key.
|
||
# Falling back would send Anthropic credentials to third-party endpoints (Fixes #1739, #minimax-401).
|
||
_is_native_anthropic = agent.provider == "anthropic"
|
||
effective_key = (api_key or resolve_anthropic_token() or "") if _is_native_anthropic else (api_key or "")
|
||
|
||
# MiniMax OAuth tokens live ~15 min and the Anthropic SDK freezes
|
||
# ``api_key`` at construction, so swap in a callable token provider:
|
||
# ``build_anthropic_client`` installs an httpx hook that mints a
|
||
# fresh bearer per request (re-reading auth.json, so refreshes from
|
||
# other processes are seen). Steady-state cost: one file read +
|
||
# timestamp compare per request.
|
||
if agent.provider == "minimax-oauth" and isinstance(effective_key, str) and effective_key:
|
||
try:
|
||
from hermes_cli.auth import build_minimax_oauth_token_provider
|
||
effective_key = build_minimax_oauth_token_provider()
|
||
except Exception as _mm_exc: # noqa: BLE001 — never block startup on this
|
||
logging.getLogger(__name__).warning(
|
||
"MiniMax OAuth: failed to install per-request token provider "
|
||
"(%s); falling back to static bearer that will expire ~15min in.",
|
||
_mm_exc,
|
||
)
|
||
|
||
agent.api_key = effective_key
|
||
agent._anthropic_api_key = effective_key
|
||
agent._anthropic_base_url = base_url
|
||
# OAuth only for native Anthropic: third-party anthropic_messages
|
||
# providers (MiniMax, Kimi, GLM, LiteLLM) must never trip OAuth
|
||
# paths — those inject Claude-Code identity headers → 401/403.
|
||
from agent.anthropic_adapter import _is_oauth_token as _is_oat
|
||
agent._is_anthropic_oauth = _is_oat(effective_key) if (_is_native_anthropic and isinstance(effective_key, str)) else False
|
||
agent._anthropic_client = build_anthropic_client(effective_key, base_url, timeout=_provider_timeout)
|
||
# No OpenAI client needed for Anthropic mode
|
||
agent.client = None
|
||
agent._client_kwargs = {}
|
||
if not agent.quiet_mode:
|
||
print(f"🤖 AI Agent initialized with model: {agent.model} (Anthropic native)")
|
||
# ``effective_key`` may be a callable Entra ID bearer provider
|
||
# (Azure Foundry) — never invoke or inspect it in the banner.
|
||
from agent.azure_identity_adapter import is_token_provider
|
||
|
||
if is_token_provider(effective_key):
|
||
print("🔑 Using credentials: Microsoft Entra ID")
|
||
elif isinstance(effective_key, str) and len(effective_key) > 12:
|
||
print(f"🔑 Using token: {effective_key[:8]}...{effective_key[-4:]}")
|
||
|
||
|
||
def _init_moa_client(agent, api_key):
|
||
"""provider == "moa": virtual Mixture-of-Agents facade, no real HTTP client."""
|
||
from agent.moa_loop import build_moa_facade
|
||
agent.api_mode = "chat_completions"
|
||
|
||
# build_moa_facade wires the reference relay ("moa.reference" /
|
||
# "moa.progress" / "moa.phase" / "moa.aggregating" events through
|
||
# tool_progress_callback) so every surface shows each reference's
|
||
# answer before the aggregator acts. Display-only, never touches
|
||
# message history; shared with fallback-restore so a restored facade
|
||
# keeps emitting.
|
||
agent.client = build_moa_facade(agent, agent.model)
|
||
agent._client_kwargs = {}
|
||
agent.api_key = api_key or "moa-virtual-provider"
|
||
agent.base_url = "moa://local"
|
||
if not agent.quiet_mode:
|
||
print(f"🤖 AI Agent initialized with MoA preset: {agent.model}")
|
||
|
||
|
||
def _init_bedrock_client(agent, base_url):
|
||
"""bedrock_converse: boto3 directly, no OpenAI client."""
|
||
agent._bedrock_region = _bedrock_region_from_url(base_url)
|
||
# Guardrail config — read from config.yaml at init time.
|
||
agent._bedrock_guardrail_config = None
|
||
try:
|
||
from hermes_cli.config import load_config_readonly as _load_br_cfg
|
||
_gr = _load_br_cfg().get("bedrock", {}).get("guardrail", {})
|
||
if _gr.get("guardrail_identifier") and _gr.get("guardrail_version"):
|
||
agent._bedrock_guardrail_config = {
|
||
"guardrailIdentifier": _gr["guardrail_identifier"],
|
||
"guardrailVersion": _gr["guardrail_version"],
|
||
}
|
||
if _gr.get("stream_processing_mode"):
|
||
agent._bedrock_guardrail_config["streamProcessingMode"] = _gr["stream_processing_mode"]
|
||
if _gr.get("trace"):
|
||
agent._bedrock_guardrail_config["trace"] = _gr["trace"]
|
||
except Exception:
|
||
pass
|
||
agent.client = None
|
||
agent._client_kwargs = {}
|
||
if not agent.quiet_mode:
|
||
_gr_label = " + Guardrails" if agent._bedrock_guardrail_config else ""
|
||
print(f"🤖 AI Agent initialized with model: {agent.model} (AWS Bedrock, {agent._bedrock_region}{_gr_label})")
|
||
|
||
|
||
def _explicit_client_kwargs(agent, api_key, base_url, _provider_timeout) -> Dict[str, Any]:
|
||
"""OpenAI-client kwargs from explicit CLI/gateway credentials (auth already resolved)."""
|
||
_parsed_url = urlparse(base_url)
|
||
if _parsed_url.query:
|
||
_clean_url = urlunparse(_parsed_url._replace(query=""))
|
||
_query_params = {
|
||
k: v[0] for k, v in parse_qs(_parsed_url.query).items()
|
||
}
|
||
client_kwargs = {
|
||
"api_key": api_key,
|
||
"base_url": _clean_url,
|
||
"default_query": _query_params,
|
||
}
|
||
else:
|
||
client_kwargs = {"api_key": api_key, "base_url": base_url}
|
||
if _provider_timeout is not None:
|
||
client_kwargs["timeout"] = _provider_timeout
|
||
if agent.provider == "copilot-acp":
|
||
client_kwargs["command"] = agent.acp_command
|
||
client_kwargs["args"] = agent.acp_args
|
||
effective_base = base_url
|
||
# OpenCode Zen free tier (*-free slugs, e.g. x-preview-f-free /
|
||
# "Ox Alpha"): the Zen relay serves these ANONYMOUSLY and 401s any
|
||
# unrecognized bearer — including our keyless placeholder. Send an
|
||
# empty Authorization header to override the SDK's "Bearer <key>".
|
||
try:
|
||
from hermes_cli.models import (
|
||
OPENCODE_ZEN_FREE_KEYLESS_PLACEHOLDER,
|
||
opencode_zen_free_headers,
|
||
)
|
||
if api_key == OPENCODE_ZEN_FREE_KEYLESS_PLACEHOLDER:
|
||
client_kwargs["default_headers"] = opencode_zen_free_headers()
|
||
except Exception:
|
||
pass
|
||
_headers_for = _host_default_headers_factory(effective_base)
|
||
if _headers_for is not None:
|
||
client_kwargs["default_headers"] = _headers_for(api_key, effective_base)
|
||
elif "default_headers" not in client_kwargs:
|
||
# Fall back to profile.default_headers for providers that
|
||
# declare custom headers (e.g. Vercel AI Gateway attribution,
|
||
# Kimi User-Agent on non-kimi.com endpoints).
|
||
try:
|
||
from providers import get_provider_profile as _gpf
|
||
_ph = _gpf(agent.provider)
|
||
if _ph and _ph.default_headers:
|
||
client_kwargs["default_headers"] = dict(_ph.default_headers)
|
||
except Exception:
|
||
pass
|
||
return client_kwargs
|
||
|
||
|
||
def _routed_client_kwargs(agent, fallback_model, _provider_timeout) -> Dict[str, Any]:
|
||
"""OpenAI-client kwargs via the centralized provider router (no explicit creds).
|
||
|
||
Falls through to the init-time fallback chain, then raises with the
|
||
missing-key / no-provider diagnostic.
|
||
"""
|
||
from agent.auxiliary_client import resolve_provider_client
|
||
_routed_client, _ = resolve_provider_client(
|
||
agent.provider or "auto", model=agent.model, raw_codex=True)
|
||
if _routed_client is not None:
|
||
client_kwargs = _client_kwargs_from_routed(_routed_client, _provider_timeout)
|
||
else:
|
||
# No credentials for the configured provider: try the
|
||
# user-configured fallback chain BEFORE failing, whichever
|
||
# provider failed (an exhausted single-entry pool must not die
|
||
# with a misleading "No LLM provider configured"). Only
|
||
# explicitly named providers keep the missing-key diagnostic.
|
||
_explicit = (agent.provider or "").strip().lower()
|
||
_fb_resolved = False
|
||
for _fb in _fallback_entries(fallback_model):
|
||
try:
|
||
from hermes_cli.fallback_config import resolve_entry_api_key
|
||
_fb_explicit_key = resolve_entry_api_key(_fb)
|
||
_fb_client, _fb_model = resolve_provider_client(
|
||
_fb["provider"], model=_fb["model"], raw_codex=True,
|
||
explicit_base_url=_fb.get("base_url"),
|
||
explicit_api_key=_fb_explicit_key,
|
||
)
|
||
except Exception as _fb_exc:
|
||
logger.debug(
|
||
"Init-time fallback entry %s failed: %s",
|
||
_fb.get("provider"), _fb_exc,
|
||
)
|
||
continue
|
||
if _fb_client is not None:
|
||
agent.provider = _fb["provider"]
|
||
agent.model = _fb_model or _fb["model"]
|
||
agent._fallback_activated = True
|
||
client_kwargs = _client_kwargs_from_routed(_fb_client, _provider_timeout)
|
||
_fb_resolved = True
|
||
break
|
||
if (
|
||
not _fb_resolved
|
||
and _explicit
|
||
and _explicit not in {"auto", "openrouter", "custom"}
|
||
):
|
||
# Explicit non-OpenRouter provider with no creds and no
|
||
# usable fallback: fail fast. Use the provider's real env
|
||
# var name (alibaba → DASHSCOPE_API_KEY, not ALIBABA_API_KEY).
|
||
_env_hint = f"{_explicit.upper()}_API_KEY"
|
||
try:
|
||
from hermes_cli.auth import PROVIDER_REGISTRY
|
||
_pcfg = PROVIDER_REGISTRY.get(_explicit)
|
||
if _pcfg and _pcfg.api_key_env_vars:
|
||
_env_hint = _pcfg.api_key_env_vars[0]
|
||
except Exception:
|
||
pass
|
||
raise RuntimeError(
|
||
f"Provider '{_explicit}' is set in config.yaml but no API key "
|
||
f"was found. Set the {_env_hint} environment "
|
||
f"variable, or switch to a different provider with `hermes model`."
|
||
)
|
||
if not getattr(agent, "_fallback_activated", False):
|
||
# No provider configured — reject with a clear message.
|
||
raise RuntimeError(
|
||
"No LLM provider configured. Run `hermes model` to "
|
||
"select a provider, or run `hermes setup` for first-time "
|
||
"configuration."
|
||
)
|
||
return client_kwargs
|
||
|
||
|
||
def _init_openai_client(agent, api_key, base_url, fallback_model, _provider_timeout):
|
||
"""OpenAI-wire client: resolve kwargs, apply header/TLS policy, construct."""
|
||
if api_key and base_url:
|
||
client_kwargs = _explicit_client_kwargs(agent, api_key, base_url, _provider_timeout)
|
||
else:
|
||
client_kwargs = _routed_client_kwargs(agent, fallback_model, _provider_timeout)
|
||
try:
|
||
from agent.bedrock_adapter import configure_bedrock_openai_client_kwargs
|
||
configure_bedrock_openai_client_kwargs(
|
||
client_kwargs,
|
||
timeout=_provider_timeout,
|
||
)
|
||
except Exception:
|
||
if agent.provider == "bedrock" and "bedrock-mantle." in str(client_kwargs.get("base_url", "")):
|
||
raise
|
||
|
||
agent._client_kwargs = client_kwargs # stored for rebuilding after interrupt
|
||
|
||
# Fine-grained tool streaming for Claude on OpenRouter: without the beta
|
||
# header Anthropic buffers the whole tool call and OpenRouter's proxy
|
||
# times out during the silence.
|
||
_effective_base = str(client_kwargs.get("base_url", "")).lower()
|
||
if base_url_host_matches(_effective_base, "openrouter.ai") and "claude" in (agent.model or "").lower():
|
||
headers = client_kwargs.get("default_headers") or {}
|
||
existing_beta = headers.get("x-anthropic-beta", "")
|
||
_FINE_GRAINED = "fine-grained-tool-streaming-2025-05-14"
|
||
if _FINE_GRAINED not in existing_beta:
|
||
if existing_beta:
|
||
headers["x-anthropic-beta"] = f"{existing_beta},{_FINE_GRAINED}"
|
||
else:
|
||
headers["x-anthropic-beta"] = _FINE_GRAINED
|
||
client_kwargs["default_headers"] = headers
|
||
|
||
# model.default_headers (config.yaml) override provider/SDK defaults (WAFs
|
||
# that reject the SDK's identifying headers). Mutates agent._client_kwargs,
|
||
# which is this same dict, so the client built below sees it.
|
||
agent._apply_user_default_headers()
|
||
|
||
try:
|
||
from hermes_cli.config import (
|
||
apply_custom_provider_extra_headers_to_client_kwargs,
|
||
apply_custom_provider_tls_to_client_kwargs,
|
||
get_compatible_custom_providers,
|
||
load_config,
|
||
)
|
||
|
||
_cp_config = load_config()
|
||
_cp_entries = get_compatible_custom_providers(_cp_config)
|
||
_cp_base_url = str(client_kwargs.get("base_url") or agent.base_url or "")
|
||
apply_custom_provider_tls_to_client_kwargs(
|
||
client_kwargs,
|
||
_cp_base_url,
|
||
_cp_entries,
|
||
)
|
||
# Per-provider extra HTTP headers (providers.<name>.extra_headers /
|
||
# custom_providers[].extra_headers) — proxies, gateways, custom
|
||
# auth. Applied last so the most specific config level wins.
|
||
# SECURITY: values may carry credentials — never log them.
|
||
apply_custom_provider_extra_headers_to_client_kwargs(
|
||
client_kwargs,
|
||
_cp_base_url,
|
||
_cp_entries,
|
||
)
|
||
except Exception:
|
||
logger.debug("custom-provider TLS resolution skipped", exc_info=True)
|
||
|
||
agent.api_key = client_kwargs.get("api_key", "")
|
||
agent.base_url = client_kwargs.get("base_url", agent.base_url)
|
||
try:
|
||
from agent.ssl_guard import verify_ca_bundle_with_fallback
|
||
|
||
verify_ca_bundle_with_fallback()
|
||
agent.client = agent._create_openai_client(client_kwargs, reason="agent_init", shared=True)
|
||
if not agent.quiet_mode:
|
||
print(f"🤖 AI Agent initialized with model: {agent.model}")
|
||
if base_url:
|
||
print(f"🔗 Using custom base URL: {base_url}")
|
||
# ``api_key`` may be a callable Entra ID bearer provider (Azure
|
||
# Foundry) — never invoke or inspect it in the banner.
|
||
from agent.azure_identity_adapter import is_token_provider
|
||
|
||
key_used = client_kwargs.get("api_key", "none")
|
||
if is_token_provider(key_used):
|
||
print("🔑 Using credentials: Microsoft Entra ID")
|
||
elif isinstance(key_used, str) and key_used and key_used != "dummy-key" and len(key_used) > 12:
|
||
print(f"🔑 Using API key: {key_used[:8]}...{key_used[-4:]}")
|
||
else:
|
||
print("⚠️ Warning: API key appears invalid or missing")
|
||
except Exception as e:
|
||
raise RuntimeError(f"Failed to initialize OpenAI client: {e}")
|
||
|
||
|
||
def _build_client(agent, api_key, base_url, fallback_model):
|
||
# LLM client per wire mode. The provider router handles auth, base URL,
|
||
# headers and Codex/Anthropic wrapping (raw_codex=True: the main agent needs
|
||
# direct responses.stream()). One provider/model timeout up front so every
|
||
# construction path applies it consistently (Bedrock Claude has its own).
|
||
agent._anthropic_client = None
|
||
agent._is_anthropic_oauth = False
|
||
_provider_timeout = get_provider_request_timeout(agent.provider, agent.model)
|
||
if agent.api_mode == "anthropic_messages":
|
||
_init_anthropic_client(agent, api_key, base_url, _provider_timeout)
|
||
elif agent.provider == "moa":
|
||
_init_moa_client(agent, api_key)
|
||
elif agent.api_mode == "bedrock_converse":
|
||
_init_bedrock_client(agent, base_url)
|
||
else:
|
||
_init_openai_client(agent, api_key, base_url, fallback_model, _provider_timeout)
|
||
|
||
|
||
def _or_headers(_key, _base):
|
||
from agent.auxiliary_client import build_or_headers
|
||
return build_or_headers()
|
||
|
||
|
||
def _nvidia_headers(_key, base):
|
||
from agent.auxiliary_client import build_nvidia_nim_headers
|
||
return build_nvidia_nim_headers(base)
|
||
|
||
|
||
def _copilot_headers(_key, _base):
|
||
from hermes_cli.models import copilot_default_headers
|
||
return copilot_default_headers()
|
||
|
||
|
||
def _codex_headers(key, base):
|
||
from agent.codex_headers import codex_cloudflare_headers
|
||
return codex_cloudflare_headers(key, base_url=base)
|
||
|
||
|
||
def _xai_headers(_key, _base):
|
||
from tools.xai_http import hermes_xai_default_headers
|
||
return hermes_xai_default_headers()
|
||
|
||
|
||
# Host → default_headers factory ``(api_key, base_url) -> dict`` for explicit
|
||
# base_url client construction. Ordered: first host match wins; no match falls
|
||
# back to the provider profile's declared headers. ``_ra()`` keeps the
|
||
# ``run_agent.*`` header helpers patchable by tests.
|
||
_HOST_DEFAULT_HEADERS: List[tuple[str, Callable[[Any, str], Dict[str, str]]]] = [
|
||
("openrouter.ai", _or_headers),
|
||
("integrate.api.nvidia.com", _nvidia_headers),
|
||
("api.routermint.com", lambda _k, _b: _ra()._routermint_headers()),
|
||
("githubcopilot.com", _copilot_headers),
|
||
("api.kimi.com", lambda _k, _b: {"User-Agent": "claude-code/0.1.0"}),
|
||
("portal.qwen.ai", lambda _k, _b: _ra()._qwen_portal_headers()),
|
||
("chatgpt.com", _codex_headers),
|
||
("x.ai", _xai_headers),
|
||
]
|
||
|
||
|
||
def _host_default_headers_factory(base_url: str):
|
||
for host, factory in _HOST_DEFAULT_HEADERS:
|
||
if base_url_host_matches(base_url, host):
|
||
return factory
|
||
return None
|
||
|
||
|
||
def _client_kwargs_from_routed(client, timeout) -> Dict[str, Any]:
|
||
"""OpenAI-client kwargs mirroring a router-resolved client.
|
||
|
||
Preserves provider-specific headers the router set: the OpenAI SDK stores
|
||
caller-provided default_headers in ``_custom_headers``; older/mocked
|
||
clients may expose ``default_headers`` / ``_default_headers`` instead.
|
||
"""
|
||
kwargs = {"api_key": client.api_key, "base_url": str(client.base_url)}
|
||
if timeout is not None:
|
||
kwargs["timeout"] = timeout
|
||
headers = (
|
||
getattr(client, "_custom_headers", None)
|
||
or getattr(client, "default_headers", None)
|
||
or getattr(client, "_default_headers", None)
|
||
)
|
||
if headers:
|
||
kwargs["default_headers"] = dict(headers)
|
||
return kwargs
|
||
|
||
|
||
def _fallback_entries(fallback_model) -> List[Dict[str, Any]]:
|
||
"""Normalize legacy single-dict ``fallback_model`` / list ``fallback_providers``."""
|
||
if isinstance(fallback_model, list):
|
||
return [
|
||
f for f in fallback_model
|
||
if isinstance(f, dict) and f.get("provider") and f.get("model")
|
||
]
|
||
if isinstance(fallback_model, dict) and fallback_model.get("provider") and fallback_model.get("model"):
|
||
return [fallback_model]
|
||
return []
|
||
|
||
|
||
def _init_fallback_chain(agent, fallback_model):
|
||
# Keep a stable identity for the pool entry that supplied this runtime.
|
||
# OAuth refreshes can replace the runtime token before a failed request is
|
||
# recovered, so the mutable API-key value alone cannot reliably attribute
|
||
# the failure to its source entry.
|
||
from agent.agent_runtime_helpers import sync_credential_pool_entry_id
|
||
sync_credential_pool_entry_id(agent)
|
||
|
||
# Provider fallback chain — ordered list of backup providers tried
|
||
# when the primary is exhausted (rate-limit, overload, connection
|
||
# failure). Supports both legacy single-dict ``fallback_model`` and
|
||
# new list ``fallback_providers`` format.
|
||
agent._fallback_chain = _fallback_entries(fallback_model)
|
||
agent._fallback_index = 0
|
||
agent._fallback_activated = getattr(agent, "_fallback_activated", False)
|
||
# Legacy attribute kept for backward compat (tests, external callers)
|
||
agent._fallback_model = agent._fallback_chain[0] if agent._fallback_chain else None
|
||
if agent._fallback_chain and not agent.quiet_mode:
|
||
if len(agent._fallback_chain) == 1:
|
||
fb = agent._fallback_chain[0]
|
||
print(f"🔄 Fallback model: {fb['model']} ({fb['provider']})")
|
||
else:
|
||
print(f"🔄 Fallback chain ({len(agent._fallback_chain)} providers): " +
|
||
" → ".join(f"{f['model']} ({f['provider']})" for f in agent._fallback_chain))
|
||
|
||
|
||
def _load_tools(agent, enabled_toolsets, disabled_toolsets):
|
||
# A multiplexed gateway may enter a different HERMES_HOME after
|
||
# ``model_tools`` was first imported. Ensure that profile's keyed plugin
|
||
# manager has discovered its registrations before taking the tool snapshot.
|
||
try:
|
||
from hermes_cli.plugins import discover_plugins
|
||
|
||
discover_plugins()
|
||
except Exception:
|
||
logger.warning("Plugin discovery failed during agent setup", exc_info=True)
|
||
|
||
# Get available tools with filtering. Capture the registry generation this
|
||
# snapshot is derived from FIRST, so a later concurrent refresh can tell
|
||
# whether it holds a newer or staler view (see refresh_agent_mcp_tools).
|
||
try:
|
||
from tools.registry import registry as _snapshot_registry
|
||
agent._tool_snapshot_generation = _snapshot_registry._generation
|
||
except Exception:
|
||
agent._tool_snapshot_generation = 0
|
||
agent.tools = _ra().get_tool_definitions(
|
||
enabled_toolsets=enabled_toolsets,
|
||
disabled_toolsets=disabled_toolsets,
|
||
quiet_mode=agent.quiet_mode,
|
||
)
|
||
|
||
# Show tool configuration and store valid tool names for validation
|
||
agent.valid_tool_names = set()
|
||
if agent.tools:
|
||
agent.valid_tool_names = {tool["function"]["name"] for tool in agent.tools}
|
||
tool_names = sorted(agent.valid_tool_names)
|
||
if not agent.quiet_mode:
|
||
print(f"🛠️ Loaded {len(agent.tools)} tools: {', '.join(tool_names)}")
|
||
# Show filtering info if applied
|
||
if enabled_toolsets:
|
||
print(f" ✅ Enabled toolsets: {', '.join(enabled_toolsets)}")
|
||
if disabled_toolsets:
|
||
print(f" ❌ Disabled toolsets: {', '.join(disabled_toolsets)}")
|
||
elif not agent.quiet_mode:
|
||
print("🛠️ No tools loaded (all tools filtered out or unavailable)")
|
||
|
||
# Kanban lifecycle guidance is session-static (kanban_show is present iff
|
||
# HERMES_KANBAN_TASK is set); resolve the ~835-token block once instead of
|
||
# on every system-prompt rebuild.
|
||
from agent.prompt_builder import KANBAN_GUIDANCE
|
||
agent._kanban_worker_guidance = (
|
||
KANBAN_GUIDANCE if "kanban_show" in agent.valid_tool_names else ""
|
||
)
|
||
|
||
# Check tool requirements
|
||
if agent.tools and not agent.quiet_mode:
|
||
requirements = _ra().check_toolset_requirements()
|
||
missing_reqs = [name for name, available in requirements.items() if not available]
|
||
if missing_reqs:
|
||
print(f"⚠️ Some tools may not work due to missing requirements: {missing_reqs}")
|
||
|
||
# Show trajectory saving status
|
||
if agent.save_trajectories and not agent.quiet_mode:
|
||
print("📝 Trajectory saving enabled")
|
||
|
||
# Show ephemeral system prompt status
|
||
if agent.ephemeral_system_prompt and not agent.quiet_mode:
|
||
prompt_preview = agent.ephemeral_system_prompt[:60] + "..." if len(agent.ephemeral_system_prompt) > 60 else agent.ephemeral_system_prompt
|
||
print(f"🔒 Ephemeral system prompt: '{prompt_preview}' (not saved to trajectories)")
|
||
|
||
# Show prompt caching status
|
||
if agent._use_prompt_caching and not agent.quiet_mode:
|
||
if agent._use_native_cache_layout and agent.provider == "anthropic":
|
||
source = "native Anthropic"
|
||
elif agent._use_native_cache_layout:
|
||
source = "Anthropic-compatible endpoint"
|
||
else:
|
||
source = "Claude via OpenRouter"
|
||
print(f"💾 Prompt caching: ENABLED ({source}, {agent._cache_ttl} TTL)")
|
||
|
||
|
||
def _init_session_state(agent, session_id, session_db, parent_session_id, reasoning_config, max_tokens,
|
||
checkpoints_enabled, checkpoint_max_snapshots, checkpoint_max_total_size_mb, checkpoint_max_file_size_mb):
|
||
# Session logging setup - auto-save conversation trajectories for debugging
|
||
agent.session_start = datetime.now()
|
||
agent.session_id = session_id or (
|
||
f"{agent.session_start.strftime('%Y%m%d_%H%M%S')}_{uuid.uuid4().hex[:6]}"
|
||
)
|
||
|
||
# Expose session ID to tools (terminal, execute_code) so agents can
|
||
# reference their own session for --resume commands, cross-session
|
||
# coordination, and logging. Keep the ContextVar and os.environ
|
||
# fallback synchronized because different tool paths still read both.
|
||
try:
|
||
from gateway.session_context import set_current_session_id
|
||
|
||
set_current_session_id(agent.session_id)
|
||
except Exception:
|
||
# Preserve the root-agent legacy fallback, but never let delegated
|
||
# construction publish a child ID process-wide even if the ContextVar
|
||
# bridge itself failed to import.
|
||
try:
|
||
from agent.delegation_context import is_delegated_child_context
|
||
|
||
delegated_child = is_delegated_child_context()
|
||
except Exception:
|
||
delegated_child = False
|
||
if not delegated_child:
|
||
os.environ["HERMES_SESSION_ID"] = agent.session_id
|
||
|
||
# Session logs go into ~/.hermes/sessions/ alongside gateway sessions
|
||
hermes_home = get_hermes_home()
|
||
agent.logs_dir = hermes_home / "sessions"
|
||
agent.logs_dir.mkdir(parents=True, exist_ok=True)
|
||
# Per-session JSON snapshot writer (~/.hermes/sessions/session_{sid}.json)
|
||
# is opt-in via sessions.write_json_snapshots (default False). state.db
|
||
# is canonical — the snapshot is only useful for external tooling that
|
||
# reads the JSON files directly. See run_agent._save_session_log.
|
||
agent._session_json_enabled = False
|
||
try:
|
||
from hermes_cli.config import load_config_readonly as _load_sess_cfg
|
||
_sess_cfg = (_load_sess_cfg().get("sessions") or {})
|
||
agent._session_json_enabled = bool(_sess_cfg.get("write_json_snapshots", False))
|
||
except Exception:
|
||
pass
|
||
# logs_dir is retained unconditionally for request_dump_*.json (debug
|
||
# breadcrumb path written by agent_runtime_helpers.dump_api_request_debug).
|
||
|
||
# Track conversation messages for session logging
|
||
agent._session_messages: List[Dict[str, Any]] = []
|
||
# Responses encrypted-reasoning replay: some routes 400 with
|
||
# ``invalid_encrypted_content`` on replay; the loop then disables it for
|
||
# the session and falls back to stateless continuity.
|
||
agent._codex_reasoning_replay_enabled = True
|
||
agent._memory_write_origin = "assistant_tool"
|
||
agent._memory_write_context = "foreground"
|
||
|
||
# Cached system prompt -- built once per session, only rebuilt on compression
|
||
agent._cached_system_prompt: Optional[str] = None
|
||
# Cross-session-stable prefix of the cached prompt. It remains separate
|
||
# from the persisted string and is used only to place an early cache marker.
|
||
agent._cached_system_prompt_static: Optional[str] = None
|
||
|
||
# Filesystem checkpoint manager (transparent — not a tool)
|
||
from tools.checkpoint_manager import CheckpointManager
|
||
agent._checkpoint_mgr = CheckpointManager(
|
||
enabled=checkpoints_enabled,
|
||
max_snapshots=checkpoint_max_snapshots,
|
||
max_total_size_mb=checkpoint_max_total_size_mb,
|
||
max_file_size_mb=checkpoint_max_file_size_mb,
|
||
)
|
||
|
||
# SQLite session store (optional -- provided by CLI or gateway)
|
||
agent._session_db = session_db
|
||
# Whether close() also closes that handle. Default False: a caller-supplied
|
||
# session_db is usually the SHARED launch handle. Callers handing over a
|
||
# DEDICATED handle (gateway per-profile state.db, the lazy self-open in
|
||
# _get_session_db_for_recall) set True so teardown releases sqlite fds.
|
||
agent._owns_session_db = False
|
||
agent._parent_session_id = parent_session_id
|
||
# A close flush and the worker's turn-start flush can overlap. The durable
|
||
# marker is attached to each in-memory message dict, so its test-and-append
|
||
# sequence must be serialized per agent rather than relying on SQLite alone.
|
||
agent._session_persist_lock = threading.RLock()
|
||
# CLI retains its just-accepted user dict until turn setup can reuse it.
|
||
# This preserves the message-local durable marker if close persistence wins
|
||
# the race before the agent's normal early turn flush.
|
||
agent._pending_cli_user_message = None
|
||
agent._last_flushed_db_idx = 0 # tracks DB-write cursor to prevent duplicate writes
|
||
agent._session_db_created = False # DB row deferred to run_conversation()
|
||
# Most agents own their session row and should finalize it on close().
|
||
# Some temporary helper agents (manual compression / session-hygiene /
|
||
# background-review forks) rotate or share the session forward to a
|
||
# continuation row that must remain open after the helper is torn down;
|
||
# those callers explicitly set this flag to False.
|
||
agent._end_session_on_close = True
|
||
# When True, this agent NEVER persists to the canonical session store
|
||
# (state.db) or the JSON snapshot, regardless of session_id. Set on the
|
||
# background skill/memory review fork so its harness turn can't leak into
|
||
# the user's real session and hijack the next live turn. Default False.
|
||
agent._persist_disabled = False
|
||
agent._session_init_model_config = {
|
||
"max_iterations": agent.max_iterations,
|
||
"reasoning_config": reasoning_config,
|
||
"max_tokens": max_tokens,
|
||
}
|
||
# Persist a process-scoped --yolo launch into the session row so a later
|
||
# `hermes --resume <id>` can restore the bypass (CLI resume paths read
|
||
# model_config.yolo_mode back via SessionDB.session_yolo_enabled).
|
||
# Session-scoped /yolo toggles persist separately through
|
||
# SessionDB.set_session_yolo at toggle time.
|
||
try:
|
||
from tools.approval import _YOLO_MODE_FROZEN
|
||
if _YOLO_MODE_FROZEN:
|
||
agent._session_init_model_config["yolo_mode"] = True
|
||
except Exception:
|
||
pass
|
||
|
||
# In-memory todo list for task planning (one per agent/session)
|
||
from tools.todo_tool import TodoStore
|
||
agent._todo_store = TodoStore()
|
||
|
||
|
||
def _apply_display_config(agent, _agent_cfg, platform):
|
||
# display.show_commentary (default true): Codex phase=commentary messages
|
||
# go to the interim message path; false routes them to the reasoning channel.
|
||
_display_section = _agent_cfg.get("display", {})
|
||
if not isinstance(_display_section, dict):
|
||
_display_section = {}
|
||
agent.show_commentary = bool(_display_section.get("show_commentary", True))
|
||
|
||
# Window (seconds) for the bounded /fast auto|cold modes (agent.fast_mode).
|
||
agent.fast_auto_seconds = (_agent_cfg.get("agent") or {}).get("fast_auto_seconds", 60)
|
||
|
||
# model.lmstudio_load_mode: "explicit" (default, preload via LM Studio's
|
||
# management API) or "jit" (LM Studio just-in-time / Auto-Evict path).
|
||
_model_section = _agent_cfg.get("model", {})
|
||
if not isinstance(_model_section, dict):
|
||
_model_section = {}
|
||
agent.lmstudio_load_mode = "explicit"
|
||
_load_mode = str(_model_section.get("lmstudio_load_mode", "explicit") or "explicit").strip().lower()
|
||
if _load_mode in {"explicit", "jit"}:
|
||
agent.lmstudio_load_mode = _load_mode
|
||
else:
|
||
logger.warning(
|
||
"Invalid model.lmstudio_load_mode=%r; expected 'explicit' or 'jit'. Using explicit.",
|
||
_model_section.get("lmstudio_load_mode"),
|
||
)
|
||
|
||
# API-transport streaming (``model.streaming``, default true). Some
|
||
# self-hosted backends have broken streaming tool-call paths (markup leaks
|
||
# into text, zero tool_calls) so ``false`` seeds ``_disable_streaming`` —
|
||
# the same non-streaming path the loop falls back to at runtime. Session-
|
||
# scoped (survives model switches); orthogonal to ``display.streaming``.
|
||
agent._disable_streaming = False
|
||
_streaming = str(_model_section.get("streaming", "true")).strip().lower()
|
||
if _streaming in {"false", "0", "no", "off"}:
|
||
agent._disable_streaming = True
|
||
elif _streaming not in {"true", "1", "yes", "on"}:
|
||
logger.warning(
|
||
"Invalid model.streaming=%r; expected a boolean. Using streaming (default).",
|
||
_model_section.get("streaming"),
|
||
)
|
||
|
||
try:
|
||
agent._tool_guardrails = ToolCallGuardrailController(
|
||
ToolCallGuardrailConfig.from_mapping(
|
||
_agent_cfg.get("tool_loop_guardrails", {}),
|
||
platform=platform,
|
||
)
|
||
)
|
||
except Exception as _tlg_err:
|
||
_ra().logger.warning("Tool loop guardrail config ignored: %s", _tlg_err)
|
||
# Cache only the derived auxiliary compression context override that is
|
||
# needed later by the startup feasibility check. Avoid exposing a
|
||
# broad pseudo-public config object on the agent instance.
|
||
agent._aux_compression_context_length_config = None
|
||
|
||
|
||
def _memory_provider_init_kwargs(agent, platform) -> Dict[str, Any]:
|
||
"""Scoping kwargs for ``MemoryManager.initialize_all``.
|
||
|
||
status_callback (deterministic retain indicator) is CLI-only — gateway
|
||
status travels a different path and the indicator no-ops without it.
|
||
"""
|
||
kwargs = {
|
||
"session_id": agent.session_id,
|
||
"platform": platform or "cli",
|
||
"hermes_home": str(get_hermes_home()),
|
||
"agent_context": "primary",
|
||
}
|
||
if kwargs["platform"] == "cli":
|
||
kwargs["warning_callback"] = agent._emit_warning
|
||
kwargs["status_callback"] = agent._emit_status
|
||
# Session title (e.g. honcho derives chat-scoped session keys from it).
|
||
if agent._session_db:
|
||
try:
|
||
_st = agent._session_db.get_session_title(agent.session_id)
|
||
if _st:
|
||
kwargs["session_title"] = _st
|
||
except Exception:
|
||
pass
|
||
# Gateway user/chat identity for per-user scoping (gateway_session_key:
|
||
# stable per-chat Honcho session isolation).
|
||
for _ident in (
|
||
"user_id", "user_id_alt", "user_name", "chat_id", "chat_name",
|
||
"chat_type", "thread_id", "gateway_session_key",
|
||
):
|
||
_val = getattr(agent, f"_{_ident}")
|
||
if _val:
|
||
kwargs[_ident] = _val
|
||
# Profile identity for per-profile provider scoping
|
||
try:
|
||
from hermes_cli.profiles import get_active_profile_name
|
||
kwargs["agent_identity"] = get_active_profile_name()
|
||
kwargs["agent_workspace"] = "hermes"
|
||
except Exception:
|
||
pass
|
||
return kwargs
|
||
|
||
|
||
def _init_memory(agent, _agent_cfg, skip_memory, platform):
|
||
# Persistent memory (MEMORY.md + USER.md) -- loaded from disk
|
||
agent._memory_store = None
|
||
agent._memory_enabled = False
|
||
agent._user_profile_enabled = False
|
||
agent._memory_nudge_interval = 10
|
||
agent._turns_since_memory = 0
|
||
agent._iters_since_skill = 0
|
||
# skip_memory=True skips the external *provider*; enabled_toolsets=["memory"]
|
||
# still gets the built-in store so the memory tool never sees store=None.
|
||
# A memory entry on disabled_toolsets is not a request.
|
||
_enabled_toolsets = agent.enabled_toolsets or []
|
||
_disabled_toolsets = agent.disabled_toolsets or []
|
||
_memory_toolset_requested = (
|
||
"memory" in _enabled_toolsets and "memory" not in _disabled_toolsets
|
||
)
|
||
if not skip_memory or _memory_toolset_requested:
|
||
try:
|
||
from tools.memory_tool import (
|
||
get_builtin_memory_config,
|
||
get_builtin_memory_store_flags,
|
||
)
|
||
|
||
mem_config = get_builtin_memory_config(_agent_cfg)
|
||
agent._memory_enabled, agent._user_profile_enabled = get_builtin_memory_store_flags(
|
||
_agent_cfg
|
||
)
|
||
agent._memory_nudge_interval = int(mem_config.get("nudge_interval", 10))
|
||
if agent._memory_enabled or agent._user_profile_enabled:
|
||
from tools.memory_tool import MemoryStore
|
||
agent._memory_store = MemoryStore(
|
||
memory_char_limit=mem_config.get("memory_char_limit", 2200),
|
||
user_char_limit=mem_config.get("user_char_limit", 1375),
|
||
memory_enabled=agent._memory_enabled,
|
||
user_profile_enabled=agent._user_profile_enabled,
|
||
)
|
||
agent._memory_store.load_from_disk()
|
||
except Exception:
|
||
pass # Memory is optional -- don't break agent init
|
||
|
||
|
||
|
||
# Memory provider plugin (external — one at a time, alongside built-in)
|
||
# Reads memory.provider from config to select which plugin to activate.
|
||
agent._memory_manager = None
|
||
if not skip_memory:
|
||
try:
|
||
_mem_provider_name = mem_config.get("provider", "") if mem_config else ""
|
||
|
||
if _mem_provider_name and _mem_provider_name.strip():
|
||
from agent.memory_manager import MemoryManager as _MemoryManager
|
||
from plugins.memory import load_memory_provider as _load_mem
|
||
agent._memory_manager = _MemoryManager()
|
||
_mp = _load_mem(_mem_provider_name)
|
||
if _mp and _mp.is_available():
|
||
agent._memory_manager.add_provider(_mp)
|
||
elif _mp is not None and _mem_provider_name not in _warned_unavailable_providers:
|
||
# unavailable_reason() reads config/probes importlib — skip
|
||
# it once warned (the gateway builds an AIAgent per message).
|
||
try:
|
||
_unavailable_reason = _mp.unavailable_reason()
|
||
except Exception:
|
||
_unavailable_reason = ""
|
||
_warn_memory_provider_unavailable(_mem_provider_name, _unavailable_reason)
|
||
if agent._memory_manager.providers:
|
||
agent._memory_manager.initialize_all(
|
||
**_memory_provider_init_kwargs(agent, platform)
|
||
)
|
||
_ra().logger.info("Memory provider '%s' activated", _mem_provider_name)
|
||
else:
|
||
_ra().logger.debug("Memory provider '%s' not found or not available", _mem_provider_name)
|
||
agent._memory_manager = None
|
||
except Exception as _mpe:
|
||
_ra().logger.warning("Memory provider plugin init failed: %s", _mpe)
|
||
agent._memory_manager = None
|
||
|
||
from agent.memory_manager import inject_memory_provider_tools as _inject_memory_provider_tools
|
||
_inject_memory_provider_tools(agent)
|
||
|
||
|
||
def _apply_agent_section(agent, _agent_cfg):
|
||
# Skills config: nudge interval for skill creation reminders
|
||
agent._skill_nudge_interval = 10
|
||
try:
|
||
skills_config = _agent_cfg.get("skills", {})
|
||
agent._skill_nudge_interval = int(skills_config.get("creation_nudge_interval", 10))
|
||
except Exception:
|
||
pass
|
||
|
||
# Tool-use enforcement config: "auto" (default — matches hardcoded
|
||
# model list), true (always), false (never), or list of substrings.
|
||
_agent_section = _agent_cfg.get("agent", {})
|
||
if not isinstance(_agent_section, dict):
|
||
_agent_section = {}
|
||
agent._tool_use_enforcement = _agent_section.get("tool_use_enforcement", "auto")
|
||
|
||
# Execution-discipline guidance gate: "auto" (default — matches
|
||
# EXECUTION_GUIDANCE_MODELS), true (always), false (never), or list of
|
||
# model-name substrings. Independent of tool_use_enforcement — see
|
||
# agent/system_prompt.py for the injection gate.
|
||
agent._execution_guidance = _agent_section.get("execution_guidance", "auto")
|
||
|
||
# Wall-clock run budget from config (agent.run_budget_seconds) — only
|
||
# consulted when the constructor arg was not given. Absent/None/invalid
|
||
# keeps the feature fully off (zero behavior change in the default path).
|
||
if agent.run_budget_seconds is None:
|
||
agent.run_budget_seconds = _normalize_run_budget_seconds(
|
||
_agent_section.get("run_budget_seconds")
|
||
)
|
||
|
||
# Empty-response retry guard config (NS-503): additive
|
||
# ``agent.empty_response_guard`` subsection. Resolution is tolerant —
|
||
# a malformed section falls back to the schema defaults (guard on,
|
||
# $0.25 threshold), matching the guard's overall fail-open posture.
|
||
from agent.empty_response_guard import resolve_guard_settings
|
||
(
|
||
agent._empty_guard_enabled,
|
||
agent._empty_guard_cost_threshold_usd,
|
||
) = resolve_guard_settings(_agent_section.get("empty_response_guard"))
|
||
|
||
# Intent-ack continuation config: "auto" (default — codex_responses only,
|
||
# the historical gate), true (all api_modes), false (never), or a list of
|
||
# model-name substrings. Resolved against the active api_mode/model in the
|
||
# conversation loop's intent-ack block.
|
||
agent._intent_ack_continuation = _agent_section.get("intent_ack_continuation", "auto")
|
||
|
||
# Runtime anti-stall guards (identical-call loop-breaker notice on tool
|
||
# results + continue-intent extension of the empty-response recovery).
|
||
# Single boolean gate, default True. Notice-only — never blocks a call.
|
||
agent._stall_guards = bool(_agent_section.get("stall_guards", True))
|
||
|
||
# Universal guidance toggles (apply to ALL models, unlike enforcement):
|
||
# task-completion, parallel-tool-call batching, and the local Python
|
||
# toolchain probe (False skips the subprocess probe + prompt line).
|
||
agent._task_completion_guidance = bool(_agent_section.get("task_completion_guidance", True))
|
||
agent._parallel_tool_call_guidance = bool(_agent_section.get("parallel_tool_call_guidance", True))
|
||
agent._environment_probe = bool(_agent_section.get("environment_probe", True))
|
||
# Warm the probe off-thread (~0.5s of python3/pip subprocesses) during
|
||
# init so the FIRST system-prompt build — on the time-to-first-token
|
||
# path — finds the line already cached.
|
||
if agent._environment_probe:
|
||
try:
|
||
from tools.env_probe import warm_environment_probe_async
|
||
warm_environment_probe_async()
|
||
except Exception:
|
||
pass
|
||
|
||
# Bot Mode teammate protocol section (tools/bot_mode_probe.py) — pure
|
||
# filesystem reads, no warm needed. Silent on non-Bot-Mode installs.
|
||
agent._bot_mode_protocol = bool(_agent_section.get("bot_mode_protocol", True))
|
||
# Session-title hint for the "Bot Chat" gate: hosts that defer the DB
|
||
# title write past the first prompt build (tui_gateway pending_title)
|
||
# set this so the gate doesn't depend on write ordering.
|
||
agent._session_title_hint = None
|
||
|
||
# Per-platform prompt-hint overrides (config.yaml → platform_hints:
|
||
# <platform>: {append: ..} | {replace: ..}). Stored verbatim; resolved in
|
||
# agent/system_prompt.py. Invalid shapes are ignored so a bad entry can
|
||
# never break prompt assembly.
|
||
_platform_hints_cfg = _agent_cfg.get("platform_hints", {})
|
||
if not isinstance(_platform_hints_cfg, dict):
|
||
_platform_hints_cfg = {}
|
||
agent._platform_hint_overrides = _platform_hints_cfg
|
||
|
||
# App-level API retry count (wraps each model API call). Default 3,
|
||
# overridable via agent.api_max_retries in config.yaml. See #11616.
|
||
try:
|
||
_raw_api_retries = _agent_section.get("api_max_retries", 3)
|
||
_api_retries = int(_raw_api_retries)
|
||
_api_retries = max(_api_retries, 1) # 1 = no retry (single attempt)
|
||
except (TypeError, ValueError):
|
||
_api_retries = 3
|
||
agent._api_max_retries = _api_retries
|
||
|
||
|
||
def _parse_compression_config(agent, _agent_cfg):
|
||
# Initialize context compressor for automatic context management
|
||
# Compresses conversation when approaching model's context limit
|
||
# Configuration via config.yaml (compression section)
|
||
_compression_cfg = _agent_cfg.get("compression", {})
|
||
if not isinstance(_compression_cfg, dict):
|
||
_compression_cfg = {}
|
||
compression_threshold = float(_compression_cfg.get("threshold", 0.50))
|
||
# Per-model compaction-threshold override: Codex gpt-5.4/5.5 raise to 85%
|
||
# (backend caps at 272K, so 50% would compact at ~136K). Opt-out flag
|
||
# restores the global threshold; when it fires a one-time notice is
|
||
# stashed for the first turn, with its own display gate.
|
||
_codex_gpt55_autoraise = _cfg_flag(_compression_cfg, "codex_gpt55_autoraise", True)
|
||
_codex_gpt55_autoraise_notice = _cfg_flag(_compression_cfg, "codex_gpt55_autoraise_notice", True)
|
||
agent._compression_threshold_autoraised = None
|
||
try:
|
||
from agent.auxiliary_client import (
|
||
_compression_threshold_for_model as _cthresh_fn,
|
||
_is_codex_gpt54_or_gpt55 as _is_codex_gpt54_or_gpt55_fn,
|
||
_is_codex_spark as _is_codex_spark_fn,
|
||
)
|
||
_model_cthresh = _cthresh_fn(
|
||
agent.model,
|
||
agent.provider,
|
||
allow_codex_gpt55_autoraise=_codex_gpt55_autoraise,
|
||
)
|
||
# Codex autoraises apply only when they RAISE (never lower a higher
|
||
# global threshold); Arcee Trinity keeps its unconditional override.
|
||
compression_threshold, agent._compression_threshold_autoraised = (
|
||
_resolve_compression_threshold(
|
||
compression_threshold,
|
||
_model_cthresh,
|
||
model=agent.model,
|
||
is_codex_autoraise=(
|
||
_is_codex_gpt54_or_gpt55_fn(agent.model, agent.provider)
|
||
or _is_codex_spark_fn(agent.model, agent.provider)
|
||
),
|
||
)
|
||
)
|
||
except Exception:
|
||
pass
|
||
compression_enabled = _cfg_flag(_compression_cfg, "enabled", True)
|
||
compression_target_ratio = float(_compression_cfg.get("target_ratio", 0.20))
|
||
compression_protect_last = int(_compression_cfg.get("protect_last_n", 20))
|
||
# compression.tail_mode: "lean" (default) keeps a clamped 2.5%/10K-25K
|
||
# verbatim tail — continuity rides the summary (digests, anchors, pointers).
|
||
# "legacy" restores the 0.20*threshold tail, which hoards 100-240K tokens
|
||
# on big windows. Unknown values fall back to lean inside the compressor.
|
||
compression_tail_mode = str(_compression_cfg.get("tail_mode", "lean")).strip().lower()
|
||
# compression.min_tail_user_messages: actionable user messages guaranteed
|
||
# to survive in the uncompressed tail (default 1, floor 1).
|
||
compression_min_tail_users = max(
|
||
1, _parse_config_int(_compression_cfg.get("min_tail_user_messages", 1), 1)
|
||
)
|
||
# compression.max_attempts: retry rounds before "max compression attempts
|
||
# reached". Some sessions legitimately need >3 (incompressible tool
|
||
# schemas keep the estimate above threshold). Default 3, floor 1, cap 10;
|
||
# bool/fractional/garbage → 3 (see _parse_config_int).
|
||
compression_max_attempts = _parse_config_int(_compression_cfg.get("max_attempts", 3), 3)
|
||
if compression_max_attempts < 1:
|
||
compression_max_attempts = 3
|
||
compression_max_attempts = min(compression_max_attempts, 10)
|
||
|
||
# Opt-in proactive tool-result prune trigger (0 = disabled — the
|
||
# default, so an unset key is behavior-neutral). Negative values are
|
||
# treated as disabled rather than erroring.
|
||
compression_proactive_prune_tokens = max(
|
||
0, _parse_config_int(_compression_cfg.get("proactive_prune_tokens", 0), 0)
|
||
)
|
||
compression_proactive_prune_min_chars = _parse_config_int(
|
||
_compression_cfg.get("proactive_prune_min_result_chars", 8000), 8000
|
||
)
|
||
compression_proactive_prune_min_reclaim = max(
|
||
0,
|
||
_parse_config_int(
|
||
_compression_cfg.get("proactive_prune_min_reclaim_tokens", 4096), 4096
|
||
),
|
||
)
|
||
# protect_first_n: non-system head messages to protect (system prompt is
|
||
# always protected). 0 = "system prompt + summary + tail" is legitimate.
|
||
compression_protect_first = max(
|
||
0, int(_compression_cfg.get("protect_first_n", 3))
|
||
)
|
||
compression_abort_on_summary_failure = _cfg_flag(_compression_cfg, "abort_on_summary_failure", False)
|
||
# Per-model threshold overrides: keys are substring-matched against the
|
||
# model name (longest match wins). Empty dict = use the global threshold
|
||
# for all models (backward compatible).
|
||
_raw_model_thresholds = _compression_cfg.get("model_thresholds", {})
|
||
if isinstance(_raw_model_thresholds, dict):
|
||
compression_model_thresholds = {
|
||
str(k): float(v) for k, v in _raw_model_thresholds.items()
|
||
if isinstance(v, (int, float)) and not isinstance(v, bool)
|
||
}
|
||
else:
|
||
compression_model_thresholds = {}
|
||
# Absolute token cap: when set, compression triggers at the lower of
|
||
# the ratio-based threshold and this absolute count. Clamped to the
|
||
# model's context length at apply-time so a cap above the window is
|
||
# a no-op (ratio-based threshold wins).
|
||
compression_threshold_tokens = _compression_cfg.get("threshold_tokens")
|
||
if compression_threshold_tokens is not None:
|
||
try:
|
||
compression_threshold_tokens = int(compression_threshold_tokens)
|
||
if compression_threshold_tokens <= 0:
|
||
compression_threshold_tokens = None
|
||
except (TypeError, ValueError):
|
||
compression_threshold_tokens = None
|
||
compression_checkpoint_required = is_truthy_value(
|
||
_compression_cfg.get("checkpoint_required"), default=False
|
||
)
|
||
_refuse_checkpoint_required_on_codex_app_server(
|
||
compression_checkpoint_required, getattr(agent, "api_mode", None)
|
||
)
|
||
# In-place compaction: compress_context() rewrites messages + system prompt
|
||
# WITHOUT rotating the session id (no parent chain / `name #N`). Rides on
|
||
# the agent, not the compressor. default=True MUST match
|
||
# DEFAULT_CONFIG["compression"]["in_place"] — a False default flipped
|
||
# agents into rotation mode whenever the merged config omitted the key.
|
||
compression_in_place = is_truthy_value(
|
||
_compression_cfg.get("in_place"), default=True
|
||
)
|
||
# Opt-in (default False): micro-compaction rewrites already-sent history
|
||
# per turn, breaking the prompt-cache prefix on a per-turn cadence — the
|
||
# very cost proactive_prune_min_reclaim_tokens exists to amortize.
|
||
compression_micro_compact = is_truthy_value(
|
||
_compression_cfg.get("micro_compact"), default=False
|
||
)
|
||
# How often a pass runs, in completed turns. Each pass rewrites
|
||
# already-sent history and costs one prompt-cache break, so this is the
|
||
# dial for how often that cost is paid: 1 = every turn (most aggressive
|
||
# reclaim), 5 = one break per five turns. Clamped to >= 1.
|
||
compression_micro_compact_every_n_turns = max(
|
||
1,
|
||
_parse_config_int(_compression_cfg.get("micro_compact_every_n_turns", 1), 1),
|
||
)
|
||
# Rolling-summary defrag threshold, in tokens. Lived on the compressor as
|
||
# a hardcoded attribute with no path from config until now.
|
||
compression_micro_compact_defrag_tokens = max(
|
||
1,
|
||
_parse_config_int(
|
||
_compression_cfg.get("micro_compact_defrag_threshold_tokens", 2000),
|
||
2000,
|
||
),
|
||
)
|
||
codex_app_server_auto_compaction = str(
|
||
_compression_cfg.get("codex_app_server_auto", "native") or "native"
|
||
).lower()
|
||
if codex_app_server_auto_compaction not in {"native", "hermes", "off"}:
|
||
_ra().logger.warning(
|
||
"Invalid compression.codex_app_server_auto=%r; using 'native'. "
|
||
"Valid values are: native, hermes, off.",
|
||
codex_app_server_auto_compaction,
|
||
)
|
||
codex_app_server_auto_compaction = "native"
|
||
# Native OpenAI Responses server-side compaction (opt-in; per-request gate
|
||
# in agent/native_compaction.py). Truthy coercion: "false"/"off" strings
|
||
# must stay disabled (bool("false") is True).
|
||
codex_responses_native_compaction = is_truthy_value(
|
||
_compression_cfg.get("codex_responses_native", False)
|
||
)
|
||
_native_threshold_raw = _compression_cfg.get("codex_responses_compact_threshold")
|
||
codex_responses_compact_threshold = None
|
||
if _native_threshold_raw is not None:
|
||
try:
|
||
if isinstance(_native_threshold_raw, (bool, float)):
|
||
raise ValueError
|
||
codex_responses_compact_threshold = int(_native_threshold_raw)
|
||
if codex_responses_compact_threshold <= 0:
|
||
raise ValueError
|
||
except (TypeError, ValueError):
|
||
_ra().logger.warning(
|
||
"Invalid compression.codex_responses_compact_threshold=%r; "
|
||
"using the automatic threshold derived from local compression.",
|
||
_native_threshold_raw,
|
||
)
|
||
codex_responses_compact_threshold = None
|
||
# Opt-in idle compaction: compact a session up front when it resumes after
|
||
# this many seconds of inactivity (0 = disabled). Time-based, so it
|
||
# complements the size-based threshold above. Consumed by build_turn_context().
|
||
compression_idle_compact_after_seconds = max(
|
||
0, int(_compression_cfg.get("idle_compact_after_seconds", 0))
|
||
)
|
||
return CompressionSettings(
|
||
threshold=compression_threshold,
|
||
autoraise_notice_enabled=_codex_gpt55_autoraise_notice,
|
||
enabled=compression_enabled,
|
||
target_ratio=compression_target_ratio,
|
||
protect_last=compression_protect_last,
|
||
tail_mode=compression_tail_mode,
|
||
min_tail_users=compression_min_tail_users,
|
||
max_attempts=compression_max_attempts,
|
||
proactive_prune_tokens=compression_proactive_prune_tokens,
|
||
proactive_prune_min_chars=compression_proactive_prune_min_chars,
|
||
proactive_prune_min_reclaim=compression_proactive_prune_min_reclaim,
|
||
protect_first=compression_protect_first,
|
||
abort_on_summary_failure=compression_abort_on_summary_failure,
|
||
model_thresholds=compression_model_thresholds,
|
||
threshold_tokens=compression_threshold_tokens,
|
||
checkpoint_required=compression_checkpoint_required,
|
||
in_place=compression_in_place,
|
||
micro_compact=compression_micro_compact,
|
||
micro_compact_every_n_turns=compression_micro_compact_every_n_turns,
|
||
micro_compact_defrag_tokens=compression_micro_compact_defrag_tokens,
|
||
codex_app_server_auto=codex_app_server_auto_compaction,
|
||
codex_responses_native=codex_responses_native_compaction,
|
||
codex_responses_compact_threshold=codex_responses_compact_threshold,
|
||
idle_compact_after_seconds=compression_idle_compact_after_seconds,
|
||
)
|
||
|
||
|
||
def _warn_invalid_config_int(
|
||
what: str, value: Any, requirement: str, fallback: str, print_fallback: str = "",
|
||
) -> None:
|
||
"""Log + stderr-print an invalid integer config value.
|
||
|
||
``print_fallback`` lets the user-facing line keep its historical wording
|
||
where it differs from the log line.
|
||
"""
|
||
_ra().logger.warning(
|
||
"Invalid %s: %r — %s. Falling back to %s.", what, value, requirement, fallback,
|
||
)
|
||
print(
|
||
f"\n⚠ Invalid {what}: {value!r}\n"
|
||
f" {requirement[0].upper() + requirement[1:]}.\n"
|
||
f" Falling back to {print_fallback or fallback}.\n",
|
||
file=sys.stderr,
|
||
)
|
||
|
||
|
||
def _custom_provider_configured_base_url(
|
||
_configured_provider: str, _agent_cfg, _custom_providers
|
||
) -> str:
|
||
"""Base URL of a named custom provider (``providers.<name>`` first, then
|
||
``custom_providers``), normalized for route comparison; "" if unknown.
|
||
Disabled ``providers.*`` entries also mask their ``custom_providers`` twin.
|
||
"""
|
||
_configured_base_url = ""
|
||
_configured_custom_provider = _normalize_custom_provider_name(
|
||
_configured_provider
|
||
)
|
||
_user_providers = _agent_cfg.get("providers")
|
||
_disabled_custom_provider_ids: set[str] = set()
|
||
if isinstance(_user_providers, dict):
|
||
from hermes_cli.config import is_provider_enabled
|
||
|
||
for _provider_key, _provider_entry in _user_providers.items():
|
||
if not isinstance(_provider_entry, dict):
|
||
continue
|
||
_entry_name = str(
|
||
_provider_entry.get("name") or ""
|
||
).strip()
|
||
_entry_provider_ids = _custom_provider_runtime_ids(
|
||
_provider_key
|
||
) | _custom_provider_runtime_ids(_entry_name)
|
||
if not is_provider_enabled(_provider_entry):
|
||
_disabled_custom_provider_ids.update(
|
||
provider_id
|
||
for provider_id in _entry_provider_ids
|
||
if provider_id
|
||
)
|
||
continue
|
||
if _configured_custom_provider not in _entry_provider_ids:
|
||
continue
|
||
_configured_base_url = _normalize_route_base_url(
|
||
_provider_entry.get("api")
|
||
or _provider_entry.get("url")
|
||
or _provider_entry.get("base_url")
|
||
)
|
||
if _configured_base_url:
|
||
break
|
||
if not _configured_base_url:
|
||
for _provider_entry in _custom_providers:
|
||
if not isinstance(_provider_entry, dict):
|
||
continue
|
||
_entry_name = str(
|
||
_provider_entry.get("name") or ""
|
||
).strip()
|
||
_entry_provider_key = str(
|
||
_provider_entry.get("provider_key") or ""
|
||
).strip().lower()
|
||
_entry_provider_ids = _custom_provider_runtime_ids(
|
||
_entry_name
|
||
) | _custom_provider_runtime_ids(_entry_provider_key)
|
||
if (
|
||
_entry_provider_key
|
||
and _custom_provider_runtime_ids(_entry_provider_key)
|
||
& _disabled_custom_provider_ids
|
||
):
|
||
continue
|
||
if _configured_custom_provider not in _entry_provider_ids:
|
||
continue
|
||
_configured_base_url = _normalize_route_base_url(
|
||
_provider_entry.get("base_url")
|
||
)
|
||
if _configured_base_url:
|
||
break
|
||
return _configured_base_url
|
||
|
||
|
||
def _scope_context_length_to_default_runtime(
|
||
agent, _agent_cfg, _model_cfg, _custom_providers, _config_context_length, base_url
|
||
) -> Optional[int]:
|
||
"""Return ``model.context_length`` only if it describes the active runtime.
|
||
|
||
``model.context_length`` describes the configured default model. A process
|
||
launched directly with ``--model`` / ``-m`` has already replaced
|
||
``agent.model`` before this initializer loads config, so carrying the
|
||
default model's explicit window into that different runtime is stale. The
|
||
live switch/fallback paths already clear this override; keep direct-start
|
||
overrides consistent with them and let provider metadata resolve the
|
||
active model's window instead.
|
||
"""
|
||
_default = _model_cfg.get("default")
|
||
if isinstance(_default, dict):
|
||
from hermes_cli.config import split_model_config_default
|
||
_default, _ = split_model_config_default(_default)
|
||
_configured_default_model = str(_default or "").strip()
|
||
_configured_default_runtime_model = _configured_default_model
|
||
_active_runtime_model = agent.model
|
||
if _configured_default_model:
|
||
try:
|
||
from hermes_cli.model_normalize import normalize_model_for_provider
|
||
|
||
_configured_default_runtime_model = normalize_model_for_provider(
|
||
_configured_default_model, agent.provider
|
||
)
|
||
_active_runtime_model = normalize_model_for_provider(
|
||
agent.model, agent.provider
|
||
)
|
||
except Exception:
|
||
pass
|
||
_configured_provider = str(_model_cfg.get("provider") or "").strip()
|
||
_configured_base_url = _normalize_route_base_url(
|
||
_model_cfg.get("base_url")
|
||
)
|
||
_configured_provider_norm = _normalize_custom_provider_name(
|
||
_configured_provider
|
||
)
|
||
_custom_provider_candidate = bool(_configured_provider_norm)
|
||
_runtime_first_provider_ids = {
|
||
"auto",
|
||
"moa",
|
||
"vertex",
|
||
"google-vertex",
|
||
"vertex-ai",
|
||
"gcp-vertex",
|
||
"vertexai",
|
||
}
|
||
if _configured_provider_norm in _runtime_first_provider_ids:
|
||
_custom_provider_candidate = False
|
||
elif (
|
||
_custom_provider_candidate
|
||
and _configured_provider_norm != "custom"
|
||
and not _configured_provider_norm.startswith("custom:")
|
||
):
|
||
try:
|
||
from hermes_cli.auth import resolve_provider as resolve_auth_provider
|
||
|
||
_resolved_auth_provider = resolve_auth_provider(
|
||
_configured_provider_norm
|
||
)
|
||
_custom_provider_candidate = (
|
||
str(_resolved_auth_provider or "").strip().lower()
|
||
!= _configured_provider_norm
|
||
)
|
||
except Exception:
|
||
pass
|
||
if not _configured_base_url and _custom_provider_candidate:
|
||
_configured_base_url = _custom_provider_configured_base_url(
|
||
_configured_provider, _agent_cfg, _custom_providers
|
||
)
|
||
_active_route_url = str(agent.base_url or "")
|
||
_requested_route_url = str(base_url or "")
|
||
if "?" in _requested_route_url.split("#", 1)[0]:
|
||
try:
|
||
_requested_parts = urlparse(_requested_route_url)
|
||
_requested_without_query = urlunparse(
|
||
_requested_parts._replace(query="")
|
||
)
|
||
if _normalize_route_base_url(
|
||
_requested_without_query
|
||
) == _normalize_route_base_url(_active_route_url):
|
||
_active_route_url = _requested_route_url
|
||
except (TypeError, ValueError):
|
||
pass
|
||
_active_base_url = _normalize_route_base_url(_active_route_url)
|
||
_route_mismatch = _context_route_mismatch(
|
||
_configured_base_url,
|
||
_active_base_url,
|
||
_configured_provider,
|
||
agent.provider,
|
||
already_normalized=True,
|
||
)
|
||
_model_mismatch = bool(
|
||
_configured_default_runtime_model
|
||
and _configured_default_runtime_model != _active_runtime_model
|
||
)
|
||
if _model_mismatch or _route_mismatch:
|
||
_ra().logger.debug(
|
||
"Ignoring model.context_length=%s for startup runtime %s at %s "
|
||
"(configured default is %s at %s)",
|
||
_config_context_length,
|
||
agent.model,
|
||
_active_base_url or agent.provider,
|
||
_configured_default_model,
|
||
_configured_base_url or _model_cfg.get("provider"),
|
||
)
|
||
return None
|
||
return _config_context_length
|
||
|
||
|
||
def _resolve_context_length(agent, _agent_cfg, base_url):
|
||
# Read optional explicit context_length override for the auxiliary
|
||
# compression model. Custom endpoints often cannot report this via
|
||
# /models, so the startup feasibility check needs the config hint.
|
||
try:
|
||
_aux_cfg = cfg_get(_agent_cfg, "auxiliary", "compression", default={})
|
||
except Exception:
|
||
_aux_cfg = {}
|
||
if isinstance(_aux_cfg, dict):
|
||
_aux_context_config = _aux_cfg.get("context_length")
|
||
else:
|
||
_aux_context_config = None
|
||
if _aux_context_config is not None:
|
||
try:
|
||
_aux_context_config = int(_aux_context_config)
|
||
except (TypeError, ValueError):
|
||
_aux_context_config = None
|
||
agent._aux_compression_context_length_config = _aux_context_config
|
||
|
||
# Read explicit model output-token override from config when the
|
||
# caller did not pass one directly.
|
||
_model_cfg = _agent_cfg.get("model", {})
|
||
if agent.max_tokens is None and isinstance(_model_cfg, dict):
|
||
_config_max_tokens = _model_cfg.get("max_tokens")
|
||
if _config_max_tokens is not None:
|
||
try:
|
||
if isinstance(_config_max_tokens, bool):
|
||
raise ValueError
|
||
_parsed_max_tokens = int(_config_max_tokens)
|
||
if _parsed_max_tokens <= 0:
|
||
raise ValueError
|
||
agent.max_tokens = _parsed_max_tokens
|
||
except (TypeError, ValueError):
|
||
_warn_invalid_config_int(
|
||
"model.max_tokens in config.yaml", _config_max_tokens,
|
||
"must be a positive integer (e.g. 4096)", "provider default",
|
||
)
|
||
agent._session_init_model_config["max_tokens"] = agent.max_tokens
|
||
|
||
# Read explicit context_length override from model config
|
||
if isinstance(_model_cfg, dict):
|
||
_config_context_length = _model_cfg.get("context_length")
|
||
else:
|
||
_config_context_length = None
|
||
if _config_context_length is not None:
|
||
try:
|
||
_config_context_length = int(_config_context_length)
|
||
except (TypeError, ValueError):
|
||
_warn_invalid_config_int(
|
||
"model.context_length in config.yaml", _config_context_length,
|
||
"must be a plain integer (e.g. 256000, not '256K')",
|
||
"auto-detection", "auto-detected context window",
|
||
)
|
||
_config_context_length = None
|
||
|
||
# Resolve custom_providers once before route-scoping a global context pin:
|
||
# a named custom provider may keep its base URL only in this list rather
|
||
# than repeating it under ``model``.
|
||
try:
|
||
from hermes_cli.config import get_compatible_custom_providers
|
||
_custom_providers = get_compatible_custom_providers(_agent_cfg)
|
||
except Exception:
|
||
_custom_providers = _agent_cfg.get("custom_providers")
|
||
if not isinstance(_custom_providers, list):
|
||
_custom_providers = []
|
||
|
||
# ``model.context_length`` describes the configured default model; drop it
|
||
# when the startup runtime (model or route) differs from that default.
|
||
if _config_context_length is not None and isinstance(_model_cfg, dict):
|
||
_config_context_length = _scope_context_length_to_default_runtime(
|
||
agent, _agent_cfg, _model_cfg, _custom_providers, _config_context_length, base_url
|
||
)
|
||
|
||
# Store for reuse by _check_compression_model_feasibility (auxiliary
|
||
# compression model context-length detection needs the same list).
|
||
agent._custom_providers = _custom_providers
|
||
_merge_custom_provider_extra_body(agent, _custom_providers)
|
||
|
||
# Check custom_providers per-model context_length
|
||
if _config_context_length is None and _custom_providers:
|
||
try:
|
||
from hermes_cli.config import get_custom_provider_context_length
|
||
_cp_ctx_resolved = get_custom_provider_context_length(
|
||
model=agent.model,
|
||
base_url=agent.base_url,
|
||
custom_providers=_custom_providers,
|
||
)
|
||
if _cp_ctx_resolved:
|
||
_config_context_length = int(_cp_ctx_resolved)
|
||
except Exception:
|
||
_cp_ctx_resolved = None
|
||
|
||
# Surface a clear warning if the user set a context_length but it
|
||
# wasn't a valid positive int — the helper silently skips those.
|
||
if _config_context_length is None:
|
||
_target = _normalize_route_base_url(agent.base_url)
|
||
for _cp_entry in _custom_providers:
|
||
if not isinstance(_cp_entry, dict):
|
||
continue
|
||
_cp_url = _normalize_route_base_url(_cp_entry.get("base_url"))
|
||
if _target and _cp_url == _target:
|
||
_cp_models = _cp_entry.get("models", {})
|
||
if isinstance(_cp_models, dict):
|
||
_cp_model_cfg = _cp_models.get(agent.model, {})
|
||
if isinstance(_cp_model_cfg, dict):
|
||
_cp_ctx = _cp_model_cfg.get("context_length")
|
||
if _cp_ctx is not None:
|
||
try:
|
||
_parsed = int(_cp_ctx)
|
||
if _parsed <= 0:
|
||
raise ValueError
|
||
except (TypeError, ValueError):
|
||
_warn_invalid_config_int(
|
||
f"context_length for model {agent.model!r} in custom_providers",
|
||
_cp_ctx,
|
||
"must be a positive integer (e.g. 256000, not '256K')",
|
||
"auto-detection", "auto-detected context window",
|
||
)
|
||
break
|
||
|
||
# Persist for reuse on switch_model / fallback activation. Must come
|
||
# AFTER the custom_providers branch so per-model overrides aren't lost.
|
||
agent._config_context_length = _config_context_length
|
||
|
||
_lmstudio_runtime_context_length = agent._ensure_lmstudio_runtime_loaded(
|
||
_config_context_length
|
||
)
|
||
if agent._lmstudio_load_was_unverified(_lmstudio_runtime_context_length):
|
||
_ra().logger.warning(
|
||
"LM Studio model activation was rejected or completed without a "
|
||
"verifiable active context length; falling back to configured context"
|
||
)
|
||
_effective_context_length = agent._effective_lmstudio_context_length(
|
||
_config_context_length,
|
||
_lmstudio_runtime_context_length,
|
||
)
|
||
return _config_context_length, _custom_providers, _effective_context_length, _model_cfg
|
||
|
||
|
||
def _build_context_engine(agent, _agent_cfg, cs, _custom_providers, _effective_context_length, session_db):
|
||
# Select context engine: config-driven (like memory providers).
|
||
# 1. Check config.yaml context.engine setting
|
||
# 2. Check plugins/context_engine/<name>/ directory (repo-shipped)
|
||
# 3. Check general plugin system (user-installed plugins)
|
||
# 4. Fall back to built-in ContextCompressor
|
||
_selected_engine = None
|
||
_copy_failed = False
|
||
_engine_name = "compressor" # default
|
||
try:
|
||
_engine_name = _agent_cfg.get("context", {}).get("engine", "compressor") or "compressor"
|
||
except Exception:
|
||
pass
|
||
|
||
if _engine_name != "compressor":
|
||
# Try loading from plugins/context_engine/<name>/
|
||
try:
|
||
from plugins.context_engine import load_context_engine
|
||
_selected_engine = load_context_engine(_engine_name)
|
||
except Exception as _ce_load_err:
|
||
_ra().logger.debug("Context engine load from plugins/context_engine/: %s", _ce_load_err)
|
||
|
||
# Try general plugin system as fallback
|
||
if _selected_engine is None:
|
||
_candidate = None
|
||
try:
|
||
from hermes_cli.plugins import get_plugin_context_engine
|
||
_candidate = get_plugin_context_engine()
|
||
except Exception:
|
||
_candidate = None
|
||
if _candidate is not None and _candidate.name == _engine_name:
|
||
# Deep-copy the shared plugin singleton so a child's
|
||
# update_model() can't mutate the parent's compressor. Uncopyable
|
||
# state (locks, DB conns) → built-in compressor with an ACCURATE
|
||
# message, not "not found".
|
||
import copy
|
||
try:
|
||
_selected_engine = copy.deepcopy(_candidate)
|
||
except Exception as _copy_err:
|
||
_copy_failed = True
|
||
_ra().logger.warning(
|
||
"Context engine '%s' could not be safely copied for this "
|
||
"agent (%s) — falling back to built-in compressor. Plugin "
|
||
"engines that hold uncopyable state (locks, DB connections) "
|
||
"should implement __deepcopy__ to copy only mutable budget "
|
||
"state.",
|
||
_engine_name, _copy_err,
|
||
)
|
||
_selected_engine = None
|
||
|
||
if _selected_engine is None and not _copy_failed:
|
||
_ra().logger.warning(
|
||
"Context engine '%s' not found — falling back to built-in compressor",
|
||
_engine_name,
|
||
)
|
||
# else: config says "compressor" — use built-in, don't auto-activate plugins
|
||
|
||
if _selected_engine is not None:
|
||
agent.context_compressor = _selected_engine
|
||
# External engines own compaction policy — the host threshold (and its
|
||
# Codex autoraise) never reaches the plugin, so drop the notice.
|
||
agent._compression_threshold_autoraised = None
|
||
# Resolve context_length for plugin engines — mirrors switch_model() path
|
||
from agent.model_metadata import get_model_context_length
|
||
_plugin_ctx_len = get_model_context_length(
|
||
agent.model,
|
||
base_url=agent.base_url,
|
||
api_key=getattr(agent, "api_key", ""),
|
||
config_context_length=_effective_context_length,
|
||
provider=agent.provider,
|
||
custom_providers=_custom_providers,
|
||
)
|
||
# Assign per-model threshold overrides BEFORE the initial update_model()
|
||
# so the first threshold resolution already sees them (assigning after
|
||
# left the initial model on the global threshold until a /model switch).
|
||
if cs.model_thresholds:
|
||
agent.context_compressor.model_thresholds = cs.model_thresholds
|
||
agent.context_compressor.update_model(
|
||
model=agent.model,
|
||
context_length=_plugin_ctx_len,
|
||
base_url=agent.base_url,
|
||
api_key=getattr(agent, "api_key", ""),
|
||
provider=agent.provider,
|
||
api_mode=agent.api_mode,
|
||
)
|
||
if not agent.quiet_mode:
|
||
_ra().logger.info("Using context engine: %s", _selected_engine.name)
|
||
else:
|
||
# Native Gemini output reservation: with model.max_tokens unset the
|
||
# generateContent adapter still sends maxOutputTokens=65,535, and the
|
||
# compressor threshold is pct×(window − max_tokens). Reserving 0 here
|
||
# while the wire reserved 65,535 let the provider 400 before compaction
|
||
# fired, so mirror the adapter's documented default (native Gemini only).
|
||
_compressor_max_tokens = agent.max_tokens
|
||
if _compressor_max_tokens is None:
|
||
try:
|
||
from agent.gemini_native_adapter import (
|
||
GEMINI_DEFAULT_MAX_OUTPUT_TOKENS,
|
||
is_native_gemini_base_url,
|
||
)
|
||
_gemini_provider = str(
|
||
getattr(agent, "provider", "") or ""
|
||
).strip().lower() in {
|
||
"gemini", "google", "google-gemini", "google-ai-studio",
|
||
}
|
||
if _gemini_provider or is_native_gemini_base_url(agent.base_url):
|
||
_compressor_max_tokens = GEMINI_DEFAULT_MAX_OUTPUT_TOKENS
|
||
except Exception:
|
||
pass
|
||
agent.context_compressor = ContextCompressor(
|
||
model=agent.model,
|
||
threshold_percent=cs.threshold,
|
||
protect_first_n=cs.protect_first,
|
||
protect_last_n=cs.protect_last,
|
||
summary_target_ratio=cs.target_ratio,
|
||
summary_model_override=None,
|
||
quiet_mode=agent.quiet_mode,
|
||
base_url=agent.base_url,
|
||
api_key=getattr(agent, "api_key", ""),
|
||
config_context_length=_effective_context_length,
|
||
provider=agent.provider,
|
||
api_mode=agent.api_mode,
|
||
abort_on_summary_failure=cs.abort_on_summary_failure,
|
||
max_tokens=_compressor_max_tokens,
|
||
model_thresholds=cs.model_thresholds,
|
||
threshold_tokens_cap=cs.threshold_tokens,
|
||
proactive_prune_tokens=cs.proactive_prune_tokens,
|
||
proactive_prune_min_result_chars=cs.proactive_prune_min_chars,
|
||
proactive_prune_min_reclaim_tokens=cs.proactive_prune_min_reclaim,
|
||
min_tail_user_messages=cs.min_tail_users,
|
||
tail_mode=cs.tail_mode,
|
||
)
|
||
_bind_session_state = getattr(agent.context_compressor, "bind_session_state", None)
|
||
if callable(_bind_session_state):
|
||
try:
|
||
_bind_session_state(session_db=session_db, session_id=agent.session_id)
|
||
except Exception:
|
||
pass
|
||
agent.compression_enabled = cs.enabled
|
||
agent.compression_in_place = cs.in_place
|
||
# Apply micro-compaction settings to the compressor (feature is opt-in)
|
||
_cc = agent.context_compressor
|
||
# checkpoint_required: micro-compaction is a lossy rewrite with no
|
||
# pre-compress checkpoint hook, so suppress it while the gate is armed
|
||
# (mirrors the native-compaction suppression in native_compaction.py).
|
||
if cs.checkpoint_required and cs.micro_compact:
|
||
logger.warning(
|
||
"compression.checkpoint_required is enabled: post-turn "
|
||
"micro-compaction is disabled for this agent so every lossy "
|
||
"rewrite passes through the checkpoint-gated compressor."
|
||
)
|
||
cs.micro_compact = False
|
||
if hasattr(_cc, "_micro_compact_enabled"):
|
||
_cc._micro_compact_enabled = cs.micro_compact
|
||
if hasattr(_cc, "_micro_compact_every_n_turns"):
|
||
_cc._micro_compact_every_n_turns = cs.micro_compact_every_n_turns
|
||
if hasattr(_cc, "_micro_compact_defrag_threshold_tokens"):
|
||
_cc._micro_compact_defrag_threshold_tokens = cs.micro_compact_defrag_tokens
|
||
agent.compression_checkpoint_required = cs.checkpoint_required
|
||
agent.codex_app_server_auto_compaction = cs.codex_app_server_auto
|
||
agent.codex_responses_native_compaction = cs.codex_responses_native
|
||
agent.codex_responses_compact_threshold = cs.codex_responses_compact_threshold
|
||
from agent.native_compaction import resolve_native_compaction_capabilities
|
||
agent.runtime_capabilities = resolve_native_compaction_capabilities(
|
||
model=agent.model,
|
||
base_url=agent.base_url,
|
||
provider=agent.provider,
|
||
is_codex_backend=(agent.provider or "").strip().lower() == "openai-codex",
|
||
)
|
||
agent.max_compression_attempts = cs.max_attempts
|
||
agent.compression_idle_compact_after_seconds = (
|
||
cs.idle_compact_after_seconds
|
||
)
|
||
|
||
|
||
def _enforce_minimum_context(agent):
|
||
# Reject models whose context window is below the minimum required
|
||
# for reliable tool-calling workflows (64K tokens).
|
||
_ctx = getattr(agent.context_compressor, "context_length", 0)
|
||
_allow_lmstudio_explicit_below_floor = (
|
||
str(getattr(agent, "provider", "") or "").strip().lower() == "lmstudio"
|
||
and isinstance(agent._config_context_length, int)
|
||
and not isinstance(agent._config_context_length, bool)
|
||
and agent._config_context_length > 0
|
||
)
|
||
if _ctx and _ctx < MINIMUM_CONTEXT_LENGTH and not _allow_lmstudio_explicit_below_floor:
|
||
raise ValueError(
|
||
f"Model {agent.model} has a context window of {_ctx:,} tokens, "
|
||
f"which is below the minimum {MINIMUM_CONTEXT_LENGTH:,} required "
|
||
f"by Hermes Agent. Choose a model with at least "
|
||
f"{MINIMUM_CONTEXT_LENGTH // 1000}K context. If your server "
|
||
f"reports a window smaller than the model's true window, set "
|
||
f"model.context_length in config.yaml to the real value "
|
||
f"(this must be at least {MINIMUM_CONTEXT_LENGTH // 1000}K)."
|
||
)
|
||
|
||
|
||
def _warn_nonagentic_hermes_model(agent):
|
||
# Nous Hermes 3/4 are chat models, not tool-call-tuned. cli.py show_banner()
|
||
# already warns on the CLI, so skip platform=="cli" to avoid a double
|
||
# warning; non-quiet non-CLI surfaces still get it.
|
||
if not agent.quiet_mode and (agent.platform or "cli") != "cli":
|
||
try:
|
||
from hermes_cli.model_switch import _check_hermes_model_warning
|
||
|
||
_hermes_warn = _check_hermes_model_warning(agent.model or "")
|
||
if _hermes_warn:
|
||
_user_msg = (
|
||
"⚠ Nous Research Hermes 3 & 4 models are NOT agentic — they "
|
||
"lack reliable tool-calling for agent workflows (delegation, "
|
||
"cron, proactive tools). Consider an agentic model instead "
|
||
"(Claude, GPT, Gemini, Qwen-Coder, etc.)."
|
||
)
|
||
if hasattr(agent, "_emit_warning"):
|
||
agent._emit_warning(_user_msg)
|
||
else:
|
||
print(f"\n{_user_msg}\n", file=sys.stderr)
|
||
_ra().logger.warning(_hermes_warn)
|
||
except Exception:
|
||
pass
|
||
|
||
|
||
def _inject_context_engine_tools(agent):
|
||
# Inject context engine tool schemas (e.g. lcm_grep, lcm_describe, lcm_expand).
|
||
# Dedup against existing names: plugin paths may register the same schemas
|
||
# via ctx.register_tool(), and a duplicate trips provider-side 'duplicate
|
||
# tool name' errors. Gate on enabled_toolsets like memory-provider tools so
|
||
# `platform_toolsets: telegram: []` cannot leak lcm_* tools.
|
||
agent._context_engine_tool_names: set = set()
|
||
if (
|
||
agent.context_compressor
|
||
and agent.tools is not None
|
||
and (
|
||
agent.enabled_toolsets is None
|
||
or "context_engine" in agent.enabled_toolsets
|
||
)
|
||
):
|
||
_existing_tool_names = {
|
||
t.get("function", {}).get("name")
|
||
for t in agent.tools
|
||
if isinstance(t, dict)
|
||
}
|
||
from agent.memory_manager import normalize_tool_schema as _normalize_tool_schema
|
||
for _raw_schema in agent.context_compressor.get_tool_schemas():
|
||
_schema = _normalize_tool_schema(_raw_schema)
|
||
if _schema is None:
|
||
# A schema with no resolvable name (e.g. an already-wrapped
|
||
# entry) would append a nameless tool that strict providers
|
||
# 400 on, disabling the whole toolset (#47707). Skip it.
|
||
_ra().logger.warning(
|
||
"Context engine returned a tool schema with no resolvable "
|
||
"name; skipping to avoid poisoning the request (%r)",
|
||
_raw_schema,
|
||
)
|
||
continue
|
||
_tname = _schema["name"]
|
||
if _tname in _existing_tool_names:
|
||
continue # already registered via plugin/cache path
|
||
_wrapped = {"type": "function", "function": _schema}
|
||
agent.tools.append(_wrapped)
|
||
agent.valid_tool_names.add(_tname)
|
||
agent._context_engine_tool_names.add(_tname)
|
||
_existing_tool_names.add(_tname)
|
||
|
||
# Notify context engine of session start
|
||
if agent.context_compressor:
|
||
try:
|
||
agent.context_compressor.on_session_start(
|
||
agent.session_id,
|
||
hermes_home=str(get_hermes_home()),
|
||
platform=agent.platform or "cli",
|
||
model=agent.model,
|
||
context_length=getattr(agent.context_compressor, "context_length", 0),
|
||
conversation_id=getattr(agent, "_gateway_session_key", None),
|
||
)
|
||
except Exception as _ce_err:
|
||
_ra().logger.debug("Context engine on_session_start: %s", _ce_err)
|
||
|
||
|
||
def _configure_ollama_num_ctx(agent, _model_cfg, _config_context_length):
|
||
# Ollama num_ctx: Ollama defaults to 2048 regardless of the model, so detect
|
||
# the max window and pass num_ctx on every request. model.ollama_num_ctx
|
||
# overrides; model.context_length caps the detected value (VRAM budget).
|
||
agent._ollama_num_ctx: int | None = None
|
||
_ollama_num_ctx_override = None
|
||
if isinstance(_model_cfg, dict):
|
||
_ollama_num_ctx_override = _model_cfg.get("ollama_num_ctx")
|
||
if _ollama_num_ctx_override is not None:
|
||
try:
|
||
agent._ollama_num_ctx = int(_ollama_num_ctx_override)
|
||
except (TypeError, ValueError):
|
||
_ra().logger.debug("Invalid ollama_num_ctx config value: %r", _ollama_num_ctx_override)
|
||
if agent._ollama_num_ctx is None and agent.base_url and is_local_endpoint(agent.base_url):
|
||
try:
|
||
# ``agent.api_key`` may be a callable (Entra token provider).
|
||
# Ollama detection makes a manual HTTP request and expects a
|
||
# string — Azure Foundry isn't a local endpoint so this branch
|
||
# never fires for Entra, but guard defensively.
|
||
_key_for_ollama = agent.api_key if isinstance(agent.api_key, str) else ""
|
||
_detected = query_ollama_num_ctx(agent.model, agent.base_url, api_key=_key_for_ollama or "")
|
||
if _detected and _detected > 0:
|
||
agent._ollama_num_ctx = _detected
|
||
except Exception as exc:
|
||
_ra().logger.debug("Ollama num_ctx detection failed: %s", exc)
|
||
# Cap auto-detected num_ctx to the explicit context_length: GGUF metadata
|
||
# can advertise 256K+ and Ollama would allocate that much VRAM.
|
||
if (
|
||
agent._ollama_num_ctx
|
||
and _config_context_length
|
||
and _ollama_num_ctx_override is None # don't override explicit ollama_num_ctx
|
||
and agent._ollama_num_ctx > _config_context_length
|
||
):
|
||
_ra().logger.info(
|
||
"Ollama num_ctx capped: %d -> %d (model.context_length override)",
|
||
agent._ollama_num_ctx, _config_context_length,
|
||
)
|
||
agent._ollama_num_ctx = _config_context_length
|
||
if agent._ollama_num_ctx and not agent.quiet_mode:
|
||
_ra().logger.info(
|
||
"Ollama num_ctx: will request %d tokens (model max from /api/show)",
|
||
agent._ollama_num_ctx,
|
||
)
|
||
# Recalibrate the compressor to the served window: it was built from the
|
||
# probed model window (GGUF may advertise 256K+) but every request runs at
|
||
# num_ctx, so with only model.ollama_num_ctx set the compaction trigger
|
||
# could sit above the real window and never fire. Clamp to num_ctx.
|
||
_cc_window = getattr(agent.context_compressor, "context_length", 0) or 0
|
||
if (
|
||
agent._ollama_num_ctx
|
||
and agent._ollama_num_ctx > 0
|
||
and _cc_window
|
||
and agent._ollama_num_ctx < _cc_window
|
||
):
|
||
_ra().logger.info(
|
||
"Compressor window clamped to Ollama num_ctx: %d -> %d",
|
||
_cc_window, agent._ollama_num_ctx,
|
||
)
|
||
agent.context_compressor.update_model(
|
||
model=agent.model,
|
||
context_length=agent._ollama_num_ctx,
|
||
base_url=agent.base_url,
|
||
api_key=getattr(agent, "api_key", ""),
|
||
provider=agent.provider,
|
||
api_mode=agent.api_mode,
|
||
)
|
||
|
||
|
||
def _emit_compression_summary(agent, cs):
|
||
# Codex gpt-5.x autoraise notice: at most once per profile/config state
|
||
# (persisted marker — the gateway rebuilds the agent per message). A
|
||
# changed threshold/model re-notifies once; the display gate suppresses
|
||
# the banner without disabling the autoraise.
|
||
_autoraise = agent._compression_threshold_autoraised or {}
|
||
_autoraise_notice = None
|
||
if (
|
||
bool(_autoraise)
|
||
and cs.enabled
|
||
and cs.autoraise_notice_enabled
|
||
and not _codex_gpt55_autoraise_notice_seen(_autoraise)
|
||
):
|
||
_autoraise_notice = _build_codex_gpt5_autoraise_notice(
|
||
_autoraise,
|
||
context_length=getattr(agent.context_compressor, "context_length", None),
|
||
)
|
||
|
||
if not agent.quiet_mode:
|
||
if cs.enabled:
|
||
# Report the active engine's own threshold — for a plugin engine
|
||
# the host cs.threshold is not in effect, and mixing the
|
||
# two printed a percent that contradicted the token count. (#44439)
|
||
_active_threshold_pct = getattr(
|
||
agent.context_compressor, "threshold_percent", cs.threshold
|
||
)
|
||
_cap_note = ""
|
||
_cap = getattr(agent.context_compressor, "threshold_tokens_cap", None)
|
||
if _cap and _cap > 0:
|
||
_cap_note = f" (capped at {_cap:,} tokens)"
|
||
print(f"📊 Context limit: {agent.context_compressor.context_length:,} tokens (compress at {int(_active_threshold_pct*100)}% = {agent.context_compressor.threshold_tokens:,}{_cap_note})")
|
||
else:
|
||
print(f"📊 Context limit: {agent.context_compressor.context_length:,} tokens (auto-compression disabled)")
|
||
# Notice with the exact opt-back-out command. Printed inline at startup
|
||
# for CLI users; gateway users get the same text replayed via
|
||
# _compression_warning on turn 1 (set below).
|
||
if _autoraise_notice:
|
||
print(_autoraise_notice)
|
||
|
||
# Gateway parity: the startup print only reaches the CLI; status_callback
|
||
# isn't wired yet, so stash the text to replay on the first run_conversation().
|
||
agent._compression_warning = _autoraise_notice
|
||
# Mark shown so repeated inits in this profile (every gateway message) stay
|
||
# silent, whether the notice went to the CLI print or the replay slot.
|
||
if _autoraise_notice:
|
||
_record_codex_gpt55_autoraise_notice(_autoraise)
|
||
# Feasibility check is deferred to the first turn near the threshold
|
||
# (eager costs ~400ms cold per init); run_conversation's preflight runs
|
||
# ensure_compression_feasibility_checked at most once per agent.
|
||
agent._compression_feasibility_checked = False
|
||
|
||
|
||
def _snapshot_primary_runtime(agent):
|
||
# Snapshot primary runtime for per-turn restoration. When fallback
|
||
# activates during a turn, the next turn restores these values so the
|
||
# preferred model gets a fresh attempt each time. Uses a single dict
|
||
# so new state fields are easy to add without N individual attributes.
|
||
_cc = agent.context_compressor
|
||
agent._primary_runtime = {
|
||
"model": agent.model,
|
||
"provider": agent.provider,
|
||
"requested_provider": agent.requested_provider,
|
||
"base_url": agent.base_url,
|
||
"api_mode": agent.api_mode,
|
||
"api_key": getattr(agent, "api_key", ""),
|
||
"request_overrides": dict(getattr(agent, "request_overrides", {}) or {}),
|
||
"client_kwargs": dict(agent._client_kwargs),
|
||
"use_prompt_caching": agent._use_prompt_caching,
|
||
"use_native_cache_layout": agent._use_native_cache_layout,
|
||
"reasoning_echo_flag": getattr(agent, "_reasoning_echo_flag", False),
|
||
# Context engine state that _try_activate_fallback() overwrites.
|
||
# Use getattr for model/base_url/api_key/provider since plugin
|
||
# engines may not have these (they're ContextCompressor-specific).
|
||
"compressor_model": getattr(_cc, "model", agent.model),
|
||
"compressor_base_url": getattr(_cc, "base_url", agent.base_url),
|
||
"compressor_api_key": getattr(_cc, "api_key", ""),
|
||
"compressor_provider": getattr(_cc, "provider", agent.provider),
|
||
"compressor_context_length": _cc.context_length,
|
||
"compressor_threshold_tokens": _cc.threshold_tokens,
|
||
}
|
||
if agent.api_mode == "anthropic_messages":
|
||
agent._primary_runtime.update({
|
||
"anthropic_api_key": agent._anthropic_api_key,
|
||
"anthropic_base_url": agent._anthropic_base_url,
|
||
"is_anthropic_oauth": agent._is_anthropic_oauth,
|
||
})
|
||
|
||
|
||
_CALLBACK_PARAMS = (
|
||
"tool_progress_callback", "tool_start_callback", "tool_complete_callback",
|
||
"thinking_callback", "reasoning_callback", "clarify_callback",
|
||
"read_terminal_callback", "read_preview_callback", "drive_preview_callback",
|
||
"read_window_below_callback", "setup_mcp_callback", "tour_callback",
|
||
"step_callback", "stream_delta_callback", "interim_assistant_callback",
|
||
"status_callback", "notice_callback", "notice_clear_callback",
|
||
"event_callback", "reaction_callback", "tool_gen_callback",
|
||
)
|
||
|
||
|
||
def init_agent(
|
||
agent,
|
||
base_url: str = None,
|
||
api_key: str = None,
|
||
provider: str = None,
|
||
api_mode: str = None,
|
||
acp_command: str = None,
|
||
acp_args: list[str] | None = None,
|
||
command: str = None,
|
||
args: list[str] | None = None,
|
||
model: str = "",
|
||
max_iterations: int = sys.maxsize, # Default: unlimited tool-calling iterations (shared with subagents)
|
||
enabled_toolsets: List[str] = None,
|
||
disabled_toolsets: List[str] = None,
|
||
save_trajectories: bool = False,
|
||
verbose_logging: bool = False,
|
||
quiet_mode: bool = False,
|
||
tool_progress_mode: str = "all",
|
||
ephemeral_system_prompt: str = None,
|
||
log_prefix_chars: int = 100,
|
||
log_prefix: str = "",
|
||
providers_allowed: List[str] = None,
|
||
providers_ignored: List[str] = None,
|
||
providers_order: List[str] = None,
|
||
provider_sort: str = None,
|
||
provider_require_parameters: bool = False,
|
||
provider_data_collection: str = None,
|
||
openrouter_min_coding_score: Optional[float] = None,
|
||
session_id: str = None,
|
||
tool_progress_callback: callable = None,
|
||
tool_start_callback: callable = None,
|
||
tool_complete_callback: callable = None,
|
||
thinking_callback: callable = None,
|
||
reasoning_callback: callable = None,
|
||
clarify_callback: callable = None,
|
||
read_terminal_callback: callable = None,
|
||
read_preview_callback: callable = None,
|
||
drive_preview_callback: callable = None,
|
||
read_window_below_callback: callable = None,
|
||
setup_mcp_callback: callable = None,
|
||
tour_callback: callable = None,
|
||
step_callback: callable = None,
|
||
stream_delta_callback: callable = None,
|
||
interim_assistant_callback: callable = None,
|
||
tool_gen_callback: callable = None,
|
||
status_callback: callable = None,
|
||
notice_callback: callable = None,
|
||
notice_clear_callback: callable = None,
|
||
event_callback: Optional[Callable[[str, dict], None]] = None,
|
||
reaction_callback: Optional[Callable[[str], None]] = None,
|
||
max_tokens: int = None,
|
||
reasoning_config: Dict[str, Any] = None,
|
||
service_tier: str = None,
|
||
request_overrides: Dict[str, Any] = None,
|
||
prefill_messages: List[Dict[str, Any]] = None,
|
||
platform: str = None,
|
||
user_id: str = None,
|
||
user_id_alt: str = None,
|
||
user_name: str = None,
|
||
chat_id: str = None,
|
||
chat_name: str = None,
|
||
chat_type: str = None,
|
||
thread_id: str = None,
|
||
gateway_session_key: str = None,
|
||
skip_context_files: bool = False,
|
||
load_soul_identity: bool = False,
|
||
skip_memory: bool = False,
|
||
skip_background_review: bool = False,
|
||
session_db=None,
|
||
parent_session_id: str = None,
|
||
iteration_budget: "IterationBudget" = None,
|
||
run_budget_seconds: Optional[float] = None,
|
||
fallback_model: Dict[str, Any] = None,
|
||
credential_pool=None,
|
||
checkpoints_enabled: bool = False,
|
||
checkpoint_max_snapshots: int = 20,
|
||
checkpoint_max_total_size_mb: int = 500,
|
||
checkpoint_max_file_size_mb: int = 10,
|
||
pass_session_id: bool = False,
|
||
requested_provider: str = None,
|
||
capabilities: Optional[Dict[str, bool]] = None,
|
||
):
|
||
"""Initialize the AI Agent (body of :meth:`AIAgent.__init__`).
|
||
|
||
Non-obvious parameters:
|
||
requested_provider: original provider identity before runtime canonicalization.
|
||
openrouter_min_coding_score: coding-score floor for ``openrouter/pareto-code``
|
||
only; None/empty lets OpenRouter pick.
|
||
clarify_callback: ``(question, choices) -> str`` from the platform layer; if
|
||
None the clarify tool returns an error.
|
||
reasoning_config: None defaults to ``{"enabled": True, "effort": "medium"}``
|
||
on OpenRouter.
|
||
prefill_messages: prepended to history as priming context. Anthropic Sonnet/
|
||
Opus 4.6+ reject a conversation ending on an assistant message (400) — use
|
||
structured outputs instead of a trailing-assistant prefill there.
|
||
skip_context_files: skip SOUL.md/.hermes.md/AGENTS.md/CLAUDE.md/.cursorrules
|
||
injection (batch/data-generation runs). load_soul_identity keeps
|
||
~/.hermes/SOUL.md as identity even when skip_context_files=True.
|
||
"""
|
||
_install_safe_stdio()
|
||
|
||
agent.model = model
|
||
agent.max_iterations = max_iterations
|
||
# Shared iteration budget — parent creates, children inherit.
|
||
# Consumed by every LLM turn across parent + all subagents.
|
||
agent.iteration_budget = iteration_budget or IterationBudget(max_iterations)
|
||
agent.save_trajectories = save_trajectories
|
||
agent.verbose_logging = verbose_logging
|
||
agent.quiet_mode = quiet_mode
|
||
agent.tool_progress_mode = tool_progress_mode
|
||
agent.ephemeral_system_prompt = ephemeral_system_prompt
|
||
agent.platform = platform # "cli", "telegram", "discord", "whatsapp", etc.
|
||
agent._user_id = user_id # Platform user identifier (gateway sessions)
|
||
agent._user_id_alt = user_id_alt # Optional stable alternate platform identifier
|
||
agent._user_name = user_name
|
||
agent._chat_id = chat_id
|
||
agent._chat_name = chat_name
|
||
agent._chat_type = chat_type
|
||
agent._thread_id = thread_id
|
||
agent._gateway_session_key = gateway_session_key # Stable per-chat key (e.g. agent:main:telegram:dm:123)
|
||
# Pluggable print function — CLI replaces this with _cprint so that
|
||
# raw ANSI status lines are routed through prompt_toolkit's renderer
|
||
# instead of going directly to stdout where patch_stdout's StdoutProxy
|
||
# would mangle the escape sequences. None = use builtins.print.
|
||
agent._print_fn = None
|
||
agent.background_review_callback = None # Optional sync callback for gateway delivery
|
||
agent.memory_notifications = "on" # Memory update notifications: "off", "on", "verbose"
|
||
agent.skip_context_files = skip_context_files
|
||
agent.load_soul_identity = load_soul_identity
|
||
# Background review (memory/skill) opt-out: skips the end-of-turn fork
|
||
# (~30K tokens/event) on cron-style sessions. Single switch for both
|
||
# review paths (skip_memory alone only disables the memory trigger).
|
||
agent.skip_background_review = bool(skip_background_review)
|
||
agent.pass_session_id = pass_session_id
|
||
agent.log_prefix_chars = log_prefix_chars
|
||
agent.log_prefix = f"{log_prefix} " if log_prefix else ""
|
||
# Store effective base URL for feature detection (prompt caching, reasoning, etc.)
|
||
agent.base_url = base_url or ""
|
||
provider_name = provider.strip().lower() if isinstance(provider, str) and provider.strip() else None
|
||
agent.provider = provider_name or ""
|
||
agent.requested_provider = (
|
||
requested_provider.strip().lower()
|
||
if isinstance(requested_provider, str) and requested_provider.strip()
|
||
else agent.provider
|
||
)
|
||
agent.capabilities = {
|
||
key: value for key, value in (capabilities or {}).items()
|
||
if isinstance(key, str) and isinstance(value, bool)
|
||
}
|
||
agent._credential_pool = credential_pool
|
||
agent.acp_command = acp_command or command
|
||
agent.acp_args = list(acp_args or args or [])
|
||
_resolve_api_mode(agent, api_mode, provider_name, base_url)
|
||
|
||
_finalize_routing(agent, api_mode, credential_pool)
|
||
|
||
# Platform callbacks are stored under their parameter names verbatim.
|
||
_params = locals()
|
||
for _cb in _CALLBACK_PARAMS:
|
||
setattr(agent, _cb, _params[_cb])
|
||
agent.suppress_status_output = False
|
||
|
||
_init_control_state(agent)
|
||
|
||
# Store OpenRouter provider preferences
|
||
agent.providers_allowed = providers_allowed
|
||
agent.providers_ignored = providers_ignored
|
||
agent.providers_order = providers_order
|
||
agent.provider_sort = provider_sort
|
||
agent.provider_require_parameters = provider_require_parameters
|
||
agent.provider_data_collection = provider_data_collection
|
||
agent.openrouter_min_coding_score = openrouter_min_coding_score
|
||
|
||
# Store toolset filtering options
|
||
agent.enabled_toolsets = enabled_toolsets
|
||
agent.disabled_toolsets = disabled_toolsets
|
||
|
||
# Model response configuration
|
||
agent.max_tokens = max_tokens # None = use model default
|
||
agent.reasoning_config = reasoning_config # None = use default (medium for OpenRouter)
|
||
# Per-provider reasoning_content echo opt-in (see _reasoning_echo_opt_in).
|
||
# Read once at init; switch_model / try_activate_fallback / restore
|
||
# keep it in sync with the active provider.
|
||
agent._reasoning_echo_flag = agent._read_reasoning_echo_from_config()
|
||
agent.service_tier = service_tier
|
||
agent.request_overrides = dict(request_overrides or {})
|
||
agent.prefill_messages = prefill_messages or [] # Prefilled conversation turns
|
||
agent._force_ascii_payload = False
|
||
|
||
_init_prompt_cache_config(agent)
|
||
|
||
_init_turn_state(agent, run_budget_seconds)
|
||
|
||
_setup_logging(agent)
|
||
|
||
_init_stream_state(agent)
|
||
|
||
_build_client(agent, api_key, base_url, fallback_model)
|
||
|
||
_init_fallback_chain(agent, fallback_model)
|
||
|
||
_load_tools(agent, enabled_toolsets, disabled_toolsets)
|
||
|
||
_init_session_state(
|
||
agent, session_id, session_db, parent_session_id, reasoning_config, max_tokens,
|
||
checkpoints_enabled, checkpoint_max_snapshots, checkpoint_max_total_size_mb, checkpoint_max_file_size_mb,
|
||
)
|
||
|
||
# Load config once for memory, skills, and compression sections
|
||
try:
|
||
from hermes_cli.config import load_config_readonly as _load_agent_config
|
||
_agent_cfg = _load_agent_config()
|
||
except Exception:
|
||
_agent_cfg = {}
|
||
|
||
_apply_display_config(agent, _agent_cfg, platform)
|
||
|
||
_init_memory(agent, _agent_cfg, skip_memory, platform)
|
||
|
||
_apply_agent_section(agent, _agent_cfg)
|
||
|
||
cs = _parse_compression_config(agent, _agent_cfg)
|
||
|
||
_config_context_length, _custom_providers, _effective_context_length, _model_cfg = _resolve_context_length(
|
||
agent, _agent_cfg, base_url
|
||
)
|
||
|
||
|
||
|
||
_build_context_engine(agent, _agent_cfg, cs, _custom_providers, _effective_context_length, session_db)
|
||
|
||
_enforce_minimum_context(agent)
|
||
|
||
_warn_nonagentic_hermes_model(agent)
|
||
|
||
_inject_context_engine_tools(agent)
|
||
|
||
from agent.runtime_cwd import scope_terminal_cwd as _scope_terminal_cwd
|
||
|
||
agent._subdirectory_hints = SubdirectoryHintTracker(
|
||
working_dir=_scope_terminal_cwd() or None,
|
||
)
|
||
agent._user_turn_count = 0
|
||
# Copilot x-initiator flag: first API call of a user turn sends "user" (#3040).
|
||
agent._is_user_initiated_turn = False
|
||
|
||
# Usage-anchored context accounting (agent/model_metadata.py): the last
|
||
# main-loop provider response's exact usage + transcript snapshot. None
|
||
# until the first response with usage; invalidated on compaction and
|
||
# session switches so stale anchors can never suppress compression.
|
||
agent._usage_anchor = None
|
||
agent._turn_base_usage_anchor = None
|
||
|
||
# Cumulative token usage for the session
|
||
agent.session_prompt_tokens = 0
|
||
agent.session_completion_tokens = 0
|
||
agent.session_total_tokens = 0
|
||
agent.session_api_calls = 0
|
||
agent.session_input_tokens = 0
|
||
agent.session_output_tokens = 0
|
||
agent.session_cache_read_tokens = 0
|
||
agent.session_cache_write_tokens = 0
|
||
agent.session_reasoning_tokens = 0
|
||
agent.session_estimated_cost_usd = 0.0
|
||
agent.session_cost_status = "unknown"
|
||
agent.session_cost_source = "none"
|
||
# Rolling history for status-bar avg latency / velocity (last 10 calls).
|
||
# Stored on the agent so both conversation_loop and codex_runtime share it
|
||
# and the CLI snapshot can read it without extra IPC.
|
||
from collections import deque as _deque
|
||
agent._api_latency_history = _deque(maxlen=10)
|
||
agent._api_output_history = _deque(maxlen=10)
|
||
|
||
_configure_ollama_num_ctx(agent, _model_cfg, _config_context_length)
|
||
|
||
_emit_compression_summary(agent, cs)
|
||
|
||
_snapshot_primary_runtime(agent)
|
||
|
||
|
||
|
||
__all__ = ["init_agent"]
|