74 lines
1.9 KiB
Python
74 lines
1.9 KiB
Python
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)
|