"""Transport bring-up for MCPServerTask: stdio spawn (OSV preflight, cached-npx swap, child PID
ledger + death-supervisor registration), Streamable HTTP / SSE connect (preflight, identity header, client certs, OAuth),
protocol negotiation and initial tool discovery. Split from tools/mcp_tool.py."""

import logging
import asyncio
import os
import urllib.parse
import urllib.request
from contextlib import asynccontextmanager
from typing import Dict, Optional, Set
from utils import normalize_proxy_url
from agent.proxy_bypass import is_loopback_host, should_bypass_proxy
from agent import runtime_cwd as _runtime_cwd
from tools.mcp_tool_errors import NonMcpEndpointError, _apply_identity_header, _describe_http_failure, _handshake_answered_with_unsupported_version, _handshake_rejected_as_modern, _is_streamable_http_rejection, _make_http_rejection_recorder, _make_mcp_body_cap_transport, _make_redirect_header_stripper, _resolve_client_cert, _unwrap_exception_group
from tools.mcp_tool_lifecycle import _filter_mcp_children, _orphan_stdio_pid_servers, _orphan_stdio_pids, _stdio_pgids, _stdio_pids
from tools.mcp_tool_common import _core
from tools.mcp_tool_node_abi import node_abi_error
from tools import mcp_tool_config as _config
from tools import mcp_tool_lifecycle as _lifecycle
from tools import mcp_tool_registration as _registration

logger = logging.getLogger("tools.mcp_tool")

_PROBE_INITIALIZE_BODY = (  # JSON-RPC ``initialize`` body for the content-type preflight POST
    '{"jsonrpc":"2.0","id":"_probe","method":"initialize","params":{"protocolVersion":"2025-03-26",'
    '"capabilities":{},"clientInfo":{"name":"hermes-probe","version":"0.1"}}}')


def _content_type_base(resp) -> str:
    """``content-type`` header of *resp* without parameters, lowercased."""
    return resp.headers.get("content-type", "").split(";")[0].strip().lower()


def _is_2xx(resp) -> bool:
    return 200 <= resp.status_code < 300


def _present(**kwargs) -> dict:
    """*kwargs* minus the ``None`` values (optional httpx client arguments)."""
    return {k: v for k, v in kwargs.items() if v is not None}


def _mcp_proxy_mounts(httpx_mod, url: str, ssl_verify, client_cert, server_name: str = "") -> Optional[dict]:
    """Proxy transports for the caller-owned MCP HTTP client, or ``None`` for a direct connect.

    httpx auto-detects proxies only when ``transport is None``
    (``allow_env_proxies = trust_env and transport is None``). The wire-body cap is exactly that
    custom transport, so HTTP_PROXY / HTTPS_PROXY and the OS (Windows-registry / macOS) proxy were
    silently ignored for every HTTP/SSE MCP server: on a network that reaches the MCP host only
    through a proxy, the connect failed with ``All connection attempts failed`` and the server was
    parked. Rebuild httpx's own behaviour as explicit ``mounts`` — same source order (environment
    first, then the OS proxy), ``NO_PROXY`` / platform bypass list respected, ``socks://``
    normalized, and TLS settings identical to the transport they accompany.

    NO_PROXY goes through ``agent.proxy_bypass.should_bypass_proxy`` — the one matcher the LLM
    transport and the gateway adapters use (CIDR ranges and ``*.host`` forms the stdlib check
    does not understand) — plus ``urllib.request.proxy_bypass`` for the OS bypass list
    (Windows ``ProxyOverride`` / macOS exceptions). Loopback is never dialed through a proxy
    (``agent.proxy_bypass.is_loopback_host``), NO_PROXY or not.

    A mount wins over ``transport=`` for the URLs it matches, so each proxy transport is wrapped in
    the same wire-body cap as the direct one. A proxy the installed httpx cannot build (e.g.
    ``socks://`` without socksio) raises here and surfaces as this server's connect error.
    """
    host = urllib.parse.urlsplit(url).hostname or ""
    if not host or is_loopback_host(host) or should_bypass_proxy(url) or urllib.request.proxy_bypass(host):
        return None
    proxies = urllib.request.getproxies()
    mounts: dict = {}
    for scheme in ("http", "https"):
        proxy_url = normalize_proxy_url(proxies.get(scheme) or proxies.get("all"))
        if not proxy_url:
            continue
        # verify/cert apply to the CONNECT+TLS leg, so the proxy transport needs its own copy.
        mounts[f"{scheme}://"] = _make_mcp_body_cap_transport(httpx_mod, httpx_mod.AsyncHTTPTransport(
            proxy=proxy_url, verify=ssl_verify, **_present(cert=client_cert)))
    return mounts or None


def _pgroup_alive(pgid: Optional[int]) -> bool:
    """Signal 0 to the group succeeds iff any member is alive (POSIX only)."""
    try:
        os.killpg(pgid, 0)  # windows-footgun: ok — guarded by AttributeError below
        return True
    except (AttributeError, TypeError, OSError):  # non-POSIX / pgid None / gone
        return False


class LiveEndpointUnavailable(ConnectionError):
    """A declared runtime file did not provide a usable live endpoint."""


def _live_endpoint(server_name: str) -> Optional[tuple[str, dict]]:
    from agent.redact import register_vault_redaction_value
    from hermes_platform import declaration
    from hermes_platform.host import facts
    from hermes_platform.resolver.app import AppResolver
    from tools.mcp_liveness import liveness_for

    live = liveness_for(server_name)
    if live.kind != "server_json":
        return None
    decl = declaration.lookup(server_name)
    definition = decl.app_for(facts.os_family()) if decl is not None else None
    endpoint = AppResolver(live.app_definition(definition)).endpoint() if definition is not None else None
    if endpoint is None:
        raise LiveEndpointUnavailable(f"MCP server '{server_name}' has no usable live endpoint")
    if endpoint.token:
        register_vault_redaction_value(endpoint.token)
    headers = {"Authorization": f"Bearer {endpoint.token}"} if endpoint.token else {}
    return endpoint.url, headers


def _http_endpoint(server_name: str, config: dict) -> tuple[str, dict]:
    """URL + configured headers, overlaid by a ``server_json`` live endpoint's URL and bearer."""
    url, headers = config["url"], dict(config.get("headers") or {})
    live = _live_endpoint(server_name)
    if live is not None:
        url, live_headers = live
        headers.update(live_headers)
    return url, headers


def _stdio_launch(config: dict) -> tuple:
    """(command, env, cwd) a stdio child is spawned with, resolved in the current profile's scope."""
    command, env = _config._resolve_stdio_command(config["command"], _config._build_safe_env(config.get("env")))
    # A stdio child inherits this process's cwd when none is configured. Hosted sessions (ACP,
    # gateway) pin a logical cwd via agent.runtime_cwd; without it the child resolves relative
    # paths against the daemon's launch dir, not the session workspace. Explicit config always
    # wins; an existing session/TERMINAL_CWD anchor becomes the default; else native (None).
    # The resolved default is fixed for the connection's lifetime.
    cwd = config.get("cwd")
    if cwd is None:
        cwd = _runtime_cwd.resolve_context_cwd() or None
    return command, env, cwd


def _connect_inputs(server_name: str, config: dict) -> tuple[list, set]:
    """What a connection is opened with, resolved in the current profile's scope: stdio
    ``[command, env, cwd]``; HTTP ``[url, headers]`` after the live endpoint and the identity
    header. The owner hashes this very list and an adopter recomputes it, so both digests are
    built from one code path. Also returns the configured header names (HTTP only), captured
    before the identity header is merged in, for the strict-redirect boundary."""
    if "url" in config:
        url, headers = _http_endpoint(server_name, config)
        configured_header_names = {key.lower() for key in headers}
        return [url, _apply_identity_header(server_name, config, headers)], configured_header_names
    return list(_stdio_launch(config)), set()


class MCPServerTransportMixin:
    """Methods of :class:`tools.mcp_tool.MCPServerTask` (mixed in; relies on its attributes)."""

    __slots__ = ()

    def _advertises_tools(self) -> bool:
        """False only when captured capabilities omit ``tools`` (prompt-/resource-only servers,
        where ``tools/list`` raises -32601); True without capability info (legacy fallback).

        Per the MCP spec, ``InitializeResult.capabilities.tools`` is non-None iff the server implements the
        ``tools/*`` request family. Prompt-only or resource-only servers omit it, and calling ``tools/list``
        against them raises ``MCPError(-32601 Method not found)`` — which previously killed the connection
        during discovery and made every keepalive fail. (Ported from anomalyco/opencode#31271.)
        """
        caps = getattr(self.initialize_result, "capabilities", None)
        return caps is None or getattr(caps, "tools", None) is not None

    def _session_kwargs(self) -> dict:
        """ClientSession kwargs: sampling, elicitation, notification + logging callbacks."""
        kwargs = {}
        for handler in (self._sampling, self._elicitation):
            if handler:
                kwargs.update(handler.session_kwargs())
        if _core._MCP_NOTIFICATION_TYPES and _core._MCP_MESSAGE_HANDLER_SUPPORTED:
            kwargs["message_handler"] = self._make_message_handler()
        if _core._MCP_LOGGING_CALLBACK_SUPPORTED:
            kwargs["logging_callback"] = self._make_logging_callback()
        return kwargs

    async def _negotiate_session(self, session, connect_timeout: float):
        """Negotiate the protocol era (``initialize`` vs ``server/discover``; both expose
        ``.capabilities``). ``auto`` tries the legacy handshake FIRST, falling back to discover only
        on a modern-only signal (-32022 / initialize -32601) — the reverse of the SDK's discover-first
        mode, so handshake-era servers pay zero extra round-trips. ``stateless`` probes discover first
        (one legacy retry on any error); ``legacy`` is handshake only. A TIMEOUT never falls back."""
        def call(method: str):
            return asyncio.wait_for(getattr(session, method)(), timeout=connect_timeout)

        async def attempt(primary, fallback, should_fallback, log_fmt, *log_extra):
            try:
                return await call(primary)
            except Exception as exc:
                if isinstance(exc, asyncio.TimeoutError) or not should_fallback(exc):
                    raise
                logger.info(log_fmt, self.name, exc, *log_extra)
                try:
                    return await call(fallback)
                except Exception as fallback_exc:
                    # #113359: the server ANSWERED ``initialize`` (200, valid result) but named a version the
                    # SDK's handshake refuses (e.g. 2026-07-28 echoed to a 2025-11-25 offer), and it has no
                    # ``server/discover`` either. The wire handshake succeeded, so complete it ourselves.
                    if (isinstance(fallback_exc, asyncio.TimeoutError) or primary != "initialize"
                            or not _handshake_answered_with_unsupported_version(exc)):
                        raise
                    logger.info("MCP server '%s': server/discover also failed (%s) — completing the handshake "
                                "at %s, the version this client offered", self.name, fallback_exc,
                                _core.LATEST_HANDSHAKE_VERSION)
                    return await asyncio.wait_for(self._complete_handshake_at_offered_version(session),
                                                  timeout=connect_timeout)
        mode = str((self._config or {}).get("protocol", "auto")).lower().strip()
        if mode in ("stateless", "modern", "2026-07-28"):
            return await attempt("discover", "initialize", lambda exc: True,
                                 "MCP server '%s': server/discover rejected (%s) despite "
                                 "protocol=%s — falling back to the legacy handshake", mode)
        if mode in ("legacy", "handshake"):
            return await call("initialize")
        if mode != "auto":
            logger.warning("MCP server '%s': unknown protocol=%r — treating as 'auto' "
                           "(valid: auto, stateless, legacy)", self.name, mode)
        # mcp 1.x has no server/discover client — nothing to fall back to.
        return await attempt(
            "initialize", "discover", lambda exc: _handshake_rejected_as_modern(exc) and hasattr(session, "discover"),
            "MCP server '%s': legacy handshake rejected (%s) — retrying via server/discover (2026-07-28 stateless server)")

    async def _complete_handshake_at_offered_version(self, session):
        """Re-run the legacy ``initialize`` exchange the SDK already proved works against this server and
        adopt its result pinned to the version WE offered (#113359). ``ClientSession.initialize()`` raises
        on a ``protocolVersion`` outside its handshake set even though the server answered 200, and a
        stateless server that echoes 2026-07-28 to every offer has no ``server/discover`` — so this is the
        only way to reach ``notifications/initialized`` and ``tools/list``. Pinning to the offered version
        keeps later requests legacy-shaped (envelope and MCP-Protocol-Version header), the form the
        handshake itself just proved the server accepts. The returned result keeps the server's own
        version for logging/diagnostics."""
        import mcp.types as types  # late: keeps the SDK import lazy
        offered = _core.LATEST_HANDSHAKE_VERSION
        build_caps = getattr(session, "_build_capabilities", None)
        capabilities = build_caps(offered) if callable(build_caps) else types.ClientCapabilities()
        client_info = getattr(session, "_client_info", None) or types.Implementation(name="hermes-agent", version="0")
        result = await session.send_request(
            types.InitializeRequest(params=types.InitializeRequestParams(
                protocolVersion=offered, capabilities=capabilities, clientInfo=client_info)),
            types.InitializeResult)
        session.adopt(result.model_copy(update={"protocol_version": offered}))
        await session.send_notification(types.InitializedNotification())
        return result

    async def _serve_session(self, session, connect_timeout: float,
                             label: str = "", mark_lifecycle: bool = False) -> str:
        """Handshake, discover, publish readiness, then serve until a lifecycle event. Clears stale
        breaker state but leaves the session UNPROVEN: flapping transports handshake fine and drop
        moments later, so only keepalive/tool-call success clears the reconnect budget."""
        self.initialize_result = await self._negotiate_session(session, connect_timeout)
        self.session = session
        if mark_lifecycle:
            self._mark_lifecycle_started()
        await self._discover_tools()
        self._ready.set()
        self._ever_connected = True
        _core._reset_server_error(self.name)
        # Session is live again: clear any breaker state from a prior outage so the first call after
        # recovery isn't gated on a stale consecutive-failure count (#16788).
        # A completed handshake alone is NOT proof of health: a flapping transport can handshake fine and
        # drop moments later, forever (#62212). The session must prove itself (keepalive success, a
        # successful tool call, or — stdio without a keepalive — surviving a full default interval
        # idle with the child alive) before the reconnect budget is cleared — see _mark_session_proven.
        # Session is live again: clear any breaker state from a prior outage so the first call after
        # recovery isn't gated on a stale consecutive-failure count (#16788).
        # Unproven until keepalive/tool-call success (#62212).
        # Session is live again: clear any breaker state from a prior outage so the first call after
        # recovery isn't gated on a stale failure count (#16788).
        # Unproven until keepalive/tool-call success (#62212).
        # Session is live again: clear any breaker state from a prior outage so the first call after
        # recovery isn't gated on a stale consecutive-failure count (#16788).
        # Unproven until keepalive/tool-call success (#62212).
        self._session_proven = False
        reason = await self._wait_for_lifecycle_event()
        if label and reason == "reconnect":
            logger.info("MCP server '%s': reconnect requested — tearing down %s session", self.name, label)
        return reason

    async def _serve_transport(self, transport_cm, label: str, connect_timeout: float) -> str:
        """Open *transport_cm*, wrap its streams in a ClientSession and serve it. Streams are indexed,
        not unpacked (mcp 1.x yields a 3-tuple, 2.x a pair); a TaskGroup drop maps to ``"reconnect"``."""
        try:
            async with transport_cm as _streams:
                async with _core.ClientSession(_streams[0], _streams[1], **self._session_kwargs()) as session:
                    return await self._serve_session(session, connect_timeout, label)
        except BaseExceptionGroup as _eg:
            return self._reconnect_or_reraise_group(_eg)

    # ------------------------------------------------------------------ stdio

    def _track_spawned_children(self, new_pids: Set[int]) -> None:
        """Ledger the freshly spawned stdio children (pids, pgids, machine spawn ledger). pgids are
        captured while alive (getpgid fails after exit; the sweep needs them for reparented descendants)."""
        new_pgids: Dict[int, int] = {}
        for pid in new_pids:
            try:
                new_pgids[pid] = os.getpgid(pid)
            except ProcessLookupError:
                # Raced and already exited. The SDK spawns with start_new_session=True, so the
                # child was its own group leader (pgid == pid): keep that group covered — any
                # descendant it left behind still has to be reaped; the prune forgets the group
                # once nothing in it is alive.
                new_pgids[pid] = pid
            except (AttributeError, OSError):  # Windows (os.getpgid is POSIX-only)
                pass
        with _core._lock:
            _stdio_pids.update(dict.fromkeys(new_pids, self.name))
            _stdio_pgids.update(new_pgids)
        # Machine spawn ledger (startup sweeps reap orphans after an unclean exit); best-effort.
        for _pid in new_pids:
            try:
                from hermes_cli.process_identity import register_child
                register_child(_pid, "mcp-helper")
            except Exception:
                logger.debug("spawn-ledger register_child failed for MCP helper pid %s", _pid, exc_info=True)
        # Hand the pgroups to the shared parent-death supervisor so an ungraceful exit of this
        # process (kill -9, crash, force-quit) can't leave this server — or its descendants, e.g.
        # mcp-remote's spawned `node` — running forever. The graceful paths (shutdown,
        # _kill_orphaned_mcp_children) still reap as before; this only covers when they never run.
        _core._update_death_supervisor("register", new_pgids.values())

    def _release_spawned_children(self, new_pids: Set[int]) -> None:
        """Drop the ledger entries; a child (or its pgroup) still alive means SDK teardown failed
        (common on mid-way cancel on Linux: setsid() children escape) — mark it orphaned for the sweep."""
        from gateway.status import _pid_exists
        # Groups with nothing left alive; the supervisor forgets them after the lock is released.
        # Groups still alive stay registered on purpose, so the supervisor still reaps them if this
        # process dies before the orphan sweep runs.
        released_pgids: list = []
        with _core._lock:
            for pid in new_pids:
                _stdio_pids.pop(pid, None)
                # Windows-safe pid probe; the child may be gone while descendants remain in its pgroup.
                if _pid_exists(pid) or _pgroup_alive(_stdio_pgids.get(pid)):
                    _orphan_stdio_pids.add(pid)
                    _orphan_stdio_pid_servers[pid] = self.name
                else:  # nothing to reap — drop the pgid so PID reuse can't surface stale pgroup state
                    dropped = _stdio_pgids.pop(pid, None)
                    if dropped is not None:
                        released_pgids.append(dropped)
        _core._update_death_supervisor("unregister", released_pgids)

    async def _run_stdio(self, config: dict):
        """Run the server using stdio transport."""
        if config.get("identity_header") is not None:  # copy-pasted HTTP block: warn, don't mislead
            logger.warning("MCP server '%s': identity_header is only supported on "
                           "HTTP/SSE transports — ignored for stdio servers", self.name)
        if not _core._ensure_mcp_sdk():
            raise ImportError(f"MCP server '{self.name}' requires the 'mcp' Python SDK, but "
                              "it is not installed. Run `hermes setup` to install MCP support, then retry.")
        command = config.get("command")
        if not command:
            raise ValueError(f"MCP server '{self.name}' has no 'command' in config")
        inputs, _ = _connect_inputs(self.name, config)
        # Hash the inputs this attempt spawns with, never a second resolution of them.
        self._resolved_identity = _registration._identity_digest(inputs)
        command, safe_env, stdio_cwd = inputs
        # OSV malware preflight, then the cached-npx swap (ordering enforced there).
        command, args = await _core._preflight_stdio_command(self.name, command, config.get("args", []))
        server_params = _core.StdioServerParameters(
            command=command, args=args, env=safe_env or None, cwd=stdio_cwd,
            # Windows pipes can split non-UTF-8 bytes at chunk boundaries; substitute, don't raise.
            encoding_error_handler="replace")
        # Windows has no POSIX parent-death supervisor / killpg safety net (#61059): when this
        # process dies ungracefully (crash, force-quit), the stdio child trees — npx.cmd →
        # node.exe — survive as orphans with ParentId=null and pile up across restarts. Attaching
        # THIS process to a KILL_ON_JOB_CLOSE job before the spawn makes every child (and
        # grandchild) created after it die with the parent at the kernel level, so no orphan can
        # outlive us. Idempotent; BREAKAWAY_OK keeps deliberate breakaway children escaping.
        # Self-guards: a cheap no-op returning False on non-Windows.
        try:
            from hermes_cli.process_identity import attach_self_to_kill_on_close_job
            attach_self_to_kill_on_close_job()
        except Exception:
            logger.debug("job-object self-attach failed before stdio spawn", exc_info=True)
        # Reap orphans of prior attempts first (else retries pile up zombie pairs); unscoped on purpose;
        # off-loop because the reaper blocks up to 2s.
        await asyncio.to_thread(_lifecycle._kill_orphaned_mcp_children)
        pids_before = _lifecycle._snapshot_child_pids()  # so the new child can be identified after spawn
        # Reap any orphaned subprocesses from prior failed connection attempts before spawning a new one.
        # Without this, each retry in the run() reconnect loop spawns a fresh process pair while the
        # previous failed pair lingers — leading to rapid zombie accumulation (see #57355, #57228). The
        # unscoped sweep also opportunistically reaps orphans left by *other* servers that never reconnect;
        # per-server filtering via ``server_name`` remains available for scoped call sites. Run in a worker
        # thread: the reaper blocks up to 2s (SIGTERM → wait → SIGKILL) when orphans exist, which would
        # otherwise stall the shared MCP event loop.
        new_pids: set = set()
        # Subprocess stderr goes to ~/.hermes/logs/mcp-stderr.log so banners can't corrupt the TUI.
        _config._write_stderr_log_header(self.name)
        stderr = _config._StderrTee(_config._get_mcp_stderr_log())
        try:
            async with _core.stdio_client(server_params, errlog=stderr.sink) as (read_stream, write_stream):
                # New PIDs for force-kill cleanup, minus non-MCP children (slash_worker, LSP) racing
                # into the window: they share the TUI's pgid — leaking them would killpg() the TUI.
                new_pids = _filter_mcp_children(_lifecycle._snapshot_child_pids() - pids_before)
                if new_pids:
                    self._track_spawned_children(new_pids)
                self._stdio_child_pids = set(new_pids)  # so in-flight calls fail fast when the child dies
                async with _core.ClientSession(read_stream, write_stream, **self._session_kwargs()) as session:
                    # Bound the handshake here (``connect_timeout`` only bounds the caller's ``.result()``):
                    # a server that never answers ``initialize`` would leak child + pipes per retry until EMFILE.
                    connect_timeout = float(config.get("connect_timeout", _core._DEFAULT_CONNECT_TIMEOUT))
                    return await self._serve_session(session, connect_timeout, mark_lifecycle=True)
        except Exception as exc:
            # The SDK only sees "Connection closed"; the child's stderr says why (#124264).
            abi_error = node_abi_error(self.name, await asyncio.to_thread(stderr.close))
            if abi_error is not None:
                raise abi_error from exc
            raise
        finally:  # clean exit, exceptions AND cancellation
            stderr.close(timeout=0)
            if new_pids:
                self._release_spawned_children(new_pids)

    # ------------------------------------------------------------------- HTTP

    async def _preflight_content_type(self, url: str, *, headers: Optional[dict] = None,
                                      ssl_verify: bool = True, client_cert=None, timeout: float = 5.0,
                                      strict_redirect_headers: bool = False) -> None:
        """Probe *url* before the SDK connects: a plain web page would make the SDK sit out the full
        ``connect_timeout`` before an opaque ``CancelledError``; this raises NonMcpEndpointError within
        ``timeout``. Allow-list based: only a 2xx with a definite non-MCP content type is rejected, and
        only after a JSON-RPC ``initialize`` POST also fails to look like MCP (some servers serve a UI
        on GET but speak MCP via POST). Anything else passes — the handshake stays the source of truth.
        Own httpx client, OUTSIDE the SDK's anyio task group, so the error isn't group-wrapped."""
        try:
            import httpx as _httpx
        except ImportError:
            return  # No httpx → skip probe; SDK import would have failed first.

        def _non_mcp_2xx(resp) -> bool:
            # Only judge 2xx (4xx/5xx may be an auth challenge); no content type advertised → trust the SDK.
            ct = _content_type_base(resp)
            return _is_2xx(resp) and bool(ct) and ct not in self._MCP_CONTENT_TYPES
        probe_headers = dict(headers) if headers else {}
        # Same route as the SDK client: TLS on an explicit transport (which also turns off httpx's own
        # env proxy auto-detection) plus the repo's proxy mounts, so the probe and the handshake agree.
        # Same redirect boundary as the transport client too: httpx strips Authorization on a
        # cross-origin hop natively, but forwards every other configured header verbatim — under
        # strict_redirect_headers those must not leave the configured origin on the probe either.
        probe_transport = _httpx.AsyncHTTPTransport(verify=ssl_verify, **_present(cert=client_cert))
        _build_client = _make_redirect_header_stripper(
            _httpx, _httpx.URL(url), strict=strict_redirect_headers,
            configured_header_names={key.lower() for key in probe_headers})
        try:
            async with _build_client(
                    follow_redirects=True, timeout=_httpx.Timeout(timeout), transport=probe_transport,
                    **_present(mounts=_mcp_proxy_mounts(_httpx, url, ssl_verify, client_cert, self.name))) as client:
                resp = await client.head(url, headers=probe_headers)  # cheapest; GET on 405/501
                if resp.status_code in (405, 501):
                    resp = await client.get(url, headers=probe_headers)
                # Non-MCP content type on HEAD/GET: try a JSON-RPC POST so POST-only servers pass.
                if _non_mcp_2xx(resp):
                    post_resp = await client.post(
                        url, content=_PROBE_INITIALIZE_BODY,
                        headers={**probe_headers, "Content-Type": "application/json",
                                 "Accept": "application/json, text/event-stream"})
                    if _is_2xx(post_resp) and _content_type_base(post_resp) in self._MCP_CONTENT_TYPES:
                        resp = post_resp
        except _httpx.HTTPError:
            return  # DNS/connect/timeout/transport error — let the SDK try.
        if not _non_mcp_2xx(resp):
            return
        ct_base = _content_type_base(resp)
        raise NonMcpEndpointError(f"MCP server '{self.name}' at {url} returned Content-Type '{ct_base}', not an MCP "
            f"response (expected one of: {', '.join(self._MCP_CONTENT_TYPES)}). The URL most likely "
            "points at a web page rather than an MCP endpoint — check it resolves to a Streamable "
            "HTTP / SSE endpoint (e.g. https://host/mcp, not https://host/).")

    def _reconnect_or_reraise_group(self, eg: BaseExceptionGroup) -> str:
        """Map an SDK transport TaskGroup failure to a clean ``"reconnect"``: HTTP/SSE stream pumps run in an anyio
        TaskGroup, so a transient drop escapes as a ``BaseExceptionGroup`` that would otherwise park the server for
        300s over a sub-second glitch. Re-raise when it is not one: shutdown in progress (``_shutdown_event`` is
        set before cancel), KeyboardInterrupt/SystemExit or a real CancelledError in the group, or no live session
        this attempt (``_ready`` unset — connect failures must back off, not hot-loop).

        Streamable-HTTP / SSE transports run their stream pump inside an anyio TaskGroup. A transient stream
        drop (idle timeout, brief backend blip, server-side TCP close) surfaces as a ``BaseExceptionGroup``
        escaping the transport context manager. Left unwrapped it reaches ``run()``'s error path, which
        applies exponential backoff and eventually *parks* the server for 300s and deregisters its tools — a
        multi-minute tool outage for what is usually a sub-second glitch while the POST path stays healthy
        (issue #66092).
        - the group carries a ``KeyboardInterrupt`` / ``SystemExit`` — fatal signals must propagate to the
        interpreter, never be converted into a reconnect; - the group carries a real ``CancelledError``
        (task cancellation must propagate to asyncio, mirroring the ``run()`` guard for #9930); - we never
        reached a live session this attempt (``_ready`` unset) — a connect/handshake failure SHOULD fall
        through to ``run()``'s backoff rather than hot-loop reconnects against a broken endpoint.
        """
        if (self._shutdown_event.is_set()
                or eg.split((KeyboardInterrupt, SystemExit))[0] is not None
                or eg.split(asyncio.CancelledError)[0] is not None
                or not self._ready.is_set()):
            raise eg
        logger.debug("MCP server '%s': transport TaskGroup exited after a live session "
                     "(%r) — reconnecting immediately instead of backing off", self.name, eg)
        return "reconnect"

    def _build_oauth_auth(self, url: str, config: dict):
        """OAuth 2.1 PKCE via the central MCPOAuthManager (one provider reused across reconnects and
        CLI paths). Setup failures re-raise (after a warning) so only this server is reported failed."""
        if self._auth_type != "oauth":
            return None
        try:
            from tools.mcp_oauth_manager import get_manager
            return get_manager().get_or_build_provider(self.name, url, config.get("oauth"))
        except Exception as exc:
            logger.warning("MCP OAuth setup failed for '%s': %s", self.name, exc)
            raise

    def _sse_transport(self, url: str, headers: dict, connect_timeout: float,
                       ssl_verify, client_cert, oauth_auth, strict_cfg_headers: bool):
        """``sse_client`` context manager for ``transport: sse`` entries."""
        if strict_cfg_headers:  # fail closed: SSE cannot enforce the redirect boundary
            raise ValueError(f"MCP server '{self.name}': strict_redirect_headers is "
                             "not supported on the SSE transport.")
        if _core.sse_client is None:
            raise ImportError(f"MCP server '{self.name}' requires SSE transport but "
                              "mcp.client.sse.sse_client is not available. "
                              "Upgrade the mcp package to get SSE support.")
        # sse_read_timeout bounds the gap between events: SSE servers idle for minutes, so 300s (the
        # Streamable HTTP read timeout), not tool_timeout. ``auth`` must be forwarded or OAuth SSE 401s silently.
        sse_kwargs: dict = {"url": url, "headers": headers or None, "timeout": float(connect_timeout),
                            "sse_read_timeout": 300.0, **_present(auth=oauth_auth)}
        # Always own the client: the httpx_client_factory forwards the SDK's (headers, auth, timeout),
        # installs the wire-body cap, layers TLS on the inner transport (client-level verify/cert are
        # inert once a custom transport= is passed) and re-adds the proxy mounts that custom transport
        # would otherwise suppress. Client MUST come from the SDK's httpx (httpx2 on mcp >= 2.0).
        _httpx_mod = _core.sdk_httpx()
        def _sse_client_factory(headers=None, timeout=None, auth=None):
            inner_transport = _httpx_mod.AsyncHTTPTransport(verify=ssl_verify, **_present(cert=client_cert))
            return _httpx_mod.AsyncClient(
                follow_redirects=True,
                timeout=timeout if timeout is not None else _httpx_mod.Timeout(30.0, read=300.0),
                transport=_make_mcp_body_cap_transport(_httpx_mod, inner_transport),
                **_present(mounts=_mcp_proxy_mounts(_httpx_mod, url, ssl_verify, client_cert, self.name),
                           headers=headers, auth=auth))
        sse_kwargs["httpx_client_factory"] = _sse_client_factory
        return _core.sse_client(**sse_kwargs)

    def _streamable_http_transport(self, url: str, headers: dict, connect_timeout: float,
                                   ssl_verify, client_cert, oauth_auth,
                                   strict_cfg_headers: bool, configured_header_names: set):
        """Streamable HTTP context manager: mcp >= 1.24.0 gets a caller-owned httpx client; on the
        deprecated API (mcp < 1.24.0) the SDK owns the client."""
        if not _core._MCP_NEW_HTTP:
            if strict_cfg_headers:  # fail closed: without an owned client redirects can't be hooked
                raise ImportError(f"MCP server '{self.name}' requires mcp >= 1.24.0 to "
                                  "enforce the portable redirect-header boundary "
                                  "(strict_redirect_headers). Upgrade the mcp package.")
            return _core.streamablehttp_client(url, headers=headers, timeout=float(connect_timeout), verify=ssl_verify,
                                               **_present(auth=oauth_auth))
        # Explicit AsyncClient matching the SDK's create_mcp_http_client defaults; MUST come from the
        # SDK's httpx (httpx2 on mcp >= 2.0) since the SDK sends its own Requests through it.
        httpx = _core.sdk_httpx()
        _build_client = _make_redirect_header_stripper(
            httpx, httpx.URL(url), strict=strict_cfg_headers, configured_header_names=configured_header_names)
        # verify/cert live on the inner transport: a custom transport= makes client-level TLS kwargs
        # inert — and suppresses httpx's own proxy auto-detection, hence the explicit mounts=.
        inner_transport = httpx.AsyncHTTPTransport(verify=ssl_verify, **_present(cert=client_cert))
        client_kwargs: dict = {"follow_redirects": True, "timeout": httpx.Timeout(float(connect_timeout), read=300.0),
                               **({"headers": headers} if headers else {}),
                               "event_hooks": {"response": [_make_http_rejection_recorder(self._http_rejection)]},
                               "transport": _make_mcp_body_cap_transport(httpx, inner_transport),
                               **_present(mounts=_mcp_proxy_mounts(httpx, url, ssl_verify, client_cert, self.name),
                                          auth=oauth_auth)}

        @asynccontextmanager
        async def _owned_client_streams():  # the SDK skips cleanup when http_client is provided
            async with _build_client(**client_kwargs) as http_client:
                async with _core.streamable_http_client(url, http_client=http_client) as streams:
                    yield streams
        return _owned_client_streams()

    async def _run_http(self, config: dict):
        """Run the server using HTTP/StreamableHTTP (or SSE) transport."""
        _core._ensure_mcp_sdk()
        if not _core._MCP_HTTP_AVAILABLE:
            raise ImportError(f"MCP server '{self.name}' requires HTTP transport but "
                              "mcp.client.streamable_http is not available. "
                              "Upgrade the mcp package to get HTTP support.")
        # Agent Plugins v1 strict_redirect_headers: configured headers MUST NOT follow a cross-origin
        # redirect — their names are captured BEFORE client-generated headers are merged in.
        inputs, configured_header_names = _connect_inputs(self.name, config)
        # Hash the endpoint this attempt connects to: a second read of the runtime file could
        # publish another endpoint's identity alongside this session.
        self._resolved_identity = _registration._identity_digest(inputs)
        url, headers = inputs
        logger.debug("MCP server '%s': connecting to %s", self.name, url)
        self._http_rejection = {}  # last 4xx/5xx the owned client saw this attempt (recorder hook)
        # Seed MCP-Protocol-Version (user override wins) from the HANDSHAKE version, not the latest: a
        # 2026-07-28 header routes the handshake-era ``initialize()`` onto the envelope ladder, which rejects it.
        if not any(key.lower() == "mcp-protocol-version" for key in headers):
            headers["mcp-protocol-version"] = _core.LATEST_HANDSHAKE_VERSION
        connect_timeout = config.get("connect_timeout", _core._DEFAULT_CONNECT_TIMEOUT)
        common = (url, headers, connect_timeout, config.get("ssl_verify", True), _resolve_client_cert(self.name, config),
                  self._build_oauth_auth(url, config), bool(config.get("strict_redirect_headers")))
        if config.get("transport") == "sse":
            return await self._serve_transport(self._sse_transport(*common), "SSE", float(connect_timeout))
        if self._sse_fallback:
            # A prior connect already proved this server SSE-only: skip the doomed Streamable
            # HTTP attempt on reconnects instead of flapping into the retry budget.
            logger.info("MCP server '%s': using latched SSE fallback transport", self.name)
            return await self._serve_transport(self._sse_transport(*common), "SSE", float(connect_timeout))
        transport = self._streamable_http_transport(*common, configured_header_names)
        label = "HTTP" if _core._MCP_NEW_HTTP else "legacy HTTP"
        try:
            return await self._serve_transport(transport, label, float(connect_timeout))
        except Exception as exc:
            # The SDK folds a non-2xx it cannot parse into ``-32603 Server returned an error response``;
            # the recorder hook kept the status/URL/body the server actually sent (#114350, #113359).
            http_detail = _describe_http_failure(exc, self._http_rejection)
            # SSE-only servers (or their load balancers) reject the Streamable HTTP chunked
            # ``initialize`` POST — with a 400-family status or an opaque SDK INTERNAL_ERROR —
            # previously a permanent failure with 0 active tools unless the user set
            # ``transport: sse`` (#53676, #104343). Retry over SSE on the initial connect, as
            # the MCP spec's transport-fallback behavior describes. Never on reconnect after a
            # proven session (``_ever_connected``: a genuine rejection on an established
            # transport must not silently switch transports), never on a timeout (not a
            # transport mismatch — ``_is_streamable_http_rejection`` matches neither), and never
            # with ``strict_redirect_headers`` (SSE cannot enforce that boundary).
            if (self._ever_connected or common[-1] or not _is_streamable_http_rejection(exc)):
                if http_detail != str(_unwrap_exception_group(exc)):  # opaque SDK error + a recorded rejection
                    raise ConnectionError(f"MCP server '{self.name}': Streamable HTTP connect failed "
                                          f"({http_detail})") from exc
                raise
            logger.warning(
                "MCP server '%s': Streamable HTTP rejected the initial connect (%s) — retrying "
                "over SSE. If this connects, set `transport: sse` for this server in config.yaml "
                "to skip the failed attempt on future startups.",
                self.name, http_detail)
            try:
                self._sse_fallback = True
                return await self._serve_transport(self._sse_transport(*common), "SSE", float(connect_timeout))
            except Exception as sse_exc:
                if self._ever_connected:  # SSE session was live and dropped: transient, keep the latch
                    raise
                self._sse_fallback = False
                raise ConnectionError(
                    f"MCP server '{self.name}': both Streamable HTTP and SSE transports failed "
                    f"(Streamable HTTP: {http_detail}; SSE: "
                    f"{_unwrap_exception_group(sse_exc)}). Check the URL points at an MCP "
                    "endpoint, or pin `transport: sse` if the server is SSE-only.") from sse_exc

    # -------------------------------------------------------------- discovery

    # Legacy Streamable-HTTP transport TaskGroup dropped: reconnect immediately instead of backoff/park
    # (#66092).
    async def _discover_tools(self):
        """Discover tools from the connected session. Capability-gated: prompt-/resource-only
        servers raise ``MCPError(-32601)`` on ``tools/list``, which would abort the connection.

        Skip the call when the server doesn't advertise the ``tools`` capability. (Ported from
        anomalyco/opencode#31271.)
        """
        self._ping_unsupported = False  # fresh transport: re-probe ``ping`` across the reconnect
        if self.session is None:
            return
        if not self._advertises_tools():
            logger.info("MCP server '%s': does not advertise 'tools' capability — "
                        "skipping tools/list (prompts/resources remain available)", self.name)
            self._tools = []
        else:
            async with self._rpc_lock:
                self._list_cache_meta = {}
                self._tools = await _core._paginate_full_list(
                    self.session.list_tools, "tools", self.name, cache_meta_out=self._list_cache_meta)
        self._register_discovered_tools_if_needed()

    def _register_discovered_tools_if_needed(self) -> None:
        """Publish freshly discovered tools when none are registered (initial registration normally happens in
        ``_discover_and_register_server``). Outage handling may clear ``_ready`` and deregister stale tools;
        ownership via ``_servers`` authorizes publishing before readiness is restored so a revival (or a server
        retained after a recoverable initial failure) never comes back with zero tools."""
        if self._registered_tool_names:
            return
        with _core._lock:
            owned = [key for key, live in _core._servers.items() if live is self]
        if not owned and not self._ready.is_set():
            return
        self._registered_tool_names = _registration._register_server_tools(self.name, self, self._config)
        with _core._lock:  # a retained initial-failure server that just published tools has recovered
            for key in owned:
                if _core._servers.get(key) is self:
                    _core._server_connect_errors.pop(key, None)
