@@ -314,7 +314,7 @@ private Flux<McpSchema.JSONRPCMessage> consumeSseStream(Flow.Publisher<List<Byte
314314 }
315315 catch (IOException e ) {
316316 return Flux .<McpSchema .JSONRPCMessage >error (
317- new McpTransportException ("Error parsing JSON-RPC message" , e ));
317+ new McpTransportException ("Error parsing JSON-RPC message: " + data , e ));
318318 }
319319 });
320320 }
@@ -368,73 +368,34 @@ private Mono<Disposable> reconnect(McpTransportStream<Disposable> stream) {
368368 return Mono .from (this .httpRequestCustomizer .customize (builder , "GET" , uri , null , transportContext ));
369369 }).flatMapMany (requestBuilder -> {
370370 var request = requestBuilder .build ();
371- // Classify the response against the session id that this very request
372- // carried, rather than the one currently held by the session, which
373- // can be established concurrently.
374- Optional <String > maybeSessionId = request .headers ().firstValue (HttpHeaders .MCP_SESSION_ID );
375-
376371 return ResponseBodyHandlers .sendAsync (this .httpClient , request ).flatMapMany (httpResponse -> {
377372 int statusCode = httpResponse .statusCode ();
378- Exception exception = null ;
379- boolean proceed = false ;
380373 if (statusCode == 401 || statusCode == 403 ) {
381374 logger .debug ("Authorization error in reconnect with code {}" , statusCode );
382375 var requestSnapshot = new HttpRequestSnapshot (request .uri (), request .method (),
383376 request .headers ());
384- exception = new McpHttpClientTransportAuthorizationException (
385- "Authorization error connecting to SSE stream" , requestSnapshot ,
386- toResponseInfo (httpResponse ));
377+ return ResponseBodyHandlers .drainThenError (httpResponse .body (), this .maxResponseSize ,
378+ new McpHttpClientTransportAuthorizationException (
379+ "Authorization error connecting to SSE stream" , requestSnapshot ,
380+ toResponseInfo (httpResponse )));
387381 }
388- else if (statusCode == METHOD_NOT_ALLOWED ) {
382+ if (statusCode == METHOD_NOT_ALLOWED ) {
389383 logger .debug ("The server does not support SSE streams, using request-response mode." );
384+ return ResponseBodyHandlers .drain (httpResponse .body (), this .maxResponseSize );
390385 }
391- else if (statusCode == NOT_FOUND ) {
392- if (maybeSessionId .isPresent ()) {
393- logger .debug ("Session not found for session ID: {}" , maybeSessionId .get ());
394- String sessionIdRepresentation = sessionIdOrPlaceholder (maybeSessionId );
395- exception = new McpTransportSessionNotFoundException (sessionIdRepresentation );
396- }
397- else {
398- exception = new McpTransportException ("Server Not Found. Status code:" + statusCode );
399- }
400- }
401- else if (statusCode == BAD_REQUEST ) {
402- // Some implementations return 400 when presented with a session
403- // id they do not know about, so the session is invalidated.
404- // https://github.com/modelcontextprotocol/typescript-sdk/issues/389
405- if (maybeSessionId .isPresent ()) {
406- String sessionIdRepresentation = sessionIdOrPlaceholder (maybeSessionId );
407- exception = new McpTransportSessionNotFoundException (
408- "Session not found for session ID: " + sessionIdRepresentation );
409- }
410- else {
411- exception = new McpTransportException ("Bad Request. Status code:" + statusCode );
412- }
386+ if (statusCode < 200 || statusCode >= 300 ) {
387+ return statusError (request , httpResponse );
413388 }
414- else if (statusCode >= 200 && statusCode < 300 ) {
415- String contentType = httpResponse .headers ()
416- .firstValue (HttpHeaders .CONTENT_TYPE )
417- .orElse ("" )
418- .toLowerCase ();
419- if (contentType .contains (TEXT_EVENT_STREAM )) {
420- logger .debug ("SSE connection established successfully" );
421- proceed = true ;
422- }
423- else {
424- exception = new McpTransportException (
425- "Unrecognized server error when connecting to SSE stream, status code: "
426- + statusCode );
427- }
389+ String contentType = httpResponse .headers ()
390+ .firstValue (HttpHeaders .CONTENT_TYPE )
391+ .orElse ("" )
392+ .toLowerCase ();
393+ if (!contentType .contains (TEXT_EVENT_STREAM )) {
394+ return ResponseBodyHandlers .readThenError (httpResponse .body (), this .maxResponseSize ,
395+ "Unrecognized server error when connecting to SSE stream, status code: " + statusCode );
428396 }
429- else {
430- exception = new McpTransportException ("Received unrecognized status code: " + statusCode );
431- }
432-
433- return proceed ? consumeSseStream (httpResponse .body (), stream , null )
434- : exception != null
435- ? ResponseBodyHandlers .drainThenError (httpResponse .body (), this .maxResponseSize ,
436- exception )
437- : ResponseBodyHandlers .drain (httpResponse .body (), this .maxResponseSize );
397+ logger .debug ("SSE connection established successfully" );
398+ return consumeSseStream (httpResponse .body (), stream , null );
438399 });
439400 })
440401 .retryWhen (authorizationErrorRetrySpec ())
@@ -551,11 +512,6 @@ public Mono<Void> sendMessage(McpSchema.JSONRPCMessage sentMessage) {
551512 .from (this .httpRequestCustomizer .customize (builder , "POST" , uri , jsonBody , transportContext ));
552513 }).flatMapMany (requestBuilder -> {
553514 var request = requestBuilder .build ();
554- // Classify the response against the session id that this very request
555- // carried, rather than the one currently held by the session, which
556- // can be established concurrently.
557- Optional <String > maybeSessionId = request .headers ().firstValue (HttpHeaders .MCP_SESSION_ID );
558-
559515 return ResponseBodyHandlers .sendAsync (this .httpClient , request ).flatMapMany (httpResponse -> {
560516 int statusCode = httpResponse .statusCode ();
561517 if (statusCode == 401 || statusCode == 403 ) {
@@ -573,76 +529,49 @@ public Mono<Void> sendMessage(McpSchema.JSONRPCMessage sentMessage) {
573529 reconnect (null ).contextWrite (deliveredSink .contextView ()).subscribe ();
574530 }
575531
576- String sessionRepresentation = sessionIdOrPlaceholder (maybeSessionId );
577-
578- if (statusCode >= 200 && statusCode < 300 ) {
579- String contentType = httpResponse .headers ()
580- .firstValue (HttpHeaders .CONTENT_TYPE )
581- .orElse ("" )
582- .toLowerCase ();
583- String contentLength = httpResponse .headers ()
584- .firstValue (HttpHeaders .CONTENT_LENGTH )
585- .orElse (null );
586-
587- if (contentType .isBlank () || "0" .equals (contentLength ) || statusCode == 202 ) {
588- logger .debug ("No body returned for POST in session {}" , sessionRepresentation );
589- markDelivered .run ();
590- return ResponseBodyHandlers .drain (httpResponse .body (), this .maxResponseSize );
591- }
592- else if (contentType .contains (TEXT_EVENT_STREAM )) {
593- return consumeSseStream (httpResponse .body (), null , markDelivered );
594- }
595- else if (contentType .contains (APPLICATION_JSON )) {
596- return ResponseBodyHandlers
597- .decodeAggregateResponse (httpResponse .body (), this .maxResponseSize )
598- .flatMapMany (data -> {
599- markDelivered .run ();
600- if (sentMessage instanceof McpSchema .JSONRPCNotification ) {
601- logger .warn ("Notification: {} received non-compliant response: {}" , sentMessage ,
602- Utils .hasText (data ) ? data : "[empty]" );
603- return Flux .empty ();
604- }
605- try {
606- return Flux .just (McpSchema .deserializeJsonRpcMessage (jsonMapper , data ));
607- }
608- catch (IOException e ) {
609- return Flux .<McpSchema .JSONRPCMessage >error (
610- new McpTransportException ("Error deserializing JSON-RPC message" , e ));
611- }
612- });
613- }
614-
615- logger .warn ("Unknown media type {} returned for POST in session {}" , contentType ,
616- sessionRepresentation );
617- return ResponseBodyHandlers .drainThenError (httpResponse .body (), this .maxResponseSize ,
618- new McpTransportException ("Unknown media type returned: " + contentType ));
532+ if (statusCode < 200 || statusCode >= 300 ) {
533+ return statusError (request , httpResponse );
619534 }
620- else if (statusCode == NOT_FOUND ) {
621- if (maybeSessionId .isPresent ()) {
622- logger .debug ("Session not found for session ID: {}" , sessionRepresentation );
623- return ResponseBodyHandlers .drainThenError (httpResponse .body (), this .maxResponseSize ,
624- new McpTransportSessionNotFoundException (
625- "Session not found for session ID: " + sessionRepresentation ));
626- }
627- return ResponseBodyHandlers .drainThenError (httpResponse .body (), this .maxResponseSize ,
628- new McpTransportException ("Server Not Found. Status code:" + statusCode ));
535+
536+ String sessionRepresentation = sessionIdOrPlaceholder (
537+ request .headers ().firstValue (HttpHeaders .MCP_SESSION_ID ));
538+ String contentType = httpResponse .headers ()
539+ .firstValue (HttpHeaders .CONTENT_TYPE )
540+ .orElse ("" )
541+ .toLowerCase ();
542+ String contentLength = httpResponse .headers ().firstValue (HttpHeaders .CONTENT_LENGTH ).orElse (null );
543+
544+ if (contentType .isBlank () || "0" .equals (contentLength ) || statusCode == 202 ) {
545+ logger .debug ("No body returned for POST in session {}" , sessionRepresentation );
546+ markDelivered .run ();
547+ return ResponseBodyHandlers .drain (httpResponse .body (), this .maxResponseSize );
629548 }
630- else if (statusCode == BAD_REQUEST ) {
631- if (maybeSessionId .isPresent ()) {
632- return ResponseBodyHandlers .drainThenError (httpResponse .body (), this .maxResponseSize ,
633- new McpTransportSessionNotFoundException (
634- "Session not found for session ID: " + sessionRepresentation ));
635- }
636- return ResponseBodyHandlers .drainThenError (httpResponse .body (), this .maxResponseSize ,
637- new McpTransportException ("Bad Request. Status code:" + statusCode ));
549+ else if (contentType .contains (TEXT_EVENT_STREAM )) {
550+ return consumeSseStream (httpResponse .body (), null , markDelivered );
638551 }
639- else if (statusCode >= 400 && statusCode < 500 ) {
640- return ResponseBodyHandlers .drainThenError (httpResponse .body (), this .maxResponseSize ,
641- new McpTransportException ("Invalid request. Status code: " + statusCode ));
552+ else if (contentType .contains (APPLICATION_JSON )) {
553+ return ResponseBodyHandlers .decodeAggregateResponse (httpResponse .body (), this .maxResponseSize )
554+ .flatMapMany (data -> {
555+ markDelivered .run ();
556+ if (sentMessage instanceof McpSchema .JSONRPCNotification ) {
557+ logger .warn ("Notification: {} received non-compliant response: {}" , sentMessage ,
558+ Utils .hasText (data ) ? data : "[empty]" );
559+ return Flux .empty ();
560+ }
561+ try {
562+ return Flux .just (McpSchema .deserializeJsonRpcMessage (jsonMapper , data ));
563+ }
564+ catch (IOException e ) {
565+ return Flux .<McpSchema .JSONRPCMessage >error (new McpTransportException (
566+ "Error deserializing JSON-RPC message: " + data , e ));
567+ }
568+ });
642569 }
643570
571+ logger .warn ("Unknown media type {} returned for POST in session {}" , contentType ,
572+ sessionRepresentation );
644573 return ResponseBodyHandlers .drainThenError (httpResponse .body (), this .maxResponseSize ,
645- new McpTransportException ("Failed to send message, status code : " + statusCode ));
574+ new McpTransportException ("Unknown media type returned : " + contentType ));
646575 });
647576 })
648577 .retryWhen (authorizationErrorRetrySpec ())
@@ -681,6 +610,30 @@ else if (statusCode >= 400 && statusCode < 500) {
681610
682611 }
683612
613+ /**
614+ * Fails the exchange over a response with an error status. A session id the server
615+ * does not recognise invalidates the session; any other failure carries the response
616+ * body, which is what the server said about it.
617+ */
618+ private <T > Flux <T > statusError (HttpRequest request , HttpResponse <Flow .Publisher <List <ByteBuffer >>> response ) {
619+ int statusCode = response .statusCode ();
620+ // Classify the response against the session id that this very request carried,
621+ // rather than the one currently held by the session, which can be established
622+ // concurrently. Some implementations return 400 rather than 404 for a session id
623+ // they do not know about.
624+ // https://github.com/modelcontextprotocol/typescript-sdk/issues/389
625+ Optional <String > sessionId = request .headers ().firstValue (HttpHeaders .MCP_SESSION_ID );
626+ if ((statusCode == NOT_FOUND || statusCode == BAD_REQUEST ) && sessionId .isPresent ()) {
627+ logger .debug ("Session not found for session ID: {}" , sessionId .get ());
628+ return ResponseBodyHandlers .drainThenError (response .body (), this .maxResponseSize ,
629+ new McpTransportSessionNotFoundException (sessionId .get ()));
630+ }
631+ String failure = statusCode == NOT_FOUND ? "Server Not Found. Status code:" + statusCode
632+ : statusCode == BAD_REQUEST ? "Bad Request. Status code:" + statusCode
633+ : "Received unexpected status code: " + statusCode ;
634+ return ResponseBodyHandlers .readThenError (response .body (), this .maxResponseSize , failure );
635+ }
636+
684637 private static String sessionIdOrPlaceholder (Optional <String > sessionId ) {
685638 return sessionId .orElse ("[missing_session_id]" );
686639 }
0 commit comments