Containerize MCP tool services
This commit is contained in:
@@ -10,12 +10,11 @@ from facts verified by crawled pages or primary APIs.
|
||||
from __future__ import annotations
|
||||
|
||||
import ipaddress
|
||||
import asyncio
|
||||
import html
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import select
|
||||
import subprocess
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
@@ -28,9 +27,9 @@ from xml.etree import ElementTree
|
||||
|
||||
|
||||
SERVER_VERSION = "2.1.0"
|
||||
TINYSEARCH_CONTAINER = os.environ.get(
|
||||
"TINYSEARCH_CONTAINER", "mike-ai-web-search-tinysearch-1"
|
||||
)
|
||||
TINYSEARCH_MCP_URL = os.environ.get(
|
||||
"TINYSEARCH_MCP_URL", "http://tinysearch:8000/mcp"
|
||||
).rstrip("/")
|
||||
SEARXNG_URL = os.environ.get("SEARXNG_URL", "").rstrip("/")
|
||||
CHILD_TIMEOUT_SECONDS = float(os.environ.get("TINYSEARCH_CHILD_TIMEOUT", "110"))
|
||||
HTTP_TIMEOUT_SECONDS = float(os.environ.get("WEB_API_TIMEOUT", "18"))
|
||||
@@ -801,90 +800,46 @@ def general_discovery(query: str, limit: int) -> tuple[list[dict[str, Any]], lis
|
||||
|
||||
|
||||
class TinySearchClient:
|
||||
"""Minimal synchronous MCP client for TinySearch's stdio server."""
|
||||
"""Synchronous facade for TinySearch's Streamable HTTP MCP endpoint.
|
||||
|
||||
TinySearch used to be reached by executing ``docker exec`` on the host.
|
||||
That required access to the Docker socket and coupled this service to a
|
||||
particular container name. The containerized tool layer instead talks to
|
||||
TinySearch over the private Docker network.
|
||||
"""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self._next_id = 1
|
||||
self._lock = threading.Lock()
|
||||
self._process = subprocess.Popen(
|
||||
[
|
||||
"/usr/bin/docker",
|
||||
"exec",
|
||||
"-e",
|
||||
"MCP_TRANSPORT=stdio",
|
||||
"-i",
|
||||
TINYSEARCH_CONTAINER,
|
||||
"tinysearch",
|
||||
"mcp",
|
||||
],
|
||||
stdin=subprocess.PIPE,
|
||||
stdout=subprocess.PIPE,
|
||||
stderr=subprocess.DEVNULL,
|
||||
text=True,
|
||||
encoding="utf-8",
|
||||
errors="replace",
|
||||
bufsize=1,
|
||||
)
|
||||
self._request(
|
||||
"initialize",
|
||||
{
|
||||
"protocolVersion": "2024-11-05",
|
||||
"capabilities": {},
|
||||
"clientInfo": {"name": "mike-ai-web-facade", "version": SERVER_VERSION},
|
||||
},
|
||||
)
|
||||
self._notify("notifications/initialized", {})
|
||||
|
||||
def _notify(self, method: str, params: dict[str, Any]) -> None:
|
||||
assert self._process.stdin is not None
|
||||
message = {"jsonrpc": "2.0", "method": method, "params": params}
|
||||
self._process.stdin.write(json.dumps(message, separators=(",", ":")) + "\n")
|
||||
self._process.stdin.flush()
|
||||
async def _call_async(self, name: str, arguments: dict[str, Any]) -> str:
|
||||
# Imported lazily so the dependency-free stdio implementation still
|
||||
# gives a useful startup error outside its production container.
|
||||
from mcp import ClientSession
|
||||
from mcp.client.streamable_http import streamablehttp_client
|
||||
|
||||
def _request(self, method: str, params: dict[str, Any]) -> dict[str, Any]:
|
||||
with self._lock:
|
||||
if self._process.poll() is not None:
|
||||
raise RuntimeError("TinySearch child process is not running")
|
||||
request_id = self._next_id
|
||||
self._next_id += 1
|
||||
assert self._process.stdin is not None
|
||||
assert self._process.stdout is not None
|
||||
message = {
|
||||
"jsonrpc": "2.0",
|
||||
"id": request_id,
|
||||
"method": method,
|
||||
"params": params,
|
||||
}
|
||||
self._process.stdin.write(json.dumps(message, separators=(",", ":")) + "\n")
|
||||
self._process.stdin.flush()
|
||||
deadline = datetime.now().timestamp() + CHILD_TIMEOUT_SECONDS
|
||||
while True:
|
||||
remaining = deadline - datetime.now().timestamp()
|
||||
if remaining <= 0:
|
||||
raise TimeoutError(f"TinySearch timed out after {CHILD_TIMEOUT_SECONDS:.0f}s")
|
||||
ready, _, _ = select.select([self._process.stdout], [], [], remaining)
|
||||
if not ready:
|
||||
raise TimeoutError(f"TinySearch timed out after {CHILD_TIMEOUT_SECONDS:.0f}s")
|
||||
line = self._process.stdout.readline()
|
||||
if not line:
|
||||
raise RuntimeError("TinySearch closed its output stream")
|
||||
payload = json.loads(line)
|
||||
if payload.get("id") != request_id:
|
||||
continue
|
||||
if "error" in payload:
|
||||
raise RuntimeError(f"TinySearch error: {payload['error']}")
|
||||
return payload.get("result") or {}
|
||||
async with streamablehttp_client(
|
||||
TINYSEARCH_MCP_URL,
|
||||
timeout=CHILD_TIMEOUT_SECONDS,
|
||||
sse_read_timeout=CHILD_TIMEOUT_SECONDS,
|
||||
) as (read_stream, write_stream, _):
|
||||
async with ClientSession(read_stream, write_stream) as session:
|
||||
await session.initialize()
|
||||
result = await session.call_tool(name, arguments)
|
||||
if result.isError:
|
||||
raise RuntimeError(clean_text(str(result.content)))
|
||||
return "\n".join(
|
||||
item.text for item in result.content
|
||||
if getattr(item, "type", None) == "text"
|
||||
)
|
||||
|
||||
def call(self, name: str, arguments: dict[str, Any]) -> str:
|
||||
result = self._request("tools/call", {"name": name, "arguments": arguments})
|
||||
if result.get("isError"):
|
||||
raise RuntimeError(clean_text(str(result.get("content"))))
|
||||
texts = [
|
||||
item.get("text", "")
|
||||
for item in result.get("content", [])
|
||||
if isinstance(item, dict) and item.get("type") == "text"
|
||||
]
|
||||
return "\n".join(texts)
|
||||
with self._lock:
|
||||
try:
|
||||
return asyncio.run(asyncio.wait_for(
|
||||
self._call_async(name, arguments), CHILD_TIMEOUT_SECONDS))
|
||||
except TimeoutError as exc:
|
||||
raise TimeoutError(
|
||||
f"TinySearch timed out after {CHILD_TIMEOUT_SECONDS:.0f}s") from exc
|
||||
|
||||
|
||||
_client: TinySearchClient | None = None
|
||||
|
||||
Reference in New Issue
Block a user