Coverage for /home/admin/Documents/AI/applications/lexigram-dev/lexigram/experimental/ai/lexigram-ai-mcp/src/lexigram/ai/mcp/transport/sse.py: 42%
33 statements
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-25 07:19 +0800
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-25 07:19 +0800
1"""HTTP+SSE transport for MCP server."""
3from __future__ import annotations
5from typing import Any
7from lexigram.ai.mcp.transport.base import AbstractTransport
8from lexigram.contracts.mcp.exceptions import MCPTransportError
9from lexigram.logging import (
10 get_logger,
11)
13logger = get_logger(__name__)
16class SSETransport(AbstractTransport):
17 """Server-Sent Events transport for MCP server.
19 This transport uses HTTP POST for requests and SSE for streaming
20 responses. Ideal for web-based integrations.
21 """
23 def __init__(
24 self,
25 server: Any | None = None,
26 ) -> None:
27 """Initialize the SSE transport.
29 Args:
30 server: Optional HTTP server instance.
31 """
32 self._server = server
33 self._running = False
34 self._message_queue: list[dict[str, Any]] = []
36 async def start(self) -> None:
37 """Start the SSE transport."""
38 if self._running:
39 return
41 self._running = True
42 logger.info("mcp_sse_transport_started")
44 async def stop(self) -> None:
45 """Stop the SSE transport."""
46 if not self._running:
47 return
49 self._running = False
50 logger.info("mcp_sse_transport_stopped")
52 async def send(self, message: dict[str, Any]) -> None:
53 """Send a message via SSE.
55 Args:
56 message: JSON-RPC message to send.
58 Raises:
59 MCPTransportError: If sending fails.
60 """
61 if not self._running:
62 raise MCPTransportError(
63 message="Transport not started",
64 transport_type="sse",
65 )
67 # Queue the message for SSE delivery
68 self._message_queue.append(message)
69 logger.debug("mcp_sse_message_queued", message_id=message.get("id"))
71 async def receive(self) -> dict[str, Any] | None:
72 """Receive is not applicable for SSE (pull-based).
74 For HTTP+SSE, use the HTTP endpoint directly instead.
76 Returns:
77 None (receiving is handled via HTTP POST).
78 """
79 # For SSE transport, receiving happens via HTTP
80 # This method is not used
81 return None
83 def get_queued_messages(self) -> list[dict[str, Any]]:
84 """Get queued messages for SSE delivery.
86 Returns:
87 List of queued JSON-RPC messages.
88 """
89 messages = self._message_queue.copy()
90 self._message_queue.clear()
91 return messages
94__all__ = ["SSETransport"]