From aa354a1bb054e7ffea37497e1abc93c5a0cc1e6a Mon Sep 17 00:00:00 2001 From: Widthdom Date: Sun, 24 May 2026 23:33:56 +0900 Subject: [PATCH 1/3] Fix MCP response write ordering (#1927 #2020 #2021) --- changelog.d/unreleased/1927.fixed.md | 19 +++++ src/CodeIndex/Mcp/McpServer.cs | 49 +++++++++++- src/CodeIndex/Mcp/StdioMcpTransport.cs | 1 + tests/CodeIndex.Tests/McpServerTests.cs | 101 ++++++++++++++++++++++++ 4 files changed, 166 insertions(+), 4 deletions(-) create mode 100644 changelog.d/unreleased/1927.fixed.md diff --git a/changelog.d/unreleased/1927.fixed.md b/changelog.d/unreleased/1927.fixed.md new file mode 100644 index 0000000000..523f439348 --- /dev/null +++ b/changelog.d/unreleased/1927.fixed.md @@ -0,0 +1,19 @@ +--- +category: fixed +issues: + - 1927 + - 2020 + - 2021 +affected: + - src/CodeIndex/Mcp/McpServer.cs + - src/CodeIndex/Mcp/StdioMcpTransport.cs + - tests/CodeIndex.Tests/McpServerTests.cs +--- + +## English + +- **MCP error responses now reach the transport before diagnostic logs (#1927, #2020, #2021)** — stdio writes explicitly flush before the server reads the next frame, and parse/error diagnostics are emitted after the response write attempt so clients receive JSON-RPC failures before operators see loop errors. + +## 日本語 + +- **MCP のエラー応答が診断ログより先に transport へ届くようになりました (#1927, #2020, #2021)** — stdio 書き込みは次のフレームを読む前に明示的に flush し、parse/error 診断は応答書き込みの試行後に出力するため、クライアントは loop error のログより先に JSON-RPC failure を受け取れます。 diff --git a/src/CodeIndex/Mcp/McpServer.cs b/src/CodeIndex/Mcp/McpServer.cs index 9346b530bf..0a23bdd8c3 100644 --- a/src/CodeIndex/Mcp/McpServer.cs +++ b/src/CodeIndex/Mcp/McpServer.cs @@ -60,6 +60,7 @@ public partial class McpServer : IDisposable // を渡せるようにするため (#1567)。 private readonly AsyncLocal _currentRequestToken = new(); private readonly AsyncLocal?> _currentOutOfBandFrameWriter = new(); + private readonly AsyncLocal?> _deferredFrameLogs = new(); private bool _running = true; // Per-session DbContext reused across MCP tool calls. Holding the connection open // avoids reopening SQLite, reapplying pragmas, and re-registering every SQL function @@ -437,6 +438,7 @@ internal async Task RunAsync(IMcpTransport transport, CancellationToken cancella _currentOutOfBandFrameWriter.Value = transport is IOutOfBandMcpTransport outOfBandTransport ? frameToWrite => outOfBandTransport.WriteOutOfBandFrameAsync(frameToWrite, loopToken).GetAwaiter().GetResult() : null; + BeginDeferredFrameLogs(); response = ProcessFrame(frame); } finally @@ -447,6 +449,7 @@ internal async Task RunAsync(IMcpTransport transport, CancellationToken cancella } await WriteFrameSafelyAsync(transport, response, loopToken).ConfigureAwait(false); + FlushDeferredFrameLogs(); // `notifications/shutdown` flips `_running` inside `HandleMessage`; exit the loop // immediately so a subsequent slow `ReadFrameAsync` does not extend the lifetime @@ -461,7 +464,9 @@ internal async Task RunAsync(IMcpTransport transport, CancellationToken cancella } catch (DecoderFallbackException ex) { + BeginDeferredFrameLogs(); await WriteFrameSafelyAsync(transport, BuildInvalidUtf8ParseErrorResponse(ex), loopToken).ConfigureAwait(false); + FlushDeferredFrameLogs(); break; } } @@ -498,7 +503,9 @@ private async Task RunConcurrentFrameLoopAsync(IMcpTransport transport, Cancella await writeGate.WaitAsync(loopToken).ConfigureAwait(false); try { + BeginDeferredFrameLogs(); await WriteFrameSafelyAsync(transport, BuildInvalidUtf8ParseErrorResponse(ex), loopToken).ConfigureAwait(false); + FlushDeferredFrameLogs(); } finally { @@ -511,11 +518,13 @@ private async Task RunConcurrentFrameLoopAsync(IMcpTransport transport, Cancella if (IsCancellationFrame(frame)) { + BeginDeferredFrameLogs(); var response = ProcessFrame(frame); await writeGate.WaitAsync(loopToken).ConfigureAwait(false); try { await WriteFrameSafelyAsync(transport, response, loopToken).ConfigureAwait(false); + FlushDeferredFrameLogs(); } finally { @@ -546,6 +555,7 @@ private async Task RunConcurrentFrameLoopAsync(IMcpTransport transport, Cancella writeGate.Release(); } }; + BeginDeferredFrameLogs(); response = ProcessFrame(frame); } finally @@ -559,6 +569,7 @@ private async Task RunConcurrentFrameLoopAsync(IMcpTransport transport, Cancella try { await WriteFrameSafelyAsync(transport, response, loopToken).ConfigureAwait(false); + FlushDeferredFrameLogs(); } finally { @@ -586,16 +597,19 @@ private async Task RunConcurrentFrameLoopAsync(IMcpTransport transport, Cancella /// internal async Task ProcessLineAsync(string line, TextWriter writer) { + BeginDeferredFrameLogs(); var response = ProcessFrame(line); if (response != null) { try { await WriteJsonLineAsync(writer, response).ConfigureAwait(false); + FlushDeferredFrameLogs(); } catch (Exception ex) when (ex is IOException or ObjectDisposedException or OperationCanceledException) { Console.Error.WriteLine(BuildResponseWriteErrorLog(ex.Message)); + FlushDeferredFrameLogs(); } } } @@ -604,6 +618,7 @@ private static async Task WriteJsonLineAsync(TextWriter writer, string response) { await writer.WriteAsync(response).ConfigureAwait(false); await writer.WriteAsync('\n').ConfigureAwait(false); + await writer.FlushAsync().ConfigureAwait(false); } private static async Task WriteFrameSafelyAsync(IMcpTransport transport, string? response, CancellationToken cancellationToken) @@ -624,7 +639,7 @@ private static async Task WriteFrameSafelyAsync(IMcpTransport transport, string? private string BuildInvalidUtf8ParseErrorResponse(DecoderFallbackException ex) { - Console.Error.WriteLine(BuildInvalidUtf8ErrorLog(ex.Message)); + DeferFrameLog(BuildInvalidUtf8ErrorLog(ex.Message)); var errorResponse = CreateErrorResponse(hasId: true, id: null, code: -32700, message: "Parse error: invalid UTF-8 input", category: McpErrorEnvelope.CategoryParseError, suggestion: "Send one JSON-RPC 2.0 object per line encoded as valid UTF-8. Reject or re-encode malformed bytes before retrying.", @@ -675,7 +690,7 @@ internal static string BuildInvalidUtf8ErrorLog(string detail) catch (JsonException ex) { // Parse error / パースエラー - Console.Error.WriteLine(BuildJsonParseErrorLog(ex.Message)); + DeferFrameLog(BuildJsonParseErrorLog(ex.Message)); var errorResponse = CreateErrorResponse(null, -32700, "Parse error", category: McpErrorEnvelope.CategoryParseError, suggestion: "Send valid JSON-RPC 2.0 framed as a single line of UTF-8 JSON.", @@ -691,7 +706,7 @@ internal static string BuildInvalidUtf8ErrorLog(string detail) // stderr には診断用に詳細を残すが、ネットワークに出るレスポンスには // 例外型のみを返し、SQLite の "near 'foo': syntax error" などを通じた // 内容漏れを防ぐ(#1530)。 - Console.Error.WriteLine(BuildUnhandledLoopErrorLog(ex.Message)); + DeferFrameLog(BuildUnhandledLoopErrorLog(ex.Message)); var classification = McpErrorEnvelope.ClassifyException(ex); var errorResponse = CreateErrorResponse(responseHasId, responseId, classification.JsonRpcCode, BuildSanitizedLoopErrorMessage(ex), @@ -710,11 +725,37 @@ private string SerializeResponseOrFallback(JsonNode response, bool hasId, JsonNo } catch (Exception ex) { - Console.Error.WriteLine(BuildResponseSerializationErrorLog(ex.Message)); + DeferFrameLog(BuildResponseSerializationErrorLog(ex.Message)); return BuildMinimalInternalErrorResponse(hasId, id, ex); } } + private void DeferFrameLog(string message) + { + var logs = _deferredFrameLogs.Value; + if (logs is null) + { + Console.Error.WriteLine(message); + return; + } + + logs.Add(message); + } + + private void BeginDeferredFrameLogs() + => _deferredFrameLogs.Value = []; + + private void FlushDeferredFrameLogs() + { + var logs = _deferredFrameLogs.Value; + if (logs is null) + return; + + _deferredFrameLogs.Value = null; + foreach (var log in logs) + Console.Error.WriteLine(log); + } + private static void ExtractResponseId(JsonNode request, out bool hasId, out JsonNode? id) { if (request is JsonObject obj) diff --git a/src/CodeIndex/Mcp/StdioMcpTransport.cs b/src/CodeIndex/Mcp/StdioMcpTransport.cs index 7e6b905ea7..d489121c2f 100644 --- a/src/CodeIndex/Mcp/StdioMcpTransport.cs +++ b/src/CodeIndex/Mcp/StdioMcpTransport.cs @@ -62,6 +62,7 @@ public async Task WriteFrameAsync(string? frame, CancellationToken cancellationT if (frame is null) return; // notifications produce no wire output on stdio. await _writer.WriteLineAsync(frame.AsMemory(), cancellationToken).ConfigureAwait(false); + await _writer.FlushAsync(cancellationToken).ConfigureAwait(false); } public ValueTask DisposeAsync() diff --git a/tests/CodeIndex.Tests/McpServerTests.cs b/tests/CodeIndex.Tests/McpServerTests.cs index 5d0a9adca0..7d43905445 100644 --- a/tests/CodeIndex.Tests/McpServerTests.cs +++ b/tests/CodeIndex.Tests/McpServerTests.cs @@ -1030,6 +1030,50 @@ await server.ProcessLineAsync( new ThrowingTextWriter()); } + [Fact] + public async Task ProcessLineAsync_ParseError_WritesResponseBeforeErrorLog() + { + using var server = new McpServer(_dbPath, ConsoleUi.LoadVersion()); + using var writer = new StringWriter(); + using var error = new StringWriter(); + var previousError = Console.Error; + Console.SetError(error); + try + { + await server.ProcessLineAsync("not json", new AssertingTextWriter(writer, () => Assert.Equal(string.Empty, error.ToString()))); + } + finally + { + Console.SetError(previousError); + } + + Assert.Contains("\"code\":-32700", writer.ToString()); + Assert.Contains("JSON parse error", error.ToString()); + } + + [Fact] + public async Task RunAsync_ParseErrorWriteFailure_LogsWriteFailureAndParseError() + { + var transport = new ShutdownProbeTransport("stdio", _ => throw new IOException("pipe closed"), "not json"); + using var server = new McpServer(_dbPath, "test"); + using var error = new StringWriter(); + var previousError = Console.Error; + Console.SetError(error); + try + { + await server.RunAsync(transport, CancellationToken.None); + } + finally + { + Console.SetError(previousError); + } + + var log = error.ToString(); + Assert.Contains("Error writing response", log); + Assert.Contains("pipe closed", log); + Assert.Contains("JSON parse error", log); + } + [Fact] public void BuildSanitizedLoopErrorMessage_CodeIndexException_EchoesStructuredFields() { @@ -1371,6 +1415,18 @@ public async Task StdioTransport_Utf16BomInput_ThrowsDecodeFailure() await Assert.ThrowsAsync(() => transport.ReadFrameAsync(CancellationToken.None)); } + [Fact] + public async Task StdioTransport_WriteFrameAsync_FlushesBeforeReturning() + { + await using var input = new MemoryStream(); + await using var output = new FlushCountingStream(); + await using var transport = new StdioMcpTransport(input, output, bufferSize: 1024); + + await transport.WriteFrameAsync("""{"jsonrpc":"2.0","id":1,"result":{}}""", CancellationToken.None); + + Assert.True(output.FlushCount > 0); + } + [Fact] public async Task RunAsync_StdioCancellationNotification_CancelsActiveRequest() { @@ -9087,7 +9143,52 @@ private sealed class ThrowingTextWriter : TextWriter { public override Encoding Encoding => Encoding.UTF8; + public override Task WriteAsync(char value) => + throw new IOException("pipe closed"); + + public override Task WriteAsync(string? value) => + throw new IOException("pipe closed"); + public override Task WriteLineAsync(string? value) => throw new IOException("pipe closed"); } + + private sealed class AssertingTextWriter : TextWriter + { + private readonly TextWriter _inner; + private readonly Action _beforeWrite; + + public AssertingTextWriter(TextWriter inner, Action beforeWrite) + { + _inner = inner; + _beforeWrite = beforeWrite; + } + + public override Encoding Encoding => _inner.Encoding; + + public override async Task WriteAsync(string? value) + { + _beforeWrite(); + await _inner.WriteAsync(value).ConfigureAwait(false); + } + + public override async Task WriteAsync(char value) + { + _beforeWrite(); + await _inner.WriteAsync(value).ConfigureAwait(false); + } + + public override Task FlushAsync() => _inner.FlushAsync(); + } + + private sealed class FlushCountingStream : MemoryStream + { + public int FlushCount { get; private set; } + + public override Task FlushAsync(CancellationToken cancellationToken) + { + FlushCount++; + return base.FlushAsync(cancellationToken); + } + } } From d5656547b10679ec0e172909e7d996f12e772806 Mon Sep 17 00:00:00 2001 From: Widthdom Date: Sun, 24 May 2026 23:48:22 +0900 Subject: [PATCH 2/3] Cover oversized MCP response ordering (#1927 #2020 #2021) --- src/CodeIndex/Mcp/McpServer.cs | 2 +- tests/CodeIndex.Tests/McpServerTests.cs | 21 +++++++++++++++++++++ 2 files changed, 22 insertions(+), 1 deletion(-) diff --git a/src/CodeIndex/Mcp/McpServer.cs b/src/CodeIndex/Mcp/McpServer.cs index 0a23bdd8c3..6404a457a2 100644 --- a/src/CodeIndex/Mcp/McpServer.cs +++ b/src/CodeIndex/Mcp/McpServer.cs @@ -666,7 +666,7 @@ internal static string BuildInvalidUtf8ErrorLog(string detail) // メモリ枯渇を防ぐため巨大メッセージを拒否 if (line.Length > MaxLineLength) { - Console.Error.WriteLine(BuildOversizedMessageLog(line.Length)); + DeferFrameLog(BuildOversizedMessageLog(line.Length)); var errorResponse = CreateErrorResponse(null, -32700, "Message too large", category: McpErrorEnvelope.CategoryMessageTooLarge, suggestion: $"JSON-RPC frame exceeds the {MaxLineLength} byte cap. Split the request into smaller calls or use `batch_query` with smaller slots.", diff --git a/tests/CodeIndex.Tests/McpServerTests.cs b/tests/CodeIndex.Tests/McpServerTests.cs index 7d43905445..53453c7aeb 100644 --- a/tests/CodeIndex.Tests/McpServerTests.cs +++ b/tests/CodeIndex.Tests/McpServerTests.cs @@ -1051,6 +1051,27 @@ public async Task ProcessLineAsync_ParseError_WritesResponseBeforeErrorLog() Assert.Contains("JSON parse error", error.ToString()); } + [Fact] + public async Task ProcessLineAsync_OversizedFrame_WritesResponseBeforeErrorLog() + { + using var server = new McpServer(_dbPath, ConsoleUi.LoadVersion()); + using var writer = new StringWriter(); + using var error = new StringWriter(); + var previousError = Console.Error; + Console.SetError(error); + try + { + await server.ProcessLineAsync(new string('x', 1_000_001), new AssertingTextWriter(writer, () => Assert.Equal(string.Empty, error.ToString()))); + } + finally + { + Console.SetError(previousError); + } + + Assert.Contains("Message too large", writer.ToString()); + Assert.Contains("Message too large", error.ToString()); + } + [Fact] public async Task RunAsync_ParseErrorWriteFailure_LogsWriteFailureAndParseError() { From 73246868afd499bc5234e7ae13a1b7c4361f1ce2 Mon Sep 17 00:00:00 2001 From: Widthdom Date: Sun, 24 May 2026 23:57:55 +0900 Subject: [PATCH 3/3] Defer responded MCP diagnostics (#1927 #2020 #2021) --- src/CodeIndex/Mcp/McpServer.cs | 26 ++++++++------ tests/CodeIndex.Tests/McpServerTests.cs | 47 +++++++++++++++++++++++++ 2 files changed, 63 insertions(+), 10 deletions(-) diff --git a/src/CodeIndex/Mcp/McpServer.cs b/src/CodeIndex/Mcp/McpServer.cs index 6404a457a2..4b047619d5 100644 --- a/src/CodeIndex/Mcp/McpServer.cs +++ b/src/CodeIndex/Mcp/McpServer.cs @@ -60,7 +60,7 @@ public partial class McpServer : IDisposable // を渡せるようにするため (#1567)。 private readonly AsyncLocal _currentRequestToken = new(); private readonly AsyncLocal?> _currentOutOfBandFrameWriter = new(); - private readonly AsyncLocal?> _deferredFrameLogs = new(); + private readonly AsyncLocal?> _deferredFrameLogs = new(); private bool _running = true; // Per-session DbContext reused across MCP tool calls. Holding the connection open // avoids reopening SQLite, reapplying pragmas, and re-registering every SQL function @@ -731,15 +731,18 @@ private string SerializeResponseOrFallback(JsonNode response, bool hasId, JsonNo } private void DeferFrameLog(string message) + => DeferFrameLog(() => Console.Error.WriteLine(message)); + + private void DeferFrameLog(Action writeLog) { var logs = _deferredFrameLogs.Value; if (logs is null) { - Console.Error.WriteLine(message); + writeLog(); return; } - logs.Add(message); + logs.Add(writeLog); } private void BeginDeferredFrameLogs() @@ -753,7 +756,7 @@ private void FlushDeferredFrameLogs() _deferredFrameLogs.Value = null; foreach (var log in logs) - Console.Error.WriteLine(log); + log(); } private static void ExtractResponseId(JsonNode request, out bool hasId, out JsonNode? id) @@ -873,7 +876,7 @@ private static string BuildMinimalInternalErrorResponse(bool hasId, JsonNode? id var authResult = _authenticator.Authenticate(request); if (!authResult.IsAuthenticated) { - Console.Error.WriteLine(BuildAuthFailureLog(method, authResult.FailureReason)); + DeferFrameLog(BuildAuthFailureLog(method, authResult.FailureReason)); return CreateErrorResponse(hasId: true, id: id, code: McpErrorEnvelope.CodeUnauthorized, message: "Unauthorized", category: McpErrorEnvelope.CategoryPermissionDenied, suggestion: "Set CDIDX_MCP_AUTH_TOKEN on the server and include a matching params.auth.token (or an `Authorization: Bearer ` header for HTTP) on each request.", @@ -1056,7 +1059,7 @@ private JsonNode HandleInitialize(JsonNode? id, JsonNode? _params) } else if (resolved != _caller && resolved != "unknown") { - Console.Error.WriteLine(BuildCallerSwapRejectionLog(_caller, resolved)); + DeferFrameLog(BuildCallerSwapRejectionLog(_caller, resolved)); } var negotiated = NegotiateProtocolVersion(_params, out var requestedVersion); if (negotiated == null) @@ -1068,7 +1071,7 @@ private JsonNode HandleInitialize(JsonNode? id, JsonNode? _params) // クライアント要求バージョンとサーバー対応集合に重なりがない場合。Issue #1554: // クライアントが分岐判定できるよう、`error.data` に要求バージョンと対応バージョン // を入れた -32602 (invalid params) を返す。 - Console.Error.WriteLine(BuildUnsupportedProtocolLog(requestedVersion)); + DeferFrameLog(BuildUnsupportedProtocolLog(requestedVersion)); return CreateUnsupportedProtocolError(id, requestedVersion); } @@ -1579,7 +1582,7 @@ private JsonNode HandleToolsCall(JsonNode? id, JsonNode? callParams) if (!decision.Allowed) { metricsError = "rate_limited"; - Console.Error.WriteLine(BuildRateLimitedLog(toolName, _caller, decision.RetryAfterMs)); + DeferFrameLog(BuildRateLimitedLog(toolName, _caller, decision.RetryAfterMs)); response = CreateRateLimitedErrorResponse(id, toolName, _caller, decision.RetryAfterMs); } else @@ -1636,8 +1639,11 @@ private JsonNode HandleToolsCall(JsonNode? id, JsonNode? callParams) // JSON-RPC のツール結果は tool 名 + 例外型のみに絞る。SQLite 例外などは // バインド値や該当リテラルを含むため、生のメッセージをクライアントに渡すと // パスや索引内容が漏れる(#1530)。 - Console.Error.WriteLine(BuildToolErrorLog(toolName, ex.Message)); - Database.DbDebug.DumpToStderr(ex); + DeferFrameLog(() => + { + Console.Error.WriteLine(BuildToolErrorLog(toolName, ex.Message)); + Database.DbDebug.DumpToStderr(ex); + }); metricsError = ex.GetType().Name; var classification = McpErrorEnvelope.ClassifyException(ex); response = CreateToolErrorResponse(true, id, BuildSanitizedToolErrorMessage(toolName, ex), diff --git a/tests/CodeIndex.Tests/McpServerTests.cs b/tests/CodeIndex.Tests/McpServerTests.cs index 53453c7aeb..1c1b40963b 100644 --- a/tests/CodeIndex.Tests/McpServerTests.cs +++ b/tests/CodeIndex.Tests/McpServerTests.cs @@ -1072,6 +1072,53 @@ public async Task ProcessLineAsync_OversizedFrame_WritesResponseBeforeErrorLog() Assert.Contains("Message too large", error.ToString()); } + [Fact] + public async Task ProcessLineAsync_UnsupportedProtocol_WritesResponseBeforeErrorLog() + { + using var server = new McpServer(_dbPath, ConsoleUi.LoadVersion()); + using var writer = new StringWriter(); + using var error = new StringWriter(); + var previousError = Console.Error; + Console.SetError(error); + try + { + await server.ProcessLineAsync( + """{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2099-01-01"}}""", + new AssertingTextWriter(writer, () => Assert.Equal(string.Empty, error.ToString()))); + } + finally + { + Console.SetError(previousError); + } + + Assert.Contains("Unsupported MCP protocolVersion", writer.ToString()); + Assert.Contains("Rejecting initialize", error.ToString()); + } + + [Fact] + public async Task ProcessLineAsync_AuthFailure_WritesResponseBeforeErrorLog() + { + using var server = new McpServer(_dbPath, ConsoleUi.LoadVersion(), false, + new TokenMcpAuthenticator("secret")); + using var writer = new StringWriter(); + using var error = new StringWriter(); + var previousError = Console.Error; + Console.SetError(error); + try + { + await server.ProcessLineAsync( + """{"jsonrpc":"2.0","id":1,"method":"tools/list"}""", + new AssertingTextWriter(writer, () => Assert.Equal(string.Empty, error.ToString()))); + } + finally + { + Console.SetError(previousError); + } + + Assert.Contains("Unauthorized", writer.ToString()); + Assert.Contains("Auth failed", error.ToString()); + } + [Fact] public async Task RunAsync_ParseErrorWriteFailure_LogsWriteFailureAndParseError() {