diff --git a/acp_adapter/server.py b/acp_adapter/server.py index 20896a905d..3586768b67 100644 --- a/acp_adapter/server.py +++ b/acp_adapter/server.py @@ -409,7 +409,7 @@ class HermesACPAgent(SlashCommandsMixin, acp.Agent): if not mcp_servers: return try: - from tools.mcp_tool import register_mcp_servers + from tools.mcp_tool_discovery import register_mcp_servers await asyncio.to_thread(register_mcp_servers, {s.name: _mcp_server_config(s) for s in mcp_servers}) except Exception: @@ -479,7 +479,7 @@ class HermesACPAgent(SlashCommandsMixin, acp.Agent): if any(int(getattr(agent, k, 0) or 0) > 0 for k in ("_user_turn_count", "_api_call_count")): return - from tools.mcp_tool import refresh_agent_mcp_tools + from tools.mcp_tool_agent import refresh_agent_mcp_tools added = refresh_agent_mcp_tools(agent, quiet_mode=True) if added: diff --git a/agent/agent_init.py b/agent/agent_init.py index 3ad54b13c2..cb393b21cc 100644 --- a/agent/agent_init.py +++ b/agent/agent_init.py @@ -72,53 +72,12 @@ def _warn_memory_provider_unavailable(name: str, reason: str = "") -> None: ) -# Canonicalize an endpoint URL for model-route identity comparisons. -_normalize_route_base_url = normalize_route_base_url - - -def _moa_reference_output_allowed(agent: Any) -> bool: - """Keep MoA display events off only the machine-readable ``-Q`` surface.""" - return not ( - getattr(agent, "platform", None) == "cli" - and getattr(agent, "tool_progress_mode", "all") == "off" - ) - - -def _relay_moa_reference_event(agent: Any, event: str, **kwargs: Any) -> None: - """Relay MoA display events while preserving the ``-Q`` stdout contract.""" - if not _moa_reference_output_allowed(agent): - return - cb = getattr(agent, "tool_progress_callback", None) - if cb is None: - return - try: - if event == "moa.reference": - cb( - "moa.reference", - str(kwargs.get("label") or ""), - str(kwargs.get("text") or ""), - None, - moa_index=kwargs.get("index"), - moa_count=kwargs.get("count"), - ) - elif event == "moa.aggregating": - cb( - "moa.aggregating", - str(kwargs.get("aggregator") or ""), - None, - None, - moa_ref_count=kwargs.get("ref_count"), - ) - except Exception: - pass - - def _provider_default_routes(provider: str) -> set[str]: """Return known exact default routes for a canonical provider id.""" routes: set[str] = set() def add(value): - route = _normalize_route_base_url(value) + route = normalize_route_base_url(value) if route: routes.add(route) @@ -151,7 +110,7 @@ def _context_route_mismatch( *, already_normalized: bool = False, ) -> bool: """Return whether a context pin's configured route differs from runtime.""" - _norm = (lambda v: str(v or "")) if already_normalized else _normalize_route_base_url + _norm = (lambda v: str(v or "")) if already_normalized else normalize_route_base_url configured_route, active_route = _norm(configured_base_url), _norm(active_base_url) if configured_route: return configured_route != active_route @@ -960,8 +919,8 @@ def _init_openai_client(agent, api_key, base_url, fallback_model, _provider_time 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() + from agent.ssl_guard import verify_ca_bundle + verify_ca_bundle() 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}") @@ -1569,7 +1528,7 @@ def _custom_provider_configured_base_url( _disabled_ids.update(_ids) continue if _wanted in _ids: - _url = _normalize_route_base_url( + _url = normalize_route_base_url( _entry.get("api") or _entry.get("url") or _entry.get("base_url") ) if _url: @@ -1581,7 +1540,7 @@ def _custom_provider_configured_base_url( if _key_ids & _disabled_ids: continue if _wanted in _key_ids | _custom_provider_runtime_ids(_entry.get("name")): - _url = _normalize_route_base_url(_entry.get("base_url")) + _url = normalize_route_base_url(_entry.get("base_url")) if _url: return _url return "" @@ -1596,7 +1555,7 @@ _RUNTIME_FIRST_PROVIDER_IDS = { def _configured_default_base_url(_agent_cfg, _model_cfg, _custom_providers) -> str: """Normalized route of the configured default model (``model.base_url``, else the named custom provider's URL when ``model.provider`` is not a first-class/auth provider).""" - _configured_base_url = _normalize_route_base_url(_model_cfg.get("base_url")) + _configured_base_url = normalize_route_base_url(_model_cfg.get("base_url")) _configured_provider = str(_model_cfg.get("provider") or "").strip() _norm = _normalize_custom_provider_name(_configured_provider) _custom_provider_candidate = bool(_norm) @@ -1622,9 +1581,9 @@ def _active_route_url(agent, base_url) -> str: if "?" in _requested.split("#", 1)[0]: with suppress(TypeError, ValueError): _without_query = urlunparse(urlparse(_requested)._replace(query="")) - if _normalize_route_base_url(_without_query) == _normalize_route_base_url(_active): + if normalize_route_base_url(_without_query) == normalize_route_base_url(_active): _active = _requested - return _normalize_route_base_url(_active) + return normalize_route_base_url(_active) def _scope_context_length_to_default_runtime( @@ -1679,13 +1638,13 @@ _CTX_LEN_REQUIREMENT = "must be a positive integer (e.g. 256000, not '256K')" def _warn_invalid_custom_provider_context_length(agent, _custom_providers) -> None: """Surface a context_length the helper silently skipped (not a positive int).""" - _target = _normalize_route_base_url(agent.base_url) + _target = normalize_route_base_url(agent.base_url) if not _target: return for _cp_entry in _custom_providers: if not isinstance(_cp_entry, dict): continue - if _normalize_route_base_url(_cp_entry.get("base_url")) != _target: + if normalize_route_base_url(_cp_entry.get("base_url")) != _target: continue _cp_models = _cp_entry.get("models", {}) _cp_model_cfg = _cp_models.get(agent.model, {}) if isinstance(_cp_models, dict) else None diff --git a/agent/auxiliary_client.py b/agent/auxiliary_client.py index 92e12a18e0..4fc976f7e3 100644 --- a/agent/auxiliary_client.py +++ b/agent/auxiliary_client.py @@ -20,7 +20,6 @@ import re import threading import time import uuid -from pathlib import Path # noqa: F401 β€” used by test mocks from types import SimpleNamespace from typing import Any, Callable, Dict, List, NamedTuple, Optional, Tuple, TYPE_CHECKING from urllib.parse import urlparse, parse_qs, urlunparse @@ -445,10 +444,6 @@ def aux_stream_deadline(deadline: Optional[float]): _aux_stream_deadline.value = previous -# Back-compat alias β€” the timing hooks were introduced with this name. -_aux_timing_hook = _aux_thread_local_hook - - def _run_protected_sync_provider_call(callback: Callable[[dict[str, Any]], Any], kwargs: dict[str, Any]) -> Any: """Run one protected provider callback in an attempt-isolated daemon thread. @@ -916,9 +911,6 @@ def _nous_extra_body() -> dict: return {"tags": _nous_portal_tags()} -# Back-compat snapshot; tests/plugins read ``NOUS_EXTRA_BODY`` directly. -NOUS_EXTRA_BODY = _nous_extra_body() - # Set at resolve time β€” True if the auxiliary client points to Nous Portal auxiliary_is_nous: bool = False @@ -3101,11 +3093,6 @@ def _is_unsupported_parameter_error(exc: Exception, param: str) -> bool: )) -def _is_unsupported_temperature_error(exc: Exception) -> bool: - """Back-compat wrapper for ``temperature``; kept as a named symbol because tests/call sites import it.""" - return _is_unsupported_parameter_error(exc, "temperature") - - def _is_structured_output_rejection(exc: Exception) -> bool: """Provider 400/422 rejecting the structured-output field, on either wire: OpenAI ``response_format`` (incl. vLLM's ``guided_grammar``/xgrammar failures) or Anthropic ``output_config.format`` ("Extra inputs @@ -4174,14 +4161,6 @@ def _resolve_auto_route( return _try_discovery_chain() -def _resolve_auto( - main_runtime: Optional[Dict[str, Any]] = None, task: Optional[str] = None -) -> Tuple[Optional[OpenAI], Optional[str]]: - """Backward-compatible auto resolver for callers that only need client/model.""" - client, model, _provider = _resolve_auto_route(main_runtime=main_runtime, task=task) - return client, model - - def _effective_provider_for_client(client: Any, fallback: str) -> str: """Return the concrete provider selected for an auto-routed client.""" effective_provider = getattr(client, "_hermes_aux_effective_provider", "") @@ -4825,7 +4804,7 @@ def resolve_provider_client( # Excluded: ``auto`` (a stale main slug could pair with any picked provider) and Nous + vision (the # Portal's tier-aware vision recommendation must win over a text-only model). if not model and provider != "auto" and not (provider == "nous" and is_vision): - # ``auto`` is intentionally excluded: `_resolve_auto(main_runtime=...)` returns the model paired + # ``auto`` is intentionally excluded: `_resolve_auto_route(main_runtime=...)` returns the model paired # with the provider it actually selected. Pre-filling an auto call from `_read_main_model()` can # leak a stale process-global runtime into a different provider (for example Claude model slug on # Codex OAuth) and override that correctly resolved model. 1. ``model`` argument (caller knew what @@ -4839,9 +4818,9 @@ def resolve_provider_client( # Each provider branch below sees a non-empty ``model`` whenever the user has *anything* configured # β€” no provider-specific empty-model guards needed. When the user has NOTHING configured (fresh # install, main_model also empty), the branches still hit their own missing-credentials returns and - # ``_resolve_auto`` falls through to the Step-2 chain as before. Do NOT pre-fill a blank ``auto`` + # ``_resolve_auto_route`` falls through to the Step-2 chain as before. Do NOT pre-fill a blank ``auto`` # request from the config/main default here. Claude model sent to Codex after the main lane fell - # back to gpt-5.5). Let _resolve_auto() return the actual current runtime model when the caller did + # back to gpt-5.5). Let _resolve_auto_route() return the actual current runtime model when the caller did # not explicitly request one. (# compression-current-model) Nous + vision is the one carve-out: the # branch below resolves its model from the Portal's tier-aware vision recommendation # (``_try_nous(vision= True)``), and ``final_model = model or default`` means anything pre-filled @@ -5417,14 +5396,14 @@ _AUX_DIRECT_API_BASE_URLS: Dict[str, str] = {"openai": "https://api.openai.com/v # MoA virtual provider: an *explicit* `provider: moa` override (either the caller-passed `provider` arg or # `auxiliary..provider` in config.yaml) reaches this function directly β€” it never goes through -# _resolve_auto(), which only unwraps the *implicit* "main provider is moa" case (#53827). Left as-is, "moa" +# _resolve_auto_route(), which only unwraps the *implicit* "main provider is moa" case (#53827). Left as-is, "moa" # is returned verbatim and resolve_provider_client() looks it up in PROVIDER_REGISTRY (which has no "moa" # entry β€” it's not a real HTTP provider), falls to the unknown-provider dead end, and call_llm surfaces a # nonsensical "MOA_API_KEY environment variable" error for a provider that was never meant to be reached # over the wire. Auxiliary tasks don't need the reference fan-out β€” resolve to the preset's aggregator slot # instead, exactly like the implicit path does (shared helper: _resolve_moa_aggregator). def _unwrap_moa_provider(prov: str, mdl: Optional[str]) -> Tuple[str, Optional[str]]: - """Resolve an *explicit* ``provider: moa`` to its preset's aggregator slot (_resolve_auto() + """Resolve an *explicit* ``provider: moa`` to its preset's aggregator slot (_resolve_auto_route() only unwraps the implicit case; "moa" isn't in PROVIDER_REGISTRY and would dead-end).""" if prov.strip().lower() != "moa": return prov, mdl @@ -5489,7 +5468,7 @@ def _resolve_task_provider_model( cfg_model = None resolved_model = model or cfg_model # Any moa:// facade endpoint belongs to the facade, not the aggregator's real provider β€” - # drop it (mirrors _resolve_auto()). + # drop it (mirrors _resolve_auto_route()). if provider and str(provider).strip().lower() == "moa": provider, resolved_model = _unwrap_moa_provider(provider, resolved_model) if provider and provider.lower() != "moa": @@ -6663,7 +6642,7 @@ def _ladder_parameter_rungs( """Rungs 1-3: retry without temperature / structured-output format / max_tokens. Returns ``(response, None, kwargs)`` or ``(None, narrowed_err, stripped_kwargs)``.""" client, task, tag = route.client, route.task, route.tag - if "temperature" in kwargs and _is_unsupported_temperature_error(first_err): + if "temperature" in kwargs and _is_unsupported_parameter_error(first_err, "temperature"): retry_kwargs = {k: v for k, v in kwargs.items() if k != "temperature"} logger.info("Auxiliary %s%s: provider rejected temperature; retrying once without it", task or "call", tag) @@ -6981,9 +6960,9 @@ def call_llm( if callable(prior_progress_hook) else ((lambda: None) if latency_info is not None else None) ), - _aux_timing_hook(_aux_dispatch, functools.partial( + _aux_thread_local_hook(_aux_dispatch, functools.partial( _stamp_latency_once, latency_info, "provider_dispatch_ms", request_started_at)), - _aux_timing_hook(_aux_provider_response, functools.partial( + _aux_thread_local_hook(_aux_provider_response, functools.partial( _stamp_latency_once, latency_info, "time_to_first_progress_ms", request_started_at)), ): response = _call_llm_impl( diff --git a/agent/azure_identity_adapter.py b/agent/azure_identity_adapter.py index a895569873..8c31bfa907 100644 --- a/agent/azure_identity_adapter.py +++ b/agent/azure_identity_adapter.py @@ -120,11 +120,10 @@ def _install_failure(allow_install: bool) -> Optional[Dict[str, Any]]: def build_token_provider(scope: Optional[str] = None, *, config: Optional[EntraIdentityConfig] = None, - base_url: Optional[str] = None, exclude_interactive_browser: bool = True, - ) -> Callable[[], str]: + exclude_interactive_browser: bool = True) -> Callable[[], str]: """Zero-arg callable minting a fresh Entra bearer JWT β€” pass as ``OpenAI(api_key=...)``. Scope precedence: - ``config.scope`` > ``scope`` kwarg > default; ``base_url`` is unused (back-compat). Not picklable: ship the - ``EntraIdentityConfig`` and rebuild in the worker.""" + ``config.scope`` > ``scope`` kwarg > default. Not picklable: ship the ``EntraIdentityConfig`` and rebuild + in the worker.""" ai = _require_azure_identity() config = _resolve_config(config, scope, exclude_interactive_browser=exclude_interactive_browser) return ai.get_bearer_token_provider(build_credential(config), config.scope) diff --git a/agent/background_review.py b/agent/background_review.py index 5d40c35995..2daf9ced40 100644 --- a/agent/background_review.py +++ b/agent/background_review.py @@ -294,8 +294,8 @@ def _digest_history(messages_snapshot: List[Dict], tail: int = 24) -> List[Dict] return [{"role": "user", "content": digest}] + keep -# Review prompts. AIAgent exposes them as class attributes -# (``_MEMORY_REVIEW_PROMPT`` etc.) for back-compat; the text lives here. +# Review prompts. AIAgent exposes them as class attributes (``_MEMORY_REVIEW_PROMPT`` etc.) so +# per-agent overrides work; the text lives here. _MEMORY_REVIEW_PROMPT = ( "Review the conversation above and consider saving to memory if appropriate.\n\n" "Focus on:\n" diff --git a/agent/bedrock_adapter.py b/agent/bedrock_adapter.py index a24513b026..f02ecffe30 100644 --- a/agent/bedrock_adapter.py +++ b/agent/bedrock_adapter.py @@ -838,65 +838,6 @@ def call_converse( return normalize_converse_response(response) -# Public API kept from main (plugins may import it): the agent loop itself streams through -# chat_completion_helpers._bedrock_converse_call, which applies the same recovery ladder. -def call_converse_stream( - region: str, - model: str, - messages: List[Dict], - tools: Optional[List[Dict]] = None, - max_tokens: Optional[int] = 4096, - temperature: Optional[float] = None, - top_p: Optional[float] = None, - stop_sequences: Optional[List[str]] = None, - guardrail_config: Optional[Dict] = None, -) -> SimpleNamespace: - """Call Bedrock ConverseStream API and return an OpenAI-compatible response. - - Consumes the full stream and returns the assembled response. For true - streaming with delta callbacks, use ``iter_converse_stream()`` instead. - """ - client = _get_bedrock_runtime_client(region) - kwargs = build_converse_kwargs( - model=model, - messages=messages, - tools=tools, - max_tokens=max_tokens, - temperature=temperature, - top_p=top_p, - stop_sequences=stop_sequences, - guardrail_config=guardrail_config, - ) - - try: - response = client.converse_stream(**kwargs) - except Exception as exc: - retry_kwargs = recover_from_cache_point_rejection(exc, kwargs) - if retry_kwargs is not None: - return normalize_converse_stream_events( - client.converse_stream(**retry_kwargs) - ) - if is_streaming_access_denied_error(exc): - # IAM allows bedrock:InvokeModel but not - # InvokeModelWithResponseStream β€” permanent for this session. - # Fall back to the non-streaming converse() path. - logger.info( - "bedrock: converse_stream denied by IAM on (region=%s, model=%s) β€” " - "falling back to non-streaming converse().", - region, model, - ) - return normalize_converse_response(client.converse(**kwargs)) - if is_stale_connection_error(exc): - logger.warning( - "bedrock: stale-connection error on converse_stream(region=%s, " - "model=%s): %s β€” evicting cached client so the next call reconnects.", - region, model, type(exc).__name__, - ) - invalidate_runtime_client(region) - raise - return normalize_converse_stream_events(response) - - # --- Model discovery --- _discovery_cache: Dict[str, Any] = {} @@ -986,62 +927,6 @@ def _extract_provider_from_arn(arn: str) -> str: return match.group(1) if match else "" -# --------------------------------------------------------------------------- -# Error classification β€” Bedrock-specific exceptions -# --------------------------------------------------------------------------- -# Mirrors OpenClaw's classifyFailoverReason() and matchesContextOverflowError() -# in extensions/amazon-bedrock/register.sync.runtime.ts. - -# Patterns that indicate the input context exceeded the model's token limit. -# Used by run_agent.py to trigger context compression instead of retrying. -CONTEXT_OVERFLOW_PATTERNS = [ - re.compile(r"ValidationException.*(?:input is too long|max input token|input token.*exceed)", re.IGNORECASE), - re.compile(r"ValidationException.*(?:exceeds? the (?:maximum|max) (?:number of )?(?:input )?tokens)", re.IGNORECASE), - re.compile(r"ModelStreamErrorException.*(?:Input is too long|too many input tokens)", re.IGNORECASE), -] - -# Patterns for throttling / rate limit errors β€” should trigger backoff + retry. -THROTTLE_PATTERNS = [ - re.compile(r"ThrottlingException", re.IGNORECASE), - re.compile(r"Too many concurrent requests", re.IGNORECASE), - re.compile(r"ServiceQuotaExceededException", re.IGNORECASE), -] - -# Patterns for transient overload β€” model is temporarily unavailable. -OVERLOAD_PATTERNS = [ - re.compile(r"ModelNotReadyException", re.IGNORECASE), - re.compile(r"ModelTimeoutException", re.IGNORECASE), - re.compile(r"InternalServerException", re.IGNORECASE), -] - - -def is_context_overflow_error(error_message: str) -> bool: - """Return True if the error indicates the input context was too large. - - When this returns True, the agent should compress context and retry - rather than treating it as a fatal error. - """ - return any(p.search(error_message) for p in CONTEXT_OVERFLOW_PATTERNS) - - -def classify_bedrock_error(error_message: str) -> str: - """Classify a Bedrock error for retry/failover decisions. - - Returns: - - ``"context_overflow"`` β€” input too long, compress and retry - - ``"rate_limit"`` β€” throttled, backoff and retry - - ``"overloaded"`` β€” model temporarily unavailable, retry with delay - - ``"unknown"`` β€” unclassified error - """ - if is_context_overflow_error(error_message): - return "context_overflow" - if any(p.search(error_message) for p in THROTTLE_PATTERNS): - return "rate_limit" - if any(p.search(error_message) for p in OVERLOAD_PATTERNS): - return "overloaded" - return "unknown" - - # --- Bedrock model context lengths --- # Static fallback when the live probe is unavailable (agent/model_metadata.py). Keys match by longest # substring, so versioned entries win over the generic "anthropic.claude-opus-4". diff --git a/agent/browser_provider.py b/agent/browser_provider.py index 84889ca242..cf9989bcd3 100644 --- a/agent/browser_provider.py +++ b/agent/browser_provider.py @@ -55,20 +55,3 @@ class BrowserProvider(ProviderBase): def emergency_cleanup(self, session_id: str) -> None: """Best-effort teardown from atexit / signal handlers. Must tolerate missing credentials and network errors; must not raise.""" - - # Legacy ``CloudBrowserProvider`` names still used by ``tools.browser_tool`` and out-of-tree subclasses. - - # ------------------------------------------------------------------ Backward-compat shims for the - # legacy CloudBrowserProvider API ------------------------------------------------------------------ The - # pre-PR-#25214 ABC exposed ``is_configured()`` and ``provider_name()``; ``tools.browser_tool`` has ~6 - # callers that still use those names. Rather than churn every callsite (and break out-of-tree downstream - # code that subclassed CloudBrowserProvider), we expose the old names as thin delegations to the new - # API. Subclasses MUST implement :meth:`is_available` and :attr:`name`; they may override - # ``is_configured`` / ``provider_name`` for compatibility with the legacy ABC but it is not required. - def is_configured(self) -> bool: - """Backward-compat alias for :meth:`is_available`.""" - return self.is_available() - - def provider_name(self) -> str: - """Backward-compat alias returning :attr:`display_name`.""" - return self.display_name diff --git a/agent/client_lifecycle.py b/agent/client_lifecycle.py index a782654afe..57384bb1c1 100644 --- a/agent/client_lifecycle.py +++ b/agent/client_lifecycle.py @@ -117,7 +117,6 @@ class ClientLifecycleMixin: _force_close_tcp_sockets = _forward_static("agent.agent_runtime_helpers", "force_close_tcp_sockets") _cleanup_dead_connections = _forward("agent.agent_runtime_helpers", "cleanup_dead_connections") _run_codex_stream = _forward("agent.codex_runtime", "run_codex_stream") - _run_codex_create_stream_fallback = _forward("agent.codex_runtime", "run_codex_create_stream_fallback") _recover_with_credential_pool = _forward("agent.agent_runtime_helpers", "recover_with_credential_pool") def _close_openai_client(self, client: Any, *, reason: str, shared: bool) -> None: diff --git a/agent/codex_responses_adapter.py b/agent/codex_responses_adapter.py index fcf58d0278..5f4395b268 100644 --- a/agent/codex_responses_adapter.py +++ b/agent/codex_responses_adapter.py @@ -209,10 +209,6 @@ def _summarize_user_message_for_log(content: Any, *, sep: str = " ") -> str: # --- ID helpers --------------------------------------------------------------- -# Deterministic call_id fallback (random ids would break the prompt-cache prefix); re-exported for run_agent/tests. -_deterministic_call_id = deterministic_call_id - - def _clamp_responses_call_id(call_id: str) -> str: """Keep ``call_id`` within the API's 64-char cap (the codex app-server namespaces MCP call ids past it). The surrogate is a pure function of the original so a ``function_call`` and its ``function_call_output`` agree.""" @@ -262,7 +258,7 @@ def _resolve_call_id( if not _nonblank(call_id) and canonicalize_fc: call_id = _canonical_call_id_from_fc(embedded_response_item_id) if not _nonblank(call_id): - call_id = _deterministic_call_id(fn_name, arguments, index) + call_id = deterministic_call_id(fn_name, arguments, index) return call_id.strip() diff --git a/agent/codex_runtime.py b/agent/codex_runtime.py index 62268dd81e..12460c2132 100644 --- a/agent/codex_runtime.py +++ b/agent/codex_runtime.py @@ -1,6 +1,6 @@ """Codex API runtime β€” App Server and Responses-API streaming paths. Every entry point takes the parent AIAgent first: ``run_codex_app_server_turn`` drives one ``codex app-server`` subprocess turn; -``run_codex_stream`` runs one streaming Codex Responses call (``run_codex_create_stream_fallback`` aliases it).""" +``run_codex_stream`` runs one streaming Codex Responses call.""" from __future__ import annotations @@ -965,12 +965,7 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta _close_event_stream(event_stream) -def run_codex_create_stream_fallback(agent, api_kwargs: dict, client: Any = None): - """Backward-compatible alias kept for tests and a few call sites.""" - return run_codex_stream(agent, api_kwargs, client=client) - - __all__ = [ - "run_codex_app_server_turn", "run_codex_stream", "run_codex_create_stream_fallback", + "run_codex_app_server_turn", "run_codex_stream", "_consume_codex_event_stream", "make_codex_app_server_event_bridge", ] diff --git a/agent/coding_context.py b/agent/coding_context.py index b35e346998..dfdf546a4e 100644 --- a/agent/coding_context.py +++ b/agent/coding_context.py @@ -345,7 +345,7 @@ class RuntimeMode: return prefix, [workspace] if workspace else [], trailing def system_blocks(self) -> list[str]: - """Posture blocks as one flat list in historical order (compat helper).""" + """Posture blocks as one flat list in historical order.""" prefix, workspace, trailing = self.system_prompt_parts() return [*prefix, *workspace, *trailing] @@ -381,7 +381,7 @@ def resolve_runtime_mode( ) -# ── Back-compat surface (thin wrappers over RuntimeMode) ──────────────────── +# ── Functional API (thin wrappers over RuntimeMode) ────────────────────────── def is_coding_context(*, platform: Optional[str] = None, cwd: Optional[str | Path] = None, config: Optional[dict[str, Any]] = None) -> bool: """Whether Hermes should operate in its coding posture right now.""" diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index 4a6aabb685..e691015e72 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -156,7 +156,7 @@ def is_compaction_progress_status(text: str | None) -> bool: def _refresh_agent_tool_definitions(agent) -> bool: """Rebuild agent.tools at the compaction commit boundary (the only moment config reaches a forever-session's frozen tool schemas; the prompt cache is already invalid). Returns True when tools were added.""" - from tools.mcp_tool import refresh_agent_mcp_tools + from tools.mcp_tool_agent import refresh_agent_mcp_tools added = refresh_agent_mcp_tools(agent, content_aware=True) if added: logger.info("Compaction tool refresh added tools: %s", sorted(added)) diff --git a/agent/conversation_loop.py b/agent/conversation_loop.py index 19ec984055..f855c28fc5 100644 --- a/agent/conversation_loop.py +++ b/agent/conversation_loop.py @@ -631,7 +631,7 @@ def _persist_system_prompt(agent, failure_message: str, *, persist_tools: bool = try: agent._session_db.update_system_prompt(agent.session_id, agent._cached_system_prompt) if persist_tools: - from tools.mcp_tool import persist_agent_tool_names + from tools.mcp_tool_agent import persist_agent_tool_names persist_agent_tool_names(agent) except Exception as exc: logger.warning(failure_message, agent.session_id, exc) @@ -692,7 +692,7 @@ def _restore_or_build_system_prompt(agent, system_message, conversation_history) try: saved_tools = session_row.get("tool_names") if session_row else None if saved_tools: - from tools.mcp_tool import restore_agent_tool_prefix + from tools.mcp_tool_agent import restore_agent_tool_prefix restore_agent_tool_prefix(agent, json.loads(saved_tools)) except Exception: logger.debug("tool prefix restore skipped", exc_info=True) diff --git a/agent/monitoring/gateway_health_export.py b/agent/monitoring/gateway_health_export.py index 616d1e5d13..e0f7a3cdf5 100644 --- a/agent/monitoring/gateway_health_export.py +++ b/agent/monitoring/gateway_health_export.py @@ -32,7 +32,6 @@ from agent.monitoring.otlp_exporter import ( _otlp_config, _resolve_headers, _runtime_resource_attributes, - _safe_resource_attributes, # noqa: F401 β€” re-exported for tests _signal_endpoint, ) from agent.monitoring.redaction import redact_bounded diff --git a/agent/reasoning_effort.py b/agent/reasoning_effort.py index ea196a68a3..8c67bbd1aa 100644 --- a/agent/reasoning_effort.py +++ b/agent/reasoning_effort.py @@ -33,7 +33,6 @@ CODEX_GPT56_EFFORTS: tuple[str, ...] = ("none", "low", "medium", "high", "xhigh" CODEX_LEGACY_EFFORTS: tuple[str, ...] = ("none", "low", "medium", "high", "xhigh") #: xAI Responses β€” Grok 4.6+ accepts xhigh; older Grok tops out at high. -# : Backward-compat alias (pre-#68365-verification name). XAI_GROK46_EFFORTS: tuple[str, ...] = ("low", "medium", "high", "xhigh") XAI_LEGACY_EFFORTS: tuple[str, ...] = ("low", "medium", "high") diff --git a/agent/ssl_guard.py b/agent/ssl_guard.py index 1135daf0fa..79b2b4e765 100644 --- a/agent/ssl_guard.py +++ b/agent/ssl_guard.py @@ -60,8 +60,3 @@ def verify_ca_bundle() -> None: except Exception as exc: raise _ssl_err(f"certifi is not importable: {exc}") from exc _validate_bundle_path("certifi", str(certifi.where()), require_substantial=True) - - -def verify_ca_bundle_with_fallback() -> None: - """Backward-compatible name for older call sites; a broken certifi bundle fails later anyway, so enforce the same check.""" - verify_ca_bundle() diff --git a/agent/stream_diag.py b/agent/stream_diag.py index 5500ee4565..f832bf976d 100644 --- a/agent/stream_diag.py +++ b/agent/stream_diag.py @@ -38,7 +38,7 @@ def stream_diag_capture_response(agent: Any, diag: Dict[str, Any], http_response try: headers = getattr(http_response, "headers", None) or {} captured: Dict[str, str] = {} - for name in getattr(agent, "_STREAM_DIAG_HEADERS", STREAM_DIAG_HEADERS): # per-agent override (back-compat) + for name in STREAM_DIAG_HEADERS: try: if val := headers.get(name): captured[name] = str(val)[:120] # keep log lines bounded diff --git a/agent/tool_dispatch_helpers.py b/agent/tool_dispatch_helpers.py index 42209f7447..f994ecfab2 100644 --- a/agent/tool_dispatch_helpers.py +++ b/agent/tool_dispatch_helpers.py @@ -75,7 +75,7 @@ def _is_destructive_command(cmd: str) -> bool: def _is_mcp_tool_parallel_safe(tool_name: str) -> bool: """Whether an MCP tool's server opted into parallel calls; False if MCP is unavailable.""" try: - from tools.mcp_tool import is_mcp_tool_parallel_safe # lazy: avoids import cycle + from tools.mcp_tool_discovery import is_mcp_tool_parallel_safe return is_mcp_tool_parallel_safe(tool_name) except Exception: return False diff --git a/agent/turn_context.py b/agent/turn_context.py index 9163a0e3a1..9f3143a8ae 100644 --- a/agent/turn_context.py +++ b/agent/turn_context.py @@ -250,13 +250,6 @@ def compression_made_progress( return new_len < orig_len or (orig_tokens > 0 and new_tokens < orig_tokens * 0.95) -# Back-compat alias: gateway callers and tests patch ``_compression_made_progress``. -# Back-compat alias: this predicate was module-private until the gateway's session-hygiene recovery gate -# needed the same semantics (#79624). Keeping the old name bound means existing callers and any test that -# patches ``_compression_made_progress`` continue to work unchanged. -_compression_made_progress = compression_made_progress - - class PreflightCompressionTimedOut(RuntimeError): """Raised when an oversized turn cannot safely finish preflight.""" diff --git a/agent/turn_context_compaction.py b/agent/turn_context_compaction.py index 415a54fefd..a9f8a9b703 100644 --- a/agent/turn_context_compaction.py +++ b/agent/turn_context_compaction.py @@ -394,7 +394,7 @@ def _run_preflight_passes( _preflight_tokens = _tc._preflight_request_tokens( agent, out.messages, out.active_system_prompt or "" ) - if not _tc._compression_made_progress( + if not _tc.compression_made_progress( _orig_len, len(out.messages), _orig_tokens, _preflight_tokens ): _tc._fail_closed_after_preflight_timeout(agent, _preflight_tokens) diff --git a/agent/turn_facade.py b/agent/turn_facade.py index 6252e3ba2e..5f147c62cf 100644 --- a/agent/turn_facade.py +++ b/agent/turn_facade.py @@ -11,7 +11,6 @@ from contextlib import suppress from typing import Any, Dict, List, Optional from agent.lazy_forward import forward as _forward -from tools.interrupt import set_interrupt as _set_interrupt # noqa: F401 (tests patch it here) # Same logger name as the origin module so log records / caplog filters are unchanged. logger = logging.getLogger("run_agent") diff --git a/agent/turn_facade_lease.py b/agent/turn_facade_lease.py index 5fb4f5c052..da33122f6f 100644 --- a/agent/turn_facade_lease.py +++ b/agent/turn_facade_lease.py @@ -159,8 +159,7 @@ class DurableTurnLease: if not message: return agent = self.agent - # Lazy via the faΓ§ade so ``patch("agent.turn_facade._set_interrupt")`` keeps intercepting. - from agent.turn_facade import _set_interrupt + from tools.interrupt import set_interrupt as _set_interrupt with getattr(agent, "_pending_redirect_lock", None) or nullcontext(): if getattr(agent, "_interrupt_message", None) != message: diff --git a/cli.py b/cli.py index 8069ce7998..ff31406bfa 100644 --- a/cli.py +++ b/cli.py @@ -783,7 +783,7 @@ def _interrupt_async_delegations() -> None: def _shutdown_mcp_servers() -> None: - from tools.mcp_tool import shutdown_mcp_servers + from tools.mcp_tool_lifecycle import shutdown_mcp_servers shutdown_mcp_servers() diff --git a/cron/scheduler.py b/cron/scheduler.py index 781d755966..3a3bc39798 100644 --- a/cron/scheduler.py +++ b/cron/scheduler.py @@ -1663,7 +1663,7 @@ def _init_cron_mcp_tools(job_id: str) -> None: # paths called discover_mcp_tools() at startup. Idempotent: subsequent ticks short-circuit on # already-connected servers inside register_mcp_servers(). Non-fatal on failure: a broken MCP server # shouldn't kill an otherwise-working cron job. See #4219. - from tools.mcp_tool import discover_mcp_tools + from tools.mcp_tool_discovery import discover_mcp_tools _mcp_tools = discover_mcp_tools() if _mcp_tools: logger.info("Job '%s': %d MCP tool(s) available", job_id, len(_mcp_tools)) @@ -3595,7 +3595,7 @@ def _sweep_mcp_orphans() -> None: """Reap MCP stdio orphans (only PIDs flagged by tools.mcp_tool._run_stdio's finally block); run AFTER jobs finish so live sessions are never touched.""" try: - from tools.mcp_tool import _kill_orphaned_mcp_children + from tools.mcp_tool_lifecycle import _kill_orphaned_mcp_children _kill_orphaned_mcp_children() except Exception as _e: logger.debug("Post-tick MCP orphan cleanup failed: %s", _e) diff --git a/gateway/run.py b/gateway/run.py index 98ce1d4a71..01cb426982 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -1719,7 +1719,7 @@ async def _discover_gateway_mcp_tools(config: object) -> None: carry the scope into the executor thread with ``copy_context()`` (the same shape as ``_run_in_executor_with_context``). See #95518. """ - from tools.mcp_tool import discover_mcp_tools + from tools.mcp_tool_discovery import discover_mcp_tools loop = asyncio.get_running_loop() if not getattr(config, "multiplex_profiles", False): await loop.run_in_executor(None, discover_mcp_tools) @@ -4592,7 +4592,7 @@ async def _shutdown_mcp_servers_nonblocking(timeout: float = 5.0) -> bool: """ def _do() -> None: try: - from tools.mcp_tool import shutdown_mcp_servers + from tools.mcp_tool_lifecycle import shutdown_mcp_servers shutdown_mcp_servers() except Exception: logger.debug("MCP shutdown raised", exc_info=True) diff --git a/gateway/run_turn.py b/gateway/run_turn.py index 4dfe8df32f..3b7ee873fc 100644 --- a/gateway/run_turn.py +++ b/gateway/run_turn.py @@ -2230,7 +2230,7 @@ class GatewayTurnMixin: a history-destroying ``/new``. Each agent keeps its build-time toolset selection EXACTLY: a session built with restricted enabled_toolsets (e.g. ["safe"]) must NOT silently gain tools.""" try: - from tools.mcp_tool import refresh_agent_mcp_tools + from tools.mcp_tool_agent import refresh_agent_mcp_tools _cache = getattr(self, "_agent_cache", None) _cache_lock = getattr(self, "_agent_cache_lock", None) if _cache_lock is None or not _cache: @@ -2262,8 +2262,10 @@ class GatewayTurnMixin: with _profile_runtime_scope(Path(profile_home)): return await self._execute_mcp_reload(event) try: - from tools.mcp_tool import shutdown_mcp_servers, discover_mcp_tools, _servers, _lock - from tools.mcp_tool import _server_scope_keys, reprobe_tool_availability + from tools.mcp_tool_lifecycle import shutdown_mcp_servers + from tools.mcp_tool_discovery import discover_mcp_tools + from tools.mcp_tool import _servers, _lock, _server_scope_keys + from tools.mcp_tool_agent import reprobe_tool_availability from tools.registry import registry reload_scope = registry.current_scope_key() if multiplex else None diff --git a/gateway/session.py b/gateway/session.py index 857f0dd57b..4cebdf9156 100644 --- a/gateway/session.py +++ b/gateway/session.py @@ -198,7 +198,7 @@ def _slack_tools_loaded() -> bool: intentionally not per-session). False on any error so a bad config never promises tools. """ try: - from tools.mcp_tool import get_registered_mcp_server_names + from tools.mcp_tool_discovery import get_registered_mcp_server_names if any("slack" in name.lower() for name in get_registered_mcp_server_names()): return True except Exception: diff --git a/hermes_cli/banner.py b/hermes_cli/banner.py index 672ea88174..62cab7fc26 100644 --- a/hermes_cli/banner.py +++ b/hermes_cli/banner.py @@ -710,7 +710,7 @@ def _mcp_configured() -> bool: def _probe_mcp_status() -> list: - from tools.mcp_tool import get_mcp_status + from tools.mcp_tool_discovery import get_mcp_status return get_mcp_status() diff --git a/hermes_cli/cli_commands_mixin.py b/hermes_cli/cli_commands_mixin.py index 4b982113d8..d760f09f46 100644 --- a/hermes_cli/cli_commands_mixin.py +++ b/hermes_cli/cli_commands_mixin.py @@ -588,7 +588,7 @@ def _browser_status() -> None: else: provider = _probe("tools.browser_tool", "_get_cloud_provider", None) if provider is not None: - print(f"🌐 Browser: {provider.provider_name()} (cloud)") + print(f"🌐 Browser: {provider.display_name} (cloud)") _print_lightpanda_engine_status() else: engine = _probe("tools.browser_tool_cloud", "_get_browser_engine", "auto") diff --git a/hermes_cli/cli_info_mixin.py b/hermes_cli/cli_info_mixin.py index 4c1ba2e3b7..cfdcb36dc4 100644 --- a/hermes_cli/cli_info_mixin.py +++ b/hermes_cli/cli_info_mixin.py @@ -856,9 +856,10 @@ class CLIInfoMixin: """Reload MCP servers: disconnect all, re-read config.yaml, reconnect, then refresh the agent's tool list so the model sees the updated tools on the next turn.""" try: - from tools.mcp_tool import ( - shutdown_mcp_servers, discover_mcp_tools, reprobe_tool_availability, _servers, _lock, - ) + from tools.mcp_tool_lifecycle import shutdown_mcp_servers + from tools.mcp_tool_discovery import discover_mcp_tools + from tools.mcp_tool_agent import reprobe_tool_availability + from tools.mcp_tool import _servers, _lock with _lock: old_servers = set(_servers.keys()) if not self._command_running: @@ -886,7 +887,7 @@ class CLIInfoMixin: # gateway reload / late-binding paths (name-diff, thread-safe, additive-preserving so # memory-provider and context-engine tools survive the rebuild). if self.agent is not None: - from tools.mcp_tool import refresh_agent_mcp_tools + from tools.mcp_tool_agent import refresh_agent_mcp_tools # Pick up servers ENABLED in config this session: enabled_toolsets was resolved at # startup, so merge now-connected names in (unless `all`/`*` is pinned) so a # freshly-added server isn't filtered out. Mirrors startup (see __init__). diff --git a/hermes_cli/doctor_platform.py b/hermes_cli/doctor_platform.py index b175c56d76..9b1f7edf2d 100644 --- a/hermes_cli/doctor_platform.py +++ b/hermes_cli/doctor_platform.py @@ -182,14 +182,14 @@ def check_certificates(should_fix: bool = False, issues: "list | None" = None) - certifi into THIS interpreter's environment and re-verifying. """ try: - from agent.ssl_guard import verify_ca_bundle_with_fallback + from agent.ssl_guard import verify_ca_bundle from agent.errors import SSLConfigurationError except Exception as e: return check_warn("SSL certificate check skipped", str(e)) if issues is None: issues = [] try: - verify_ca_bundle_with_fallback() + verify_ca_bundle() return check_ok("SSL CA certificate bundle is valid") except SSLConfigurationError as e: first_error = str(e) @@ -215,7 +215,7 @@ def check_certificates(should_fix: bool = False, issues: "list | None" = None) - sys.modules.pop(mod_name, None) importlib.invalidate_caches() try: - verify_ca_bundle_with_fallback() + verify_ca_bundle() check_ok("SSL CA certificate bundle repaired (certifi reinstalled)") except SSLConfigurationError as e: _fail_and_issue("SSL CA certificate bundle still broken after reinstall", str(e), diff --git a/hermes_cli/main.py b/hermes_cli/main.py index 918a3e0e6c..99171b98e8 100644 --- a/hermes_cli/main.py +++ b/hermes_cli/main.py @@ -105,8 +105,8 @@ _oneshot_cleanup_done = False _ONESHOT_CLEANUPS = ( ("tools.terminal_tool", "cleanup_all_environments", {}, Exception), ("tools.async_delegation", "interrupt_all", {"reason": "oneshot shutdown"}, Exception), - ("tools.browser_tool", "_emergency_cleanup_all_sessions", {}, Exception), - ("tools.mcp_tool", "shutdown_mcp_servers", {}, BaseException), + ("tools.browser_tool_lifecycle", "_emergency_cleanup_all_sessions", {}, Exception), + ("tools.mcp_tool_lifecycle", "shutdown_mcp_servers", {}, BaseException), ("agent.auxiliary_client", "shutdown_cached_clients", {}, Exception), ) @@ -2966,7 +2966,7 @@ def _prepare_agent_startup(args) -> None: if _run_inline_mcp_discovery: try: # synchronous for entrypoints without a later bounded startup path from hermes_cli.mcp_startup import get_mcp_server_filter - from tools.mcp_tool import discover_mcp_tools + from tools.mcp_tool_discovery import discover_mcp_tools _mcp_filter = get_mcp_server_filter() if _mcp_filter is None: diff --git a/hermes_cli/mcp_config.py b/hermes_cli/mcp_config.py index 7e0d956f00..c2ba4df047 100644 --- a/hermes_cli/mcp_config.py +++ b/hermes_cli/mcp_config.py @@ -18,7 +18,8 @@ from hermes_cli.config import ( from hermes_cli.colors import Colors, color from hermes_constants import display_hermes_home from hermes_cli.mcp_security import validate_mcp_server_entry -from tools.mcp_tool import _ENV_VAR_PATTERN, _env_ref_name +from tools.mcp_tool_config import _ENV_VAR_PATTERN +from tools.mcp_tool_common import _env_ref_name logger = logging.getLogger(__name__) @@ -240,7 +241,7 @@ def _resolve_mcp_server_config(config: dict) -> dict: probe sent the literal placeholder and auth-requiring servers (e.g. n8n) returned 401 β€” while runtime tool loading worked because it interpolates. (#37792) """ - from tools.mcp_tool import _interpolate_env_vars + from tools.mcp_tool_config import _interpolate_env_vars from agent.secret_scope import current_secret_scope if current_secret_scope() is None: @@ -264,8 +265,10 @@ def _probe_single_server( if issues: raise ValueError("; ".join(issues)) - from tools.mcp_tool import ( - _ensure_mcp_loop, _run_on_mcp_loop, _connect_server, _stop_mcp_loop_if_idle, _parse_boolish) + from tools.mcp_tool_loop import _ensure_mcp_loop, _run_on_mcp_loop + from tools.mcp_tool_discovery import _connect_server + from tools.mcp_tool_lifecycle import _stop_mcp_loop_if_idle + from tools.mcp_tool_common import _parse_boolish config = _resolve_mcp_server_config(config) if connect_timeout is None: @@ -290,7 +293,7 @@ def _probe_single_server( # the desktop can estimate per-call token cost. Best-effort, absent on failure. try: import json as _json - from tools.mcp_tool import _convert_mcp_schema + from tools.mcp_tool_schema import _convert_mcp_schema details["schema_chars"] = { t.name: len(_json.dumps(_convert_mcp_schema(name, t), separators=(",", ":"), default=str)) @@ -796,7 +799,7 @@ def cmd_mcp_configure(args): # Same matching semantics as runtime registration (tools/mcp_tool.py): exact names or globs. try: - from tools.mcp_tool import matches_name_filter + from tools.mcp_tool_schema import matches_name_filter except ImportError: # pragma: no cover β€” defensive fallback def matches_name_filter(tool_name, patterns): return tool_name in patterns diff --git a/hermes_cli/mcp_startup.py b/hermes_cli/mcp_startup.py index 11620e359c..b5ab4282a2 100644 --- a/hermes_cli/mcp_startup.py +++ b/hermes_cli/mcp_startup.py @@ -61,7 +61,7 @@ def _has_configured_mcp_servers() -> bool: def _any_mcp_connected() -> bool: - from tools.mcp_tool import get_mcp_status + from tools.mcp_tool_discovery import get_mcp_status return any(entry.get("connected") for entry in (get_mcp_status() or [])) @@ -155,7 +155,7 @@ def _discover_mcp_tools_without_interactive_oauth() -> None: suppress_interactive_oauth = nullcontext with suppress_interactive_oauth(): - from tools.mcp_tool import discover_mcp_tools + from tools.mcp_tool_discovery import discover_mcp_tools # Only pass the kwarg when a filter is set: many tests (and any # out-of-tree caller) stub discover_mcp_tools with a zero-arg diff --git a/hermes_cli/plugins.py b/hermes_cli/plugins.py index 34c11150db..684f41a912 100644 --- a/hermes_cli/plugins.py +++ b/hermes_cli/plugins.py @@ -524,7 +524,7 @@ class PluginContext: except (TypeError, ValueError): timeout = 30.0 timeout = max(1.0, min(timeout, 600.0)) - from tools.mcp_tool import _make_tool_handler + from tools.mcp_tool_handlers import _make_tool_handler raw = _make_tool_handler(server, tool, timeout)(dict(arguments or {})) logger.debug("Plugin %s called MCP %s/%s (timeout=%ss, %d chars returned)", self.manifest.name, server, tool, timeout, len(raw or "")) diff --git a/hermes_cli/tools_config_mcp.py b/hermes_cli/tools_config_mcp.py index c2123400e1..67e9a5250e 100644 --- a/hermes_cli/tools_config_mcp.py +++ b/hermes_cli/tools_config_mcp.py @@ -18,7 +18,7 @@ def _mcp_match_filter(): excludes (e.g. ``*team_member*`` from catalog default_excluded manifests) as if nothing were excluded.""" try: - from tools.mcp_tool import matches_name_filter + from tools.mcp_tool_schema import matches_name_filter return matches_name_filter except ImportError: # pragma: no cover β€” defensive fallback return lambda tool_name, patterns: tool_name in patterns @@ -93,7 +93,7 @@ def _configure_mcp_tools_interactive(config: dict): print(color(f" Connecting to {len(enabled_names)} server(s): {', '.join(enabled_names)}", Colors.DIM)) try: - from tools.mcp_tool import probe_mcp_server_tools + from tools.mcp_tool_discovery import probe_mcp_server_tools server_tools = probe_mcp_server_tools() except Exception as exc: _print_error(f"Failed to probe MCP servers: {exc}") diff --git a/hermes_cli/web_server_mcp.py b/hermes_cli/web_server_mcp.py index 5308a35e1a..2d28d907d4 100644 --- a/hermes_cli/web_server_mcp.py +++ b/hermes_cli/web_server_mcp.py @@ -149,7 +149,7 @@ def _run_dashboard_mcp_oauth(flow, cfg: dict) -> None: flow.tools = [{"name": t, "description": d} for t, d in tools] flow.mark_approved() if flow.reconnect_live: - from tools.mcp_tool import reconnect_mcp_server + from tools.mcp_tool_loop import reconnect_mcp_server reconnect_mcp_server(flow.server_name) except Exception: diff --git a/run_agent.py b/run_agent.py index 1766c4555d..76ac04e091 100644 --- a/run_agent.py +++ b/run_agent.py @@ -124,11 +124,11 @@ from agent.session_activity import ActivityProvenance from agent.model_metadata import is_local_endpoint from agent.message_sanitization import ( coalesce_tool_call_id as _sanitize_coalesce_tool_call_id, + deterministic_call_id as _codex_deterministic_call_id, uniquify_tool_call_ids as _sanitize_uniquify_tool_call_ids, ) from agent.codex_responses_adapter import ( _derive_responses_function_call_id as _codex_derive_responses_function_call_id, - _deterministic_call_id as _codex_deterministic_call_id, _split_responses_tool_id as _codex_split_responses_tool_id, _summarize_user_message_for_log, ) diff --git a/tests/agent/test_aux_progress_streaming.py b/tests/agent/test_aux_progress_streaming.py index fcc5f7d81d..6d30f17f0b 100644 --- a/tests/agent/test_aux_progress_streaming.py +++ b/tests/agent/test_aux_progress_streaming.py @@ -494,7 +494,7 @@ class TestContentBearingProgress: any kind (transport liveness), not only on the first token.""" from agent.auxiliary_client import ( _aux_provider_response, - _aux_timing_hook, + _aux_thread_local_hook, _notify_aux_timing_response, ) @@ -507,7 +507,7 @@ class TestContentBearingProgress: accumulator = _ChatStreamAccumulator() with ( - _aux_timing_hook(_aux_provider_response, _timed_response), + _aux_thread_local_hook(_aux_provider_response, _timed_response), aux_progress_hook(lambda: None), ): accumulator.feed(keepalive) diff --git a/tests/agent/test_auxiliary_client.py b/tests/agent/test_auxiliary_client.py index f2930a0604..7bcb5f3fba 100644 --- a/tests/agent/test_auxiliary_client.py +++ b/tests/agent/test_auxiliary_client.py @@ -32,7 +32,7 @@ from agent.auxiliary_client import ( _try_openrouter, _OPENROUTER_MODEL, OPENROUTER_BASE_URL, - _resolve_auto, + _resolve_auto_route, _resolve_task_provider_model, _resolve_xai_oauth_for_aux, _CodexCompletionsAdapter, @@ -72,7 +72,7 @@ def _clean_env(monkeypatch): monkeypatch.delenv(key, raising=False) # Module-level unhealthy cache (10-min TTL) leaks between tests; # earlier tests that call _mark_provider_unhealthy() poison the - # cache for later ones, causing _resolve_auto to skip providers + # cache for later ones, causing _resolve_auto_route to skip providers # that the test patched to return valid clients. import agent.auxiliary_client as _aux_mod _aux_mod._aux_unhealthy_until.clear() @@ -166,7 +166,7 @@ class TestResolveTaskProviderModel: """An *explicit* `provider="moa"` arg (e.g. a per-task model override naming a MoA preset) must resolve to the preset's aggregator, not the literal "moa" string β€” mirrors #53827's fix for the implicit - "main provider is moa" case in _resolve_auto(), which this function + "main provider is moa" case in _resolve_auto_route(), which this function never went through.""" preset = { "aggregator": {"provider": "openrouter", "model": "anthropic/claude-opus-4.8"}, @@ -796,7 +796,7 @@ class TestResolveProviderClientUniversalModelFallback: Pre-fix the OAuth providers (xai-oauth, openai-codex) returned ``(None, None)`` on an empty model β€” both lack a catalog default because their accepted-model lists drift on the backend. That - silent failure caused ``_resolve_auto`` to drop to its Step-2 + silent failure caused ``_resolve_auto_route`` to drop to its Step-2 fallback chain (OpenRouter / Nous / etc.), so aux tasks billed against the wrong subscription. """ @@ -890,8 +890,8 @@ class TestExpiredCodexFallback: monkeypatch.setenv("ANTHROPIC_TOKEN", "sk-ant-oat01-test-fallback") with patch("agent.anthropic_adapter.build_anthropic_client") as mock_build: mock_build.return_value = MagicMock() - from agent.auxiliary_client import _resolve_auto - client, model = _resolve_auto() + from agent.auxiliary_client import _resolve_auto_route + client, model, _provider = _resolve_auto_route() # Should NOT be Codex, should be Anthropic (or another available provider) assert not isinstance(client, type(None)), "Should find a provider after expired Codex" @@ -932,8 +932,8 @@ class TestExpiredCodexFallback: with patch("agent.auxiliary_client.OpenAI") as mock_openai: mock_openai.return_value = MagicMock() - from agent.auxiliary_client import _resolve_auto - client, model = _resolve_auto() + from agent.auxiliary_client import _resolve_auto_route + client, model, _provider = _resolve_auto_route() assert client is not None # OpenRouter is 1st in chain, should win mock_openai.assert_called() @@ -2427,7 +2427,7 @@ class TestKimiTemperatureOmitted: class TestStaleBaseUrlWarning: - """_resolve_auto() warns when OPENAI_BASE_URL conflicts with config provider (#5161).""" + """_resolve_auto_route() warns when OPENAI_BASE_URL conflicts with config provider (#5161).""" def test_warns_when_openai_base_url_set_with_named_provider(self, monkeypatch, caplog): """Warning fires when OPENAI_BASE_URL is set but provider is a named provider.""" @@ -2440,7 +2440,7 @@ class TestStaleBaseUrlWarning: with patch("agent.auxiliary_client._read_main_provider", return_value="openrouter"), \ patch("agent.auxiliary_client._read_main_model", return_value="google/gemini-flash"), \ caplog.at_level(logging.WARNING, logger="agent.auxiliary_client"): - _resolve_auto() + _resolve_auto_route() assert any("OPENAI_BASE_URL is set" in rec.message for rec in caplog.records), \ "Expected a warning about stale OPENAI_BASE_URL" @@ -2611,7 +2611,7 @@ class TestAuxiliaryTaskExtraBody: patch("agent.auxiliary_client.OpenAI") as mock_openai, \ caplog.at_level(logging.WARNING, logger="agent.auxiliary_client"): mock_openai.return_value = MagicMock() - _resolve_auto() + _resolve_auto_route() assert not any("OPENAI_BASE_URL is set" in rec.message for rec in caplog.records), \ "Should NOT warn when provider is 'custom'" @@ -3318,7 +3318,7 @@ class TestCodexAdapterGithubResponsesMessageIdDrop: class TestVisionAutoSkipsKimiCoding: - """_resolve_auto vision branch skips providers that have no vision on + """_resolve_auto_route vision branch skips providers that have no vision on their main endpoint (e.g. Kimi Coding Plan /coding) and falls through to the aggregator chain instead of handing back a client that will 404 on every request (#17076). diff --git a/tests/agent/test_auxiliary_explicit_base_anthropic.py b/tests/agent/test_auxiliary_explicit_base_anthropic.py index c882e6b30f..9ad35975d1 100644 --- a/tests/agent/test_auxiliary_explicit_base_anthropic.py +++ b/tests/agent/test_auxiliary_explicit_base_anthropic.py @@ -5,7 +5,7 @@ When the main provider is ``custom`` and its ``base_url`` ends in ``/anthropic`` (a proxied Anthropic gateway β€” MiniMax, Zhipu GLM, LiteLLM, or a self-hosted LLM proxy), auxiliary tasks reach ``resolve_provider_client("custom", explicit_base_url=..., api_mode="anthropic_messages")`` β€” directly for a -per-task ``auxiliary.`` override, or via ``_resolve_auto`` Step 1 which +per-task ``auxiliary.`` override, or via ``_resolve_auto_route`` Step 1 which forwards the main runtime's ``api_mode``. The bug (issue #16254): this branch called ``_to_openai_base_url()`` diff --git a/tests/agent/test_auxiliary_main_first.py b/tests/agent/test_auxiliary_main_first.py index e999364cee..47923aedff 100644 --- a/tests/agent/test_auxiliary_main_first.py +++ b/tests/agent/test_auxiliary_main_first.py @@ -17,11 +17,11 @@ from unittest.mock import MagicMock, patch -# ── Text aux tasks β€” _resolve_auto ────────────────────────────────────────── +# ── Text aux tasks β€” _resolve_auto_route ────────────────────────────────────────── class TestResolveAutoMainFirst: - """_resolve_auto() must prefer main provider + main model for every user.""" + """_resolve_auto_route() must prefer main provider + main model for every user.""" def test_title_generation_auto_honors_main_model(self): """The default auto title route must not replace the selected main model.""" @@ -37,9 +37,9 @@ class TestResolveAutoMainFirst: ) as mock_resolve, patch( "agent.auxiliary_client._is_provider_unhealthy", return_value=False ): - from agent.auxiliary_client import _resolve_auto + from agent.auxiliary_client import _resolve_auto_route - client, model = _resolve_auto( + client, model, _provider = _resolve_auto_route( main_runtime={ "provider": "opencode-zen", "model": main_model, @@ -71,9 +71,9 @@ class TestResolveAutoMainFirst: ), patch( "agent.auxiliary_client._is_provider_unhealthy", return_value=False ): - from agent.auxiliary_client import _resolve_auto + from agent.auxiliary_client import _resolve_auto_route - client, model = _resolve_auto( + client, model, _provider = _resolve_auto_route( main_runtime={ "provider": "opencode-zen", "model": "deepseek-v4-flash-free", @@ -124,9 +124,9 @@ class TestResolveAutoMainFirst: mock_client = MagicMock() mock_resolve.return_value = (mock_client, "anthropic/claude-opus-4.8") - from agent.auxiliary_client import _resolve_auto + from agent.auxiliary_client import _resolve_auto_route - client, model = _resolve_auto( + client, model, _provider = _resolve_auto_route( main_runtime={ "provider": "moa", "model": "opus-gpt", @@ -166,9 +166,9 @@ class TestResolveAutoMainFirst: ) as mock_main_chain, patch( "agent.auxiliary_client._try_openrouter", ) as mock_openrouter: - from agent.auxiliary_client import _resolve_auto + from agent.auxiliary_client import _resolve_auto_route - client, model = _resolve_auto(task="title_generation") + client, model, _provider = _resolve_auto_route(task="title_generation") assert client is task_client assert model == "task-free-model" @@ -221,9 +221,9 @@ class TestResolveAutoMainFirst: ) as mock_resolve: mock_resolve.return_value = (MagicMock(), "mimo-v2.5-pro") - from agent.auxiliary_client import _resolve_auto + from agent.auxiliary_client import _resolve_auto_route - _resolve_auto(main_runtime={ + _resolve_auto_route(main_runtime={ "provider": "xiaomi", "model": "mimo-v2.5-pro", "base_url": token_plan_url, @@ -614,7 +614,7 @@ def test_aggregator_providers_constant_removed(): import agent.auxiliary_client as aux_mod assert not hasattr(aux_mod, "_AGGREGATOR_PROVIDERS"), ( - "_AGGREGATOR_PROVIDERS was removed when _resolve_auto stopped " + "_AGGREGATOR_PROVIDERS was removed when _resolve_auto_route stopped " "treating aggregators specially. If you re-added it, the main-first " "policy may have regressed." ) diff --git a/tests/agent/test_bedrock_adapter.py b/tests/agent/test_bedrock_adapter.py index 1546369787..93c3b6fc89 100644 --- a/tests/agent/test_bedrock_adapter.py +++ b/tests/agent/test_bedrock_adapter.py @@ -90,7 +90,6 @@ class TestResolveAwsAuthEnvVar: """ - def test_requires_both_access_key_and_secret(self): from agent.bedrock_adapter import resolve_aws_auth_env_var # Only access key, no secret β†’ should not match @@ -98,8 +97,6 @@ class TestResolveAwsAuthEnvVar: assert resolve_aws_auth_env_var(env) != "AWS_ACCESS_KEY_ID" - - def test_returns_none_when_no_aws_auth(self): from agent.bedrock_adapter import resolve_aws_auth_env_var # Mock botocore to return no credentials (covers EC2 IMDS fallback) @@ -111,7 +108,6 @@ class TestResolveAwsAuthEnvVar: assert resolve_aws_auth_env_var({}) is None - class TestHasAwsCredentials: def test_true_with_profile(self): from agent.bedrock_adapter import has_aws_credentials @@ -143,8 +139,6 @@ class TestResolveBedrocRegion: assert resolve_bedrock_region({}) == "us-east-1" - - # --------------------------------------------------------------------------- # Tool conversion # --------------------------------------------------------------------------- @@ -183,7 +177,6 @@ class TestConvertToolsToConverse: assert convert_tools_to_converse(None) == [] - # --------------------------------------------------------------------------- # Message conversion: OpenAI β†’ Converse # --------------------------------------------------------------------------- @@ -255,9 +248,6 @@ class TestConvertMessagesToConverse: assert tr["toolResult"]["content"][0]["text"] == "file contents here" - - - def test_empty_content_gets_placeholder(self): from agent.bedrock_adapter import convert_messages_to_converse messages = [{"role": "user", "content": ""}] @@ -266,8 +256,6 @@ class TestConvertMessagesToConverse: assert msgs[0]["content"][0]["text"].strip() != "" or msgs[0]["content"][0]["text"] == " " - - # --------------------------------------------------------------------------- # Response normalization: Converse β†’ OpenAI # --------------------------------------------------------------------------- @@ -428,10 +416,6 @@ class TestNormalizeConverseResponse: assert blocks[2]["reasoningContent"]["redactedContent"] == b"r2" - - - - # --------------------------------------------------------------------------- # Streaming response normalization # --------------------------------------------------------------------------- @@ -503,8 +487,6 @@ class TestNormalizeConverseStreamEvents: assert json.loads(tc[0].function.arguments) == {"path": "/tmp/f"} - - # --------------------------------------------------------------------------- # build_converse_kwargs # --------------------------------------------------------------------------- @@ -574,27 +556,6 @@ class TestBuildConverseKwargs: ) assert "inferenceConfig" not in kwargs - def test_call_converse_stream_omits_cap_for_none(self): - """The streaming entry point funnels through the same builder β€” pin - that max_tokens=None omits the cap there too.""" - from unittest.mock import MagicMock, patch as mock_patch - from agent.bedrock_adapter import call_converse_stream - boto3_client = MagicMock() - boto3_client.converse_stream.return_value = {"stream": []} - with mock_patch( - "agent.bedrock_adapter._get_bedrock_runtime_client", - return_value=boto3_client, - ): - call_converse_stream( - region="us-east-1", - model="test-model", - messages=[{"role": "user", "content": "Hi"}], - max_tokens=None, - temperature=0.2, - ) - wire_kwargs = boto3_client.converse_stream.call_args.kwargs - assert "maxTokens" not in wire_kwargs.get("inferenceConfig", {}) - def test_cache_point_added_for_supported_model(self): """Claude and Nova on the Converse path get cachePoint markers on system, tools, and the message before the newest turn.""" @@ -813,8 +774,6 @@ class TestDiscoverBedrockModels: """Test Bedrock model discovery with mocked AWS API calls.""" - - def test_provider_filter(self): from agent.bedrock_adapter import discover_bedrock_models, reset_discovery_cache reset_discovery_cache() @@ -877,7 +836,6 @@ class TestDiscoverBedrockModels: assert first == second - def test_handles_api_error_gracefully(self): from agent.bedrock_adapter import discover_bedrock_models, reset_discovery_cache reset_discovery_cache() @@ -980,9 +938,6 @@ class TestStreamConverseWithCallbacks: assert len(result.choices[0].message.tool_calls) == 1 - - - # --------------------------------------------------------------------------- # Guardrail config in build_converse_kwargs # --------------------------------------------------------------------------- @@ -1029,35 +984,15 @@ class TestGuardrailConfig: # Error classification # --------------------------------------------------------------------------- -class TestBedrockErrorClassification: - """Test Bedrock-specific error classification.""" - - def test_context_overflow_validation_exception(self): - from agent.bedrock_adapter import classify_bedrock_error - assert classify_bedrock_error( - "ValidationException: input is too long for model" - ) == "context_overflow" - - class TestBedrockContextLength: """Test Bedrock model context length lookup.""" - - - - - - - - - def test_unknown_model_gets_default(self): from agent.bedrock_adapter import get_bedrock_context_length, BEDROCK_DEFAULT_CONTEXT_LENGTH assert get_bedrock_context_length("unknown.model-v1:0") == BEDROCK_DEFAULT_CONTEXT_LENGTH - def test_no_region_skips_probe_uses_table(self): # Default call (no region) must NOT hit the network β€” returns the # static table value. Guards backward compatibility for callers that @@ -1078,7 +1013,6 @@ class TestBedrockContextProbe: return client - def test_probe_returns_none_when_client_unavailable(self): from agent.bedrock_adapter import probe_bedrock_context_length with patch("agent.bedrock_adapter._get_bedrock_runtime_client", @@ -1097,7 +1031,6 @@ class TestBedrockContextProbe: region="eu-central-1") == 1_000_000 - # --------------------------------------------------------------------------- # Tool-calling capability detection # --------------------------------------------------------------------------- @@ -1110,17 +1043,11 @@ class TestModelSupportsToolUse: assert _model_supports_tool_use("us.anthropic.claude-sonnet-4-6") is True - - def test_deepseek_r1_no_tools(self): from agent.bedrock_adapter import _model_supports_tool_use assert _model_supports_tool_use("us.deepseek.r1-v1:0") is False - - - - class TestBuildConverseKwargsToolStripping: """Test that tools are stripped for non-tool-calling models.""" @@ -1157,22 +1084,17 @@ class TestIsAnthropicBedrockModel: assert is_anthropic_bedrock_model("us.anthropic.claude-sonnet-4-6") is True - def test_nova_is_not_anthropic(self): from agent.bedrock_adapter import is_anthropic_bedrock_model assert is_anthropic_bedrock_model("us.amazon.nova-pro-v1:0") is False - - - def test_au_inference_profile(self): from agent.bedrock_adapter import is_anthropic_bedrock_model assert is_anthropic_bedrock_model("au.anthropic.claude-haiku-4-5-20251001-v1:0") is True assert is_anthropic_bedrock_model("au.anthropic.claude-sonnet-4-6") is True - class TestEmptyTextBlockFix: """Test that empty/whitespace-only text blocks are replaced with a non-whitespace placeholder (not a literal space, which is itself @@ -1185,15 +1107,12 @@ class TestEmptyTextBlockFix: assert blocks[0]["text"].strip() - def test_real_text_preserved(self): from agent.bedrock_adapter import _convert_content_to_converse blocks = _convert_content_to_converse("Hello") assert blocks[0]["text"] == "Hello" - - # --------------------------------------------------------------------------- # Stale-connection detection and per-region client invalidation # --------------------------------------------------------------------------- @@ -1227,7 +1146,6 @@ class TestIsStaleConnectionError: """Classifier that decides whether an exception warrants client eviction.""" - def test_detects_botocore_read_timeout(self): pytest.importorskip("botocore.exceptions", reason="botocore (with working exceptions module) required") from agent.bedrock_adapter import is_stale_connection_error @@ -1254,7 +1172,6 @@ class TestIsStaleConnectionError: pytest.fail("AssertionError not raised") - def test_ignores_unrelated_exceptions(self): from agent.bedrock_adapter import is_stale_connection_error assert is_stale_connection_error(ValueError("bad input")) is False @@ -1262,35 +1179,9 @@ class TestIsStaleConnectionError: class TestCallConverseInvalidatesOnStaleError: - """call_converse / call_converse_stream evict the cached client when the - boto3 call raises a stale-connection error β€” so the next invocation - reconnects instead of reusing the dead socket.""" - - - def test_converse_stream_evicts_client_on_stale_error(self): - pytest.importorskip("botocore.exceptions", reason="botocore (with working exceptions module) required") - from agent.bedrock_adapter import ( - _bedrock_runtime_client_cache, - call_converse_stream, - reset_client_cache, - ) - from botocore.exceptions import ConnectionClosedError - - reset_client_cache() - dead_client = MagicMock() - dead_client.converse_stream.side_effect = ConnectionClosedError( - endpoint_url="https://bedrock.example", - ) - _bedrock_runtime_client_cache["us-east-1"] = dead_client - - with pytest.raises(ConnectionClosedError): - call_converse_stream( - region="us-east-1", - model="anthropic.claude-3-sonnet-20240229-v1:0", - messages=[{"role": "user", "content": "hi"}], - ) - - assert "us-east-1" not in _bedrock_runtime_client_cache + """call_converse evicts the cached client only on a stale-connection error β€” so the + next invocation reconnects instead of reusing the dead socket (the agent's streaming + path is pinned in ``TestAgentBedrockStreamRecovery``).""" def test_converse_does_not_evict_on_non_stale_error(self): """Non-stale errors (e.g. ValidationException) leave the client cache alone.""" @@ -1322,7 +1213,6 @@ class TestCallConverseInvalidatesOnStaleError: ) - class TestStreamingAccessDeniedDetection: """is_streaming_access_denied_error() recognizes IAM denials of bedrock:InvokeModelWithResponseStream (InvokeModel-only policies).""" @@ -1351,8 +1241,6 @@ class TestStreamingAccessDeniedDetection: assert is_streaming_access_denied_error(self._denied_client_error()) is True - - def test_ignores_unrelated_errors(self): from agent.bedrock_adapter import is_streaming_access_denied_error assert is_streaming_access_denied_error(ValueError("boom")) is False @@ -1361,55 +1249,9 @@ class TestStreamingAccessDeniedDetection: ) is False -class TestCallConverseStreamIamFallback: - """call_converse_stream() falls back to converse() when IAM denies the - streaming action β€” InvokeModel-only policies keep working.""" - - def test_falls_back_to_converse_on_streaming_denial(self): - pytest.importorskip("botocore.exceptions", reason="botocore (with working exceptions module) required") - from agent.bedrock_adapter import ( - _bedrock_runtime_client_cache, - call_converse_stream, - reset_client_cache, - ) - from botocore.exceptions import ClientError - - reset_client_cache() - client = MagicMock() - client.converse_stream.side_effect = ClientError( - error_response={ - "Error": { - "Code": "AccessDeniedException", - "Message": ( - "User is not authorized to perform: " - "bedrock:InvokeModelWithResponseStream" - ), - } - }, - operation_name="ConverseStream", - ) - client.converse.return_value = { - "output": {"message": {"role": "assistant", "content": [{"text": "hi"}]}}, - "stopReason": "end_turn", - "usage": {"inputTokens": 1, "outputTokens": 1, "totalTokens": 2}, - } - _bedrock_runtime_client_cache["us-east-1"] = client - - result = call_converse_stream( - region="us-east-1", - model="anthropic.claude-3-sonnet-20240229-v1:0", - messages=[{"role": "user", "content": "hi"}], - ) - - client.converse.assert_called_once() - assert result.choices[0].message.content == "hi" - # Not a stale connection β€” client stays cached. - assert _bedrock_runtime_client_cache.get("us-east-1") is client - - class TestAgentBedrockStreamRecovery: - """The agent loop streams through ``chat_completion_helpers._bedrock_converse_call`` - (not ``call_converse_stream``); pin the same recovery ladder on that live path: + """The agent loop streams through ``chat_completion_helpers._bedrock_converse_call``; + pin the recovery ladder on that live path: IAM streaming denial β†’ ``_BedrockStream._fall_back_to_converse`` (non-streaming converse, streaming disabled for the session, client kept), stale connection β†’ cached client evicted so the outer retry reconnects.""" @@ -1498,7 +1340,6 @@ class TestRequireBoto3VersionCheck: assert result is fake_boto3 - class TestImageBase64Decoding: """Image data URLs must be decoded to raw bytes before passing to Converse API. diff --git a/tests/agent/test_compression_progress.py b/tests/agent/test_compression_progress.py index 85799318f7..0dc8e09e9c 100644 --- a/tests/agent/test_compression_progress.py +++ b/tests/agent/test_compression_progress.py @@ -7,7 +7,7 @@ estimated request token count without removing any rows β€” and surfaces a spurious ``Context length exceeded`` failure followed by an auto-reset of an otherwise healthy session. -These tests pin the contract of ``_compression_made_progress``: a +These tests pin the contract of ``compression_made_progress``: a row-count reduction OR a *material* (>5%) token-count reduction counts as progress. """ @@ -15,7 +15,7 @@ progress. from __future__ import annotations from agent.turn_context import ( - _compression_made_progress, + compression_made_progress, _compression_warrants_another_preflight_pass, ) @@ -23,7 +23,7 @@ from agent.turn_context import ( class TestCompressionMadeProgress: def test_rows_reduced_counts_as_progress(self): """Removing message rows is the obvious progress signal.""" - assert _compression_made_progress( + assert compression_made_progress( orig_len=10, new_len=5, orig_tokens=1000, new_tokens=1000 ) is True @@ -31,7 +31,7 @@ class TestCompressionMadeProgress: def test_neither_moved_means_no_progress(self): """The genuine "stuck" case β€” same rows, same tokens, give up.""" - assert _compression_made_progress( + assert compression_made_progress( orig_len=10, new_len=10, orig_tokens=1000, new_tokens=1000 ) is False @@ -43,11 +43,11 @@ class TestCompressionMadeProgress: progress β€” matching the overflow-handler retry path (#39550) so a marginal wobble can't keep the multi-pass loop spinning.""" # 1000 -> 970 is a 3% drop, below the 5% floor. - assert _compression_made_progress( + assert compression_made_progress( orig_len=10, new_len=10, orig_tokens=1000, new_tokens=970 ) is False # 1000 -> 940 is a 6% drop, above the floor. - assert _compression_made_progress( + assert compression_made_progress( orig_len=10, new_len=10, orig_tokens=1000, new_tokens=940 ) is True diff --git a/tests/agent/test_fast_compression_lane.py b/tests/agent/test_fast_compression_lane.py index 920e248053..e3491a1173 100644 --- a/tests/agent/test_fast_compression_lane.py +++ b/tests/agent/test_fast_compression_lane.py @@ -412,7 +412,7 @@ def test_timing_hooks_propagate_to_protected_call_worker_thread(): source both active). """ from agent.auxiliary_client import ( - _aux_timing_hook, + _aux_thread_local_hook, _aux_dispatch, _aux_provider_response, _notify_aux_dispatch, @@ -431,8 +431,8 @@ def test_timing_hooks_propagate_to_protected_call_worker_thread(): return "ok" with ( - _aux_timing_hook(_aux_dispatch, lambda: seen.append("dispatch")), - _aux_timing_hook(_aux_provider_response, lambda: seen.append("response")), + _aux_thread_local_hook(_aux_dispatch, lambda: seen.append("dispatch")), + _aux_thread_local_hook(_aux_provider_response, lambda: seen.append("response")), aux_interrupt_protection(cancel_check=lambda: False), ): result = _run_protected_sync_provider_call(_callback, {}) diff --git a/tests/agent/test_message_sanitization_policy.py b/tests/agent/test_message_sanitization_policy.py index 683d2d8df0..74db041d98 100644 --- a/tests/agent/test_message_sanitization_policy.py +++ b/tests/agent/test_message_sanitization_policy.py @@ -52,11 +52,6 @@ class TestDeterministicCallId: out = deterministic_call_id("t", "bad \ud800 arg", 0) assert out.startswith("call_") - def test_codex_adapter_wrapper_delegates(self): - from agent.codex_responses_adapter import _deterministic_call_id - assert _deterministic_call_id("terminal", '{"command":"ls"}', 0) == \ - deterministic_call_id("terminal", '{"command":"ls"}', 0) - def test_run_agent_static_delegates(self): from run_agent import AIAgent assert AIAgent._deterministic_call_id("terminal", '{"command":"ls"}', 0) == \ diff --git a/tests/agent/test_protected_tail_pressure_61932.py b/tests/agent/test_protected_tail_pressure_61932.py index 339b45b529..7515849e7a 100644 --- a/tests/agent/test_protected_tail_pressure_61932.py +++ b/tests/agent/test_protected_tail_pressure_61932.py @@ -23,7 +23,7 @@ from agent.context_compressor import ( _PRESSURE_KEEP_RECENT_MESSAGES, ) from agent.model_metadata import estimate_messages_tokens_rough -from agent.turn_context import _compression_made_progress +from agent.turn_context import compression_made_progress def _unique_tool_pair(i: int, chars: int) -> list[dict]: @@ -129,7 +129,7 @@ class TestProtectedTailPressure61932: o_len, o_tok = len(cur), tok out = c.compress(list(cur), current_tokens=tok) n_tok = estimate_messages_tokens_rough(out) - last_progress = _compression_made_progress( + last_progress = compression_made_progress( o_len, len(out), o_tok, n_tok ) cur, tok = out, n_tok diff --git a/tests/agent/test_set_runtime_main_custom_provider.py b/tests/agent/test_set_runtime_main_custom_provider.py index d84dfc5d62..032e33dad0 100644 --- a/tests/agent/test_set_runtime_main_custom_provider.py +++ b/tests/agent/test_set_runtime_main_custom_provider.py @@ -1,5 +1,5 @@ """Regression test: set_runtime_main() must pass base_url/api_key/api_mode -so that _resolve_auto() can route custom: providers in Step 1. +so that _resolve_auto_route() can route custom: providers in Step 1. Fixes https://github.com/NousResearch/hermes-agent/issues/34777 """ @@ -38,7 +38,7 @@ class TestSetRuntimeMainCustomProvider: assert v == "", f"Expected empty, got {v!r}" def test_resolve_auto_uses_globals_for_custom_provider(self): - """_resolve_auto reads base_url/api_key from globals when main_runtime is None.""" + """_resolve_auto_route reads base_url/api_key from globals when main_runtime is None.""" import agent.auxiliary_client as mod mod.clear_runtime_main() @@ -52,7 +52,7 @@ class TestSetRuntimeMainCustomProvider: with patch.object(mod, "resolve_provider_client") as mock_resolve: mock_resolve.return_value = (MagicMock(), "test-model") - client, resolved = mod._resolve_auto(main_runtime=None) + client, resolved, _provider = mod._resolve_auto_route(main_runtime=None) mock_resolve.assert_called_once() call_args = mock_resolve.call_args @@ -218,12 +218,12 @@ class TestResolveAutoCustomEndToEnd: # The original /anthropic URL must survive β€” no /v1 rewrite. assert getattr(client, "base_url", "").rstrip("/") == proxy_base - # Wiring check: _resolve_auto must hand the FULL custom: + # Wiring check: _resolve_auto_route must hand the FULL custom: # string to resolve_provider_client, with no explicit_base_url # override (the named arm reads base_url/api_key from config). with patch.object(mod, "resolve_provider_client") as mock_resolve: mock_resolve.return_value = (MagicMock(), "claude-4-6-opus") - mod._resolve_auto(main_runtime=None) + mod._resolve_auto_route(main_runtime=None) mock_resolve.assert_called_once() assert mock_resolve.call_args.args[0] == "custom:palantir" assert mock_resolve.call_args.kwargs["explicit_base_url"] is None diff --git a/tests/agent/test_ssl_ca_guard.py b/tests/agent/test_ssl_ca_guard.py index a381352a45..065edd2c46 100644 --- a/tests/agent/test_ssl_ca_guard.py +++ b/tests/agent/test_ssl_ca_guard.py @@ -6,7 +6,7 @@ import certifi import pytest from agent.errors import SSLConfigurationError -from agent.ssl_guard import verify_ca_bundle, verify_ca_bundle_with_fallback +from agent.ssl_guard import verify_ca_bundle def test_healthy_bundle_passes(monkeypatch): @@ -71,4 +71,3 @@ def test_truststore_get_ca_certs_not_implemented_is_accepted(monkeypatch, tmp_pa # Must not raise on the explicit env bundle nor the certifi check. verify_ca_bundle() - verify_ca_bundle_with_fallback() diff --git a/tests/agent/test_unsupported_parameter_retry.py b/tests/agent/test_unsupported_parameter_retry.py index 160d545b4d..b14b2254f8 100644 --- a/tests/agent/test_unsupported_parameter_retry.py +++ b/tests/agent/test_unsupported_parameter_retry.py @@ -9,7 +9,6 @@ pattern. These tests lock in: * ``_is_unsupported_parameter_error(exc, param)`` across common phrasings - * the back-compat wrapper ``_is_unsupported_temperature_error`` still works * the max_tokens retry branch no longer pops a key that was never set (``max_tokens is None`` gate) * the max_tokens retry branch matches via the generic helper on top of the @@ -24,7 +23,6 @@ from agent.auxiliary_client import ( call_llm, async_call_llm, _is_unsupported_parameter_error, - _is_unsupported_temperature_error, ) @@ -49,13 +47,12 @@ class TestIsUnsupportedParameterError: - def test_temperature_wrapper_delegates_to_generic(self): - """Back-compat: ``_is_unsupported_temperature_error`` still routes through.""" + def test_temperature_param_routes_through_generic(self): msg = "HTTP 400: Unsupported parameter: temperature" - assert _is_unsupported_temperature_error(RuntimeError(msg)) is True + assert _is_unsupported_parameter_error(RuntimeError(msg), "temperature") is True # And the unrelated-case still holds - assert _is_unsupported_temperature_error( - RuntimeError("max_tokens is too large")) is False + assert _is_unsupported_parameter_error( + RuntimeError("max_tokens is too large"), "temperature") is False def _dummy_response(): diff --git a/tests/agent/test_unsupported_temperature_retry.py b/tests/agent/test_unsupported_temperature_retry.py index fd9c72ef97..e60b835d71 100644 --- a/tests/agent/test_unsupported_temperature_retry.py +++ b/tests/agent/test_unsupported_temperature_retry.py @@ -30,7 +30,7 @@ import pytest from agent.auxiliary_client import ( call_llm, async_call_llm, - _is_unsupported_temperature_error, + _is_unsupported_parameter_error, ) @@ -52,7 +52,7 @@ class TestIsUnsupportedTemperatureError: "unrecognized request argument supplied: temperature", ]) def test_matches_real_provider_messages(self, message): - assert _is_unsupported_temperature_error(RuntimeError(message)) is True + assert _is_unsupported_parameter_error(RuntimeError(message), "temperature") is True @pytest.mark.parametrize("message", [ # Unrelated 400s must NOT trigger a silent-retry @@ -64,7 +64,7 @@ class TestIsUnsupportedTemperatureError: "temperature must be between 0 and 2", ]) def test_does_not_match_unrelated_errors(self, message): - assert _is_unsupported_temperature_error(RuntimeError(message)) is False + assert _is_unsupported_parameter_error(RuntimeError(message), "temperature") is False def _dummy_response(): diff --git a/tests/hermes_cli/test_certifi_repair.py b/tests/hermes_cli/test_certifi_repair.py index 6ca002a926..442b9b2a8d 100644 --- a/tests/hermes_cli/test_certifi_repair.py +++ b/tests/hermes_cli/test_certifi_repair.py @@ -167,7 +167,7 @@ class TestDoctorCertificates: return _R() monkeypatch.setattr( - "agent.ssl_guard.verify_ca_bundle_with_fallback", fake_verify + "agent.ssl_guard.verify_ca_bundle", fake_verify ) monkeypatch.setattr(subprocess, "run", fake_run) diff --git a/tests/monitoring/test_gateway_health_export.py b/tests/monitoring/test_gateway_health_export.py index 8c8190f7f1..5e24671e5d 100644 --- a/tests/monitoring/test_gateway_health_export.py +++ b/tests/monitoring/test_gateway_health_export.py @@ -41,7 +41,7 @@ def test_otlp_attrs_redact_strings_and_never_export_profile(): def test_resource_attributes_are_allowlisted_and_sanitized(): - from agent.monitoring.gateway_health_export import _safe_resource_attributes + from agent.monitoring.otlp_exporter import _safe_resource_attributes attrs = _safe_resource_attributes({ "service.name": "hermes-gateway", diff --git a/tests/plugins/browser/check_parity_vs_main.py b/tests/plugins/browser/check_parity_vs_main.py index b706ce3e9c..e0ee6fcffd 100644 --- a/tests/plugins/browser/check_parity_vs_main.py +++ b/tests/plugins/browser/check_parity_vs_main.py @@ -78,20 +78,20 @@ from tools.browser_tool import _get_cloud_provider, _is_local_mode provider = _get_cloud_provider() -# Pull the human-readable backend name via the API that exists on BOTH -# legacy (origin/main: CloudBrowserProvider.provider_name()) and the new -# ABC (BrowserProvider exposes provider_name() as a backward-compat alias -# returning display_name). Both shapes resolve to the same string β€” -# 'Browserbase' / 'Browser Use' / 'Firecrawl' β€” so we can compare safely. +# Pull the human-readable backend name: legacy (origin/main: CloudBrowserProvider.provider_name()) +# or the new ABC (BrowserProvider.display_name / is_available()). Both shapes resolve to the same +# string β€” 'Browserbase' / 'Browser Use' / 'Firecrawl' β€” so we can compare safely. provider_name = None is_available = None if provider is not None: - pn = getattr(provider, "provider_name", None) + pn = getattr(provider, "display_name", None) + if pn is None: + pn = getattr(provider, "provider_name", None) if callable(pn): provider_name = pn() elif isinstance(pn, str): provider_name = pn - is_conf = getattr(provider, "is_configured", None) + is_conf = getattr(provider, "is_available", None) or getattr(provider, "is_configured", None) if callable(is_conf): is_available = bool(is_conf()) diff --git a/tests/plugins/browser/test_browser_provider_plugins.py b/tests/plugins/browser/test_browser_provider_plugins.py index f29645fc07..7bbfafe755 100644 --- a/tests/plugins/browser/test_browser_provider_plugins.py +++ b/tests/plugins/browser/test_browser_provider_plugins.py @@ -239,45 +239,6 @@ class TestRegistryResolution: assert provider.name == "browser-use" -# --------------------------------------------------------------------------- -# Legacy ABC backward-compat aliases (is_configured / provider_name) -# --------------------------------------------------------------------------- - - -class TestLegacyAbcAliases: - """is_configured() and provider_name() delegate to the new API.""" - - @pytest.mark.parametrize( - "plugin_name", - ["browserbase", "browser-use", "firecrawl"], - ) - def test_is_configured_delegates_to_is_available(self, plugin_name: str) -> None: - _ensure_plugins_loaded() - from agent.browser_registry import get_provider - - p = get_provider(plugin_name) - assert p is not None - assert p.is_configured() is p.is_available() - - @pytest.mark.parametrize( - "plugin_name,expected_label", - [ - ("browserbase", "Browserbase"), - ("browser-use", "Browser Use"), - ("firecrawl", "Firecrawl"), - ], - ) - def test_provider_name_returns_display_name( - self, plugin_name: str, expected_label: str - ) -> None: - _ensure_plugins_loaded() - from agent.browser_registry import get_provider - - p = get_provider(plugin_name) - assert p is not None - assert p.provider_name() == expected_label - - # --------------------------------------------------------------------------- # Picker integration # --------------------------------------------------------------------------- diff --git a/tests/run_agent/test_codex_xai_oauth_recovery.py b/tests/run_agent/test_codex_xai_oauth_recovery.py index 57dc20efa9..ecbe658eeb 100644 --- a/tests/run_agent/test_codex_xai_oauth_recovery.py +++ b/tests/run_agent/test_codex_xai_oauth_recovery.py @@ -232,7 +232,7 @@ def test_summarize_api_error_does_not_accuse_subscribers(): # --------------------------------------------------------------------------- # Fix D: _StreamErrorEvent xAI entitlement classified as auth, not retryable # -# run_codex_create_stream_fallback raises _StreamErrorEvent (status_code=None) +# run_codex_stream raises _StreamErrorEvent (status_code=None) # when the Responses stream emits a ``type=error`` SSE frame. Before this # fix, classify_api_error had no match for "grok subscription" in its pattern # lists, so it returned FailoverReason.unknown (retryable=True) β€” burning diff --git a/tests/run_agent/test_run_agent.py b/tests/run_agent/test_run_agent.py index 90cf8d6798..0da74057bd 100644 --- a/tests/run_agent/test_run_agent.py +++ b/tests/run_agent/test_run_agent.py @@ -1780,7 +1780,6 @@ class TestExecuteToolCalls: patch("model_tools.handle_function_call", side_effect=KeyboardInterrupt), patch("run_agent._set_interrupt"), patch("agent.interrupt_control._set_interrupt"), - patch("agent.turn_facade._set_interrupt"), pytest.raises(KeyboardInterrupt), ): agent._execute_tool_calls_sequential(mock_msg, messages, "task-1") @@ -3411,7 +3410,6 @@ class TestRunConversation: patch.object(agent, "_cleanup_task_resources"), patch("run_agent._set_interrupt"), patch("agent.interrupt_control._set_interrupt"), - patch("agent.turn_facade._set_interrupt"), patch.object( agent, "_interruptible_api_call", side_effect=interrupt_side_effect ), diff --git a/tests/run_agent/test_streaming.py b/tests/run_agent/test_streaming.py index 040f374e40..654db7e164 100644 --- a/tests/run_agent/test_streaming.py +++ b/tests/run_agent/test_streaming.py @@ -962,7 +962,7 @@ class TestCodexStreamCallbacks: mock_client = MagicMock() mock_client.responses.create.return_value = mock_stream - agent._run_codex_create_stream_fallback( + agent._run_codex_stream( {"model": "test/model", "instructions": "hi", "input": []}, client=mock_client, ) diff --git a/tests/run_agent/test_switch_model_context.py b/tests/run_agent/test_switch_model_context.py index e768a63868..79e7b0812f 100644 --- a/tests/run_agent/test_switch_model_context.py +++ b/tests/run_agent/test_switch_model_context.py @@ -6,7 +6,7 @@ import pytest from hermes_cli.models import LMStudioLoadResult from run_agent import AIAgent -from agent.agent_init import _normalize_route_base_url +from hermes_cli.route_identity import normalize_route_base_url from agent.context_compressor import ContextCompressor @@ -26,9 +26,9 @@ class _StubStartupCompressor: def test_route_url_normalization_preserves_path_slash_before_query(): """A path slash before a query changes OpenAI SDK URL joining.""" - assert _normalize_route_base_url( + assert normalize_route_base_url( "https://example.com/v1/?tenant=large" - ) != _normalize_route_base_url("https://example.com/v1?tenant=large") + ) != normalize_route_base_url("https://example.com/v1?tenant=large") diff --git a/tests/tools/test_browser_cloud_provider_cache.py b/tests/tools/test_browser_cloud_provider_cache.py index d7b9671000..0df17f3372 100644 --- a/tests/tools/test_browser_cloud_provider_cache.py +++ b/tests/tools/test_browser_cloud_provider_cache.py @@ -224,9 +224,9 @@ class TestCloudProviderCachePolicy: ) bu_unconfigured = Mock() - bu_unconfigured.is_configured.return_value = False + bu_unconfigured.is_available.return_value = False bb_unconfigured = Mock() - bb_unconfigured.is_configured.return_value = False + bb_unconfigured.is_available.return_value = False monkeypatch.setattr( browser_tool, "BrowserUseProvider", lambda: bu_unconfigured ) @@ -239,7 +239,7 @@ class TestCloudProviderCachePolicy: # Credentials self-heal β€” next call must retry and pick up the provider. healed = Mock(name="healed-provider") - healed.is_configured.return_value = True + healed.is_available.return_value = True monkeypatch.setattr(browser_tool, "BrowserUseProvider", lambda: healed) assert browser_tool._get_cloud_provider() is healed diff --git a/tests/tools/test_browser_lightpanda.py b/tests/tools/test_browser_lightpanda.py index d27b6f95b6..48d774b5b4 100644 --- a/tests/tools/test_browser_lightpanda.py +++ b/tests/tools/test_browser_lightpanda.py @@ -556,7 +556,7 @@ class TestLightpandaEngineStatus: def test_shadowed_by_cloud_provider(self, monkeypatch): provider = MagicMock() - provider.provider_name.return_value = "Browserbase" + provider.display_name = "Browserbase" bt = self._gates(monkeypatch, _get_cloud_provider=lambda: provider) used, reason = bt.lightpanda_engine_status() assert used is False and "Browserbase" in reason @@ -578,7 +578,7 @@ class TestLightpandaEngineStatus: """browser_exec resolves real-profile before the backend, so with both set the real-profile toggle is the actual shadow.""" provider = MagicMock() - provider.provider_name.return_value = "Browserbase" + provider.display_name = "Browserbase" bt = self._gates( monkeypatch, _use_real_profile=lambda: True, diff --git a/tools/browser_tool_cloud.py b/tools/browser_tool_cloud.py index c5946a4288..8b92b0572e 100644 --- a/tools/browser_tool_cloud.py +++ b/tools/browser_tool_cloud.py @@ -110,7 +110,7 @@ def _autodetect_cloud_provider() -> Optional[CloudBrowserProvider]: try: for cls in (_bt.BrowserUseProvider, _bt.BrowserbaseProvider): fallback_provider = cls() - if fallback_provider.is_configured(): + if fallback_provider.is_available(): return fallback_provider except Exception: # pragma: no cover - defensive: never poison cache _bt.logger.debug("Cloud provider auto-detect failed", exc_info=True) diff --git a/tools/browser_tool_install.py b/tools/browser_tool_install.py index 6661f45466..dfceddd629 100644 --- a/tools/browser_tool_install.py +++ b/tools/browser_tool_install.py @@ -315,7 +315,7 @@ def check_browser_requirements() -> bool: # Cloud mode also requires provider credentials; no local Chromium needed. provider = _bt._get_cloud_provider() if provider is not None: - return provider.is_configured() + return provider.is_available() # Lightpanda provides text/navigation tools without Chromium; screenshots/vision still return install errors. if _bt._using_lightpanda_engine(): return True diff --git a/tools/browser_tool_lightpanda_fallback.py b/tools/browser_tool_lightpanda_fallback.py index 55fb80c7a9..82f0c5817c 100644 --- a/tools/browser_tool_lightpanda_fallback.py +++ b/tools/browser_tool_lightpanda_fallback.py @@ -47,7 +47,7 @@ def lightpanda_engine_status() -> Tuple[bool, str]: provider = None if provider is not None: try: - name = provider.provider_name() + name = provider.display_name except Exception: name = type(provider).__name__ return False, f"cloud provider {name} is selected (browser.cloud_provider, or auto-detected from credentials)" diff --git a/tools/browser_tool_vision.py b/tools/browser_tool_vision.py index a89736d276..8c01a2552b 100644 --- a/tools/browser_tool_vision.py +++ b/tools/browser_tool_vision.py @@ -16,7 +16,7 @@ from tools.browser_tool_origin import origin as _bt def _vision_mode_label() -> str: _cp = _bt._get_cloud_provider() - return "local" if _cp is None else f"cloud ({_cp.provider_name()})" + return "local" if _cp is None else f"cloud ({_cp.display_name})" def _lightpanda_vision_preroute( diff --git a/tui_gateway/mcp_oauth_sessions.py b/tui_gateway/mcp_oauth_sessions.py index 650575fd0d..69384d28b0 100644 --- a/tui_gateway/mcp_oauth_sessions.py +++ b/tui_gateway/mcp_oauth_sessions.py @@ -104,7 +104,7 @@ def _probe_with_rollback( flow.tools = [{"name": t, "description": d} for t, d in tools] flow.mark_approved() if reconnect_live: - from tools.mcp_tool import reconnect_mcp_server + from tools.mcp_tool_loop import reconnect_mcp_server reconnect_mcp_server(server_name) except Exception: storage.restore(backup, only_if_absent=True) diff --git a/tui_gateway/server.py b/tui_gateway/server.py index 01fee96b9c..22a29eaeda 100644 --- a/tui_gateway/server.py +++ b/tui_gateway/server.py @@ -2089,7 +2089,7 @@ def _session_info(agent, session: dict | None = None) -> dict: info["skills"] = get_available_skills() info["mcp_servers"] = [] with contextlib.suppress(Exception): - from tools.mcp_tool import get_mcp_status + from tools.mcp_tool_discovery import get_mcp_status info["mcp_servers"] = get_mcp_status() with contextlib.suppress(Exception): info["system_prompt"] = ( @@ -2153,7 +2153,7 @@ def _schedule_mcp_late_refresh(sid: str, agent) -> None: if int(getattr(agent, "_user_turn_count", 0) or 0) > 0 or int(getattr(agent, "_api_call_count", 0) or 0) > 0: return # conversation started: a rebuild would invalidate the cached prompt prefix try: - from tools.mcp_tool import refresh_agent_mcp_tools + from tools.mcp_tool_agent import refresh_agent_mcp_tools added = refresh_agent_mcp_tools(agent, quiet_mode=True) except Exception as exc: logger.warning("Late MCP refresh: tool snapshot rebuild failed for %s: %s", sid, exc)