Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@
import java.util.function.Consumer;
import java.util.function.Function;

import io.modelcontextprotocol.client.transport.ResponseSubscribers.ResponseEvent;
import io.modelcontextprotocol.client.transport.ResponseSubscribers.SseEvent;
import io.modelcontextprotocol.client.transport.customizer.McpAsyncHttpClientRequestCustomizer;
import io.modelcontextprotocol.client.transport.customizer.McpSyncHttpClientRequestCustomizer;
import io.modelcontextprotocol.common.McpTransportContext;
Expand Down Expand Up @@ -353,59 +353,59 @@ public Mono<Void> connect(Function<Mono<JSONRPCMessage>, Mono<JSONRPCMessage>> h
var transportContext = ctx.getOrDefault(McpTransportContext.KEY, McpTransportContext.EMPTY);
return Mono.from(this.httpRequestCustomizer.customize(builder, "GET", uri, null, transportContext));
}).flatMap(requestBuilder -> Mono.create(sink -> {
Disposable connection = Flux.<ResponseEvent>create(sseSink -> this.httpClient
.sendAsync(requestBuilder.build(),
responseInfo -> ResponseSubscribers.sseToBodySubscriber(responseInfo, sseSink))
.exceptionallyCompose(e -> {
sseSink.error(e);
return CompletableFuture.failedFuture(e);
}))
.map(responseEvent -> (ResponseSubscribers.SseResponseEvent) responseEvent)
.flatMap(responseEvent -> {
Disposable connection = Mono
.fromFuture(() -> this.httpClient.sendAsync(requestBuilder.build(),
HttpResponse.BodyHandlers.ofPublisher()))
.flatMapMany(response -> {
if (isClosing) {
return Mono.empty();
return Flux.empty();
}

int statusCode = responseEvent.responseInfo().statusCode();
int statusCode = response.statusCode();

if (statusCode >= 200 && statusCode < 300) {
try {
if (ENDPOINT_EVENT_TYPE.equals(responseEvent.sseEvent().event())) {
String messageEndpointUri = responseEvent.sseEvent().data();
try {
messageEndpointValidator.validate(uri, messageEndpointUri);
}
catch (InvalidSseMessageEndpointException e) {
sink.error(e);
this.messageEndpointSink.tryEmitError(e);
return Flux.error(e);
}
if (this.messageEndpointSink.tryEmitValue(messageEndpointUri).isSuccess()) {
sink.success();
return Flux.empty(); // No further processing needed
}
else {
sink.error(new RuntimeException("Failed to handle SSE endpoint event"));
}
Flux<String> lines = ResponseSubscribers.decodeLines(response.body());
return ResponseSubscribers.decodeSseResponse(lines);
}
else {
return ResponseSubscribers.drainThenError(response.body(),
new RuntimeException("Failed to connect to SSE stream: " + statusCode));
}
})
.flatMap(sseEvent -> {
try {
if (ENDPOINT_EVENT_TYPE.equals(sseEvent.event())) {
String messageEndpointUri = sseEvent.data();
try {
messageEndpointValidator.validate(uri, messageEndpointUri);
}
else if (MESSAGE_EVENT_TYPE.equals(responseEvent.sseEvent().event())) {
JSONRPCMessage message = McpSchema.deserializeJsonRpcMessage(jsonMapper,
responseEvent.sseEvent().data());
catch (InvalidSseMessageEndpointException e) {
sink.error(e);
this.messageEndpointSink.tryEmitError(e);
return Flux.error(e);
}
if (this.messageEndpointSink.tryEmitValue(messageEndpointUri).isSuccess()) {
sink.success();
return Flux.just(message);
return Flux.empty(); // No further processing needed
}
else {
logger.debug("Received unrecognized SSE event type: {}", responseEvent.sseEvent());
sink.success();
sink.error(new RuntimeException("Failed to handle SSE endpoint event"));
}
}
catch (IOException e) {
sink.error(new McpTransportException("Error processing SSE event", e));
else if (MESSAGE_EVENT_TYPE.equals(sseEvent.event())) {
JSONRPCMessage message = McpSchema.deserializeJsonRpcMessage(jsonMapper, sseEvent.data());
sink.success();
return Flux.just(message);
}
else {
logger.debug("Received unrecognized SSE event type: {}", sseEvent);
sink.success();
}
}
return Flux.<McpSchema.JSONRPCMessage>error(
new RuntimeException("Failed to send message: " + responseEvent));

catch (IOException e) {
sink.error(new McpTransportException("Error processing SSE event", e));
}
return Flux.<McpSchema.JSONRPCMessage>empty();
})
.flatMap(jsonRpcMessage -> handler.apply(Mono.just(jsonRpcMessage)))
.onErrorComplete(t -> {
Expand Down Expand Up @@ -490,7 +490,7 @@ private Mono<HttpResponse<String>> sendHttpPost(final String endpoint, final Str
return Mono.from(this.httpRequestCustomizer.customize(builder, "POST", requestUri, body, transportContext));
}).flatMap(customizedBuilder -> {
var request = customizedBuilder.build();
return Mono.fromFuture(httpClient.sendAsync(request, HttpResponse.BodyHandlers.ofString()));
return Mono.fromFuture(this.httpClient.sendAsync(request, HttpResponse.BodyHandlers.ofString()));
});
}

Expand Down
Loading
Loading