|
11 | 11 | import java.net.http.HttpResponse; |
12 | 12 | import java.time.Duration; |
13 | 13 | import java.util.List; |
14 | | -import java.util.concurrent.CompletableFuture; |
15 | 14 | import java.util.concurrent.atomic.AtomicReference; |
16 | 15 | import java.util.function.Consumer; |
17 | 16 | import java.util.function.Function; |
18 | 17 |
|
19 | | -import io.modelcontextprotocol.client.transport.ResponseSubscribers.ResponseEvent; |
20 | 18 | import io.modelcontextprotocol.client.transport.customizer.McpAsyncHttpClientRequestCustomizer; |
21 | 19 | import io.modelcontextprotocol.client.transport.customizer.McpSyncHttpClientRequestCustomizer; |
22 | 20 | import io.modelcontextprotocol.common.McpTransportContext; |
@@ -390,61 +388,59 @@ public Mono<Void> connect(Function<Mono<JSONRPCMessage>, Mono<JSONRPCMessage>> h |
390 | 388 | var transportContext = ctx.getOrDefault(McpTransportContext.KEY, McpTransportContext.EMPTY); |
391 | 389 | return Mono.from(this.httpRequestCustomizer.customize(builder, "GET", uri, null, transportContext)); |
392 | 390 | }).flatMap(requestBuilder -> Mono.create(sink -> { |
393 | | - Disposable connection = Flux.<ResponseEvent>create( |
394 | | - sseSink -> this.httpClient |
395 | | - .sendAsync(requestBuilder.build(), |
396 | | - responseInfo -> ResponseSubscribers.sseToBodySubscriber(responseInfo, sseSink, |
397 | | - this.maxResponseSize)) |
398 | | - .exceptionallyCompose(e -> { |
399 | | - sseSink.error(e); |
400 | | - return CompletableFuture.failedFuture(e); |
401 | | - })) |
402 | | - .map(responseEvent -> (ResponseSubscribers.SseResponseEvent) responseEvent) |
403 | | - .flatMap(responseEvent -> { |
| 391 | + Disposable connection = Mono |
| 392 | + .fromFuture(() -> this.httpClient.sendAsync(requestBuilder.build(), |
| 393 | + ResponseSubscribers.boundedPublisherBodyHandler(this.maxResponseSize))) |
| 394 | + .flatMapMany(response -> { |
404 | 395 | if (isClosing) { |
405 | | - return Mono.empty(); |
| 396 | + return Flux.empty(); |
406 | 397 | } |
407 | 398 |
|
408 | | - int statusCode = responseEvent.responseInfo().statusCode(); |
| 399 | + int statusCode = response.statusCode(); |
409 | 400 |
|
410 | 401 | if (statusCode >= 200 && statusCode < 300) { |
411 | | - try { |
412 | | - if (ENDPOINT_EVENT_TYPE.equals(responseEvent.sseEvent().event())) { |
413 | | - String messageEndpointUri = responseEvent.sseEvent().data(); |
414 | | - try { |
415 | | - messageEndpointValidator.validate(uri, messageEndpointUri); |
416 | | - } |
417 | | - catch (InvalidSseMessageEndpointException e) { |
418 | | - sink.error(e); |
419 | | - this.messageEndpointSink.tryEmitError(e); |
420 | | - return Flux.error(e); |
421 | | - } |
422 | | - if (this.messageEndpointSink.tryEmitValue(messageEndpointUri).isSuccess()) { |
423 | | - sink.success(); |
424 | | - return Flux.empty(); // No further processing needed |
425 | | - } |
426 | | - else { |
427 | | - sink.error(new RuntimeException("Failed to handle SSE endpoint event")); |
428 | | - } |
| 402 | + Flux<String> lines = ResponseSubscribers.decodeLines(response.body()); |
| 403 | + return ResponseSubscribers.decodeSseResponse(lines, this.maxResponseSize); |
| 404 | + } |
| 405 | + else { |
| 406 | + return ResponseSubscribers.drainThenError(response.body(), |
| 407 | + new RuntimeException("Failed to connect to SSE stream: " + statusCode)); |
| 408 | + } |
| 409 | + }) |
| 410 | + .flatMap(sseEvent -> { |
| 411 | + try { |
| 412 | + if (ENDPOINT_EVENT_TYPE.equals(sseEvent.event())) { |
| 413 | + String messageEndpointUri = sseEvent.data(); |
| 414 | + try { |
| 415 | + messageEndpointValidator.validate(uri, messageEndpointUri); |
| 416 | + } |
| 417 | + catch (InvalidSseMessageEndpointException e) { |
| 418 | + sink.error(e); |
| 419 | + this.messageEndpointSink.tryEmitError(e); |
| 420 | + return Flux.error(e); |
429 | 421 | } |
430 | | - else if (MESSAGE_EVENT_TYPE.equals(responseEvent.sseEvent().event())) { |
431 | | - JSONRPCMessage message = McpSchema.deserializeJsonRpcMessage(jsonMapper, |
432 | | - responseEvent.sseEvent().data()); |
| 422 | + if (this.messageEndpointSink.tryEmitValue(messageEndpointUri).isSuccess()) { |
433 | 423 | sink.success(); |
434 | | - return Flux.just(message); |
| 424 | + return Flux.empty(); // No further processing needed |
435 | 425 | } |
436 | 426 | else { |
437 | | - logger.debug("Received unrecognized SSE event type: {}", responseEvent.sseEvent()); |
438 | | - sink.success(); |
| 427 | + sink.error(new RuntimeException("Failed to handle SSE endpoint event")); |
439 | 428 | } |
440 | 429 | } |
441 | | - catch (IOException e) { |
442 | | - sink.error(new McpTransportException("Error processing SSE event", e)); |
| 430 | + else if (MESSAGE_EVENT_TYPE.equals(sseEvent.event())) { |
| 431 | + JSONRPCMessage message = McpSchema.deserializeJsonRpcMessage(jsonMapper, sseEvent.data()); |
| 432 | + sink.success(); |
| 433 | + return Flux.just(message); |
| 434 | + } |
| 435 | + else { |
| 436 | + logger.debug("Received unrecognized SSE event type: {}", sseEvent); |
| 437 | + sink.success(); |
443 | 438 | } |
444 | 439 | } |
445 | | - return Flux.<McpSchema.JSONRPCMessage>error( |
446 | | - new RuntimeException("Failed to send message: " + responseEvent)); |
447 | | - |
| 440 | + catch (IOException e) { |
| 441 | + sink.error(new McpTransportException("Error processing SSE event", e)); |
| 442 | + } |
| 443 | + return Flux.<McpSchema.JSONRPCMessage>empty(); |
448 | 444 | }) |
449 | 445 | .flatMap(jsonRpcMessage -> handler.apply(Mono.just(jsonRpcMessage))) |
450 | 446 | .onErrorComplete(t -> { |
@@ -529,8 +525,8 @@ private Mono<HttpResponse<String>> sendHttpPost(final String endpoint, final Str |
529 | 525 | return Mono.from(this.httpRequestCustomizer.customize(builder, "POST", requestUri, body, transportContext)); |
530 | 526 | }).flatMap(customizedBuilder -> { |
531 | 527 | var request = customizedBuilder.build(); |
532 | | - return Mono.fromFuture( |
533 | | - httpClient.sendAsync(request, ResponseSubscribers.boundedStringBodyHandler(this.maxResponseSize))); |
| 528 | + return Mono.fromFuture(this.httpClient.sendAsync(request, |
| 529 | + ResponseSubscribers.boundedStringBodyHandler(this.maxResponseSize))); |
534 | 530 | }); |
535 | 531 | } |
536 | 532 |
|
|
0 commit comments