|
4 | 4 |
|
5 | 5 | package io.modelcontextprotocol.client.transport; |
6 | 6 |
|
| 7 | +import io.modelcontextprotocol.client.transport.ResponseSubscribers.ResponseEvent; |
| 8 | +import io.modelcontextprotocol.client.transport.customizer.McpAsyncHttpClientRequestCustomizer; |
| 9 | +import io.modelcontextprotocol.client.transport.customizer.McpSyncHttpClientRequestCustomizer; |
| 10 | +import io.modelcontextprotocol.common.McpTransportContext; |
| 11 | +import io.modelcontextprotocol.json.McpJsonMapper; |
| 12 | +import io.modelcontextprotocol.json.TypeRef; |
| 13 | +import io.modelcontextprotocol.spec.*; |
| 14 | +import io.modelcontextprotocol.util.Assert; |
| 15 | +import io.modelcontextprotocol.util.Utils; |
| 16 | +import org.reactivestreams.Publisher; |
| 17 | +import org.slf4j.Logger; |
| 18 | +import org.slf4j.LoggerFactory; |
| 19 | +import reactor.core.Disposable; |
| 20 | +import reactor.core.publisher.Flux; |
| 21 | +import reactor.core.publisher.FluxSink; |
| 22 | +import reactor.core.publisher.Mono; |
| 23 | +import reactor.util.function.Tuple2; |
| 24 | +import reactor.util.function.Tuples; |
| 25 | + |
7 | 26 | import java.io.IOException; |
8 | 27 | import java.net.URI; |
9 | 28 | import java.net.http.HttpClient; |
|
12 | 31 | import java.net.http.HttpResponse.BodyHandler; |
13 | 32 | import java.time.Duration; |
14 | 33 | import java.util.List; |
| 34 | +import java.util.Objects; |
15 | 35 | import java.util.Optional; |
16 | 36 | import java.util.concurrent.CompletionException; |
17 | 37 | import java.util.concurrent.atomic.AtomicReference; |
18 | 38 | import java.util.function.Consumer; |
19 | 39 | import java.util.function.Function; |
20 | 40 |
|
21 | | -import org.reactivestreams.Publisher; |
22 | | -import org.slf4j.Logger; |
23 | | -import org.slf4j.LoggerFactory; |
24 | | - |
25 | | -import io.modelcontextprotocol.json.TypeRef; |
26 | | -import io.modelcontextprotocol.json.McpJsonMapper; |
27 | | - |
28 | | -import io.modelcontextprotocol.client.transport.customizer.McpAsyncHttpClientRequestCustomizer; |
29 | | -import io.modelcontextprotocol.client.transport.customizer.McpSyncHttpClientRequestCustomizer; |
30 | | -import io.modelcontextprotocol.client.transport.ResponseSubscribers.ResponseEvent; |
31 | | -import io.modelcontextprotocol.common.McpTransportContext; |
32 | | -import io.modelcontextprotocol.spec.DefaultMcpTransportSession; |
33 | | -import io.modelcontextprotocol.spec.DefaultMcpTransportStream; |
34 | | -import io.modelcontextprotocol.spec.HttpHeaders; |
35 | | -import io.modelcontextprotocol.spec.McpClientTransport; |
36 | | -import io.modelcontextprotocol.spec.McpSchema; |
37 | | -import io.modelcontextprotocol.spec.McpTransportException; |
38 | | -import io.modelcontextprotocol.spec.McpTransportSession; |
39 | | -import io.modelcontextprotocol.spec.McpTransportSessionNotFoundException; |
40 | | -import io.modelcontextprotocol.spec.McpTransportStream; |
41 | | -import io.modelcontextprotocol.spec.ProtocolVersions; |
42 | | -import io.modelcontextprotocol.util.Assert; |
43 | | -import io.modelcontextprotocol.util.Utils; |
44 | | -import reactor.core.Disposable; |
45 | | -import reactor.core.publisher.Flux; |
46 | | -import reactor.core.publisher.FluxSink; |
47 | | -import reactor.core.publisher.Mono; |
48 | | -import reactor.util.function.Tuple2; |
49 | | -import reactor.util.function.Tuples; |
50 | | - |
51 | 41 | /** |
52 | 42 | * An implementation of the Streamable HTTP protocol as defined by the |
53 | 43 | * <code>2025-03-26</code> version of the MCP specification. |
@@ -87,7 +77,9 @@ public class HttpClientStreamableHttpTransport implements McpClientTransport { |
87 | 77 | */ |
88 | 78 | private final HttpClient httpClient; |
89 | 79 |
|
90 | | - /** HTTP request builder for building requests to send messages to the server */ |
| 80 | + /** |
| 81 | + * HTTP request builder for building requests to send messages to the server |
| 82 | + */ |
91 | 83 | private final HttpRequest.Builder requestBuilder; |
92 | 84 |
|
93 | 85 | /** |
@@ -442,8 +434,11 @@ public Mono<Void> sendMessage(McpSchema.JSONRPCMessage sentMessage) { |
442 | 434 | })).onErrorMap(CompletionException.class, t -> t.getCause()).onErrorComplete().subscribe(); |
443 | 435 |
|
444 | 436 | })).flatMap(responseEvent -> { |
445 | | - if (transportSession.markInitialized( |
446 | | - responseEvent.responseInfo().headers().firstValue("mcp-session-id").orElseGet(() -> null))) { |
| 437 | + String mcpSessionId = responseEvent.responseInfo() |
| 438 | + .headers() |
| 439 | + .firstValue("mcp-session-id") |
| 440 | + .orElseGet(() -> null); |
| 441 | + if (Objects.nonNull(mcpSessionId) && transportSession.markInitialized(mcpSessionId)) { |
447 | 442 | // Once we have a session, we try to open an async stream for |
448 | 443 | // the server to send notifications and requests out-of-band. |
449 | 444 |
|
|
0 commit comments