Fix new Sonar smells
This commit is contained in:
@@ -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.
|
||||
* <p>
|
||||
* 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<Object, RSocketRequester> clientRSocketRequesters = new HashMap<>();
|
||||
|
||||
private BiFunction<String, DataBuffer, Object> 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())));
|
||||
|
||||
@@ -157,7 +157,7 @@ class ChannelSendOperator<T> extends Mono<Void> 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<T> extends Mono<Void> 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<T> extends Mono<Void> implements Scannable {
|
||||
}
|
||||
|
||||
private Subscriber<? super T> requiredWriteSubscriber() {
|
||||
Assert.state(this.writeSubscriber != null, "No write subscriber");
|
||||
return this.writeSubscriber;
|
||||
Subscriber<? super T> writeSubscriberToReturn = this.writeSubscriber;
|
||||
Assert.state(writeSubscriberToReturn != null, "No write subscriber");
|
||||
return writeSubscriberToReturn;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -255,7 +256,7 @@ class ChannelSendOperator<T> extends Mono<Void> 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<T> extends Mono<Void> 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<T> extends Mono<Void> 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<T> extends Mono<Void> implements Scannable {
|
||||
}
|
||||
}
|
||||
}
|
||||
s.request(n);
|
||||
s.request(requests);
|
||||
}
|
||||
|
||||
private boolean emitCachedSignals() {
|
||||
@@ -319,10 +321,13 @@ class ChannelSendOperator<T> extends Mono<Void> 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<T> extends Mono<Void> 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<T> extends Mono<Void> 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();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user