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
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-25 07:19 +0800
1"""Stdio transport for MCP server."""
3from __future__ import annotations
5import asyncio
6import sys
7from typing import Any
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
16logger = get_logger(__name__)
19class StdioTransport(AbstractTransport):
20 """Stdio-based transport for MCP server.
22 This transport reads JSON-RPC messages from stdin and writes
23 responses to stdout. Ideal for CLI/desktop integration.
24 """
26 def __init__(
27 self,
28 reader: asyncio.StreamReader | None = None,
29 writer: asyncio.StreamWriter | None = None,
30 ) -> None:
31 """Initialize the stdio transport.
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
41 async def start(self) -> None:
42 """Start the stdio transport."""
43 if self._running:
44 return
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)
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 )
58 self._running = True
59 logger.info("mcp_stdio_transport_started")
61 async def stop(self) -> None:
62 """Stop the stdio transport."""
63 if not self._running:
64 return
66 if self._writer:
67 self._writer.close()
68 await self._writer.wait_closed()
70 self._running = False
71 logger.info("mcp_stdio_transport_stopped")
73 async def send(self, message: dict[str, Any]) -> None:
74 """Send a message through stdio.
76 Args:
77 message: JSON-RPC message to send.
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 )
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
98 async def receive(self) -> dict[str, Any] | None:
99 """Receive a message from stdio.
101 Returns:
102 Parsed JSON-RPC message, or None if no data available.
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 )
113 try:
114 # Read a line from stdin
115 line = await self._reader.readline()
116 if not line:
117 return None
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)
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
147__all__ = ["StdioTransport"]