jarvis_mcp/client.py aktualisiert
This commit is contained in:
+21
-37
@@ -1,5 +1,7 @@
|
|||||||
# jarvis_mcp/client.py
|
# jarvis_mcp/client.py
|
||||||
import asyncio
|
import asyncio
|
||||||
|
import os
|
||||||
|
from contextlib import AsyncExitStack
|
||||||
from mcp import ClientSession
|
from mcp import ClientSession
|
||||||
from mcp.client.streamable_http import streamablehttp_client
|
from mcp.client.streamable_http import streamablehttp_client
|
||||||
from util.logger import get_logger
|
from util.logger import get_logger
|
||||||
@@ -8,44 +10,32 @@ logger = get_logger("MCP")
|
|||||||
|
|
||||||
|
|
||||||
class MCPClient:
|
class MCPClient:
|
||||||
def __init__(self, url: str):
|
def __init__(self, url: str = "http://localhost:8081"):
|
||||||
# Entferne alles nach /mcp und füge exakt /mcp/ hinzu
|
self.base_url = url.rstrip('/')
|
||||||
base = url.rstrip('/').split('/mcp')[0]
|
|
||||||
self.url = base + "/mcp/" # ← WICHTIG: Slash am Ende
|
|
||||||
|
|
||||||
self.session = None
|
self.session = None
|
||||||
self._client_context = None
|
self._exit_stack = AsyncExitStack()
|
||||||
|
|
||||||
async def connect(self):
|
async def connect(self):
|
||||||
if self.session is not None:
|
if self.session is not None:
|
||||||
return
|
return
|
||||||
|
|
||||||
logger.info(f"Verbinde mit OpenHAB MCP Server: {self.url}")
|
mcp_url = f"{self.base_url}/mcp/"
|
||||||
|
logger.info(f"Verbinde mit OpenHAB MCP: {mcp_url}")
|
||||||
|
|
||||||
try:
|
try:
|
||||||
self._client_context = streamablehttp_client(
|
# Context Manager wie im offiziellen Beispiel
|
||||||
self.url,
|
self._client_cm = streamablehttp_client(mcp_url)
|
||||||
terminate_on_close=True
|
read, write, _ = await self._exit_stack.enter_async_context(self._client_cm)
|
||||||
|
|
||||||
|
self.session = await self._exit_stack.enter_async_context(
|
||||||
|
ClientSession(read, write)
|
||||||
)
|
)
|
||||||
|
|
||||||
streams = await self._client_context.__aenter__()
|
await self.session.initialize()
|
||||||
|
logger.info("✅ OpenHAB MCP Verbindung erfolgreich hergestellt")
|
||||||
|
|
||||||
# Flexibles Unpacking
|
|
||||||
if isinstance(streams, tuple):
|
|
||||||
read, write = streams[:2] # nimm nur die ersten zwei
|
|
||||||
else:
|
|
||||||
raise ValueError("Unerwartetes Response-Format vom Client")
|
|
||||||
|
|
||||||
self.session = ClientSession(read, write)
|
|
||||||
await asyncio.wait_for(self.session.initialize(), timeout=12.0)
|
|
||||||
|
|
||||||
logger.info("✅ OpenHAB MCP Verbindung erfolgreich hergestellt!")
|
|
||||||
|
|
||||||
except asyncio.TimeoutError:
|
|
||||||
logger.error("❌ Timeout - MCP Server antwortet nicht rechtzeitig")
|
|
||||||
raise
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"❌ MCP Verbindungsfehler: {e}")
|
logger.error(f"❌ MCP Connect Fehler: {e}")
|
||||||
raise
|
raise
|
||||||
|
|
||||||
async def list_tools(self):
|
async def list_tools(self):
|
||||||
@@ -57,15 +47,9 @@ class MCPClient:
|
|||||||
return await self.session.call_tool(name, arguments)
|
return await self.session.call_tool(name, arguments)
|
||||||
|
|
||||||
async def close(self):
|
async def close(self):
|
||||||
for obj in [self.session, self._client_context]:
|
try:
|
||||||
if obj:
|
await self._exit_stack.aclose()
|
||||||
try:
|
except Exception as e:
|
||||||
if hasattr(obj, "close"):
|
logger.warning(f"Close Warning: {e}")
|
||||||
await obj.close()
|
|
||||||
elif hasattr(obj, "__aexit__"):
|
|
||||||
await obj.__aexit__(None, None, None)
|
|
||||||
except Exception as e:
|
|
||||||
logger.warning(f"Close Warning: {e}")
|
|
||||||
|
|
||||||
self.session = None
|
self.session = None
|
||||||
self._client_context = None
|
logger.info("MCP Verbindung geschlossen")
|
||||||
Reference in New Issue
Block a user