Fix new Sonar smells
This commit is contained in:
@@ -613,13 +613,14 @@ public abstract class MessagingGatewaySupport extends AbstractEndpoint
|
||||
|
||||
private Mono<Message<?>> doSendAndReceiveMessageReactive(MessageChannel requestChannel, Object object,
|
||||
boolean error) {
|
||||
final Message<?> requestMessage;
|
||||
|
||||
Message<?> requestMessage;
|
||||
try {
|
||||
Message<?> message =
|
||||
object instanceof Message<?>
|
||||
? (Message<?>) object
|
||||
: this.requestMapper.toMessage(object);
|
||||
|
||||
Assert.state(message != null, () -> "request mapper resulted in no message for " + object);
|
||||
message = this.historyWritingPostProcessor.postProcessMessage(message);
|
||||
requestMessage = message;
|
||||
}
|
||||
@@ -628,7 +629,6 @@ public abstract class MessagingGatewaySupport extends AbstractEndpoint
|
||||
}
|
||||
|
||||
return Mono.defer(() -> {
|
||||
|
||||
Object originalReplyChannelHeader = requestMessage.getHeaders().getReplyChannel();
|
||||
Object originalErrorChannelHeader = requestMessage.getHeaders().getErrorChannel();
|
||||
|
||||
|
||||
@@ -204,35 +204,39 @@ public abstract class AbstractRemoteFileStreamingMessageSource<F>
|
||||
break;
|
||||
}
|
||||
}
|
||||
try {
|
||||
String remotePath = remotePath(file);
|
||||
Session<?> session = this.remoteFileTemplate.getSession();
|
||||
if (maxFetchSize > 0) {
|
||||
this.fetched.incrementAndGet();
|
||||
}
|
||||
try {
|
||||
return getMessageBuilderFactory()
|
||||
.withPayload(session.readRaw(remotePath))
|
||||
.setHeader(IntegrationMessageHeaderAccessor.CLOSEABLE_RESOURCE, session)
|
||||
.setHeader(FileHeaders.REMOTE_DIRECTORY, file.getRemoteDirectory())
|
||||
.setHeader(FileHeaders.REMOTE_FILE, file.getFilename())
|
||||
.setHeader(FileHeaders.REMOTE_HOST_PORT, session.getHostPort())
|
||||
.setHeader(FileHeaders.REMOTE_FILE_INFO,
|
||||
this.fileInfoJson ? file.toJson() : file);
|
||||
}
|
||||
catch (IOException e) {
|
||||
session.close();
|
||||
throw new UncheckedIOException("IOException when retrieving " + remotePath, e);
|
||||
}
|
||||
}
|
||||
catch (RuntimeException e) {
|
||||
resetFilterIfNecessary(file);
|
||||
throw e;
|
||||
if (maxFetchSize > 0) {
|
||||
this.fetched.incrementAndGet();
|
||||
}
|
||||
return remoteFileToMessage(file);
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
private Object remoteFileToMessage(AbstractFileInfo<F> file) {
|
||||
try {
|
||||
String remotePath = remotePath(file);
|
||||
Session<?> session = this.remoteFileTemplate.getSession();
|
||||
try {
|
||||
return getMessageBuilderFactory()
|
||||
.withPayload(session.readRaw(remotePath))
|
||||
.setHeader(IntegrationMessageHeaderAccessor.CLOSEABLE_RESOURCE, session)
|
||||
.setHeader(FileHeaders.REMOTE_DIRECTORY, file.getRemoteDirectory())
|
||||
.setHeader(FileHeaders.REMOTE_FILE, file.getFilename())
|
||||
.setHeader(FileHeaders.REMOTE_HOST_PORT, session.getHostPort())
|
||||
.setHeader(FileHeaders.REMOTE_FILE_INFO,
|
||||
this.fileInfoJson ? file.toJson() : file);
|
||||
}
|
||||
catch (IOException e) {
|
||||
session.close();
|
||||
throw new UncheckedIOException("IOException when retrieving " + remotePath, e);
|
||||
}
|
||||
}
|
||||
catch (RuntimeException e) {
|
||||
resetFilterIfNecessary(file);
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
|
||||
protected AbstractFileInfo<F> poll() {
|
||||
if (this.toBeReceived.size() == 0) {
|
||||
listFiles();
|
||||
|
||||
Reference in New Issue
Block a user