diff --git a/agent/agent_init.py b/agent/agent_init.py index c94e25911e..956d0e7315 100644 --- a/agent/agent_init.py +++ b/agent/agent_init.py @@ -44,8 +44,7 @@ from utils import base_url_host_matches, is_truthy_value logger = logging.getLogger("run_agent") -# Memory providers already warned unavailable — the gateway builds a fresh AIAgent per -# message, so an un-deduped warning would fire every turn. +# Deduped: the gateway builds a fresh AIAgent per message, so it would warn every turn. _warned_unavailable_providers: set[str] = set() @@ -198,8 +197,7 @@ def _build_codex_gpt5_autoraise_notice( if isinstance(context_length, int) and context_length > 0: cap = f"{round(context_length / 1000)}K" else: - # Static fallback: 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. + # Static fallback: codex-spark is natively 128K; gpt-5.4/5.5/5.6 are capped at 272K. 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)) @@ -470,8 +468,8 @@ def _resolve_api_mode(agent, api_mode, provider_name, base_url): ): agent.api_mode = "bedrock_converse" elif agent.provider in {"nous", "nous-portal", "nousresearch"}: - # Portal is dual-wire: anthropic/* → Messages, everything else → chat_completions. - # Covers direct AIAgent construction without a resolved runtime. + # Portal is dual-wire (anthropic/* → Messages, else chat_completions); covers direct + # AIAgent construction without a resolved runtime. from hermes_cli.providers import nous_api_mode agent.api_mode = nous_api_mode(agent.model) @@ -502,8 +500,7 @@ def _finalize_routing(agent, api_mode, credential_pool): except Exception: agent._credential_pool = None - # Eagerly warm the transport cache so import errors surface at init, not - # mid-conversation. Non-fatal — transport may not exist for all modes yet. + # Warm the transport cache so import errors surface at init (non-fatal: some modes lack one). try: agent._get_transport() except Exception: @@ -519,11 +516,10 @@ def _finalize_routing(agent, api_mode, credential_pool): except Exception: pass - # Auto-upgrade to Responses for GPT-5.x-style models and direct OpenAI URLs, unless: - # api_mode was explicit, the runtime is ACP (`acp://` — ACP clients route themselves - # and lack the Responses surface), or the URL is Azure OpenAI (gpt-5.x on - # /chat/completions only). Provider exceptions live in - # _provider_model_requires_responses_api. + # Auto-upgrade to Responses for GPT-5.x-style models and direct OpenAI URLs, unless + # api_mode was explicit, the runtime is ACP (`acp://` clients route themselves, no + # Responses surface) or Azure OpenAI (gpt-5.x on /chat/completions only). Provider + # exceptions live in _provider_model_requires_responses_api. _base_lower = str(agent.base_url or "").lower() if ( api_mode is None @@ -541,9 +537,8 @@ def _finalize_routing(agent, api_mode, credential_pool): 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. + # Pre-warm the OpenRouter metadata cache (1h TTL) off-thread so the first pricing estimate + # doesn't block. Process-level Event guard: an unguarded spawn leaks a thread per message. if (agent.provider == "openrouter" or agent._is_openrouter_url()) and \ not _ra()._openrouter_prewarm_done.is_set(): _ra()._openrouter_prewarm_done.set() @@ -553,14 +548,13 @@ def _finalize_routing(agent, api_mode, credential_pool): def _init_control_state(agent): - # Tool execution state — allows _vprint during tool execution even when stream - # consumers are registered (no tokens streaming then). + # Lets _vprint print during tool execution even with stream consumers registered. agent._executing_tools = False agent._tool_guardrails = ToolCallGuardrailController() agent._tool_guardrail_halt_decision: ToolGuardrailDecision | None = None - # Interrupt mechanism for breaking out of tool loops. Hard cancellation is separate - # from redirect/message state; a thread-safe Event makes the cause atomic for pollers. + # Interrupts. Hard cancellation is separate from redirect/message state; the Event makes + # the cause atomic for auxiliary stream pollers. agent._interrupt_requested = False agent._interrupt_message = None # Optional message that triggered interrupt agent._hard_interrupt_requested = threading.Event() @@ -570,15 +564,13 @@ def _init_control_state(agent): 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, preserving - # role alternation (no new user turn). + # /steer: the drain hook appends the note to the last tool result after the current + # batch — no interrupt, no new user turn (role alternation preserved). agent._pending_steer: Optional[str] = None agent._pending_steer_lock = threading.Lock() - # Active-turn redirect: unlike a hard /stop, preserve the valid turn prefix, cancel - # only the in-flight request and rebuild its tail with the correction. Drained at a - # role-safe boundary. + # Active-turn redirect: keep the valid turn prefix, cancel only the in-flight request, + # rebuild the tail with the correction. Drained at a role-safe boundary. agent._pending_redirect: Optional[str] = None agent._pending_redirect_lock = threading.Lock() @@ -592,27 +584,24 @@ def _init_control_state(agent): agent._active_children = [] # Running child AIAgents (for interrupt propagation) agent._active_children_lock = threading.Lock() - # Background memory/skill review state (agent/background_review.py). The run is - # installed before the worker starts and fences its first provider-capable phase; the - # direct agent pointer keeps interrupt propagation available once the fork exists. + # Background review (agent/background_review.py): the run is installed before the worker + # starts and fences its first provider phase; the agent pointer enables interrupt fan-out. 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``. + # Anthropic prompt caching (~75% input savings): auto-enabled for Claude on native + # Anthropic, OpenRouter and anthropic_messages gateways. See _anthropic_prompt_cache_policy. 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 — OAuth plans billing cache - # writes, or proxies adding their own cache_control. The disable survives /model - # switches and fallback re-derivation via anthropic_prompt_cache_policy(). + # cache_ttl: "5m" (default) or "1h" (2x write cost; pays off with >5-minute pauses); + # unknown values keep "5m". A falsy/off value disables caching entirely (OAuth plans + # billing cache writes, proxies adding their own cache_control); the disable survives + # /model switches and fallback re-derivation. agent._cache_ttl = "5m" try: from hermes_cli.config import load_config_readonly as _load_pc_cfg @@ -633,44 +622,39 @@ def _init_prompt_cache_config(agent): 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). Intermediate pressure warnings made models - # give up early on complex tasks. + # Iteration budget: notify the LLM only on exhaustion (one message, one grace call, then + # a forced summary) — intermediate pressure warnings made models give up early. agent._budget_exhausted_injected = False agent._budget_grace_call = False - # Wall-clock run budget (seconds per run_conversation turn). Explicit constructor arg - # wins; else resolved from config.yaml (agent.run_budget_seconds) in - # _apply_agent_section. None = fully off: no clock reads, no injection, no capping. + # Wall-clock run budget per turn: constructor arg wins, else agent.run_budget_seconds + # (in _apply_agent_section). None = fully off (no clock reads, injection, or capping). agent.run_budget_seconds = _normalize_run_budget_seconds(run_budget_seconds) # 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. Read - # by the gateway timeout handler and the "still working" notifications. + # Activity tracking (API call / tool / stream chunk) for the gateway timeout handler and + # "still working" notifications. agent._last_activity_ts: float = time.time() agent._last_activity_desc: str = "initializing" - # Default paths and _touch_activity stamp unknown; named provenances are stamped by - # compression writers (heartbeat / timeout / cooldown). + # Named provenances are stamped only by compression writers (heartbeat/timeout/cooldown). agent._last_activity_provenance = ActivityProvenance.UNKNOWN # Rate-limit durable SessionDB activity stamps from _touch_activity. agent._session_activity_last_persist_mono: float = 0.0 agent._current_tool: str | None = None agent._api_call_count: int = 0 - # Opt-out for the between-turns MCP tool refresh (build_turn_context). Set on internal - # forks (background_review) that must keep ``tools[]`` byte-identical for cache parity. + # Opt-out for the between-turns MCP refresh; set on forks that need byte-identical tools[]. agent._skip_mcp_refresh = False - # Registry generation the tool snapshot was derived from: lets a late/concurrent - # refresh reject a stale rebuild instead of clobbering a newer one (set in _load_tools). + # Registry generation of the tool snapshot (set in _load_tools): a late refresh rejects + # a stale rebuild instead of clobbering a newer one. agent._tool_snapshot_generation = 0 # Rate limit tracking from x-ratelimit-* response headers; read by /usage. agent._rate_limit_state = None - # Credits tracking (dev-only, behind HERMES_DEV_CREDITS) from x-nous-credits-* headers. - # Session-start remaining is latched the first time a header is seen so cumulative - # micros spent can be reported. Threshold-notice latch: sticky-notice keys + gates. + # Credits tracking (dev-only, HERMES_DEV_CREDITS) from x-nous-credits-* headers; session + # start is latched on the first header so cumulative spend can be reported. agent._credits_state = None agent._credits_session_start_micros = None from agent.credits_tracker import new_credits_latch @@ -682,53 +666,46 @@ def _init_turn_state(agent, run_budget_seconds): def _setup_logging(agent): - # agent.log (INFO+) and errors.log (WARNING+) under ~/.hermes/logs/. Idempotent, so - # gateway mode (new AIAgent per message) won't duplicate handlers. + # agent.log (INFO+) + errors.log (WARNING+); idempotent so per-message gateway agents + # don'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 deliberately does NOT raise per-logger levels: that would starve the root - # file handlers (isEnabledFor() is checked before propagation). setup_logging() - # installs no console handler in quiet mode; noise reduction belongs in hermes_logging. + # Quiet mode must NOT raise per-logger levels: isEnabledFor() runs before propagation and + # would starve the root file handlers. Noise reduction belongs in hermes_logging. def _init_stream_state(agent): # Internal stream callback (streaming TTS); set here so _vprint can reference it early. agent._stream_callback = None - # Deferred paragraph break — set after tool iterations so one "\n\n" precedes the next - # real text delta. + # Set after tool iterations so one "\n\n" precedes the next real text delta. agent._stream_needs_break = False - # Stateful scrubbers for / thinking spans split across stream deltas: - # per-delta regexes can't survive chunk boundaries (both tags needed in one string). + # Stateful scrubbers: / thinking spans split across deltas defeat + # per-delta regexes (both tags must be in one string). agent._stream_context_scrubber = StreamingContextScrubber() agent._stream_think_scrubber = StreamingThinkScrubber() - # Visible assistant text already delivered via live token callbacks this response — - # avoids re-sending commentary the provider later returns as a completed interim. + # Text already streamed this response, so a later completed interim isn't re-sent. agent._current_streamed_assistant_text = "" - # Completed interim messages delivered this user turn; spans Codex continuation/tool - # calls so repeated commentary is not re-sent before normalization dedups it. + # Interims delivered this user turn (spans Codex continuations) so repeats aren't re-sent. 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. + # Single-writer guard for the delta sink: each attempt claims a monotonic writer token and + # the sink drops chunks from threads holding a stale one, so a superseded stream can't + # interleave with the retry's. Threads that never claimed are never fenced. agent._stream_writer_lock = threading.Lock() agent._stream_writer_token = 0 agent._stream_writer_tls = threading.local() agent._stream_writer_dropped = 0 - # Current-turn user-message override when the API-facing message intentionally differs - # from the persisted transcript (e.g. CLI voice mode's temporary prefix). + # API-facing user message override when it differs from the persisted transcript (voice). agent._persist_user_message_idx = None agent._persist_user_message_override = None agent._persist_user_message_timestamp = None - # Anthropic image-to-text fallbacks cached per image payload/URL so one tool loop - # doesn't repeatedly re-run auxiliary vision on the same image history. + # Image-to-text fallbacks cached per payload/URL so one tool loop doesn't re-run vision. agent._anthropic_image_fallback_cache: Dict[str, str] = {} @@ -768,16 +745,14 @@ def _init_anthropic_client(agent, api_key, base_url, _provider_timeout): if not agent.quiet_mode: print(f"🤖 AI Agent initialized with model: {agent.model} (AWS Bedrock + AnthropicBedrock SDK, {_br_region})") return - # Only fall back to ANTHROPIC_TOKEN when the provider is actually Anthropic. Other - # anthropic_messages providers (MiniMax, Alibaba, …) must use their own key — falling - # back would send Anthropic credentials to third-party endpoints. + # ANTHROPIC_TOKEN fallback only for native Anthropic — other anthropic_messages providers + # must use their own key or Anthropic credentials leak to third-party endpoints. _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). Cost: one file read per request. + # MiniMax OAuth tokens live ~15 min and the SDK freezes api_key at construction, so use a + # callable provider: build_anthropic_client mints a fresh bearer per request (re-reading + # auth.json, so other processes' refreshes are seen). if agent.provider == "minimax-oauth" and isinstance(effective_key, str) and effective_key: try: from hermes_cli.auth import build_minimax_oauth_token_provider @@ -806,10 +781,9 @@ def _init_moa_client(agent, api_key): 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; - # shared with fallback-restore so a restored facade keeps emitting. + # build_moa_facade relays "moa.*" events through tool_progress_callback so every surface + # shows each reference's answer before the aggregator acts. Display-only; 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" @@ -856,9 +830,8 @@ def _explicit_client_kwargs(agent, api_key, base_url, _provider_timeout) -> Dict if agent.provider == "copilot-acp": client_kwargs["command"] = agent.acp_command client_kwargs["args"] = agent.acp_args - # OpenCode Zen free tier (*-free slugs): the relay serves these ANONYMOUSLY and 401s any - # unrecognized bearer — including our keyless placeholder. Send an empty Authorization - # header to override the SDK's "Bearer ". + # OpenCode Zen free tier is served ANONYMOUSLY and 401s any bearer (incl. our keyless + # placeholder): send an empty Authorization header to override the SDK's "Bearer ". try: from hermes_cli.models import ( OPENCODE_ZEN_FREE_KEYLESS_PLACEHOLDER, opencode_zen_free_headers @@ -894,10 +867,9 @@ def _routed_client_kwargs(agent, fallback_model, _provider_timeout) -> Dict[str, agent.provider or "auto", model=agent.model, raw_codex=True) if _routed_client is not None: return _client_kwargs_from_routed(_routed_client, _provider_timeout) - # 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. + # No credentials: try the fallback chain BEFORE failing (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() for _fb in _fallback_entries(fallback_model): try: @@ -966,8 +938,8 @@ def _init_openai_client(agent, api_key, base_url, fallback_model, _provider_time headers["x-anthropic-beta"] = ",".join(filter(None, (existing_beta, _FINE_GRAINED_BETA))) 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 — this same dict. + # model.default_headers override provider/SDK defaults (WAFs rejecting SDK headers); + # mutates agent._client_kwargs in place. agent._apply_user_default_headers() try: @@ -980,9 +952,8 @@ def _init_openai_client(agent, api_key, base_url, fallback_model, _provider_time _cp_entries = get_compatible_custom_providers(load_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..extra_headers / - # custom_providers[].extra_headers). Applied last so the most specific config level - # wins. SECURITY: values may carry credentials — never log them. + # Per-provider extra_headers 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) @@ -1004,10 +975,8 @@ def _init_openai_client(agent, api_key, base_url, fallback_model, _provider_time 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). + # LLM client per wire mode (raw_codex=True: the main agent needs direct + # responses.stream()). One provider/model timeout up front so every path applies it. agent._anthropic_client = None agent._is_anthropic_oauth = False _provider_timeout = get_provider_request_timeout(agent.provider, agent.model) @@ -1091,14 +1060,12 @@ def _fallback_entries(fallback_model) -> List[Dict[str, Any]]: def _init_fallback_chain(agent, fallback_model): - # Stable identity for the pool entry that supplied this runtime: OAuth refreshes can - # replace the token before a failed request is recovered, so the mutable API-key value - # alone cannot attribute the failure to its source entry. + # Stable pool-entry identity: OAuth refreshes can replace the token before a failed + # request is recovered, so the key value alone can't attribute the failure. from agent.agent_runtime_helpers import sync_credential_pool_entry_id sync_credential_pool_entry_id(agent) - # Provider fallback chain — ordered backups tried when the primary is exhausted - # (rate-limit, overload, connection failure). Legacy single-dict or list format. + # Ordered backups tried when the primary is exhausted (legacy single-dict or list). agent._fallback_chain = _fallback_entries(fallback_model) agent._fallback_index = 0 agent._fallback_activated = getattr(agent, "_fallback_activated", False) @@ -1114,8 +1081,8 @@ def _init_fallback_chain(agent, fallback_model): 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 plugin manager has discovered its registrations first. + # A multiplexed gateway may have switched HERMES_HOME since model_tools was imported; + # make sure this profile's plugins are discovered before the tool snapshot. try: from hermes_cli.plugins import discover_plugins @@ -1123,8 +1090,7 @@ def _load_tools(agent, enabled_toolsets, disabled_toolsets): except Exception: logger.warning("Plugin discovery failed during agent setup", exc_info=True) - # Capture the registry generation FIRST so a later concurrent refresh can tell whether - # it holds a newer or staler view (see refresh_agent_mcp_tools). + # Capture the registry generation FIRST so a concurrent refresh can detect staleness. try: from tools.registry import registry as _snapshot_registry agent._tool_snapshot_generation = _snapshot_registry._generation @@ -1148,8 +1114,7 @@ def _load_tools(agent, enabled_toolsets, 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, not per prompt rebuild. + # Kanban guidance is session-static (kanban_show iff HERMES_KANBAN_TASK); resolve once. from agent.prompt_builder import KANBAN_GUIDANCE agent._kanban_worker_guidance = ( KANBAN_GUIDANCE if "kanban_show" in agent.valid_tool_names else "" @@ -1277,15 +1242,13 @@ def _init_session_state(agent, session_id, session_db, parent_session_id, reason 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. + # show_commentary: Codex phase=commentary → interim path (true) or reasoning channel. agent.show_commentary = bool(_cfg_dict(_agent_cfg, "display").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). + # lmstudio_load_mode: "explicit" (preload via management API) or "jit" (Auto-Evict path). _model_section = _cfg_dict(_agent_cfg, "model") agent.lmstudio_load_mode = "explicit" _load_mode = str(_model_section.get("lmstudio_load_mode", "explicit") or "explicit").strip().lower() @@ -1297,10 +1260,8 @@ def _apply_display_config(agent, _agent_cfg, platform): _model_section.get("lmstudio_load_mode"), ) - # API-transport streaming (``model.streaming``, default true). Some self-hosted backends - # have broken streaming tool-call paths, 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``. + # model.streaming=false seeds _disable_streaming (the loop's runtime fallback) for + # backends with broken streaming tool calls. Session-scoped; orthogonal to display.streaming. agent._disable_streaming = False _streaming = str(_model_section.get("streaming", "true")).strip().lower() if _streaming in {"false", "0", "no", "off"}: @@ -1319,8 +1280,7 @@ def _apply_display_config(agent, _agent_cfg, platform): ) except Exception as _tlg_err: _ra().logger.warning("Tool loop guardrail config ignored: %s", _tlg_err) - # Only the derived auxiliary compression context override is cached (needed by the - # startup feasibility check) — no broad pseudo-public config object on the agent. + # Only this derived override is cached (startup feasibility check) — no config object. agent._aux_compression_context_length_config = None @@ -1374,9 +1334,8 @@ def _init_memory(agent, _agent_cfg, skip_memory, platform): 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. + # skip_memory skips the external *provider*; enabled_toolsets=["memory"] still gets the + # built-in store so the memory tool never sees store=None. _memory_toolset_requested = ( "memory" in (agent.enabled_toolsets or []) and "memory" not in (agent.disabled_toolsets or []) @@ -1404,8 +1363,7 @@ def _init_memory(agent, _agent_cfg, skip_memory, platform): except Exception: pass # Memory is optional — don't break agent init - # Memory provider plugin (external — one at a time, alongside built-in), selected by - # memory.provider. + # External memory provider plugin (one at a time, alongside built-in): memory.provider. agent._memory_manager = None if not skip_memory: try: @@ -1450,9 +1408,8 @@ def _apply_agent_section(agent, _agent_cfg): pass _agent_section = _cfg_dict(_agent_cfg, "agent") - # Tool-use enforcement: "auto" (default — hardcoded model list), true, false, or list of - # substrings. Execution-discipline guidance: same shape against EXECUTION_GUIDANCE_MODELS, - # independent of enforcement (injection gate in agent/system_prompt.py). + # Both: "auto" (model-list match), true, false, or list of model substrings; independent + # of each other (gates in agent/system_prompt.py). agent._tool_use_enforcement = _agent_section.get("tool_use_enforcement", "auto") agent._execution_guidance = _agent_section.get("execution_guidance", "auto") @@ -1462,28 +1419,23 @@ def _apply_agent_section(agent, _agent_cfg): _agent_section.get("run_budget_seconds") ) - # Empty-response retry guard (``agent.empty_response_guard``): tolerant resolution — a - # malformed section falls back to schema defaults (guard on, $0.25 threshold). + # Empty-response guard: a malformed section falls back to schema defaults (on, $0.25). 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: "auto" (default — codex_responses only), true (all api_modes), - # false, or a list of model-name substrings; resolved in the loop's intent-ack block. + # "auto" (codex_responses only), true (all api_modes), false, or model substrings. agent._intent_ack_continuation = _agent_section.get("intent_ack_continuation", "auto") - # Runtime anti-stall guards (identical-call loop-breaker notice + continue-intent - # extension of empty-response recovery). Notice-only — never blocks a call. + # Anti-stall guards (identical-call notice + continue-intent recovery); notice-only. agent._stall_guards = bool(_agent_section.get("stall_guards", True)) - # Universal guidance toggles (ALL models, unlike enforcement): task-completion, - # parallel-tool-call batching, and the local Python toolchain probe. + # Universal guidance toggles (ALL models, unlike enforcement). 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 subprocesses) so the FIRST system-prompt build — on - # the time-to-first-token path — finds the line already cached. + # Warm the probe (~0.5s of subprocesses) off-thread so the first prompt build finds it cached. if agent._environment_probe: try: from tools.env_probe import warm_environment_probe_async @@ -1493,12 +1445,10 @@ def _apply_agent_section(agent, _agent_cfg): # Bot Mode teammate protocol section (tools/bot_mode_probe.py) — pure filesystem reads. 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. + # "Bot Chat" gate hint for hosts that defer the DB title write past the first prompt build. agent._session_title_hint = None - # Per-platform prompt-hint overrides (platform_hints: : {append|replace}), - # stored verbatim; resolved in agent/system_prompt.py. Invalid shapes are ignored. + # platform_hints: : {append|replace}, stored verbatim (agent/system_prompt.py). agent._platform_hint_overrides = _cfg_dict(_agent_cfg, "platform_hints") # App-level API retry count (wraps each model API call). Default 3; 1 = single attempt. @@ -1566,8 +1516,8 @@ def _compression_codex_settings(cfg: Dict[str, Any]) -> tuple[str, bool, Optiona app_server_auto, ) app_server_auto = "native" - # Native OpenAI Responses server-side compaction (opt-in; per-request gate in - # agent/native_compaction.py). Truthy coercion: "false"/"off" strings stay disabled. + # Native Responses server-side compaction (opt-in; gate in agent/native_compaction.py). + # Truthy coercion so "false"/"off" strings stay disabled. responses_native = is_truthy_value(cfg.get("codex_responses_native", False)) _raw = cfg.get("codex_responses_compact_threshold") compact_threshold = None @@ -1594,8 +1544,8 @@ def _parse_compression_config(agent, _agent_cfg) -> CompressionSettings: max_attempts = _parse_config_int(cfg.get("max_attempts", 3), 3) if max_attempts < 1: max_attempts = 3 - # threshold_tokens: absolute cap — compression triggers at the lower of the ratio - # threshold and this count; clamped to the window at apply-time (cap above window = no-op). + # threshold_tokens: absolute cap (lower of ratio threshold and this); clamped to the + # window at apply-time. threshold_tokens = cfg.get("threshold_tokens") if threshold_tokens is not None: threshold_tokens = _positive_int(threshold_tokens) @@ -1616,9 +1566,8 @@ def _parse_compression_config(agent, _agent_cfg) -> CompressionSettings: enabled=_cfg_flag(cfg, "enabled", True), target_ratio=target_ratio, protect_last=protect_last, - # tail_mode: "lean" (default) keeps a clamped 2.5%/10K-25K verbatim tail — continuity - # rides the summary. "legacy" restores the 0.20*threshold tail (hoards 100-240K on - # big windows). Unknown values fall back to lean inside the compressor. + # "lean" keeps a clamped 2.5%/10K-25K verbatim tail (continuity rides the summary); + # "legacy" restores the 0.20*threshold tail. Unknown → lean inside the compressor. tail_mode=str(cfg.get("tail_mode", "lean")).strip().lower(), # Actionable user messages guaranteed to survive in the tail (default 1, floor 1). min_tail_users=max(1, _parse_config_int(cfg.get("min_tail_user_messages", 1), 1)), @@ -1641,12 +1590,10 @@ def _parse_compression_config(agent, _agent_cfg) -> CompressionSettings: }, threshold_tokens=threshold_tokens, checkpoint_required=checkpoint_required, - # In-place compaction rewrites messages + system prompt WITHOUT rotating the session - # id. default=True MUST match DEFAULT_CONFIG — a False default flipped agents into - # rotation mode whenever the merged config omitted the key. + # In-place compaction: no session-id rotation. default=True MUST match DEFAULT_CONFIG + # (a False default flipped agents into rotation mode when the key was omitted). in_place=is_truthy_value(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. + # Opt-in: micro-compaction rewrites sent history per turn (breaks the cache prefix). micro_compact=is_truthy_value(cfg.get("micro_compact"), default=False), # Pass cadence in completed turns; each pass costs one prompt-cache break (>= 1). micro_compact_every_n_turns=max( @@ -1847,8 +1794,7 @@ def _warn_invalid_custom_provider_context_length(agent, _custom_providers) -> No def _resolve_context_length(agent, _agent_cfg, base_url): - # Explicit context_length for the auxiliary compression model: custom endpoints often - # can't report it via /models, so the startup feasibility check needs the hint. + # Aux compression model context_length hint (custom endpoints often can't report it). try: _aux_cfg = cfg_get(_agent_cfg, "auxiliary", "compression", default={}) except Exception: @@ -1884,8 +1830,7 @@ def _resolve_context_length(agent, _agent_cfg, base_url): ) _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. + # Resolve custom_providers before route-scoping: a named provider may keep its URL here. try: from hermes_cli.config import get_compatible_custom_providers _custom_providers = get_compatible_custom_providers(_agent_cfg) @@ -1918,8 +1863,7 @@ def _resolve_context_length(agent, _agent_cfg, base_url): if _config_context_length is None: _warn_invalid_custom_provider_context_length(agent, _custom_providers) - # Persist for switch_model / fallback activation — AFTER the custom_providers branch so - # per-model overrides aren't lost. + # Persisted for switch_model / fallback AFTER the custom_providers branch (per-model overrides). agent._config_context_length = _config_context_length _lmstudio_runtime_context_length = agent._ensure_lmstudio_runtime_loaded(_config_context_length) @@ -1959,9 +1903,8 @@ def _select_context_engine(_agent_cfg): 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". + # Deep-copy the shared singleton so a child's update_model() can't mutate the + # parent's. Uncopyable state (locks, DB conns) → built-in with an ACCURATE message. import copy try: _selected_engine = copy.deepcopy(_candidate) @@ -2053,8 +1996,8 @@ def _build_context_engine(agent, _agent_cfg, cs, _custom_providers, _effective_c agent.compression_enabled = cs.enabled agent.compression_in_place = cs.in_place _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 native_compaction.py). + # Micro-compaction has no pre-compress checkpoint hook; suppress it while the gate is + # armed (mirrors native_compaction.py). if cs.checkpoint_required and cs.micro_compact: logger.warning( "compression.checkpoint_required is enabled: post-turn " @@ -2129,10 +2072,9 @@ def _warn_nonagentic_hermes_model(agent): def _inject_context_engine_tools(agent): - # Inject context engine tool schemas (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. Gated on enabled_toolsets - # like memory-provider tools so `platform_toolsets: telegram: []` can't leak lcm_* tools. + # Context engine tool schemas (lcm_*), deduped against existing names (plugins may + # register the same schemas; duplicates 400 provider-side) and gated on enabled_toolsets + # so `platform_toolsets: telegram: []` can't leak them. agent._context_engine_tool_names: set = set() if ( agent.context_compressor @@ -2174,9 +2116,8 @@ def _inject_context_engine_tools(agent): def _configure_ollama_num_ctx(agent, _model_cfg, _config_context_length): - # Ollama defaults num_ctx to 2048 regardless of 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). + # Ollama defaults num_ctx to 2048, so detect the max window and send num_ctx per request. + # model.ollama_num_ctx overrides; model.context_length caps the detected value (VRAM). agent._ollama_num_ctx: int | None = None _override = _model_cfg.get("ollama_num_ctx") if isinstance(_model_cfg, dict) else None if _override is not None: @@ -2226,9 +2167,8 @@ def _configure_ollama_num_ctx(agent, _model_cfg, _config_context_length): 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. + # Codex autoraise notice: once per profile/config state (persisted marker; the gateway + # rebuilds the agent per message). The display gate hides the banner, not the autoraise. _autoraise = agent._compression_threshold_autoraised or {} _autoraise_notice = None if ( @@ -2243,8 +2183,7 @@ def _emit_compression_summary(agent, cs): 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 the percent would contradict the token count. + # The active engine's own threshold — a plugin's differs from cs.threshold. _active_threshold_pct = getattr( agent.context_compressor, "threshold_percent", cs.threshold ) @@ -2255,25 +2194,22 @@ def _emit_compression_summary(agent, cs): 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)") - # Printed inline for CLI users; gateway users get the same text replayed via - # _compression_warning on turn 1. + # Gateway users get the same text via _compression_warning on turn 1. if _autoraise_notice: print(_autoraise_notice) - # Gateway parity: status_callback isn't wired yet, so stash the text to replay on the - # first run_conversation(). Mark shown so repeated inits in this profile stay silent. + # status_callback isn't wired yet: stash for replay on the first turn; mark shown so + # repeated inits stay silent. agent._compression_warning = _autoraise_notice 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 it at most once per agent. + # Feasibility check deferred to the first turn near threshold (eager costs ~400ms cold). 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 so the preferred model gets a fresh attempt. One - # dict so new state fields are easy to add. + # Per-turn restoration snapshot: after a fallback, the next turn restores these so the + # preferred model gets a fresh attempt. _cc = agent.context_compressor agent._primary_runtime = { "model": agent.model, @@ -2287,8 +2223,7 @@ def _snapshot_primary_runtime(agent): "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. getattr because - # plugin engines may lack these ContextCompressor-specific attrs. + # Engine state _try_activate_fallback() overwrites (getattr: plugin engines may lack them). "compressor_model": getattr(_cc, "model", agent.model), "compressor_base_url": getattr(_cc, "base_url", agent.base_url), "compressor_api_key": getattr(_cc, "api_key", ""), @@ -2312,10 +2247,8 @@ def _init_usage_state(agent): # Copilot x-initiator flag: first API call of a user turn sends "user". agent._is_user_initiated_turn = False - # Usage-anchored context accounting (agent/model_metadata.py): last provider response's - # exact usage + transcript snapshot. None until the first response with usage; - # invalidated on compaction and session switches so stale anchors never suppress - # compression. + # Usage anchors (agent/model_metadata.py): last response's exact usage + transcript + # snapshot; invalidated on compaction/session switch so stale anchors never suppress compression. agent._usage_anchor = None agent._turn_base_usage_anchor = None @@ -2329,8 +2262,7 @@ def _init_usage_state(agent): 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), shared by - # conversation_loop and codex_runtime and readable by the CLI snapshot without IPC. + # Status-bar latency/velocity history (last 10 calls), shared by loop + codex_runtime. from collections import deque as _deque agent._api_latency_history = _deque(maxlen=10) agent._api_output_history = _deque(maxlen=10) @@ -2472,17 +2404,14 @@ def init_agent( setattr(agent, _name, _params[_name]) for _name in _GATEWAY_IDENTITY_PARAMS: setattr(agent, f"_{_name}", _params[_name]) - # Shared iteration budget — parent creates, children inherit; consumed by every LLM - # turn across parent + all subagents. + # Shared iteration budget: parent creates, children inherit. agent.iteration_budget = iteration_budget or IterationBudget(max_iterations) - # Pluggable print function — CLI replaces this with _cprint so raw ANSI status lines go - # through prompt_toolkit's renderer (patch_stdout's StdoutProxy would mangle them). - # None = builtins.print. + # CLI replaces this with _cprint so raw ANSI status lines go through prompt_toolkit's + # renderer (StdoutProxy would mangle them). None = 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" - # 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. + # Skips the end-of-turn review fork (~30K tokens/event); one switch for both review paths. agent.skip_background_review = bool(skip_background_review) agent.log_prefix = f"{log_prefix} " if log_prefix else "" # Effective base URL for feature detection (prompt caching, reasoning, etc.) @@ -2511,8 +2440,7 @@ def init_agent( _init_control_state(agent) - # 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. + # reasoning_content echo opt-in; switch_model / fallback / restore keep it in sync. agent._reasoning_echo_flag = agent._read_reasoning_echo_from_config() agent.request_overrides = dict(request_overrides or {}) agent.prefill_messages = prefill_messages or [] # Prefilled conversation turns