diff --git a/spring-messaging/src/main/java/org/springframework/messaging/rsocket/annotation/support/MessagingRSocket.java b/spring-messaging/src/main/java/org/springframework/messaging/rsocket/annotation/support/MessagingRSocket.java index f2f1da621dd4..23ad32208752 100644 --- a/spring-messaging/src/main/java/org/springframework/messaging/rsocket/annotation/support/MessagingRSocket.java +++ b/spring-messaging/src/main/java/org/springframework/messaging/rsocket/annotation/support/MessagingRSocket.java @@ -103,9 +103,6 @@ class MessagingRSocket implements RSocket { * @return completion handle for success or error */ public Mono handleConnectionSetupPayload(ConnectionSetupPayload payload) { - // frameDecoder does not apply to connectionSetupPayload - // so retain here since handle expects it. - payload.retain(); return handle(payload, FrameType.SETUP); } diff --git a/spring-messaging/src/main/java/org/springframework/messaging/rsocket/annotation/support/RSocketMessageHandler.java b/spring-messaging/src/main/java/org/springframework/messaging/rsocket/annotation/support/RSocketMessageHandler.java index 3e6f7f035813..d60dd836329f 100644 --- a/spring-messaging/src/main/java/org/springframework/messaging/rsocket/annotation/support/RSocketMessageHandler.java +++ b/spring-messaging/src/main/java/org/springframework/messaging/rsocket/annotation/support/RSocketMessageHandler.java @@ -425,6 +425,7 @@ public SocketAcceptor responder() { responder = createResponder(setupPayload, sendingRSocket); } catch (Throwable ex) { + setupPayload.release(); return Mono.error(ex); } return responder.handleConnectionSetupPayload(setupPayload).then(Mono.just(responder));