diff --git a/Sources/MCP/Base/Transports/StdioTransport.swift b/Sources/MCP/Base/Transports/StdioTransport.swift index 45522fac..17baa952 100644 --- a/Sources/MCP/Base/Transports/StdioTransport.swift +++ b/Sources/MCP/Base/Transports/StdioTransport.swift @@ -56,6 +56,7 @@ import struct Foundation.Data private var isConnected = false private let messageStream: AsyncThrowingStream private let messageContinuation: AsyncThrowingStream.Continuation + private var lastSend: Task? /// Creates a new stdio transport with the specified file descriptors /// @@ -195,6 +196,18 @@ import struct Foundation.Data /// - Parameter message: The message data to send (without a trailing newline) /// - Throws: Error if the message cannot be sent public func send(_ message: Data) async throws { + let previousSend = lastSend + let currentSend = Task { + await previousSend?.value + try await write(message) + } + lastSend = Task { + try? await currentSend.value + } + try await currentSend.value + } + + private func write(_ message: Data) async throws { guard isConnected else { throw MCPError.transportError(Errno(rawValue: ENOTCONN)) } diff --git a/Tests/MCPTests/StdioTransportTests.swift b/Tests/MCPTests/StdioTransportTests.swift index e395fbe0..5fe989f1 100644 --- a/Tests/MCPTests/StdioTransportTests.swift +++ b/Tests/MCPTests/StdioTransportTests.swift @@ -43,6 +43,50 @@ struct StdioTransportTests { await transport.disconnect() } + @Test("Concurrent sends preserve message framing under backpressure") + func testConcurrentSendsPreserveMessageFramingUnderBackpressure() async throws { + let (reader, output) = try FileDescriptor.pipe() + let (input, _) = try FileDescriptor.pipe() + let transport = StdioTransport(input: input, output: output, logger: nil) + try await transport.connect() + + let largeMessage = Data(repeating: UInt8(ascii: "a"), count: 512 * 1024) + let smallMessage = Data(#"{"id":2}"#.utf8) + let expected = + largeMessage + Data([UInt8(ascii: "\n")]) + + smallMessage + Data([UInt8(ascii: "\n")]) + + let firstSend = Task { + try await transport.send(largeMessage) + } + + // Leave the pipe undrained until the first send has filled its buffer + // and suspended in the EAGAIN retry path. + try await Task.sleep(for: .milliseconds(50)) + + let secondSend = Task { + try await transport.send(smallMessage) + } + + let received = try await Task.detached { + var received = Data() + var buffer = [UInt8](repeating: 0, count: 4096) + while received.count < expected.count { + let count = try buffer.withUnsafeMutableBufferPointer { pointer in + try reader.read(into: UnsafeMutableRawBufferPointer(pointer)) + } + received.append(contentsOf: buffer[..