diff --git a/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/ServerRSocketConnector.java b/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/ServerRSocketConnector.java index e89685f4ac..abd070fa48 100644 --- a/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/ServerRSocketConnector.java +++ b/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/ServerRSocketConnector.java @@ -16,6 +16,7 @@ package org.springframework.integration.rsocket; +import java.lang.reflect.Method; import java.util.Collections; import java.util.HashMap; import java.util.Map; @@ -38,7 +39,6 @@ import org.springframework.util.ReflectionUtils; import org.springframework.util.RouteMatcher; import io.rsocket.RSocketFactory; -import io.rsocket.SocketAcceptor; import io.rsocket.transport.ServerTransport; import io.rsocket.transport.netty.server.CloseableChannel; import io.rsocket.transport.netty.server.TcpServerTransport; @@ -50,7 +50,7 @@ import reactor.netty.http.server.HttpServer; /** * A server {@link AbstractRSocketConnector} extension to accept and manage client RSocket connections. *

- * Note: the {@link RSocketFactory.ServerRSocketFactory#acceptor(SocketAcceptor)} + * Note: the {@link RSocketFactory.ServerRSocketFactory#acceptor(io.rsocket.SocketAcceptor)} * in the provided {@link #factoryConfigurer} is overridden with an internal * {@link ServerRSocketMessageHandler#serverAcceptor()} * for the proper Spring Integration channel adapter mappings. @@ -175,6 +175,10 @@ public class ServerRSocketConnector extends AbstractRSocketConnector private static class ServerRSocketMessageHandler extends IntegrationRSocketMessageHandler { + private static final Method HANDLE_CONNECTION_SETUP_METHOD = + ReflectionUtils.findMethod(ServerRSocketMessageHandler.class, "handleConnectionSetup", Message.class); + + private final Map clientRSocketRequesters = new HashMap<>(); private BiFunction clientRSocketKeyStrategy = (destination, data) -> destination; @@ -182,9 +186,7 @@ public class ServerRSocketConnector extends AbstractRSocketConnector private ApplicationEventPublisher applicationEventPublisher; private void registerHandleConnectionSetupMethod() { - registerHandlerMethod(this, - ReflectionUtils.findMethod(ServerRSocketMessageHandler.class, "handleConnectionSetup", // NOSONAR - Message.class), + registerHandlerMethod(this, HANDLE_CONNECTION_SETUP_METHOD, new CompositeMessageCondition( RSocketFrameTypeMessageCondition.CONNECT_CONDITION, new DestinationPatternsMessageCondition(new String[] { "*" }, getRouteMatcher()))); diff --git a/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/inbound/ChannelSendOperator.java b/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/inbound/ChannelSendOperator.java index fd97256c4c..23cd0fa23c 100644 --- a/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/inbound/ChannelSendOperator.java +++ b/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/inbound/ChannelSendOperator.java @@ -157,7 +157,7 @@ class ChannelSendOperator extends Mono implements Scannable { private long demandBeforeReadyToWrite; /** Current state. */ - private State state = State.NEW; + private volatile State state = State.NEW; /** The actual writeSubscriber from the HTTP server adapter. */ @Nullable @@ -198,7 +198,7 @@ class ChannelSendOperator extends Mono implements Scannable { try { result = ChannelSendOperator.this.writeFunction.apply(this); } - catch (Throwable ex) { + catch (Throwable ex) { // NOSONAR this.writeCompletionBarrier.onError(ex); return; } @@ -214,8 +214,9 @@ class ChannelSendOperator extends Mono implements Scannable { } private Subscriber requiredWriteSubscriber() { - Assert.state(this.writeSubscriber != null, "No write subscriber"); - return this.writeSubscriber; + Subscriber writeSubscriberToReturn = this.writeSubscriber; + Assert.state(writeSubscriberToReturn != null, "No write subscriber"); + return writeSubscriberToReturn; } @Override @@ -255,7 +256,7 @@ class ChannelSendOperator extends Mono implements Scannable { try { result = ChannelSendOperator.this.writeFunction.apply(this); } - catch (Throwable ex) { + catch (Throwable ex) { // NOSONAR this.writeCompletionBarrier.onError(ex); return; } @@ -277,18 +278,19 @@ class ChannelSendOperator extends Mono implements Scannable { @Override public void request(long n) { + long requests = n; Subscription s = this.subscription; if (s == null) { return; } if (this.state == State.READY_TO_WRITE) { - s.request(n); + s.request(requests); return; } synchronized (this) { if (this.writeSubscriber != null) { if (this.state == State.EMITTING_CACHED_SIGNALS) { - this.demandBeforeReadyToWrite = n; + this.demandBeforeReadyToWrite = requests; return; } try { @@ -296,8 +298,8 @@ class ChannelSendOperator extends Mono implements Scannable { if (emitCachedSignals()) { return; } - n = n + this.demandBeforeReadyToWrite - 1; - if (n == 0) { + requests = requests + this.demandBeforeReadyToWrite - 1; + if (requests == 0) { return; } } @@ -306,7 +308,7 @@ class ChannelSendOperator extends Mono implements Scannable { } } } - s.request(n); + s.request(requests); } private boolean emitCachedSignals() { @@ -319,10 +321,13 @@ class ChannelSendOperator extends Mono implements Scannable { } return true; } - T item = this.item; - this.item = null; - if (item != null) { - requiredWriteSubscriber().onNext(item); + T itemToUse; + synchronized (this) { + itemToUse = this.item; + this.item = null; + } + if (itemToUse != null) { + requiredWriteSubscriber().onNext(itemToUse); } if (this.completed) { requiredWriteSubscriber().onComplete(); @@ -347,9 +352,9 @@ class ChannelSendOperator extends Mono implements Scannable { private void releaseCachedItem() { synchronized (this) { - Object item = this.item; - if (item instanceof DataBuffer) { - DataBufferUtils.release((DataBuffer) item); + Object itemToRelease = this.item; + if (itemToRelease instanceof DataBuffer) { + DataBufferUtils.release((DataBuffer) itemToRelease); } this.item = null; } @@ -451,9 +456,9 @@ class ChannelSendOperator extends Mono implements Scannable { @Override public void cancel() { this.writeBarrier.cancel(); - Subscription subscription = this.subscription; - if (subscription != null) { - subscription.cancel(); + Subscription subscriptionToCancel = this.subscription; + if (subscriptionToCancel != null) { + subscriptionToCancel.cancel(); } }