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

1"""HTTP+SSE transport for MCP server.""" 

2 

3from __future__ import annotations 

4 

5from typing import Any 

6 

7from lexigram.ai.mcp.transport.base import AbstractTransport 

8from lexigram.contracts.mcp.exceptions import MCPTransportError 

9from lexigram.logging import ( 

10 get_logger, 

11) 

12 

13logger = get_logger(__name__) 

14 

15 

16class SSETransport(AbstractTransport): 

17 """Server-Sent Events transport for MCP server. 

18 

19 This transport uses HTTP POST for requests and SSE for streaming 

20 responses. Ideal for web-based integrations. 

21 """ 

22 

23 def __init__( 

24 self, 

25 server: Any | None = None, 

26 ) -> None: 

27 """Initialize the SSE transport. 

28 

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]] = [] 

35 

36 async def start(self) -> None: 

37 """Start the SSE transport.""" 

38 if self._running: 

39 return 

40 

41 self._running = True 

42 logger.info("mcp_sse_transport_started") 

43 

44 async def stop(self) -> None: 

45 """Stop the SSE transport.""" 

46 if not self._running: 

47 return 

48 

49 self._running = False 

50 logger.info("mcp_sse_transport_stopped") 

51 

52 async def send(self, message: dict[str, Any]) -> None: 

53 """Send a message via SSE. 

54 

55 Args: 

56 message: JSON-RPC message to send. 

57 

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 ) 

66 

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")) 

70 

71 async def receive(self) -> dict[str, Any] | None: 

72 """Receive is not applicable for SSE (pull-based). 

73 

74 For HTTP+SSE, use the HTTP endpoint directly instead. 

75 

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 

82 

83 def get_queued_messages(self) -> list[dict[str, Any]]: 

84 """Get queued messages for SSE delivery. 

85 

86 Returns: 

87 List of queued JSON-RPC messages. 

88 """ 

89 messages = self._message_queue.copy() 

90 self._message_queue.clear() 

91 return messages 

92 

93 

94__all__ = ["SSETransport"]