1111import java .net .http .HttpResponse ;
1212import java .time .Duration ;
1313import java .util .List ;
14- import java .util .concurrent . atomic . AtomicBoolean ;
14+ import java .util .Optional ;
1515import java .util .concurrent .atomic .AtomicReference ;
1616import java .util .function .Consumer ;
1717import java .util .function .Function ;
@@ -389,14 +389,6 @@ public Mono<Void> connect(Function<Mono<JSONRPCMessage>, Mono<JSONRPCMessage>> h
389389 var transportContext = ctx .getOrDefault (McpTransportContext .KEY , McpTransportContext .EMPTY );
390390 return Mono .from (this .httpRequestCustomizer .customize (builder , "GET" , uri , null , transportContext ));
391391 }).flatMap (requestBuilder -> Mono .create (sink -> {
392- // Once connect() has completed, a later failure can no longer be reported
393- // through its sink: signalling it there would only have Reactor drop it.
394- AtomicBoolean connected = new AtomicBoolean ();
395- Runnable markConnected = () -> {
396- if (connected .compareAndSet (false , true )) {
397- sink .success ();
398- }
399- };
400392 Disposable connection = ResponseBodyHandlers .sendAsync (this .httpClient , requestBuilder .build ())
401393 .flatMapMany (response -> {
402394 if (isClosing ) {
@@ -417,7 +409,10 @@ public Mono<Void> connect(Function<Mono<JSONRPCMessage>, Mono<JSONRPCMessage>> h
417409 "Failed to connect to SSE stream: " + statusCode );
418410 }
419411 })
420- .flatMap (sseEvent -> {
412+ // Every successfully processed event yields exactly one element, empty
413+ // when it carries no message, so that the first one can mark the
414+ // connection as established.
415+ .<Optional <JSONRPCMessage >>handle ((sseEvent , events ) -> {
421416 try {
422417 if (ENDPOINT_EVENT_TYPE .equals (sseEvent .event ())) {
423418 String messageEndpointUri = sseEvent .data ();
@@ -426,46 +421,61 @@ public Mono<Void> connect(Function<Mono<JSONRPCMessage>, Mono<JSONRPCMessage>> h
426421 }
427422 catch (InvalidSseMessageEndpointException e ) {
428423 this .messageEndpointSink .tryEmitError (e );
429- return Flux .error (e );
424+ events .error (e );
425+ return ;
430426 }
431427 if (this .messageEndpointSink .tryEmitValue (messageEndpointUri ).isSuccess ()) {
432- markConnected .run ();
433- return Flux .empty (); // No further processing needed
428+ events .next (Optional .empty ());
429+ }
430+ else {
431+ events .error (new McpTransportException ("Failed to handle SSE endpoint event" ));
434432 }
435- return Flux .error (new McpTransportException ("Failed to handle SSE endpoint event" ));
436433 }
437434 else if (MESSAGE_EVENT_TYPE .equals (sseEvent .event ())) {
438435 String data = sseEvent .data ();
439436 if (data == null || data .isBlank ()) {
440437 logger .debug ("Skipping SSE event with empty data (stream primer)" );
441- markConnected .run ();
442- return Flux .<McpSchema .JSONRPCMessage >empty ();
438+ events .next (Optional .empty ());
439+ }
440+ else {
441+ events .next (Optional .of (McpSchema .deserializeJsonRpcMessage (jsonMapper , data )));
443442 }
444- JSONRPCMessage message = McpSchema .deserializeJsonRpcMessage (jsonMapper , data );
445- markConnected .run ();
446- return Flux .just (message );
447443 }
448444 else {
449445 logger .debug ("Received unrecognized SSE event type: {}" , sseEvent );
450- markConnected .run ();
451- return Flux .<McpSchema .JSONRPCMessage >empty ();
446+ events .next (Optional .empty ());
452447 }
453448 }
454449 catch (IOException e ) {
455- return Flux .<McpSchema .JSONRPCMessage >error (
456- new McpTransportException ("Error processing SSE event" , e ));
450+ events .error (new McpTransportException ("Error processing SSE event" , e ));
451+ }
452+ })
453+ // connect() is resolved by the first signal only: any later failure is
454+ // merely logged below, as connect() has already completed by then.
455+ .switchOnFirst ((first , events ) -> {
456+ if (first .hasValue ()) {
457+ sink .success ();
458+ }
459+ else if (first .isOnError ()) {
460+ sink .error (first .getThrowable ());
457461 }
462+ else if (first .isOnComplete ()) {
463+ sink .error (new McpTransportException ("SSE stream closed before any event was received" ));
464+ }
465+ return events ;
458466 })
459- .flatMap (jsonRpcMessage -> handler .apply (Mono .just (jsonRpcMessage )))
467+ .<JSONRPCMessage >handle ((message , messages ) -> message .ifPresent (messages ::next ))
468+ .flatMap (message -> handler .apply (Mono .just (message )))
460469 .onErrorComplete (t -> {
461470 if (!isClosing ) {
462471 logger .warn ("SSE stream observed an error" , t );
463- if (connected .compareAndSet (false , true )) {
464- sink .error (t );
465- }
466472 }
467473 return true ;
468474 })
475+ // A closeGracefully() before the first signal cancels the stream:
476+ // complete
477+ // connect() instead of leaving it pending. A no-op once it has resolved.
478+ .doOnCancel (sink ::success )
469479 .doFinally (s -> {
470480 Disposable ref = this .sseSubscription .getAndSet (null );
471481 if (ref != null && !ref .isDisposed ()) {
0 commit comments