From 0c7f21b4341bbf317fce9ccf3e199913fbcc43ad Mon Sep 17 00:00:00 2001 From: Algis Dumbris Date: Sat, 3 Oct 2026 20:41:38 +0300 Subject: [PATCH 1/3] fix(upstream): detect a dead stdio transport and respawn the process (RC4-STDIO-001) Cause: when a stdio upstream's child process dies, mcp-go's stdio reader hits EOF and closes the transport; every later request returns transport.ErrTransportClosed ("transport error: transport closed"). The managed client's isConnectionError matched none of that text, so: - the health-loop ping logged it as "timeout (high activity), ignoring" and never flipped the state machine to Error; - a failed tools/call took the "not a connection error" branch. The server therefore kept reporting Ready/healthy with every call failing, and the backoff reconnect (which only runs from StateError) never fired. Fix: - add isDeadTransportError (ErrTransportClosed, io.EOF/ErrUnexpectedEOF, io.ErrClosedPipe, os.ErrClosed via errors.Is, plus text fallbacks for flattened chains; bare "EOF" text deliberately not matched) and treat it as a hard connection error in isConnectionError and as non-transient in isTransientHealthCheckError, so one failed ping or call flips the server to Error and the existing ShouldRetry backoff + tryReconnect respawns it. - route the call-path verdict through recordDeadConnection, which applies SetError only while the client is still Ready on the same connection generation (under epochMu, like the ambiguous-call probe): a call racing a deliberate Disconnect cannot flip it to Error, and a burst of calls on one dead transport charges one retry instead of inflating the backoff. Tests: classification table, health-ping and tools/call flips to Error, burst/one-retry and post-disconnect guards, and an end-to-end test that SIGKILLs a real cmd/mcpfixture stdio child and asserts the client leaves Ready and reconnects to a new fixture instance (fails without the fix). --- internal/upstream/managed/client.go | 79 ++++++++++- .../managed/dead_transport_fixture_test.go | 110 +++++++++++++++ .../upstream/managed/dead_transport_test.go | 131 ++++++++++++++++++ 3 files changed, 319 insertions(+), 1 deletion(-) create mode 100644 internal/upstream/managed/dead_transport_fixture_test.go create mode 100644 internal/upstream/managed/dead_transport_test.go diff --git a/internal/upstream/managed/client.go b/internal/upstream/managed/client.go index 96d6b34f4..d8c0db744 100644 --- a/internal/upstream/managed/client.go +++ b/internal/upstream/managed/client.go @@ -4,8 +4,10 @@ import ( "context" "errors" "fmt" + "io" "net" "net/url" + "os" "strings" "sync" "sync/atomic" @@ -22,6 +24,7 @@ import ( "github.com/smart-mcp-proxy/mcpproxy-go/internal/upstream/core" "github.com/smart-mcp-proxy/mcpproxy-go/internal/upstream/types" + mcptransport "github.com/mark3labs/mcp-go/client/transport" "github.com/mark3labs/mcp-go/mcp" "go.uber.org/zap" ) @@ -1138,6 +1141,9 @@ func (mc *Client) callTool(ctx context.Context, toolName string, args map[string if !mc.IsConnected() { return nil, fmt.Errorf("client not connected (state: %s)", mc.StateManager.GetState().String()) } + // The connection generation this call runs on, so a transport failure is + // only ever charged to the session that produced it (RC4-STDIO-001). + callEpoch := mc.connectionEpoch.Load() // #1317 round 2: counted from HERE, before admission control, not just // around the transport call below — a call already queued in @@ -1247,7 +1253,7 @@ func (mc *Client) callTool(ctx context.Context, toolName string, args map[string zap.String("tool", toolName), zap.Error(err)) } - mc.StateManager.SetError(err) + mc.recordDeadConnection(toolName, err, callEpoch) } return nil, err } @@ -1262,6 +1268,33 @@ func (mc *Client) callTool(ctx context.Context, toolName string, args map[string return result, nil } +// recordDeadConnection flips the server to Error after a tools/call produced +// hard evidence that its connection is broken (refused/reset, broken pipe, or +// a dead stdio transport). The health loop's next tick then reconnects through +// the normal ShouldRetry backoff. +// +// The verdict is applied only while the client is still Ready on the SAME +// connection generation the call ran on, serialized with Connect/Disconnect +// through epochMu (the same pairing probeAfterAmbiguousCallError uses): +// - a call that was on the wire when the user disconnected the server sees +// "transport closed" too; that belongs to the closed generation and must +// not flip the Disconnected client to Error; +// - a burst of concurrent calls hitting the same dead transport records ONE +// failure, not one per call — every extra SetError would bump retryCount +// and stretch the reconnect backoff for a single process death. +func (mc *Client) recordDeadConnection(toolName string, err error, callEpoch int64) { + mc.epochMu.Lock() + defer mc.epochMu.Unlock() + if !mc.IsConnected() || mc.connectionEpoch.Load() != callEpoch { + mc.logger.Debug("Tool call connection failure belongs to a connection that is already gone; not re-marking", + zap.String("server", mc.GetConfig().Name), + zap.String("tool", toolName), + zap.Error(err)) + return + } + mc.StateManager.SetError(err) +} + // recordCallToolOAuthSignal inspects a failed tools/call error and, when it is an // "authorization required" / 401 from an otherwise-connected server (NOT a // connection error), flags the server as needing OAuth sign-in. This covers @@ -1836,6 +1869,12 @@ func isTransientHealthCheckError(err error) bool { if err == nil { return false } + // A dead transport (the stdio child exited, its pipes closed) never + // recovers by waiting: every later request fails the same way until the + // process is respawned (RC4-STDIO-001). + if isDeadTransportError(err) { + return false + } msg := strings.ToLower(err.Error()) // Hard failures: short-circuit to "not transient" so the caller flips // Error on the first miss. Order matters — check these BEFORE the @@ -2133,11 +2172,49 @@ func (mc *Client) waitForReconnectCompletion(ctx context.Context) error { } } +// isDeadTransportError reports whether err proves the transport itself is +// closed, as opposed to one request failing on a live connection. The stdio +// case is RC4-STDIO-001: when the upstream's child process dies, mcp-go's +// stdio reader hits EOF and closes the transport, and every later request +// returns transport.ErrTransportClosed ("transport closed") — which matched +// none of isConnectionError's markers, so the health loop ignored it as a +// "high activity timeout" and the server stayed Ready with every call failing. +// +// Sentinels are matched through the error chain; the text fallbacks cover +// chains flattened before they reach us. Bare "EOF" is deliberately NOT +// matched as text: a live server's JSON-RPC error can mention it. +func isDeadTransportError(err error) bool { + if err == nil { + return false + } + if errors.Is(err, mcptransport.ErrTransportClosed) || + errors.Is(err, io.EOF) || + errors.Is(err, io.ErrUnexpectedEOF) || + errors.Is(err, io.ErrClosedPipe) || + errors.Is(err, os.ErrClosed) { + return true + } + msg := err.Error() + for _, marker := range []string{ + mcptransport.ErrTransportClosed.Error(), // "transport closed" + os.ErrClosed.Error(), // "file already closed" + io.ErrClosedPipe.Error(), // "io: read/write on closed pipe" + } { + if containsString(msg, marker) { + return true + } + } + return false +} + // isConnectionError checks if an error indicates a connection problem func (mc *Client) isConnectionError(err error) bool { if err == nil { return false } + if isDeadTransportError(err) { + return true + } errStr := err.Error() connectionErrors := []string{ diff --git a/internal/upstream/managed/dead_transport_fixture_test.go b/internal/upstream/managed/dead_transport_fixture_test.go new file mode 100644 index 000000000..e63446ff7 --- /dev/null +++ b/internal/upstream/managed/dead_transport_fixture_test.go @@ -0,0 +1,110 @@ +package managed + +import ( + "context" + "encoding/json" + "os/exec" + "path/filepath" + "runtime" + "testing" + "time" + + "github.com/mark3labs/mcp-go/mcp" + "github.com/stretchr/testify/require" + "go.uber.org/zap" + + "github.com/smart-mcp-proxy/mcpproxy-go/internal/config" + "github.com/smart-mcp-proxy/mcpproxy-go/internal/secret" +) + +// pingInstanceID calls the fixture's deterministic `ping` tool and returns the +// per-process instance_id, which changes on every process start. +func pingInstanceID(ctx context.Context, mc *Client) (string, error) { + res, err := mc.CallTool(ctx, "ping", nil) + if err != nil { + return "", err + } + for _, c := range res.Content { + if tc, ok := c.(mcp.TextContent); ok { + var body struct { + InstanceID string `json:"instance_id"` + } + if jerr := json.Unmarshal([]byte(tc.Text), &body); jerr == nil && body.InstanceID != "" { + return body.InstanceID, nil + } + } + } + return "", nil +} + +// TestStdioChildKilled_ClientLeavesReadyAndRespawns is the end-to-end +// RC4-STDIO-001 regression: an admitted stdio upstream (cmd/mcpfixture) whose +// child process is SIGKILLed must stop reporting Ready on the first failed +// call ("transport closed") and be respawned by the managed client's own +// health loop + backoff, without any explicit restart. +func TestStdioChildKilled_ClientLeavesReadyAndRespawns(t *testing.T) { + if testing.Short() { + t.Skip("builds and spawns the mcpfixture stdio binary") + } + if runtime.GOOS == "windows" { + t.Skip("uses pkill to SIGKILL the child") + } + if _, err := exec.LookPath("pkill"); err != nil { + t.Skip("pkill not available") + } + + bin := filepath.Join(t.TempDir(), "mcpfixture-deadtransport") + build := exec.Command("go", "build", "-o", bin, "github.com/smart-mcp-proxy/mcpproxy-go/cmd/mcpfixture") + out, err := build.CombinedOutput() + require.NoError(t, err, "build mcpfixture: %s", out) + + interval := config.Duration(200 * time.Millisecond) + cfg := &config.ServerConfig{ + Name: "stdio-dead-transport", + Command: bin, + Args: []string{"--transport", "stdio"}, + Protocol: "stdio", + Enabled: true, + HealthCheckInterval: &interval, + } + mc, err := NewClient(cfg.Name, cfg, zap.NewNop(), nil, &config.Config{}, nil, secret.NewResolver()) + require.NoError(t, err) + + ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second) + defer cancel() + require.NoError(t, mc.Connect(ctx)) + t.Cleanup(func() { _ = mc.Disconnect() }) + + firstID, err := pingInstanceID(ctx, mc) + require.NoError(t, err) + require.NotEmpty(t, firstID) + epochBefore := mc.ConnectionEpoch() + + // SIGKILL only the child: the proxy-side client is left untouched. + require.NoError(t, exec.Command("pkill", "-KILL", "-f", bin).Run()) + + // The dead transport must surface as a failed call, and that call must + // stop the client from reporting Ready on the dead generation. + require.Eventually(t, func() bool { + _, callErr := mc.CallTool(ctx, "ping", nil) + return callErr != nil + }, 5*time.Second, 20*time.Millisecond, "calls must start failing once the child is dead") + require.True(t, !mc.IsConnected() || mc.ConnectionEpoch() != epochBefore, + "after a call failed on the dead transport the client must not still report Ready on that connection (state=%s)", + mc.GetState()) + + // ...and the managed client must respawn the process on its own. + var newID string + require.Eventually(t, func() bool { + if !mc.IsConnected() { + return false + } + id, callErr := pingInstanceID(ctx, mc) + if callErr != nil || id == "" { + return false + } + newID = id + return true + }, 20*time.Second, 100*time.Millisecond, "the client must reconnect a fresh stdio process (state=%s)", mc.GetState()) + require.NotEqual(t, firstID, newID, "the reconnected upstream must be a new process") +} diff --git a/internal/upstream/managed/dead_transport_test.go b/internal/upstream/managed/dead_transport_test.go new file mode 100644 index 000000000..fb155828b --- /dev/null +++ b/internal/upstream/managed/dead_transport_test.go @@ -0,0 +1,131 @@ +package managed + +import ( + "context" + "errors" + "fmt" + "io" + "os" + "testing" + + "github.com/mark3labs/mcp-go/client/transport" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/smart-mcp-proxy/mcpproxy-go/internal/upstream/types" +) + +// RC4-STDIO-001: when a stdio upstream's child process dies, mcp-go's stdio +// transport closes its done channel and every later request fails with +// transport.ErrTransportClosed ("transport closed"). The managed layer used to +// classify that as "not a connection error", so the health loop logged it as +// "timeout (high activity), ignoring" and tool calls failed forever while the +// server kept reporting Ready/healthy. These tests pin the classification and +// the two state-machine entry points (health ping, failed tool call). + +// coreShapedTransportClosed reproduces the exact error chain the core client +// returns for a tools/call on a dead stdio transport: mcp-go wraps the +// sentinel in transport.Error, and internal/upstream/core/client.go wraps +// that with %w. +func coreShapedTransportClosed() error { + return fmt.Errorf("CallTool failed for 'echo': %w", transport.NewError(transport.ErrTransportClosed)) +} + +func TestIsDeadTransportError(t *testing.T) { + cases := []struct { + name string + err error + want bool + }{ + {"nil", nil, false}, + {"sentinel transport closed", transport.ErrTransportClosed, true}, + {"core-wrapped transport closed", coreShapedTransportClosed(), true}, + {"flattened transport closed text", errors.New("transport error: transport closed"), true}, + {"wrapped io.EOF", fmt.Errorf("failed to read: %w", io.EOF), true}, + {"wrapped io.ErrUnexpectedEOF", fmt.Errorf("read: %w", io.ErrUnexpectedEOF), true}, + {"wrapped io.ErrClosedPipe", fmt.Errorf("write: %w", io.ErrClosedPipe), true}, + {"wrapped os.ErrClosed", fmt.Errorf("failed to write request: %w", os.ErrClosed), true}, + {"flattened file already closed", errors.New("failed to write request: write |1: file already closed"), true}, + // A JSON-RPC error answered by a LIVE server whose message merely + // mentions EOF is a tool failure, not a dead transport. + {"tool error mentioning EOF", errors.New("tool error: unexpected EOF while parsing input"), false}, + {"caller canceled", context.Canceled, false}, + {"deadline", context.DeadlineExceeded, false}, + {"generic", errors.New("invalid params"), false}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + assert.Equal(t, tc.want, isDeadTransportError(tc.err)) + }) + } +} + +func TestDeadTransportIsHardConnectionError(t *testing.T) { + mc := newTestClientForHealth(t) + err := coreShapedTransportClosed() + assert.True(t, mc.isConnectionError(err), "a closed transport is a connection error") + assert.False(t, isTransientHealthCheckError(err), "a closed transport is hard evidence, not a transient miss") +} + +// TestPerformHealthCheck_TransportClosedFlipsToError is the health-loop half +// of RC4-STDIO-001: one ping on a dead stdio transport must flip the server to +// Error so the next tick's ShouldRetry→tryReconnect respawns the process. +func TestPerformHealthCheck_TransportClosedFlipsToError(t *testing.T) { + mc := newTestClientForHealth(t) + mc.healthProbe = &fakeProber{pingErr: transport.NewError(transport.ErrTransportClosed)} + + mc.performHealthCheck() + + assert.Equal(t, types.StateError, mc.StateManager.GetState(), + "a ping that fails with 'transport closed' must flip the server to Error on the first miss") + info := mc.StateManager.GetConnectionInfo() + assert.Equal(t, 1, info.RetryCount, + "one death charges one retry, so the existing backoff starts at its shortest step") + assert.False(t, info.Terminal, "a dead transport is not a parked/permanent failure") +} + +// TestCallTool_TransportClosedMarksServerError is the call-path half: the +// first tool call that hits the dead transport stops the server reporting +// Ready, instead of every call failing while status stays healthy. +func TestCallTool_TransportClosedMarksServerError(t *testing.T) { + mc, fake := newTestClientForCallTool(t, coreShapedTransportClosed()) + + _, err := mc.CallTool(context.Background(), "echo", nil) + require.Error(t, err) + assert.Equal(t, 1, fake.callCount()) + assert.Equal(t, types.StateError, mc.StateManager.GetState(), + "a tool call failing with 'transport closed' must flip the server to Error") +} + +// TestCallTool_DeadTransportBurstCountsOneRetry: a burst of concurrent calls +// that all hit the same dead transport must record ONE failure, not one per +// call — otherwise each extra SetError bumps retryCount and stretches the +// reconnect backoff (or exhausts MaxConnectionRetries) for a single death. +func TestCallTool_DeadTransportBurstCountsOneRetry(t *testing.T) { + mc, _ := newTestClientForCallTool(t, coreShapedTransportClosed()) + + _, err := mc.CallTool(context.Background(), "echo", nil) + require.Error(t, err) + require.Equal(t, types.StateError, mc.StateManager.GetState()) + first := mc.StateManager.GetConnectionInfo().RetryCount + + // A call that raced the first (already past its connectivity checks when + // the state flipped) reports the same dead transport afterwards. + mc.recordDeadConnection("echo", coreShapedTransportClosed(), mc.ConnectionEpoch()) + assert.Equal(t, first, mc.StateManager.GetConnectionInfo().RetryCount, + "a second report of the same dead connection must not bump retryCount") +} + +// TestCallTool_DeadTransportAfterDisconnectKeepsDisconnected: a call that was +// on the wire when the user disconnected the server sees "transport closed" +// too; that verdict belongs to the closed generation and must not flip the +// deliberately Disconnected client to Error. +func TestCallTool_DeadTransportAfterDisconnectKeepsDisconnected(t *testing.T) { + mc, _ := newTestClientForCallTool(t, nil) + staleEpoch := mc.ConnectionEpoch() + mc.connectionEpoch.Store(nextConnectionEpoch()) + mc.StateManager.Reset() + + mc.recordDeadConnection("echo", coreShapedTransportClosed(), staleEpoch) + assert.Equal(t, types.StateDisconnected, mc.StateManager.GetState()) +} From 603614463a2db7ca351b4ce5876d494f11fef1c8 Mon Sep 17 00:00:00 2001 From: Algis Dumbris Date: Sat, 3 Oct 2026 20:54:54 +0300 Subject: [PATCH 2/3] fix(upstream): narrow dead-transport classification and guard health verdict by epoch Follow-ups from the zcode review of the RC4-STDIO-001 fix: - isDeadTransportError no longer matches io.EOF / io.ErrUnexpectedEOF. The classifier is shared by all protocols and net/http wraps those for a single truncated/reset response on a live HTTP/SSE upstream, which must keep the normal flap tolerance. A stdio child's death surfaces as ErrTransportClosed (mcp-go's reader converts its EOF), so stdio coverage is unchanged. - The text fallbacks ("transport closed", "file already closed", closed pipe) now only count after mcp-go's "transport error: " prefix, so a live server's JSON-RPC error message echoing errno text cannot evict it. - performHealthCheck applies its SetError through the same epoch/epochMu guard as the call path (setErrorIfCurrentConnection): a ping that failed because Disconnect closed the transport no longer flips the freshly Reset client to Error. Tests: HTTP-shaped EOF chains and JSON-RPC errors echoing the markers are not dead transports; HTTP EOF ping stays Ready; tool error echoing "file already closed" stays Ready; ping racing Disconnect stays Disconnected with retryCount 0. All fail on the previous commit. --- internal/upstream/managed/client.go | 66 +++++++++++---- .../upstream/managed/dead_transport_test.go | 83 ++++++++++++++++++- 2 files changed, 131 insertions(+), 18 deletions(-) diff --git a/internal/upstream/managed/client.go b/internal/upstream/managed/client.go index d8c0db744..edc2f62bc 100644 --- a/internal/upstream/managed/client.go +++ b/internal/upstream/managed/client.go @@ -1283,16 +1283,30 @@ func (mc *Client) callTool(ctx context.Context, toolName string, args map[string // failure, not one per call — every extra SetError would bump retryCount // and stretch the reconnect backoff for a single process death. func (mc *Client) recordDeadConnection(toolName string, err error, callEpoch int64) { - mc.epochMu.Lock() - defer mc.epochMu.Unlock() - if !mc.IsConnected() || mc.connectionEpoch.Load() != callEpoch { + if !mc.setErrorIfCurrentConnection(err, callEpoch) { mc.logger.Debug("Tool call connection failure belongs to a connection that is already gone; not re-marking", zap.String("server", mc.GetConfig().Name), zap.String("tool", toolName), zap.Error(err)) - return + } +} + +// setErrorIfCurrentConnection applies SetError(err) only while the client is +// still Ready on connection generation epoch, serialized with Connect and +// Disconnect through epochMu. It reports whether the error was recorded. +// +// Both the tools/call path and the health loop use it: Disconnect closes the +// transport (every in-flight request then fails with "transport closed") +// BEFORE it bumps the epoch and Resets the state, so a verdict computed from +// that failure must not land on the freshly Disconnected client. +func (mc *Client) setErrorIfCurrentConnection(err error, epoch int64) bool { + mc.epochMu.Lock() + defer mc.epochMu.Unlock() + if !mc.IsConnected() || mc.connectionEpoch.Load() != epoch { + return false } mc.StateManager.SetError(err) + return true } // recordCallToolOAuthSignal inspects a failed tools/call error and, when it is an @@ -1731,6 +1745,10 @@ func (mc *Client) performHealthCheck() { if prober == nil { prober = mc.coreClient } + // The generation this ping runs on: a Disconnect racing the ping closes + // the transport (the ping fails with "transport closed") and only then + // bumps the epoch, so the verdict below must land on this generation only. + pingEpoch := mc.connectionEpoch.Load() err := prober.Ping(ctx) if err != nil { @@ -1780,11 +1798,16 @@ func (mc *Client) performHealthCheck() { mc.resetInFlightSuppression() } if mc.recordHealthCheckFailure(err) { - mc.logger.Warn("Health check failed repeatedly, marking as error", - zap.String("server", mc.GetConfig().Name), - zap.Int("consecutive_failures", mc.consecutiveHealthFailures), - zap.Error(err)) - mc.StateManager.SetError(err) + if mc.setErrorIfCurrentConnection(err, pingEpoch) { + mc.logger.Warn("Health check failed repeatedly, marking as error", + zap.String("server", mc.GetConfig().Name), + zap.Int("consecutive_failures", mc.consecutiveHealthFailures), + zap.Error(err)) + } else { + mc.logger.Debug("Health check failure belongs to a connection that is already gone; not marking", + zap.String("server", mc.GetConfig().Name), + zap.Error(err)) + } } else { mc.logger.Info("Health check failed transiently, tolerating below threshold", zap.String("server", mc.GetConfig().Name), @@ -2180,27 +2203,40 @@ func (mc *Client) waitForReconnectCompletion(ctx context.Context) error { // none of isConnectionError's markers, so the health loop ignored it as a // "high activity timeout" and the server stayed Ready with every call failing. // -// Sentinels are matched through the error chain; the text fallbacks cover -// chains flattened before they reach us. Bare "EOF" is deliberately NOT -// matched as text: a live server's JSON-RPC error can mention it. +// The classifier is shared by every protocol, so it only accepts evidence that +// the LOCAL end of the transport is closed: +// - io.EOF / io.ErrUnexpectedEOF are deliberately NOT matched. net/http wraps +// them for a single truncated or reset response on a live HTTP/SSE +// upstream (LB idle timeout, keep-alive race), which must stay subject to +// the normal flap tolerance. A stdio child's death never surfaces as them: +// mcp-go's stdio reader turns its EOF into ErrTransportClosed. +// - The text fallbacks (for chains flattened before they reach us) only +// count after mcp-go's "transport error: " prefix, i.e. text produced by +// transport.Error — never a live server's JSON-RPC error message, which +// mcp-go returns unwrapped and which may echo errno text such as +// "file already closed". func isDeadTransportError(err error) bool { if err == nil { return false } if errors.Is(err, mcptransport.ErrTransportClosed) || - errors.Is(err, io.EOF) || - errors.Is(err, io.ErrUnexpectedEOF) || errors.Is(err, io.ErrClosedPipe) || errors.Is(err, os.ErrClosed) { return true } msg := err.Error() + const transportPrefix = "transport error: " + idx := strings.Index(msg, transportPrefix) + if idx < 0 { + return false + } + transportMsg := msg[idx+len(transportPrefix):] for _, marker := range []string{ mcptransport.ErrTransportClosed.Error(), // "transport closed" os.ErrClosed.Error(), // "file already closed" io.ErrClosedPipe.Error(), // "io: read/write on closed pipe" } { - if containsString(msg, marker) { + if strings.Contains(transportMsg, marker) { return true } } diff --git a/internal/upstream/managed/dead_transport_test.go b/internal/upstream/managed/dead_transport_test.go index fb155828b..66df53748 100644 --- a/internal/upstream/managed/dead_transport_test.go +++ b/internal/upstream/managed/dead_transport_test.go @@ -5,7 +5,9 @@ import ( "errors" "fmt" "io" + "net/url" "os" + "sync/atomic" "testing" "github.com/mark3labs/mcp-go/client/transport" @@ -31,6 +33,15 @@ func coreShapedTransportClosed() error { return fmt.Errorf("CallTool failed for 'echo': %w", transport.NewError(transport.ErrTransportClosed)) } +// httpShapedEOF reproduces the chain mcp-go's streamable-HTTP transport hands +// back when net/http's POST fails mid-response: transport.Error wrapping +// "failed to send request: %w" wrapping *url.Error wrapping the io error. +func httpShapedEOF(ioErr error) error { + urlErr := &url.Error{Op: "Post", URL: "https://upstream.example/mcp", Err: ioErr} + return fmt.Errorf("CallTool failed for 'echo': %w", + transport.NewError(fmt.Errorf("failed to send request: %w", urlErr))) +} + func TestIsDeadTransportError(t *testing.T) { cases := []struct { name string @@ -41,11 +52,21 @@ func TestIsDeadTransportError(t *testing.T) { {"sentinel transport closed", transport.ErrTransportClosed, true}, {"core-wrapped transport closed", coreShapedTransportClosed(), true}, {"flattened transport closed text", errors.New("transport error: transport closed"), true}, - {"wrapped io.EOF", fmt.Errorf("failed to read: %w", io.EOF), true}, - {"wrapped io.ErrUnexpectedEOF", fmt.Errorf("read: %w", io.ErrUnexpectedEOF), true}, {"wrapped io.ErrClosedPipe", fmt.Errorf("write: %w", io.ErrClosedPipe), true}, {"wrapped os.ErrClosed", fmt.Errorf("failed to write request: %w", os.ErrClosed), true}, - {"flattened file already closed", errors.New("failed to write request: write |1: file already closed"), true}, + {"flattened file already closed", errors.New("transport error: failed to write request: write |1: file already closed"), true}, + // One truncated/reset HTTP response on a LIVE streamable-HTTP or SSE + // upstream (LB idle timeout, keep-alive race): net/http wraps io.EOF / + // io.ErrUnexpectedEOF. That is not a dead transport and must keep the + // normal flap tolerance instead of evicting on the first miss. + {"http POST unexpected EOF", httpShapedEOF(io.ErrUnexpectedEOF), false}, + {"http POST EOF", httpShapedEOF(io.EOF), false}, + {"wrapped io.EOF", fmt.Errorf("failed to read: %w", io.EOF), false}, + // A live server's JSON-RPC error (mcp-go returns it unwrapped, not as + // transport.Error) can echo errno text; that is a tool failure. + {"jsonrpc error echoing file already closed", errors.New("CallTool failed for 'rm': close /x: file already closed"), false}, + {"jsonrpc error echoing transport closed", errors.New("CallTool failed for 'ws': upstream websocket transport closed"), false}, + {"jsonrpc error echoing closed pipe", errors.New("CallTool failed for 'exec': io: read/write on closed pipe"), false}, // A JSON-RPC error answered by a LIVE server whose message merely // mentions EOF is a tool failure, not a dead transport. {"tool error mentioning EOF", errors.New("tool error: unexpected EOF while parsing input"), false}, @@ -129,3 +150,59 @@ func TestCallTool_DeadTransportAfterDisconnectKeepsDisconnected(t *testing.T) { mc.recordDeadConnection("echo", coreShapedTransportClosed(), staleEpoch) assert.Equal(t, types.StateDisconnected, mc.StateManager.GetState()) } + +// TestPerformHealthCheck_HTTPEOFIsTolerated: an HTTP-shaped EOF from a live +// upstream stays below the eviction threshold on the first miss. +func TestPerformHealthCheck_HTTPEOFIsTolerated(t *testing.T) { + mc := newTestClientForHealth(t) + mc.healthProbe = &fakeProber{pingErr: httpShapedEOF(io.ErrUnexpectedEOF)} + + mc.performHealthCheck() + + assert.Equal(t, types.StateReady, mc.StateManager.GetState(), + "one truncated HTTP response must not evict a live upstream") +} + +// TestCallTool_JSONRPCErrorEchoingClosedFileKeepsReady: a healthy server's own +// tool error that mentions "file already closed" must not flip it to Error. +func TestCallTool_JSONRPCErrorEchoingClosedFileKeepsReady(t *testing.T) { + mc, _ := newTestClientForCallTool(t, errors.New("CallTool failed for 'echo': close /x: file already closed")) + + _, err := mc.CallTool(context.Background(), "echo", nil) + require.Error(t, err) + assert.Equal(t, types.StateReady, mc.StateManager.GetState()) +} + +// disconnectDuringPingProber simulates Disconnect running while the health +// ping is on the wire: the transport closes (the ping fails with "transport +// closed"), Disconnect bumps the epoch and Resets the state, and only then +// does the health goroutine reach its verdict. +type disconnectDuringPingProber struct { + mc *Client + calls atomic.Int32 +} + +func (p *disconnectDuringPingProber) Ping(_ context.Context) error { + p.calls.Add(1) + p.mc.epochMu.Lock() + p.mc.connectionEpoch.Store(nextConnectionEpoch()) + p.mc.epochMu.Unlock() + p.mc.StateManager.Reset() + return transport.NewError(transport.ErrTransportClosed) +} + +// TestPerformHealthCheck_DeadTransportAfterDisconnectKeepsDisconnected is the +// health-loop twin of the call-path guard: a ping that failed because +// Disconnect closed the transport must not flip the Reset client to Error. +func TestPerformHealthCheck_DeadTransportAfterDisconnectKeepsDisconnected(t *testing.T) { + mc := newTestClientForHealth(t) + p := &disconnectDuringPingProber{mc: mc} + mc.healthProbe = p + + mc.performHealthCheck() + + require.Equal(t, int32(1), p.calls.Load()) + assert.Equal(t, types.StateDisconnected, mc.StateManager.GetState(), + "a ping failure from the closed generation must not re-mark the Disconnected client") + assert.Equal(t, 0, mc.StateManager.GetConnectionInfo().RetryCount) +} From 0045fb9137ef69c48697627ad152ba76bfd0d1a8 Mon Sep 17 00:00:00 2001 From: Algis Dumbris Date: Sat, 3 Oct 2026 21:20:13 +0300 Subject: [PATCH 3/3] fix(upstream): guard every dead-transport verdict by connection epoch opencode review of RC4-STDIO-001: - ListTools, the tool-count refresh and GetPrompt now record a connection error through setErrorIfCurrentConnection with the epoch captured before the request, like the call and health paths: concurrent failures on one dead transport mark it once, and a failure from a connection Disconnect already replaced never flips the reset client to Error. - tools/call reads its connection epoch after the admission wait, so a call queued across a reconnect is charged to the connection it actually ran on instead of being discarded as stale. - the text fallback now matches only mcp-go's typed *transport.Error or a message that starts with its prefix, not server-supplied JSON-RPC text that merely contains the phrase. - test the epoch arm of the guard on its own. --- internal/upstream/managed/client.go | 42 +++++++++++++------ .../upstream/managed/dead_transport_test.go | 24 +++++++++++ internal/upstream/managed/prompts.go | 5 ++- 3 files changed, 57 insertions(+), 14 deletions(-) diff --git a/internal/upstream/managed/client.go b/internal/upstream/managed/client.go index edc2f62bc..4fb839cf1 100644 --- a/internal/upstream/managed/client.go +++ b/internal/upstream/managed/client.go @@ -1061,6 +1061,7 @@ func (mc *Client) runListToolsAsLeader(listCtx context.Context, release func() b } }() + listEpoch := mc.connectionEpoch.Load() tools, err := mc.coreClient.ListTools(listCtx) mc.publishListToolsResult(tools, err) @@ -1077,7 +1078,9 @@ func (mc *Client) runListToolsAsLeader(listCtx context.Context, release func() b mc.logger.Warn("Connection error detected during ListTools, updating server state", zap.String("server", mc.GetConfig().Name), zap.Error(err)) - mc.StateManager.SetError(err) + // Guarded like the call and health paths: only the connection + // that produced the failure may be marked, once. + mc.setErrorIfCurrentConnection(err, listEpoch) } return nil, fmt.Errorf("ListTools failed: %w", err) } @@ -1141,10 +1144,6 @@ func (mc *Client) callTool(ctx context.Context, toolName string, args map[string if !mc.IsConnected() { return nil, fmt.Errorf("client not connected (state: %s)", mc.StateManager.GetState().String()) } - // The connection generation this call runs on, so a transport failure is - // only ever charged to the session that produced it (RC4-STDIO-001). - callEpoch := mc.connectionEpoch.Load() - // #1317 round 2: counted from HERE, before admission control, not just // around the transport call below — a call already queued in // acquireAdmission's wait is just as "in flight" from tryReconnect()'s @@ -1198,6 +1197,10 @@ func (mc *Client) callTool(ctx context.Context, toolName string, args map[string invoker = mc.coreClient } + // The connection generation this call runs on, read AFTER the admission + // wait (a queued call may outlive a reconnect), so a transport failure is + // only ever charged to the session that produced it (RC4-STDIO-001). + callEpoch := mc.connectionEpoch.Load() result, err := invoker.CallTool(ctx, toolName, args) if err != nil { mc.recordCallToolOAuthSignal(toolName, err) @@ -2224,13 +2227,24 @@ func isDeadTransportError(err error) bool { errors.Is(err, os.ErrClosed) { return true } - msg := err.Error() - const transportPrefix = "transport error: " - idx := strings.Index(msg, transportPrefix) - if idx < 0 { - return false + // Text fallback, for chains that were flattened into a string: only the + // payload of mcp-go's own transport wrapper counts — either the typed + // *transport.Error, or a message that STARTS with its prefix. A server's + // JSON-RPC error is surfaced unwrapped (errors.New(message)), so text it + // supplies can only match if it is the entire message, not by containing + // the phrase somewhere. + var transportMsg string + var tErr *mcptransport.Error + if errors.As(err, &tErr) && tErr.Err != nil { + transportMsg = tErr.Err.Error() + } else { + const transportPrefix = "transport error: " + msg := err.Error() + if !strings.HasPrefix(msg, transportPrefix) { + return false + } + transportMsg = strings.TrimPrefix(msg, transportPrefix) } - transportMsg := msg[idx+len(transportPrefix):] for _, marker := range []string{ mcptransport.ErrTransportClosed.Error(), // "transport closed" os.ErrClosed.Error(), // "file already closed" @@ -2465,6 +2479,7 @@ func (mc *Client) GetCachedToolCount(ctx context.Context) (int, error) { // Fetch fresh tool count with timeout. Publish the result so any concurrent // ListTools waiter coalesced behind us receives the real tools list. + countEpoch := mc.connectionEpoch.Load() tools, err := mc.coreClient.ListTools(listCtx) mc.publishListToolsResult(tools, err) if err != nil { @@ -2473,9 +2488,10 @@ func (mc *Client) GetCachedToolCount(ctx context.Context) (int, error) { zap.Error(err), zap.Int("cached_count", cachedCount)) - // Check if it's a connection error and update state + // Check if it's a connection error and update state (guarded: only + // the connection that produced the failure may be marked, once). if mc.isConnectionError(err) { - mc.StateManager.SetError(err) + mc.setErrorIfCurrentConnection(err, countEpoch) } // Return cached count if available, even if stale diff --git a/internal/upstream/managed/dead_transport_test.go b/internal/upstream/managed/dead_transport_test.go index 66df53748..48bd15b15 100644 --- a/internal/upstream/managed/dead_transport_test.go +++ b/internal/upstream/managed/dead_transport_test.go @@ -67,6 +67,11 @@ func TestIsDeadTransportError(t *testing.T) { {"jsonrpc error echoing file already closed", errors.New("CallTool failed for 'rm': close /x: file already closed"), false}, {"jsonrpc error echoing transport closed", errors.New("CallTool failed for 'ws': upstream websocket transport closed"), false}, {"jsonrpc error echoing closed pipe", errors.New("CallTool failed for 'exec': io: read/write on closed pipe"), false}, + // Even the full transport phrase, when it is server-supplied text + // inside a JSON-RPC error rather than mcp-go's own wrapper, is a + // tool failure (opencode review: match only the typed wrapper or a + // message that starts with its prefix). + {"jsonrpc error echoing the transport wrapper text", errors.New("CallTool failed for 'proxy': upstream said transport error: transport closed"), false}, // A JSON-RPC error answered by a LIVE server whose message merely // mentions EOF is a tool failure, not a dead transport. {"tool error mentioning EOF", errors.New("tool error: unexpected EOF while parsing input"), false}, @@ -206,3 +211,22 @@ func TestPerformHealthCheck_DeadTransportAfterDisconnectKeepsDisconnected(t *tes "a ping failure from the closed generation must not re-mark the Disconnected client") assert.Equal(t, 0, mc.StateManager.GetConnectionInfo().RetryCount) } + +// TestSetErrorIfCurrentConnection_EpochArm pins the generation half of the +// guard on its own: the client is still Ready, but on a NEWER connection than +// the one that produced the failure (a reconnect landed while the request was +// in flight). The stale verdict must be dropped, not charged to the new +// session. The !IsConnected() half is covered by the after-Disconnect tests. +func TestSetErrorIfCurrentConnection_EpochArm(t *testing.T) { + mc, _ := newTestClientForCallTool(t, nil) + require.Equal(t, types.StateReady, mc.StateManager.GetState()) + staleEpoch := mc.ConnectionEpoch() + mc.connectionEpoch.Store(nextConnectionEpoch()) // reconnected, still Ready + + assert.False(t, mc.setErrorIfCurrentConnection(coreShapedTransportClosed(), staleEpoch)) + assert.Equal(t, types.StateReady, mc.StateManager.GetState(), + "a failure from a replaced connection must not mark the current one") + + assert.True(t, mc.setErrorIfCurrentConnection(coreShapedTransportClosed(), mc.ConnectionEpoch())) + assert.Equal(t, types.StateError, mc.StateManager.GetState()) +} diff --git a/internal/upstream/managed/prompts.go b/internal/upstream/managed/prompts.go index 97ec4ad3c..34b51ed18 100644 --- a/internal/upstream/managed/prompts.go +++ b/internal/upstream/managed/prompts.go @@ -32,10 +32,13 @@ func (mc *Client) GetPrompt(ctx context.Context, name string, args map[string]st return nil, fmt.Errorf("client not connected (state: %s)", mc.StateManager.GetState().String()) } + promptEpoch := mc.connectionEpoch.Load() result, err := mc.coreClient.GetPrompt(ctx, name, args) if err != nil { if mc.isConnectionError(err) { - mc.StateManager.SetError(err) + // Guarded: concurrent failures on one dead transport mark it + // once, and never a connection Disconnect already replaced. + mc.setErrorIfCurrentConnection(err, promptEpoch) } mc.logger.Error("GetPrompt failed", zap.String("server", mc.GetConfig().Name),