Optimize ServerRSocketConnector connection
When RSocket client connects to the server there is no reason to wrap a `ConnectionSetupPayload` into a `Message` since we are not going to send it downstream * Refactor `IntegrationRSocket.handleConnectionSetupPayload()` just return a `Mono<DataBuffer>` for converted `ConnectionSetupPayload` * Ask for a `destination` and `RSocketRequester` from the `IntegrationRSocket` instead of message headers
This commit is contained in:
@@ -110,18 +110,20 @@ class IntegrationRSocket extends AbstractRSocket {
|
||||
this.bufferFactory = bufferFactory;
|
||||
}
|
||||
|
||||
RSocketRequester getRequester() {
|
||||
return this.requester;
|
||||
}
|
||||
|
||||
/**
|
||||
* Wrap the {@link ConnectionSetupPayload} with a {@link Message} and
|
||||
* delegate to {@link #handle(Payload)} for handling.
|
||||
* @param payload the connection payload
|
||||
* @return completion handle for success or error
|
||||
*/
|
||||
Mono<Message<DataBuffer>> handleConnectionSetupPayload(ConnectionSetupPayload payload) {
|
||||
String destination = getDestination(payload);
|
||||
MessageHeaders headers = createHeaders(destination, null);
|
||||
Mono<DataBuffer> handleConnectionSetupPayload(ConnectionSetupPayload payload) {
|
||||
DataBuffer dataBuffer = retainDataAndReleasePayload(payload);
|
||||
int refCount = refCount(dataBuffer);
|
||||
return Mono.just(MessageBuilder.createMessage(dataBuffer, headers))
|
||||
return Mono.just(dataBuffer)
|
||||
.doFinally(s -> {
|
||||
if (refCount(dataBuffer) == refCount) {
|
||||
DataBufferUtils.release(dataBuffer);
|
||||
|
||||
@@ -29,12 +29,8 @@ import org.springframework.context.ApplicationEventPublisher;
|
||||
import org.springframework.context.ApplicationEventPublisherAware;
|
||||
import org.springframework.core.io.buffer.DataBuffer;
|
||||
import org.springframework.lang.Nullable;
|
||||
import org.springframework.messaging.MessageHeaders;
|
||||
import org.springframework.messaging.handler.DestinationPatternsMessageCondition;
|
||||
import org.springframework.messaging.rsocket.RSocketRequester;
|
||||
import org.springframework.messaging.rsocket.annotation.support.RSocketRequesterMethodArgumentResolver;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.RouteMatcher;
|
||||
|
||||
import io.rsocket.RSocketFactory;
|
||||
import io.rsocket.SocketAcceptor;
|
||||
@@ -181,17 +177,10 @@ public class ServerRSocketConnector extends AbstractRSocketConnector
|
||||
return (setupPayload, sendingRSocket) -> {
|
||||
IntegrationRSocket rsocket = createRSocket(setupPayload, sendingRSocket);
|
||||
return rsocket.handleConnectionSetupPayload(setupPayload)
|
||||
.doOnNext((message) -> {
|
||||
MessageHeaders messageHeaders = message.getHeaders();
|
||||
DataBuffer dataBuffer = message.getPayload();
|
||||
String destination =
|
||||
messageHeaders.get(DestinationPatternsMessageCondition.LOOKUP_DESTINATION_HEADER,
|
||||
RouteMatcher.Route.class)
|
||||
.value();
|
||||
.doOnNext((dataBuffer) -> {
|
||||
String destination = rsocket.getDestination(setupPayload);
|
||||
Object rsocketRequesterKey = this.clientRSocketKeyStrategy.apply(destination, dataBuffer);
|
||||
RSocketRequester rsocketRequester =
|
||||
messageHeaders.get(RSocketRequesterMethodArgumentResolver.RSOCKET_REQUESTER_HEADER,
|
||||
RSocketRequester.class);
|
||||
RSocketRequester rsocketRequester = rsocket.getRequester();
|
||||
this.clientRSocketRequesters.put(rsocketRequesterKey, rsocketRequester);
|
||||
RSocketConnectedEvent rSocketConnectedEvent =
|
||||
new RSocketConnectedEvent(rsocket, destination, dataBuffer, rsocketRequester);
|
||||
|
||||
Reference in New Issue
Block a user