Optimized toolchain configuration. Rewrite of tool-router.
This commit is contained in:
+1
-1
@@ -20,7 +20,7 @@ DNS_SERVER=8.8.8.8
|
||||
|
||||
# --- HOST MOUNT VOLUMES ---
|
||||
# Local Path to Your openclaw.json
|
||||
HOST_OPENCLAW_JSON_PATH=/home/ghost/source/repos/ghostnet-openclaw/config/openclaw.json
|
||||
HOST_OPENCLAW_JSON_PATH=/home/ghost/source/repos/ghostnet-openclaw/openclaw/config/openclaw.json
|
||||
# Local Path for the AI Workspace (persistent storage for agent's work)
|
||||
HOST_WORKSPACE_PATH=./storage/workspace
|
||||
# Local Path for Additional Repositories or Code Directories for the Agent
|
||||
|
||||
+45
-34
@@ -25,19 +25,33 @@ services:
|
||||
- OLLAMA_KEEP_ALIVE=${OLLAMA_KEEP_ALIVE}
|
||||
- HSA_OVERRIDE_GFX_VERSION=11.0.0
|
||||
- OLLAMA_FLASH_ATTENTION=1
|
||||
- http_proxy=http://ollama-proxy:3128
|
||||
- https_proxy=http://ollama-proxy:3128
|
||||
- no_proxy=localhost,127.0.0.1,0.0.0.0,ollama,openclaw-agent,tool-router,openclaw-ollama-bridge
|
||||
ports:
|
||||
- "11434:11434"
|
||||
dns:
|
||||
- ${DNS_SERVER}
|
||||
networks:
|
||||
- net.ghost.openclaw
|
||||
depends_on:
|
||||
- ollama-proxy
|
||||
|
||||
ollama-proxy:
|
||||
image: docker.io/ubuntu/squid
|
||||
container_name: ollama-proxy
|
||||
volumes:
|
||||
- ./ollama/proxy/squid.conf:/etc/squid/squid.conf:Z
|
||||
networks:
|
||||
- net.ghost.openclaw # Erreichbar für Ollama
|
||||
- net.ghost.tools # Hat Internetzugriff für Registry-Pulls
|
||||
restart: unless-stopped
|
||||
stop_grace_period: 3s
|
||||
|
||||
# --- CORE AGENT (GATEWAY) ---
|
||||
agent:
|
||||
openclaw:
|
||||
image: ghcr.io/openclaw/openclaw:latest
|
||||
container_name: openclaw-agent
|
||||
depends_on:
|
||||
- ollama
|
||||
container_name: openclaw
|
||||
userns_mode: "keep-id"
|
||||
environment:
|
||||
- OPENCLAW_GATEWAY_MODE=local
|
||||
@@ -52,34 +66,32 @@ services:
|
||||
- ${HOST_REPOS_PATH}:/repos:Z,U
|
||||
networks:
|
||||
- net.ghost.openclaw
|
||||
depends_on:
|
||||
- ollama
|
||||
- tool-router
|
||||
|
||||
# --- NETWORK TUNNEL (SIDECAR) ---
|
||||
ollama-bridge:
|
||||
image: docker.io/alpine/socat:latest
|
||||
container_name: openclaw-ollama-bridge
|
||||
depends_on:
|
||||
- agent
|
||||
network_mode: "service:agent"
|
||||
network_mode: "service:openclaw"
|
||||
command: TCP-LISTEN:11434,fork TCP:ollama:11434
|
||||
stop_grace_period: 1s
|
||||
restart: always
|
||||
depends_on:
|
||||
- openclaw
|
||||
|
||||
# --- TOOL ROUTER (GATEWAY) ---
|
||||
tool-router:
|
||||
image: docker.io/python:3.12-slim
|
||||
build:
|
||||
context: ./tools/router
|
||||
dockerfile: Containerfile
|
||||
container_name: tool-router
|
||||
depends_on:
|
||||
- tool-searchfetch
|
||||
init: true
|
||||
working_dir: /app
|
||||
environment:
|
||||
- PORT=3000
|
||||
command: >
|
||||
sh -c "pip install --no-cache-dir fastapi uvicorn 'modelcontextprotocol[python-sdk]' &&
|
||||
uvicorn router:app --host 0.0.0.0 --port 3000"
|
||||
volumes:
|
||||
- ./tools/router/src/router.py:/app/router.py:ro,Z
|
||||
- ./tools/router/config:/app/config:ro,Z
|
||||
- LOG_LEVEL=debug
|
||||
expose:
|
||||
- "3000"
|
||||
dns:
|
||||
@@ -87,35 +99,31 @@ services:
|
||||
networks:
|
||||
- net.ghost.openclaw
|
||||
- net.ghost.tools
|
||||
depends_on:
|
||||
- tool-searchfetch
|
||||
|
||||
# --- SEARCHFETCH TOOL ---
|
||||
tool-searchfetch:
|
||||
image: docker.io/nikolaik/python-nodejs:python3.14-nodejs26-slim
|
||||
build:
|
||||
context: ./tools/searchfetch
|
||||
dockerfile: Containerfile
|
||||
container_name: tool-searchfetch
|
||||
depends_on:
|
||||
- tool-searchfetch-searxng
|
||||
networks:
|
||||
- net.ghost.tools
|
||||
environment:
|
||||
- SEARXNG_URL=http://tool-searchfetch-searxng:8080
|
||||
- MCP_HOST=0.0.0.0
|
||||
- MCP_PORT=3000
|
||||
command: >
|
||||
sh -c "pip install --no-cache-dir mcp-proxy &&
|
||||
npm install -g mcp-searxng &&
|
||||
mcp-proxy --host $$MCP_HOST --port $$MCP_PORT --pass-environment -- mcp-searxng"
|
||||
networks:
|
||||
- net.ghost.tools
|
||||
depends_on:
|
||||
- tool-searchfetch-searxng
|
||||
healthcheck:
|
||||
test: ["CMD", "curl", "-f", "http://localhost:$$MCP_PORT/sse"]
|
||||
interval: 30s
|
||||
timeout: 10s
|
||||
test: ["CMD-SHELL", "timeout 2s nc -z localhost 3000 || exit 1"]
|
||||
interval: 10s
|
||||
timeout: 5s
|
||||
retries: 3
|
||||
|
||||
# --- SEARXNG SERVICE FOR SEARCHFETCH TOOL ---
|
||||
tool-searchfetch-searxng:
|
||||
image: docker.io/searxng/searxng:latest
|
||||
container_name: tool-searchfetch-searxng
|
||||
depends_on:
|
||||
- tool-searchfetch-valkey
|
||||
restart: always
|
||||
volumes:
|
||||
- ./tools/searchfetch/config/searxng.settings.yml:/etc/searxng/settings.yml:ro,Z
|
||||
@@ -129,20 +137,23 @@ services:
|
||||
- ${DNS_SERVER}
|
||||
networks:
|
||||
- net.ghost.tools
|
||||
depends_on:
|
||||
- tool-searchfetch-valkey
|
||||
|
||||
tool-searchfetch-valkey:
|
||||
image: docker.io/valkey/valkey:8-alpine
|
||||
container_name: tool-searchfetch-valkey
|
||||
restart: always
|
||||
user: "1000:1000"
|
||||
volumes:
|
||||
- ./storage/searxng_valkey:/data:U,Z
|
||||
- ./storage/searxng_valkey:/data:Z,U
|
||||
networks:
|
||||
- net.ghost.tools
|
||||
|
||||
tool-test:
|
||||
image: docker.io/nikolaik/python-nodejs:python3.14-nodejs26-slim
|
||||
container_name: tool-test
|
||||
command: ["sh", "-c", "while true; do sleep 1000; done"]
|
||||
command: ["sh", "-c", "exec tail -f /dev/null"]
|
||||
dns:
|
||||
- ${DNS_SERVER}
|
||||
networks:
|
||||
|
||||
@@ -0,0 +1,28 @@
|
||||
# Definition der Ziele (Ollama Registry)
|
||||
acl ollama_domains dstdomain .ollama.ai
|
||||
acl ollama_domains dstdomain .ollama.com
|
||||
acl ollama_domains dstdomain .cloudflare.com
|
||||
acl ollama_domains dstdomain .github.com
|
||||
acl ollama_domains dstdomain .githubusercontent.com
|
||||
|
||||
# Erlaubt Zugriff von allen gängigen privaten Docker/Podman Subnetzen
|
||||
acl private_nets src 10.0.0.0/8
|
||||
acl private_nets src 172.16.0.0/12
|
||||
acl private_nets src 192.168.0.0/16
|
||||
http_access allow private_nets
|
||||
http_access allow localhost
|
||||
|
||||
# Standard-Ports erlauben
|
||||
acl Safe_ports port 80
|
||||
acl Safe_ports port 443
|
||||
http_access deny !Safe_ports
|
||||
|
||||
# Nur die Whitelist erlauben, alles andere blockieren
|
||||
http_access allow ollama_domains
|
||||
http_access deny all
|
||||
|
||||
# Proxy-Port
|
||||
http_port 3128
|
||||
|
||||
acl local_network src 10.89.0.0/16 # Adjust to your Podman subnet if different
|
||||
http_access allow local_network
|
||||
@@ -9,13 +9,13 @@
|
||||
},
|
||||
"agents": {
|
||||
"defaults": {
|
||||
"model": "ollama/gemma4:e4b"
|
||||
"model": "ollama/qwen3.6:35b-a3b-q4_K_M"
|
||||
}
|
||||
},
|
||||
"mcp": {
|
||||
"servers": {
|
||||
"router": {
|
||||
"url": "http://mcp-router:3000/sse",
|
||||
"url": "http://tool-router:3000/sse",
|
||||
"transport": "sse"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,17 @@
|
||||
FROM docker.io/library/python:3.12-slim
|
||||
|
||||
# System dependencies
|
||||
RUN apt-get update && apt-get install -y curl && rm -rf /var/lib/apt/lists/*
|
||||
|
||||
# Application dependencies
|
||||
WORKDIR /app
|
||||
COPY app/requirements.txt .
|
||||
RUN pip install --no-cache-dir -r requirements.txt
|
||||
|
||||
COPY app/ .
|
||||
COPY config/ /app/config/
|
||||
|
||||
# Default environment variables
|
||||
ENV CONFIG_PATH=/app/config/gateway.json
|
||||
|
||||
CMD ["uvicorn", "router:app", "--host", "0.0.0.0", "--port", "3000", "--log-level", "info", "--timeout-keep-alive", "120", "--limit-concurrency", "10"]
|
||||
@@ -0,0 +1,75 @@
|
||||
# Local Development Setup für GhostNet MCP Router
|
||||
|
||||
## Voraussetzungen
|
||||
- Python 3.12+
|
||||
- VSCode mit Python/Pylance Extension
|
||||
|
||||
## Einrichtung
|
||||
|
||||
### 1. Virtual Environment anlegen
|
||||
```bash
|
||||
cd /home/ghost/source/repos/ghostnet-openclaw/tools/router
|
||||
python3 -m venv .venv
|
||||
source .venv/bin/activate
|
||||
pip install -r app/requirements.txt
|
||||
```
|
||||
|
||||
### 2. VSCode Konfiguration
|
||||
Öffne den Ordner `tools/router/` in VSCode und installiere folgende Extensions:
|
||||
- Python (Microsoft)
|
||||
- Pylance
|
||||
- Black Formatter (optional)
|
||||
|
||||
### 3. Launch Configuration (.vscode/launch.json)
|
||||
```json
|
||||
{
|
||||
"version": "0.2.0",
|
||||
"configurations": [
|
||||
{
|
||||
"name": "Router debuggen",
|
||||
"type": "debugpy",
|
||||
"request": "launch",
|
||||
"module": "uvicorn",
|
||||
"args": [
|
||||
"app.main:app",
|
||||
"--host", "0.0.0.0",
|
||||
"--port", "3000",
|
||||
"--reload"
|
||||
],
|
||||
"jinja": true,
|
||||
"env": {
|
||||
"CONFIG_PATH": "config/gateway.json"
|
||||
},
|
||||
"cwd": "${workspaceFolder}"
|
||||
}
|
||||
]
|
||||
}
|
||||
```
|
||||
|
||||
### 4. Pylance Konfiguration (.vscode/settings.json)
|
||||
```json
|
||||
{
|
||||
"python.analysis.typeCheckingMode": "basic",
|
||||
"python.pythonPath": ".venv/bin/python",
|
||||
"python.defaultInterpreterPath": ".venv/bin/python"
|
||||
}
|
||||
```
|
||||
|
||||
### 5. Testen
|
||||
```bash
|
||||
# Server starten
|
||||
uvicorn app.main:app --host 0.0.0.0 --port 3000 --reload
|
||||
|
||||
# In anderem Terminal testen:
|
||||
curl http://localhost:3000/sse
|
||||
```
|
||||
|
||||
## Wichtige Pfade
|
||||
- Router Source: `app/`
|
||||
- Config: `config/gateway.json`
|
||||
- Log-Output: stdout (INFO level default)
|
||||
|
||||
## Bekannte Issues
|
||||
- Der Router benötigt laufende MCP-Skills (z.B. searchfetch), damit er Tools registrieren kann
|
||||
- Ohne config/mcpServers ist der Router nutzlos (keine Tools registriert)
|
||||
- `config/gateway.json` muss URL zu einem MCP-fähigen Skill enthalten
|
||||
@@ -0,0 +1,36 @@
|
||||
from fastapi import Request
|
||||
|
||||
|
||||
def register_routes(
|
||||
app,
|
||||
transport,
|
||||
server
|
||||
):
|
||||
|
||||
@app.get("/sse")
|
||||
async def sse(
|
||||
request: Request
|
||||
):
|
||||
|
||||
async with transport.connect_sse(
|
||||
request.scope,
|
||||
request.receive,
|
||||
request._send
|
||||
) as (read, write):
|
||||
|
||||
await server.run(
|
||||
read,
|
||||
write,
|
||||
server.create_initialization_options()
|
||||
)
|
||||
|
||||
@app.post("/messages")
|
||||
async def messages(
|
||||
request: Request
|
||||
):
|
||||
|
||||
await transport.handle_post_message(
|
||||
request.scope,
|
||||
request.receive,
|
||||
request._send
|
||||
)
|
||||
@@ -0,0 +1,24 @@
|
||||
import json
|
||||
import logging
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def load_gateway_config():
|
||||
|
||||
try:
|
||||
|
||||
with open(
|
||||
"/app/config/gateway.json",
|
||||
"r"
|
||||
) as f:
|
||||
|
||||
return json.load(f)
|
||||
|
||||
except Exception:
|
||||
|
||||
logger.exception(
|
||||
"Gateway Config konnte nicht geladen werden"
|
||||
)
|
||||
|
||||
return {"mcpServers": {}}
|
||||
@@ -0,0 +1,19 @@
|
||||
import logging
|
||||
import os
|
||||
import sys
|
||||
|
||||
|
||||
def configure_logging():
|
||||
|
||||
level = os.getenv(
|
||||
"LOG_LEVEL",
|
||||
"INFO"
|
||||
).upper()
|
||||
|
||||
logging.basicConfig(
|
||||
level=getattr(logging, level),
|
||||
format="%(asctime)s - %(levelname)s - %(message)s",
|
||||
handlers=[
|
||||
logging.StreamHandler(sys.stdout)
|
||||
]
|
||||
)
|
||||
@@ -0,0 +1,82 @@
|
||||
import asyncio
|
||||
from contextlib import asynccontextmanager
|
||||
|
||||
from fastapi import FastAPI
|
||||
|
||||
from mcp.server.sse import SseServerTransport
|
||||
|
||||
from app.config.loader import load_gateway_config
|
||||
from app.logging.setup import configure_logging
|
||||
|
||||
from app.registry.skill_registry import SkillRegistry
|
||||
|
||||
from app.services.skill_manager import SkillManager
|
||||
|
||||
from app.mcp.server import master_server
|
||||
from app.mcp.handlers import register_handlers
|
||||
|
||||
from app.api.routes import register_routes
|
||||
|
||||
configure_logging()
|
||||
|
||||
registry = SkillRegistry()
|
||||
|
||||
register_handlers(
|
||||
master_server,
|
||||
registry
|
||||
)
|
||||
|
||||
skill_manager = SkillManager(
|
||||
registry
|
||||
)
|
||||
|
||||
|
||||
@asynccontextmanager
|
||||
async def lifespan(app):
|
||||
|
||||
config = load_gateway_config()
|
||||
|
||||
tasks = []
|
||||
|
||||
for name, cfg in config.get(
|
||||
"mcpServers",
|
||||
{}
|
||||
).items():
|
||||
|
||||
if "url" not in cfg:
|
||||
continue
|
||||
|
||||
tasks.append(
|
||||
asyncio.create_task(
|
||||
skill_manager.connect_skill(
|
||||
name,
|
||||
cfg["url"]
|
||||
)
|
||||
)
|
||||
)
|
||||
|
||||
yield
|
||||
|
||||
for task in tasks:
|
||||
task.cancel()
|
||||
|
||||
await asyncio.gather(
|
||||
*tasks,
|
||||
return_exceptions=True
|
||||
)
|
||||
|
||||
|
||||
app = FastAPI(
|
||||
title="GhostNet-Router",
|
||||
lifespan=lifespan
|
||||
)
|
||||
|
||||
transport = SseServerTransport(
|
||||
"/messages"
|
||||
)
|
||||
|
||||
register_routes(
|
||||
app,
|
||||
transport,
|
||||
master_server
|
||||
)
|
||||
@@ -0,0 +1,43 @@
|
||||
import mcp.types as types
|
||||
|
||||
|
||||
def register_handlers(
|
||||
server,
|
||||
registry
|
||||
):
|
||||
|
||||
@server.list_tools()
|
||||
async def list_tools():
|
||||
|
||||
tools = []
|
||||
|
||||
for skill_tools in registry.cached_tools.values():
|
||||
tools.extend(skill_tools)
|
||||
|
||||
return tools
|
||||
|
||||
@server.call_tool()
|
||||
async def call_tool(
|
||||
name,
|
||||
arguments
|
||||
):
|
||||
|
||||
session = await registry.get_session_for_tool(
|
||||
name
|
||||
)
|
||||
|
||||
if not session:
|
||||
|
||||
return [
|
||||
types.TextContent(
|
||||
type="text",
|
||||
text=f"Tool '{name}' offline"
|
||||
)
|
||||
]
|
||||
|
||||
result = await session.call_tool(
|
||||
name,
|
||||
arguments or {}
|
||||
)
|
||||
|
||||
return result.content
|
||||
@@ -0,0 +1,5 @@
|
||||
from mcp.server import Server
|
||||
|
||||
master_server = Server(
|
||||
"GhostNet-Router"
|
||||
)
|
||||
@@ -0,0 +1,56 @@
|
||||
import asyncio
|
||||
|
||||
class SkillRegistry:
|
||||
|
||||
def __init__(self):
|
||||
|
||||
self.sessions = {}
|
||||
self.cached_tools = {}
|
||||
self.tool_map = {}
|
||||
|
||||
self.lock = asyncio.Lock()
|
||||
|
||||
async def register_skill(
|
||||
self,
|
||||
name,
|
||||
session,
|
||||
tools
|
||||
):
|
||||
|
||||
async with self.lock:
|
||||
|
||||
self.sessions[name] = session
|
||||
self.cached_tools[name] = tools
|
||||
|
||||
for tool in tools:
|
||||
self.tool_map[tool.name] = name
|
||||
|
||||
async def unregister_skill(
|
||||
self,
|
||||
name
|
||||
):
|
||||
|
||||
async with self.lock:
|
||||
|
||||
self.sessions.pop(name, None)
|
||||
self.cached_tools.pop(name, None)
|
||||
|
||||
self.tool_map = {
|
||||
tool: skill
|
||||
for tool, skill in self.tool_map.items()
|
||||
if skill != name
|
||||
}
|
||||
|
||||
async def get_session_for_tool(
|
||||
self,
|
||||
tool_name
|
||||
):
|
||||
|
||||
async with self.lock:
|
||||
|
||||
skill = self.tool_map.get(tool_name)
|
||||
|
||||
if not skill:
|
||||
return None
|
||||
|
||||
return self.sessions.get(skill)
|
||||
@@ -0,0 +1,45 @@
|
||||
annotated-doc==0.0.4
|
||||
annotated-types==0.7.0
|
||||
anyio==4.13.0
|
||||
attrs==26.1.0
|
||||
backoff==2.2.1
|
||||
certifi==2026.5.20
|
||||
cffi==2.0.0
|
||||
charset-normalizer==3.4.7
|
||||
click==8.4.1
|
||||
cryptography==48.0.0
|
||||
distro==1.9.0
|
||||
fastapi==0.136.3
|
||||
h11==0.16.0
|
||||
httpcore==1.0.9
|
||||
httpx==0.28.1
|
||||
httpx-sse==0.4.3
|
||||
idna==3.18
|
||||
Jinja2==3.1.6
|
||||
jsonschema==4.26.0
|
||||
jsonschema-specifications==2025.9.1
|
||||
loguru==0.7.3
|
||||
markdown-it-py==4.2.0
|
||||
MarkupSafe==3.0.3
|
||||
mcp==1.27.2
|
||||
mdurl==0.1.2
|
||||
modelcontextprotocol==1.0.1
|
||||
posthog==7.17.0
|
||||
pycparser==3.0
|
||||
pydantic==2.13.4
|
||||
pydantic-settings==2.14.1
|
||||
pydantic_core==2.46.4
|
||||
Pygments==2.20.0
|
||||
PyJWT==2.13.0
|
||||
python-dotenv==1.2.2
|
||||
python-multipart==0.0.32
|
||||
referencing==0.37.0
|
||||
requests==2.34.2
|
||||
rich==15.0.0
|
||||
rpds-py==2026.5.1
|
||||
sse-starlette==3.4.4
|
||||
starlette==1.2.1
|
||||
typing-inspection==0.4.2
|
||||
typing_extensions==4.15.0
|
||||
urllib3==2.7.0
|
||||
uvicorn==0.49.0
|
||||
@@ -0,0 +1,73 @@
|
||||
from mcp import ClientSession
|
||||
from mcp.client.sse import sse_client
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class SkillManager:
|
||||
|
||||
def __init__(self, registry):
|
||||
|
||||
self.registry = registry
|
||||
|
||||
async def connect_skill(
|
||||
self,
|
||||
name,
|
||||
url
|
||||
):
|
||||
|
||||
while True:
|
||||
|
||||
try:
|
||||
|
||||
async with sse_client(url) as (
|
||||
read_stream,
|
||||
write_stream
|
||||
):
|
||||
|
||||
async with ClientSession(
|
||||
read_stream,
|
||||
write_stream
|
||||
) as session:
|
||||
|
||||
await session.initialize()
|
||||
|
||||
result = await session.list_tools()
|
||||
|
||||
await self.registry.register_skill(
|
||||
name,
|
||||
session,
|
||||
result.tools
|
||||
)
|
||||
|
||||
logger.info(
|
||||
f"{name}: {len(result.tools)} Tools"
|
||||
)
|
||||
|
||||
# Der sse_client-Context-Manager hält die Verbindung offen.
|
||||
# Solange dieser Block nicht verlassen wird, lebt die SSE-Verbindung.
|
||||
# Bei einem Disconnect (z.B. Server-Error) wird der Block verlassen.
|
||||
task = asyncio.current_task()
|
||||
while task is not None and not task.cancelled():
|
||||
await asyncio.sleep(30)
|
||||
task = asyncio.current_task()
|
||||
|
||||
except asyncio.CancelledError:
|
||||
raise
|
||||
|
||||
except Exception:
|
||||
|
||||
logger.exception(
|
||||
f"{name} disconnected"
|
||||
)
|
||||
|
||||
finally:
|
||||
|
||||
await self.registry.unregister_skill(
|
||||
name
|
||||
)
|
||||
|
||||
await asyncio.sleep(5)
|
||||
@@ -1,135 +0,0 @@
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import sys
|
||||
from contextlib import asynccontextmanager
|
||||
|
||||
from fastapi import FastAPI, Request
|
||||
from mcp import ClientSession
|
||||
from mcp.client.sse import sse_client
|
||||
from mcp.server import Server
|
||||
from mcp.server.sse import SseServerTransport
|
||||
import mcp.types as types
|
||||
|
||||
# --- LOGGING ---
|
||||
log_level_str = os.getenv("LOG_LEVEL", "INFO").upper()
|
||||
log_level = getattr(logging, log_level_str, logging.INFO)
|
||||
|
||||
logging.basicConfig(
|
||||
level=log_level,
|
||||
format="%(asctime)s - %(name)s - %(levelname)s - %(message)s",
|
||||
handlers=[logging.StreamHandler(sys.stdout)]
|
||||
)
|
||||
logger = logging.getLogger("ghostnet-router")
|
||||
|
||||
# --- GLOBAL STATE ---
|
||||
# Hier speichern wir die aktiven Sessions zu den Microservices (z.B. searchfetch)
|
||||
skill_sessions: dict[str, ClientSession] = {}
|
||||
master_server = Server("GhostNet-Router")
|
||||
|
||||
# --- LIFESPAN (Connection Management) ---
|
||||
@asynccontextmanager
|
||||
async def lifespan(app: FastAPI):
|
||||
"""Verwaltet den Lebenszyklus der Verbindungen zu den Microservices."""
|
||||
config_path = "/app/config/gateway.json"
|
||||
|
||||
try:
|
||||
with open(config_path, "r") as f:
|
||||
config = json.load(f)
|
||||
except Exception as e:
|
||||
logger.error(f"❌ Failed to load config from {config_path}: {e}")
|
||||
config = {"mcpServers": {}}
|
||||
|
||||
# Hintergrund-Tasks für jeden konfigurierten Service starten
|
||||
tasks = []
|
||||
for name, info in config.get("mcpServers", {}).items():
|
||||
url = info.get("url")
|
||||
if url:
|
||||
tasks.append(asyncio.create_task(maintain_connection(name, url)))
|
||||
|
||||
yield # Der Router ist jetzt bereit und bedient Anfragen
|
||||
|
||||
# Cleanup beim Herunterfahren
|
||||
logger.info("Closing all skill connections...")
|
||||
for task in tasks:
|
||||
task.cancel()
|
||||
logger.info("Router shutdown complete.")
|
||||
|
||||
async def maintain_connection(name: str, url: str):
|
||||
while True:
|
||||
try:
|
||||
logger.info(f"🔗 Attempting to connect to Skill '{name}' at {url}")
|
||||
# Wir nutzen einen ContextManager, der die Verbindung offen hält
|
||||
async with sse_client(url) as (read_stream, write_stream):
|
||||
async with ClientSession(read_stream, write_stream) as session:
|
||||
# 1. Handshake
|
||||
await session.initialize()
|
||||
|
||||
# 2. In den globalen State schreiben
|
||||
skill_sessions[name] = session
|
||||
logger.info(f"✅ Skill '{name}' successfully integrated.")
|
||||
|
||||
# 3. WICHTIG: Wir müssen hier blockieren, damit wir den
|
||||
# Context (async with) NICHT verlassen.
|
||||
# wait_until_disconnected() sollte eigentlich funktionieren,
|
||||
# aber wir können es mit einem asyncio.Event sicherer machen:
|
||||
stop_event = asyncio.Event()
|
||||
|
||||
# Optional: Ein kleiner Loop, der prüft, ob die Session noch lebt
|
||||
while not stop_event.is_set():
|
||||
await asyncio.sleep(1)
|
||||
# Hier könnte man einen Heartbeat prüfen, falls das SDK das bietet
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"⚠️ Connection lost to '{name}': {e}. Retrying in 5s...")
|
||||
finally:
|
||||
skill_sessions.pop(name, None)
|
||||
await asyncio.sleep(5)
|
||||
|
||||
# --- MCP SERVER LOGIC ---
|
||||
|
||||
@master_server.list_tools()
|
||||
async def handle_list_tools() -> list[types.Tool]:
|
||||
"""Bündelt die Tool-Listen aller aktuell verbundenen Microservices."""
|
||||
all_tools = []
|
||||
for name, session in skill_sessions.items():
|
||||
try:
|
||||
result = await session.list_tools()
|
||||
all_tools.extend(result.tools)
|
||||
except Exception as e:
|
||||
logger.error(f"Could not list tools for {name}: {e}")
|
||||
return all_tools
|
||||
|
||||
@master_server.call_tool()
|
||||
async def handle_call_tool(name: str, arguments: dict | None) -> list[types.TextContent]:
|
||||
"""Sucht das angeforderte Tool in den verbundenen Services und leitet den Call weiter."""
|
||||
for skill_name, session in skill_sessions.items():
|
||||
# Wir prüfen, ob dieser Service das gesuchte Tool anbietet
|
||||
tools_result = await session.list_tools()
|
||||
if any(t.name == name for t in tools_result.tools):
|
||||
logger.info(f"⚡ Relaying tool call '{name}' to Skill '{skill_name}'")
|
||||
result = await session.call_tool(name, arguments or {})
|
||||
return result.content
|
||||
|
||||
raise ValueError(f"Tool '{name}' not found in any connected GhostNet Skill.")
|
||||
|
||||
# --- FASTAPI & SSE TRANSPORT ---
|
||||
|
||||
app = FastAPI(title="GhostNet-Router", lifespan=lifespan)
|
||||
sse = SseServerTransport("/messages")
|
||||
|
||||
@app.get("/sse")
|
||||
async def handle_sse(request: Request):
|
||||
"""SSE-Endpunkt für den Agenten (OpenClaw)."""
|
||||
async with sse.connect_sse(request.scope, request.receive, request._send) as (read_stream, write_stream):
|
||||
await master_server.run(
|
||||
read_stream,
|
||||
write_stream,
|
||||
master_server.create_initialization_options()
|
||||
)
|
||||
|
||||
@app.post("/messages")
|
||||
async def handle_messages(request: Request):
|
||||
"""Post-Messages-Endpunkt für den bidirektionalen MCP-Austausch."""
|
||||
await sse.handle_post_message(request.scope, request.receive, request._send)
|
||||
@@ -0,0 +1,17 @@
|
||||
FROM docker.io/nikolaik/python-nodejs:python3.14-nodejs26-slim
|
||||
|
||||
# System dependencies
|
||||
RUN apt-get update && apt-get install -y curl && rm -rf /var/lib/apt/lists/*
|
||||
|
||||
# Application dependencies
|
||||
RUN pip install --no-cache-dir mcp-proxy
|
||||
RUN npm install -g mcp-searxng
|
||||
|
||||
# Default environment variables
|
||||
ENV SEARXNG_URL=http://tool-searchfetch-searxng:8080
|
||||
ENV MCP_HOST=0.0.0.0
|
||||
ENV MCP_PORT=3000
|
||||
|
||||
ENTRYPOINT ["mcp-proxy"]
|
||||
|
||||
CMD ["--host", "0.0.0.0", "--port", "3000", "--pass-environment", "--", "mcp-searxng"]
|
||||
@@ -4,7 +4,8 @@ server:
|
||||
bind_address: 0.0.0.0
|
||||
port: 8080
|
||||
secret_key: "48d28819fa08ae64546d354e72dbc6c7a1e8fc5d2ef69be0099dd38d2f22d8d7"
|
||||
limiter: false
|
||||
|
||||
limiter: false
|
||||
|
||||
valkey:
|
||||
url: "redis://tool-searchfetch-valkey:6379/0"
|
||||
|
||||
Reference in New Issue
Block a user