Spaces:
Sleeping
Sleeping
| """MCP client backend — executes actions by calling the in-process MCP server. | |
| WHAT THIS IS | |
| ------------ | |
| This is the fourth execution backend the dispatcher can route an authorized | |
| action to. Unlike Direct-API and Function-Call (which do the work in plain Python | |
| inside this process), this backend does its work by acting as an **MCP client**: | |
| it connects to the FastMCP server in ``tools/mcp_server.py`` and calls that | |
| server's tools over the Model Context Protocol. That mirrors how a real agent | |
| would reach an external tool server, but here both ends live in this process and | |
| talk over an **in-memory transport** — there is no network endpoint and, because | |
| nothing is spawned as a subprocess, there is no orphan process to clean up. | |
| WHERE IT SITS IN THE TRUST MODEL (same as every backend) | |
| -------------------------------------------------------- | |
| gate decides ─► dispatcher verifies grant ─► THIS backend ─► MCP server runs | |
| By the time :meth:`run` is called the gate has already said ``ALLOW`` and the | |
| dispatcher has verified the signed grant. This backend just executes — plus it | |
| keeps the same defense-in-depth **allow-list** guard as the others: an action | |
| whose name is not one of the server's published tools is refused here, before any | |
| call is made. | |
| THE ONE AWKWARD BIT: SYNC DISPATCHER, ASYNC MCP CLIENT | |
| ------------------------------------------------------ | |
| The dispatcher calls backends synchronously (``run(action) -> result``), but the | |
| FastMCP client is asynchronous. We bridge the two with a tiny helper that runs a | |
| single asyncio event loop on a background thread (:class:`_EventLoopThread`). The | |
| backend opens **one** persistent client connection at construction and reuses it | |
| for every call, then tears it down cleanly in :meth:`close`. This is why the | |
| backend is a context manager: ``with McpClientBackend() as backend: ...`` makes | |
| the clean shutdown automatic. | |
| """ | |
| from __future__ import annotations | |
| import asyncio | |
| import threading | |
| from typing import Any | |
| from fastmcp import Client, FastMCP | |
| from control_plane.schema import ProposedAction | |
| from tools.mcp_server import MCP_TOOL_NAMES, build_server | |
| class McpBackendError(Exception): | |
| """Base error for anything this backend refuses or cannot do.""" | |
| class McpUnsupportedActionError(McpBackendError): | |
| """The requested action is not one of the MCP server's published tools. | |
| This is the defense-in-depth guard: anything that isn't in the server's tool | |
| allow-list is refused here rather than sent to the server. | |
| """ | |
| class _EventLoopThread: | |
| """Runs one asyncio event loop on a daemon thread. | |
| This lets synchronous code (the dispatcher → :meth:`McpClientBackend.run`) | |
| drive asynchronous coroutines (the FastMCP client) without each call spinning | |
| up and tearing down its own event loop. We submit a coroutine with | |
| :func:`asyncio.run_coroutine_threadsafe` and block on its result. | |
| """ | |
| def __init__(self) -> None: | |
| self._loop = asyncio.new_event_loop() | |
| self._thread = threading.Thread( | |
| target=self._loop.run_forever, name="mcp-client-loop", daemon=True | |
| ) | |
| self._thread.start() | |
| def run(self, coro: Any) -> Any: | |
| """Run *coro* on the background loop and block until it returns.""" | |
| return asyncio.run_coroutine_threadsafe(coro, self._loop).result() | |
| def close(self) -> None: | |
| """Stop the loop and join the thread — no background thread is left running.""" | |
| self._loop.call_soon_threadsafe(self._loop.stop) | |
| self._thread.join() | |
| self._loop.close() | |
| class McpClientBackend: | |
| """MCP-client adapter satisfying the dispatcher's ``BackendAdapter`` contract. | |
| Construct it once (it builds the in-process server and opens a persistent | |
| client connection), let the dispatcher call :meth:`run` per authorized action, | |
| and call :meth:`close` — or use it as a context manager — to shut down. | |
| """ | |
| def __init__(self, server: FastMCP | None = None) -> None: | |
| # Build (or accept) the in-process MCP server. A fresh server has its own | |
| # private ticket store, so this backend instance is isolated from others. | |
| self._server = server or build_server() | |
| # The allow-list of executable actions = the server's published tool names. | |
| # Anything not in here is refused by :meth:`run` before any call is made. | |
| self._allowed: frozenset[str] = frozenset(MCP_TOOL_NAMES) | |
| # Bridge async↔sync, then open ONE persistent client connection to the | |
| # server over the in-memory transport. ``__aenter__`` is the connect step | |
| # of FastMCP's ``async with Client(...)`` pattern; we drive it manually so | |
| # the connection can outlive a single call and be reused. | |
| self._loop = _EventLoopThread() | |
| self._client = Client(self._server) | |
| self._loop.run(self._client.__aenter__()) | |
| # -- lifecycle (context manager + explicit close) -------------------------- | |
| def __enter__(self) -> "McpClientBackend": | |
| return self | |
| def __exit__(self, *exc: object) -> None: | |
| self.close() | |
| def close(self) -> None: | |
| """Disconnect the client and stop the background loop — clean, no orphans. | |
| Safe to call more than once. Because the server runs in-memory (not as a | |
| subprocess), closing the client connection and stopping the loop is all the | |
| teardown there is. | |
| """ | |
| if self._client is not None: | |
| self._loop.run(self._client.__aexit__(None, None, None)) | |
| self._client = None | |
| self._loop.close() | |
| # -- introspection --------------------------------------------------------- | |
| def list_tools(self) -> list[str]: | |
| """Return the names of the tools the connected server publishes (via MCP).""" | |
| tools = self._loop.run(self._client.list_tools()) | |
| return [tool.name for tool in tools] | |
| # -- dispatcher entry point ------------------------------------------------ | |
| def run(self, action: ProposedAction) -> Any: | |
| """Execute *action* by calling the matching MCP tool and return its result. | |
| Refuses (raising :class:`McpUnsupportedActionError`) any action whose name | |
| is not one of the server's tools, before contacting the server at all. | |
| Otherwise it forwards the action's arguments to the tool and returns the | |
| tool's structured result (``result.data``). | |
| """ | |
| if action.action_name not in self._allowed: | |
| raise McpUnsupportedActionError( | |
| f"mcp_client backend cannot perform {action.action_name!r} " | |
| f"(not in the server's tool allow-list)" | |
| ) | |
| result = self._loop.run( | |
| self._client.call_tool(action.action_name, action.arguments) | |
| ) | |
| # ``result.data`` is the structured value the tool returned (a dict here). | |
| return result.data | |