Skip to content

Commit 73cfe88

Browse files
committed
fix: propagate stdio process exit during initialization
Signed-off-by: Dongliang Xie <dragonfsky@gmail.com>
1 parent fb029c8 commit 73cfe88

6 files changed

Lines changed: 193 additions & 5 deletions

File tree

‎disclosure.txt‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
This change was submitted despite me reading the rules and understanding AI contribution guidelines.

‎mcp-core/src/main/java/io/modelcontextprotocol/client/LifecycleInitializer.java‎

Lines changed: 12 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@
1111
import java.util.concurrent.atomic.AtomicReference;
1212
import java.util.function.Function;
1313

14+
import io.modelcontextprotocol.client.transport.McpStdioServerProcessExitException;
1415
import io.modelcontextprotocol.spec.McpClientSession;
1516
import io.modelcontextprotocol.spec.McpError;
1617
import io.modelcontextprotocol.spec.McpSchema;
@@ -209,7 +210,7 @@ private Mono<McpSchema.InitializeResult> await() {
209210

210211
private void complete(McpSchema.InitializeResult initializeResult) {
211212
// inform all the subscribers waiting for the initialization
212-
this.initSink.emitValue(initializeResult, Sinks.EmitFailureHandler.FAIL_FAST);
213+
this.initSink.tryEmitValue(initializeResult);
213214
}
214215

215216
private void cacheResult(McpSchema.InitializeResult initializeResult) {
@@ -218,7 +219,7 @@ private void cacheResult(McpSchema.InitializeResult initializeResult) {
218219
}
219220

220221
private void error(Throwable t) {
221-
this.initSink.emitError(t, Sinks.EmitFailureHandler.FAIL_FAST);
222+
this.initSink.tryEmitError(t);
222223
}
223224

224225
private void close() {
@@ -259,6 +260,12 @@ public void handleException(Throwable t) {
259260
// the implicit initialization step.
260261
this.withInitialization("re-initializing", result -> Mono.empty()).subscribe();
261262
}
263+
else if (t instanceof McpStdioServerProcessExitException) {
264+
DefaultInitialization current = this.initializationRef.get();
265+
if (current != null && current.initializeResult() == null) {
266+
current.error(t);
267+
}
268+
}
262269
}
263270

264271
/**
@@ -277,8 +284,8 @@ public <T> Mono<T> withInitialization(String actionName, Function<Initialization
277284
boolean needsToInitialize = previous == null;
278285
logger.debug(needsToInitialize ? "Initialization process started" : "Joining previous initialization");
279286

280-
Mono<McpSchema.InitializeResult> initializationJob = needsToInitialize
281-
? this.doInitialize(newInit, this.postInitializationHook, ctx) : previous.await();
287+
Mono<McpSchema.InitializeResult> initializationJob = needsToInitialize ? Mono.firstWithSignal(
288+
newInit.await(), this.doInitialize(newInit, this.postInitializationHook, ctx)) : previous.await();
282289

283290
return initializationJob.map(initializeResult -> this.initializationRef.get())
284291
.timeout(this.initializationTimeout)
@@ -355,4 +362,4 @@ public Mono<?> closeGracefully() {
355362
});
356363
}
357364

358-
}
365+
}
Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,42 @@
1+
/*
2+
* Copyright 2026-2026 the original author or authors.
3+
*/
4+
5+
package io.modelcontextprotocol.client.transport;
6+
7+
import io.modelcontextprotocol.spec.McpTransportException;
8+
import io.modelcontextprotocol.util.Assert;
9+
10+
/**
11+
* Thrown when an MCP stdio server process exits unexpectedly.
12+
*
13+
* @author Dongliang Xie
14+
*/
15+
public class McpStdioServerProcessExitException extends McpTransportException {
16+
17+
private static final long serialVersionUID = 1L;
18+
19+
private final int exitCode;
20+
21+
private final String command;
22+
23+
public McpStdioServerProcessExitException(int exitCode, String command) {
24+
super(message(exitCode, command));
25+
this.exitCode = exitCode;
26+
this.command = command;
27+
}
28+
29+
public int getExitCode() {
30+
return this.exitCode;
31+
}
32+
33+
public String getCommand() {
34+
return this.command;
35+
}
36+
37+
private static String message(int exitCode, String command) {
38+
Assert.hasText(command, "The command can not be empty");
39+
return "MCP server process exited unexpectedly with code " + exitCode + " for command: " + command;
40+
}
41+
42+
}

‎mcp-core/src/main/java/io/modelcontextprotocol/client/transport/StdioClientTransport.java‎

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@
1414
import java.util.List;
1515
import java.util.Set;
1616
import java.util.concurrent.Executors;
17+
import java.util.concurrent.atomic.AtomicReference;
1718
import java.util.function.Consumer;
1819
import java.util.function.Function;
1920
import java.util.stream.IntStream;
@@ -62,6 +63,10 @@ public class StdioClientTransport implements McpClientTransport {
6263
/** The server process being communicated with */
6364
private Process process;
6465

66+
private final AtomicReference<McpStdioServerProcessExitException> unexpectedExitException = new AtomicReference<>();
67+
68+
private final AtomicReference<Consumer<Throwable>> exceptionHandler = new AtomicReference<>();
69+
6570
private McpJsonMapper jsonMapper;
6671

6772
/** Scheduler for handling inbound messages from the server process */
@@ -82,6 +87,8 @@ public class StdioClientTransport implements McpClientTransport {
8287

8388
private volatile boolean isClosing = false;
8489

90+
private volatile boolean closeRequested = false;
91+
8592
// visible for tests
8693
private Consumer<String> stdErrorHandler = error -> logger.info("STDERR Message received: {}", error);
8794

@@ -165,6 +172,7 @@ public Mono<Void> connect(Function<Mono<JSONRPCMessage>, Mono<JSONRPCMessage>> h
165172
startInboundProcessing();
166173
startOutboundProcessing();
167174
startErrorProcessing();
175+
startExitMonitoring();
168176
logger.info("MCP server started");
169177
}).subscribeOn(Schedulers.boundedElastic());
170178
}
@@ -191,6 +199,11 @@ public void setStdErrorHandler(Consumer<String> errorHandler) {
191199
this.stdErrorHandler = errorHandler;
192200
}
193201

202+
@Override
203+
public void setExceptionHandler(Consumer<Throwable> handler) {
204+
this.exceptionHandler.set(handler);
205+
}
206+
194207
/**
195208
* Waits for the server process to exit.
196209
* @throws RuntimeException if the process is interrupted while waiting
@@ -258,6 +271,14 @@ private void handleIncomingErrors() {
258271

259272
@Override
260273
public Mono<Void> sendMessage(JSONRPCMessage message) {
274+
McpStdioServerProcessExitException exitException = this.unexpectedExitException.get();
275+
if (exitException != null) {
276+
return Mono.error(exitException);
277+
}
278+
if (!this.closeRequested && this.process != null && !this.process.isAlive()) {
279+
exitException = signalUnexpectedProcessExit(this.process.exitValue());
280+
return Mono.error(exitException);
281+
}
261282
if (this.outboundSink.tryEmitNext(message).isSuccess()) {
262283
// TODO: essentially we could reschedule ourselves in some time and make
263284
// another attempt with the already read data but pause reading until
@@ -271,6 +292,32 @@ public Mono<Void> sendMessage(JSONRPCMessage message) {
271292
}
272293
}
273294

295+
private void startExitMonitoring() {
296+
this.process.onExit().thenAccept(process -> {
297+
if (!closeRequested) {
298+
signalUnexpectedProcessExit(process.exitValue());
299+
}
300+
});
301+
}
302+
303+
private McpStdioServerProcessExitException signalUnexpectedProcessExit(int exitCode) {
304+
McpStdioServerProcessExitException exception = new McpStdioServerProcessExitException(exitCode,
305+
this.params.getCommand());
306+
if (this.unexpectedExitException.compareAndSet(null, exception)) {
307+
logger.warn(exception.getMessage());
308+
isClosing = true;
309+
inboundSink.tryEmitComplete();
310+
outboundSink.tryEmitComplete();
311+
errorSink.tryEmitComplete();
312+
313+
Consumer<Throwable> handler = this.exceptionHandler.get();
314+
if (handler != null) {
315+
handler.accept(exception);
316+
}
317+
}
318+
return this.unexpectedExitException.get();
319+
}
320+
274321
/**
275322
* Starts the inbound processing thread that reads JSON-RPC messages from the
276323
* process's input stream. Messages are deserialized and emitted to the inbound sink.
@@ -402,6 +449,7 @@ protected void handleOutbound(Function<Flux<JSONRPCMessage>, Flux<JSONRPCMessage
402449
@Override
403450
public Mono<Void> closeGracefully() {
404451
return Mono.fromRunnable(() -> {
452+
closeRequested = true;
405453
isClosing = true;
406454
logger.debug("Initiating graceful shutdown");
407455
}).then(Mono.<Void>defer(() -> {
Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,17 @@
1+
/*
2+
* Copyright 2024-2026 the original author or authors.
3+
*/
4+
5+
package io.modelcontextprotocol.client;
6+
7+
final class FailingStdioServer {
8+
9+
private FailingStdioServer() {
10+
}
11+
12+
public static void main(String[] args) {
13+
System.err.println("Exiting before MCP initialization with code 127");
14+
System.exit(127);
15+
}
16+
17+
}
Lines changed: 73 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,73 @@
1+
/*
2+
* Copyright 2024-2026 the original author or authors.
3+
*/
4+
5+
package io.modelcontextprotocol.client;
6+
7+
import java.nio.file.Path;
8+
import java.time.Duration;
9+
import java.util.concurrent.TimeUnit;
10+
11+
import io.modelcontextprotocol.client.transport.McpStdioServerProcessExitException;
12+
import io.modelcontextprotocol.client.transport.ServerParameters;
13+
import io.modelcontextprotocol.client.transport.StdioClientTransport;
14+
import org.junit.jupiter.api.Test;
15+
import org.junit.jupiter.api.Timeout;
16+
17+
import static io.modelcontextprotocol.util.McpJsonMapperUtils.JSON_MAPPER;
18+
import static org.assertj.core.api.Assertions.assertThat;
19+
import static org.assertj.core.api.Assertions.catchThrowable;
20+
21+
/**
22+
* Tests for initialization failures reported by {@link StdioClientTransport}.
23+
*
24+
* @author Dongliang Xie
25+
*/
26+
@Timeout(10)
27+
class StdioMcpClientInitializationFailureTests {
28+
29+
@Test
30+
void initializeShouldFailWithProcessExitInsteadOfRequestTimeout() {
31+
Duration requestTimeout = Duration.ofSeconds(3);
32+
String classpath = System.getProperty("java.class.path");
33+
ServerParameters stdioParams = ServerParameters.builder(javaExecutable())
34+
.args("-cp", classpath, FailingStdioServer.class.getName())
35+
.build();
36+
StdioClientTransport transport = new StdioClientTransport(stdioParams, JSON_MAPPER);
37+
McpSyncClient client = McpClient.sync(transport)
38+
.requestTimeout(requestTimeout)
39+
.initializationTimeout(Duration.ofSeconds(5))
40+
.build();
41+
42+
Throwable failure;
43+
long elapsedMillis;
44+
try {
45+
long startNanos = System.nanoTime();
46+
failure = catchThrowable(client::initialize);
47+
elapsedMillis = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNanos);
48+
}
49+
finally {
50+
client.closeGracefully();
51+
}
52+
53+
assertThat(failure).isNotNull();
54+
assertThat(elapsedMillis).isLessThan(requestTimeout.toMillis());
55+
assertThat(rootCause(failure)).isInstanceOfSatisfying(McpStdioServerProcessExitException.class, processExit -> {
56+
assertThat(processExit.getExitCode()).isEqualTo(127);
57+
assertThat(processExit.getCommand()).isEqualTo(javaExecutable());
58+
});
59+
}
60+
61+
private String javaExecutable() {
62+
String executable = System.getProperty("os.name").toLowerCase().contains("win") ? "java.exe" : "java";
63+
return Path.of(System.getProperty("java.home"), "bin", executable).toString();
64+
}
65+
66+
private Throwable rootCause(Throwable failure) {
67+
while (failure.getCause() != null) {
68+
failure = failure.getCause();
69+
}
70+
return failure;
71+
}
72+
73+
}

0 commit comments

Comments
 (0)