Skip to content

Commit 9df07a5

Browse files
committed
Minor race conition polish
Signed-off-by: Daniel Garnier-Moiroux <git@garnier.wf>
1 parent a252295 commit 9df07a5

1 file changed

Lines changed: 14 additions & 10 deletions

File tree

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

Lines changed: 14 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -342,15 +342,13 @@ private Mono<Disposable> reconnect(McpTransportStream<Disposable> stream) {
342342

343343
final AtomicReference<Disposable> disposableRef = new AtomicReference<>();
344344

345-
Optional<String> maybeSessionId = transportSession == null ? Optional.empty()
346-
: transportSession.sessionId();
347-
348345
Disposable connection = Mono.deferContextual(connectionCtx -> {
349346
var uri = Utils.resolveUri(this.baseUri, this.endpoint);
350347
HttpRequest.Builder requestBuilder = this.requestBuilder.copy();
351348

352-
if (maybeSessionId.isPresent()) {
353-
requestBuilder = requestBuilder.header(HttpHeaders.MCP_SESSION_ID, maybeSessionId.get());
349+
if (transportSession != null && transportSession.sessionId().isPresent()) {
350+
requestBuilder = requestBuilder.header(HttpHeaders.MCP_SESSION_ID,
351+
transportSession.sessionId().get());
354352
}
355353

356354
if (stream != null && stream.lastId().isPresent()) {
@@ -366,17 +364,22 @@ private Mono<Disposable> reconnect(McpTransportStream<Disposable> stream) {
366364
.GET();
367365
var transportContext = connectionCtx.getOrDefault(McpTransportContext.KEY, McpTransportContext.EMPTY);
368366
return Mono.from(this.httpRequestCustomizer.customize(builder, "GET", uri, null, transportContext));
369-
})
370-
.flatMapMany(requestBuilder -> Mono
371-
.fromFuture(() -> this.httpClient.sendAsync(requestBuilder.build(),
367+
}).flatMapMany(requestBuilder -> {
368+
var request = requestBuilder.build();
369+
// Classify the response against the session id that this very request
370+
// carried, rather than the one currently held by the session, which
371+
// can be established concurrently.
372+
Optional<String> maybeSessionId = request.headers().firstValue(HttpHeaders.MCP_SESSION_ID);
373+
374+
return Mono
375+
.fromFuture(() -> this.httpClient.sendAsync(request,
372376
ResponseSubscribers.boundedPublisherBodyHandler(this.maxResponseSize)))
373377
.flatMapMany(httpResponse -> {
374378
int statusCode = httpResponse.statusCode();
375379
Exception exception = null;
376380
boolean proceed = false;
377381
if (statusCode == 401 || statusCode == 403) {
378382
logger.debug("Authorization error in reconnect with code {}", statusCode);
379-
var request = requestBuilder.build();
380383
var requestSnapshot = new HttpRequestSnapshot(request.uri(), request.method(),
381384
request.headers());
382385
exception = new McpHttpClientTransportAuthorizationException(
@@ -429,7 +432,8 @@ else if (statusCode >= 200 && statusCode < 300) {
429432
return proceed ? consumeSseStream(httpResponse.body(), stream, null)
430433
: exception != null ? ResponseSubscribers.drainThenError(httpResponse.body(), exception)
431434
: ResponseSubscribers.drain(httpResponse.body());
432-
}))
435+
});
436+
})
433437
.retryWhen(authorizationErrorRetrySpec())
434438
.flatMap(jsonrpcMessage -> requestHandler.apply(Mono.just(jsonrpcMessage)))
435439
.onErrorComplete(t -> {

0 commit comments

Comments
 (0)