Coverage for /home/admin/Documents/AI/applications/lexigram-dev/lexigram/experimental/ai/lexigram-ai-mcp/src/lexigram/ai/mcp/transport/stdio.py: 27%

60 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-25 07:19 +0800

1"""Stdio transport for MCP server.""" 

2 

3from __future__ import annotations 

4 

5import asyncio 

6import sys 

7from typing import Any 

8 

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

10from lexigram.contracts.mcp.exceptions import MCPTransportError 

11from lexigram.logging import ( 

12 get_logger, 

13) 

14from lexigram.serialization import JSONDecodeError, dumps_str, loads 

15 

16logger = get_logger(__name__) 

17 

18 

19class StdioTransport(AbstractTransport): 

20 """Stdio-based transport for MCP server. 

21 

22 This transport reads JSON-RPC messages from stdin and writes 

23 responses to stdout. Ideal for CLI/desktop integration. 

24 """ 

25 

26 def __init__( 

27 self, 

28 reader: asyncio.StreamReader | None = None, 

29 writer: asyncio.StreamWriter | None = None, 

30 ) -> None: 

31 """Initialize the stdio transport. 

32 

33 Args: 

34 reader: Optional stream reader (for testing). 

35 writer: Optional stream writer (for testing). 

36 """ 

37 self._reader = reader 

38 self._writer = writer 

39 self._running = False 

40 

41 async def start(self) -> None: 

42 """Start the stdio transport.""" 

43 if self._running: 

44 return 

45 

46 if self._reader is None: 

47 self._reader = asyncio.StreamReader() 

48 loop = asyncio.get_event_loop() 

49 protocol = asyncio.StreamReaderProtocol(self._reader) 

50 await loop.connect_read_pipe(lambda: protocol, sys.stdin) 

51 

52 if self._writer is None: 

53 _transport, self._writer = await asyncio.open_connection( 

54 sys.stdin.fileno(), # type: ignore[arg-type] 

55 sys.stdout.fileno(), 

56 ) 

57 

58 self._running = True 

59 logger.info("mcp_stdio_transport_started") 

60 

61 async def stop(self) -> None: 

62 """Stop the stdio transport.""" 

63 if not self._running: 

64 return 

65 

66 if self._writer: 

67 self._writer.close() 

68 await self._writer.wait_closed() 

69 

70 self._running = False 

71 logger.info("mcp_stdio_transport_stopped") 

72 

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

74 """Send a message through stdio. 

75 

76 Args: 

77 message: JSON-RPC message to send. 

78 

79 Raises: 

80 MCPTransportError: If sending fails. 

81 """ 

82 if not self._running or self._writer is None: 

83 raise MCPTransportError( 

84 message="Transport not started", 

85 transport_type="stdio", 

86 ) 

87 

88 try: 

89 data = dumps_str(message) + "\n" 

90 self._writer.write(data.encode()) 

91 await self._writer.drain() 

92 except (OSError, RuntimeError, AttributeError, TypeError, ConnectionError) as e: 

93 raise MCPTransportError( 

94 message=f"Failed to send message: {e!s}", 

95 transport_type="stdio", 

96 ) from e 

97 

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

99 """Receive a message from stdio. 

100 

101 Returns: 

102 Parsed JSON-RPC message, or None if no data available. 

103 

104 Raises: 

105 MCPTransportError: If receiving fails. 

106 """ 

107 if not self._running or self._reader is None: 

108 raise MCPTransportError( 

109 message="Transport not started", 

110 transport_type="stdio", 

111 ) 

112 

113 try: 

114 # Read a line from stdin 

115 line = await self._reader.readline() 

116 if not line: 

117 return None 

118 

119 decoded = line.decode() 

120 stripped = decoded.lstrip() 

121 if not stripped.startswith("{"): 

122 raise MCPTransportError( 

123 message="Invalid JSON: expected object payload", 

124 transport_type="stdio", 

125 ) 

126 return loads(decoded) 

127 

128 except JSONDecodeError as e: 

129 raise MCPTransportError( 

130 message=f"Invalid JSON: {e!s}", 

131 transport_type="stdio", 

132 ) from e 

133 except ( 

134 OSError, 

135 RuntimeError, 

136 AttributeError, 

137 TypeError, 

138 ConnectionError, 

139 UnicodeDecodeError, 

140 ) as e: 

141 raise MCPTransportError( 

142 message=f"Failed to receive message: {e!s}", 

143 transport_type="stdio", 

144 ) from e 

145 

146 

147__all__ = ["StdioTransport"]