Fix Sonar Smells for AmqpOutEndpoint & MessPubEH
This commit is contained in:
@@ -301,7 +301,7 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin
|
||||
* @param delayExpression the expression.
|
||||
* @since 4.3.5
|
||||
*/
|
||||
public void setDelayExpressionString(String delayExpression) {
|
||||
public void setDelayExpressionString(@Nullable String delayExpression) {
|
||||
if (delayExpression == null) {
|
||||
this.delayExpression = null;
|
||||
}
|
||||
@@ -329,7 +329,7 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin
|
||||
* @see #setConfirmNackChannel(MessageChannel)
|
||||
*/
|
||||
public void setConfirmTimeout(long confirmTimeout) {
|
||||
this.confirmTimeout = Duration.ofMillis(confirmTimeout);
|
||||
this.confirmTimeout = Duration.ofMillis(confirmTimeout); // NOSONAR sync inconsistency
|
||||
}
|
||||
|
||||
protected final synchronized void setConnectionFactory(ConnectionFactory connectionFactory) {
|
||||
@@ -408,28 +408,43 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin
|
||||
|
||||
@Override
|
||||
protected final void doInit() {
|
||||
BeanFactory beanFactory = getBeanFactory();
|
||||
configureExchangeNameGenerator(beanFactory);
|
||||
configureRoutingKeyGenerator(beanFactory);
|
||||
configureCorrelationDataGenerator(beanFactory);
|
||||
configureDelayGenerator(beanFactory);
|
||||
|
||||
endpointInit();
|
||||
}
|
||||
|
||||
private void configureExchangeNameGenerator(BeanFactory beanFactory) {
|
||||
Assert.state(this.exchangeNameExpression == null || this.exchangeName == null,
|
||||
"Either an exchangeName or an exchangeNameExpression can be provided, but not both");
|
||||
BeanFactory beanFactory = getBeanFactory();
|
||||
if (this.exchangeNameExpression != null) {
|
||||
this.exchangeNameGenerator = new ExpressionEvaluatingMessageProcessor<String>(this.exchangeNameExpression,
|
||||
String.class);
|
||||
this.exchangeNameGenerator =
|
||||
new ExpressionEvaluatingMessageProcessor<>(this.exchangeNameExpression, String.class);
|
||||
if (beanFactory != null) {
|
||||
this.exchangeNameGenerator.setBeanFactory(beanFactory);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void configureRoutingKeyGenerator(BeanFactory beanFactory) {
|
||||
Assert.state(this.routingKeyExpression == null || this.routingKey == null,
|
||||
"Either a routingKey or a routingKeyExpression can be provided, but not both");
|
||||
if (this.routingKeyExpression != null) {
|
||||
this.routingKeyGenerator = new ExpressionEvaluatingMessageProcessor<String>(this.routingKeyExpression,
|
||||
String.class);
|
||||
this.routingKeyGenerator =
|
||||
new ExpressionEvaluatingMessageProcessor<>(this.routingKeyExpression, String.class);
|
||||
if (beanFactory != null) {
|
||||
this.routingKeyGenerator.setBeanFactory(beanFactory);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void configureCorrelationDataGenerator(BeanFactory beanFactory) {
|
||||
if (this.confirmCorrelationExpression != null) {
|
||||
this.correlationDataGenerator =
|
||||
new ExpressionEvaluatingMessageProcessor<Object>(this.confirmCorrelationExpression, Object.class);
|
||||
new ExpressionEvaluatingMessageProcessor<>(this.confirmCorrelationExpression, Object.class);
|
||||
if (beanFactory != null) {
|
||||
this.correlationDataGenerator.setBeanFactory(beanFactory);
|
||||
}
|
||||
@@ -444,14 +459,15 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin
|
||||
(this.confirmNackChannel == null || nullChannel != null) && this.confirmNackChannelName == null,
|
||||
"A 'confirmCorrelationExpression' is required when specifying a 'confirmNackChannel'");
|
||||
}
|
||||
}
|
||||
|
||||
private void configureDelayGenerator(BeanFactory beanFactory) {
|
||||
if (this.delayExpression != null) {
|
||||
this.delayGenerator = new ExpressionEvaluatingMessageProcessor<Integer>(this.delayExpression,
|
||||
Integer.class);
|
||||
this.delayGenerator = new ExpressionEvaluatingMessageProcessor<>(this.delayExpression, Integer.class);
|
||||
if (beanFactory != null) {
|
||||
this.delayGenerator.setBeanFactory(beanFactory);
|
||||
}
|
||||
}
|
||||
endpointInit();
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -486,9 +502,14 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin
|
||||
|
||||
private Runnable checkUnconfirmed() {
|
||||
return () -> {
|
||||
Collection<CorrelationData> unconfirmed =
|
||||
getRabbitTemplate().getUnconfirmed(getConfirmTimeout().toMillis());
|
||||
unconfirmed.forEach(correlation -> handleConfirm(correlation, false, "Confirm timed out"));
|
||||
RabbitTemplate rabbitTemplate = getRabbitTemplate();
|
||||
if (rabbitTemplate != null) {
|
||||
Collection<CorrelationData> unconfirmed =
|
||||
rabbitTemplate.getUnconfirmed(getConfirmTimeout().toMillis());
|
||||
if (unconfirmed != null) {
|
||||
unconfirmed.forEach(correlation -> handleConfirm(correlation, false, "Confirm timed out"));
|
||||
}
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
@@ -561,21 +582,25 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin
|
||||
|
||||
protected AbstractIntegrationMessageBuilder<?> buildReply(MessageConverter converter,
|
||||
org.springframework.amqp.core.Message amqpReplyMessage) {
|
||||
|
||||
Object replyObject = converter.fromMessage(amqpReplyMessage);
|
||||
AbstractIntegrationMessageBuilder<?> builder = (replyObject instanceof Message)
|
||||
? this.getMessageBuilderFactory().fromMessage((Message<?>) replyObject)
|
||||
: this.getMessageBuilderFactory().withPayload(replyObject);
|
||||
AbstractIntegrationMessageBuilder<?> builder = prepareMessageBuilder(replyObject);
|
||||
Map<String, ?> headers = getHeaderMapper().toHeadersFromReply(amqpReplyMessage.getMessageProperties());
|
||||
builder.copyHeadersIfAbsent(headers);
|
||||
return builder;
|
||||
}
|
||||
|
||||
private AbstractIntegrationMessageBuilder<?> prepareMessageBuilder(Object replyObject) {
|
||||
return replyObject instanceof Message
|
||||
? getMessageBuilderFactory().fromMessage((Message<?>) replyObject)
|
||||
: getMessageBuilderFactory().withPayload(replyObject);
|
||||
}
|
||||
|
||||
protected Message<?> buildReturnedMessage(org.springframework.amqp.core.Message message,
|
||||
int replyCode, String replyText, String exchange, String returnedRoutingKey, MessageConverter converter) {
|
||||
|
||||
Object returnedObject = converter.fromMessage(message);
|
||||
AbstractIntegrationMessageBuilder<?> builder = (returnedObject instanceof Message)
|
||||
? this.getMessageBuilderFactory().fromMessage((Message<?>) returnedObject)
|
||||
: this.getMessageBuilderFactory().withPayload(returnedObject);
|
||||
AbstractIntegrationMessageBuilder<?> builder = prepareMessageBuilder(returnedObject);
|
||||
Map<String, ?> headers = getHeaderMapper().toHeadersFromReply(message.getMessageProperties());
|
||||
if (this.errorMessageStrategy == null) {
|
||||
builder.copyHeadersIfAbsent(headers)
|
||||
@@ -602,25 +627,7 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin
|
||||
}
|
||||
Object userCorrelationData = wrapper.getUserData();
|
||||
Message<?> confirmMessage;
|
||||
if (this.errorMessageStrategy == null || ack) {
|
||||
Map<String, Object> headers = new HashMap<String, Object>();
|
||||
headers.put(AmqpHeaders.PUBLISH_CONFIRM, ack);
|
||||
if (!ack && StringUtils.hasText(cause)) {
|
||||
headers.put(AmqpHeaders.PUBLISH_CONFIRM_NACK_CAUSE, cause);
|
||||
}
|
||||
|
||||
AbstractIntegrationMessageBuilder<?> builder = userCorrelationData instanceof Message
|
||||
? this.getMessageBuilderFactory().fromMessage((Message<?>) userCorrelationData)
|
||||
: this.getMessageBuilderFactory().withPayload(userCorrelationData);
|
||||
|
||||
confirmMessage = builder
|
||||
.copyHeaders(headers)
|
||||
.build();
|
||||
}
|
||||
else {
|
||||
confirmMessage = this.errorMessageStrategy.buildErrorMessage(
|
||||
new NackedAmqpMessageException(wrapper.getMessage(), wrapper.getUserData(), cause), null);
|
||||
}
|
||||
confirmMessage = buildConfirmMessage(ack, cause, wrapper, userCorrelationData);
|
||||
if (ack && getConfirmAckChannel() != null) {
|
||||
sendOutput(confirmMessage, getConfirmAckChannel(), true);
|
||||
}
|
||||
@@ -636,6 +643,26 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin
|
||||
}
|
||||
}
|
||||
|
||||
private Message<?> buildConfirmMessage(boolean ack, String cause, CorrelationDataWrapper wrapper,
|
||||
Object userCorrelationData) {
|
||||
|
||||
if (this.errorMessageStrategy == null || ack) {
|
||||
Map<String, Object> headers = new HashMap<>();
|
||||
headers.put(AmqpHeaders.PUBLISH_CONFIRM, ack);
|
||||
if (!ack && StringUtils.hasText(cause)) {
|
||||
headers.put(AmqpHeaders.PUBLISH_CONFIRM_NACK_CAUSE, cause);
|
||||
}
|
||||
|
||||
return prepareMessageBuilder(userCorrelationData)
|
||||
.copyHeaders(headers)
|
||||
.build();
|
||||
}
|
||||
else {
|
||||
return this.errorMessageStrategy.buildErrorMessage(new NackedAmqpMessageException(wrapper.getMessage(),
|
||||
wrapper.getUserData(), cause), null);
|
||||
}
|
||||
}
|
||||
|
||||
protected static final class CorrelationDataWrapper extends CorrelationData {
|
||||
|
||||
private final Object userData;
|
||||
@@ -673,6 +700,7 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin
|
||||
}
|
||||
super.setReturnedMessage(returnedMessage);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -96,8 +96,11 @@ public class MessagePublishingErrorHandler extends ErrorMessagePublisher impleme
|
||||
getMessagingTemplate().send(errorChannel, getErrorMessageStrategy().buildErrorMessage(ex, null));
|
||||
sent = true;
|
||||
}
|
||||
catch (Throwable errorDeliveryError) {
|
||||
handleDeliveryError(errorDeliveryError);
|
||||
catch (Exception errorDeliveryError) {
|
||||
// message will be logged only
|
||||
if (this.logger.isWarnEnabled()) {
|
||||
this.logger.warn("Error message was not delivered.", errorDeliveryError);
|
||||
}
|
||||
}
|
||||
}
|
||||
if (!sent && this.logger.isErrorEnabled()) {
|
||||
@@ -111,16 +114,6 @@ public class MessagePublishingErrorHandler extends ErrorMessagePublisher impleme
|
||||
}
|
||||
}
|
||||
|
||||
private void handleDeliveryError(Throwable errorDeliveryError) {
|
||||
// message will be logged only
|
||||
if (this.logger.isWarnEnabled()) {
|
||||
this.logger.warn("Error message was not delivered.", errorDeliveryError);
|
||||
}
|
||||
if (errorDeliveryError instanceof Error) {
|
||||
throw ((Error) errorDeliveryError);
|
||||
}
|
||||
}
|
||||
|
||||
@Nullable
|
||||
private MessageChannel resolveErrorChannel(Throwable t) {
|
||||
DestinationResolver<MessageChannel> channelResolver = getChannelResolver();
|
||||
|
||||
Reference in New Issue
Block a user