null)
+ * @param resultArg the result object to handle
* @param request the original request message
* @param channel the Rabbit channel to operate on (maybe null)
* @param source the source data for the method invocation - e.g.
@@ -383,7 +387,7 @@ public abstract class AbstractAdaptableMessageListener implements ChannelAwareMe
protected void handleResult(@Nullable InvocationResult resultArg, Message request,
@Nullable Channel channel, @Nullable Object source) {
- if (channel != null && resultArg != null) {
+ if (resultArg != null) {
if (resultArg.getReturnValue() instanceof CompletableFuture> completable) {
if (!this.isManualAck) {
this.logger.warn("Container AcknowledgeMode must be MANUAL for a Future> return type; "
@@ -413,13 +417,9 @@ public abstract class AbstractAdaptableMessageListener implements ChannelAwareMe
doHandleResult(resultArg, request, channel, source);
}
}
- else if (this.logger.isWarnEnabled()) {
- this.logger.warn("Listener method returned result [" + resultArg
- + "]: not generating response message for it because no Rabbit Channel given");
- }
}
- private void asyncSuccess(InvocationResult resultArg, Message request, Channel channel,
+ private void asyncSuccess(InvocationResult resultArg, Message request, @Nullable Channel channel,
@Nullable Object source, @Nullable Object deferredResult) {
if (deferredResult == null) {
@@ -458,8 +458,9 @@ public abstract class AbstractAdaptableMessageListener implements ChannelAwareMe
}
}
- protected void asyncFailure(Message request, Channel channel, Throwable t, @Nullable Object source) {
+ protected void asyncFailure(Message request, @Nullable Channel channel, Throwable t, @Nullable Object source) {
this.logger.error("Future, Mono, or suspend function was completed with an exception for " + request, t);
+ Assert.notNull(channel, "'channel' must not be null.");
try {
channel.basicNack(request.getMessageProperties().getDeliveryTag(), false,
ContainerUtils.shouldRequeue(this.defaultRequeueRejected, t, this.logger));
@@ -469,7 +470,7 @@ public abstract class AbstractAdaptableMessageListener implements ChannelAwareMe
}
}
- protected void doHandleResult(InvocationResult resultArg, Message request, Channel channel,
+ protected void doHandleResult(InvocationResult resultArg, Message request, @Nullable Channel channel,
@Nullable Object source) {
if (this.logger.isDebugEnabled()) {
@@ -500,12 +501,13 @@ public abstract class AbstractAdaptableMessageListener implements ChannelAwareMe
/**
* Build a Rabbit message to be sent as response based on the given result object.
* @param channel the Rabbit Channel to operate on.
+ * Can be null if implementation does not support AMQP 0.9.1.
* @param result the content of the message, as returned from the listener method.
* @param genericType the generic type to populate type headers.
* @return the Rabbit Message (never null).
* @see #setMessageConverter
*/
- protected Message buildMessage(Channel channel, @Nullable Object result, @Nullable Type genericType) {
+ protected Message buildMessage(@Nullable Channel channel, @Nullable Object result, @Nullable Type genericType) {
MessageConverter converter = getMessageConverter();
if (converter != null && !(result instanceof Message)) {
return convert(result, genericType, converter);
@@ -633,7 +635,8 @@ public abstract class AbstractAdaptableMessageListener implements ChannelAwareMe
* @see #postProcessResponse(Message, Message)
* @see #setReplyPostProcessor(ReplyPostProcessor)
*/
- protected void sendResponse(Channel channel, Address replyTo, Message messageIn) {
+ protected void sendResponse(@Nullable Channel channel, Address replyTo, Message messageIn) {
+ Assert.notNull(channel, "'channel' must not be null.");
Message message = messageIn;
if (this.beforeSendReplyPostProcessors != null) {
for (MessagePostProcessor postProcessor : this.beforeSendReplyPostProcessors) {
diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MessagingMessageListenerAdapter.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MessagingMessageListenerAdapter.java
index a0a99834..20ce47ea 100644
--- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MessagingMessageListenerAdapter.java
+++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MessagingMessageListenerAdapter.java
@@ -169,7 +169,7 @@ public class MessagingMessageListenerAdapter extends AbstractAdaptableMessageLis
}
@Override
- protected void asyncFailure(org.springframework.amqp.core.Message request, Channel channel, Throwable t,
+ protected void asyncFailure(org.springframework.amqp.core.Message request, @Nullable Channel channel, Throwable t,
@Nullable Object source) {
try {
@@ -183,7 +183,7 @@ public class MessagingMessageListenerAdapter extends AbstractAdaptableMessageLis
super.asyncFailure(request, channel, t, source);
}
- private void handleException(org.springframework.amqp.core.Message amqpMessage, @Nullable Channel channel,
+ protected void handleException(org.springframework.amqp.core.Message amqpMessage, @Nullable Channel channel,
@Nullable Message> message, ListenerExecutionFailedException e) throws Exception { // NOSONAR
if (this.errorHandler != null) {
@@ -307,7 +307,7 @@ public class MessagingMessageListenerAdapter extends AbstractAdaptableMessageLis
* @see #setMessageConverter
*/
@Override
- protected org.springframework.amqp.core.Message buildMessage(Channel channel, @Nullable Object result,
+ protected org.springframework.amqp.core.Message buildMessage(@Nullable Channel channel, @Nullable Object result,
@Nullable Type genericType) {
MessageConverter converter = getMessageConverter();
diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/adapter/MessageListenerAdapterTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/adapter/MessageListenerAdapterTests.java
index 58b1c5d9..4fc1d3c4 100644
--- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/adapter/MessageListenerAdapterTests.java
+++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/adapter/MessageListenerAdapterTests.java
@@ -24,6 +24,7 @@ import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
import com.rabbitmq.client.Channel;
+import org.jspecify.annotations.Nullable;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import reactor.core.publisher.Mono;
@@ -33,7 +34,6 @@ import org.springframework.amqp.core.Address;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageProperties;
import org.springframework.amqp.support.SendRetryContextAccessor;
-import org.springframework.amqp.support.converter.SimpleMessageConverter;
import org.springframework.aop.framework.ProxyFactory;
import org.springframework.retry.RetryPolicy;
import org.springframework.retry.policy.SimpleRetryPolicy;
@@ -67,8 +67,15 @@ public class MessageListenerAdapterTests {
public void init() {
this.messageProperties = new MessageProperties();
this.messageProperties.setContentType(MessageProperties.CONTENT_TYPE_TEXT_PLAIN);
- this.adapter = new MessageListenerAdapter();
- this.adapter.setMessageConverter(new SimpleMessageConverter());
+ this.adapter = new MessageListenerAdapter() {
+
+ @Override
+ protected void doHandleResult(InvocationResult resultArg, Message request, @Nullable Channel channel,
+ @Nullable Object source) {
+
+ }
+
+ };
}
@Test
@@ -77,7 +84,7 @@ public class MessageListenerAdapterTests {
@Override
protected Object[] buildListenerArguments(Object extractedMessage, Channel channel, Message message) {
- return new Object[] { extractedMessage, channel, message };
+ return new Object[] {extractedMessage, channel, message};
}
}
@@ -131,7 +138,15 @@ public class MessageListenerAdapterTests {
}
}
- this.adapter = new MessageListenerAdapter(new Delegate(), "myPojoMessageMethod");
+ this.adapter = new MessageListenerAdapter(new Delegate(), "myPojoMessageMethod") {
+
+ @Override
+ protected void doHandleResult(InvocationResult resultArg, Message request, @Nullable Channel channel,
+ @Nullable Object source) {
+
+ }
+
+ };
this.adapter.onMessage(new Message("foo".getBytes(), messageProperties), null);
assertThat(called.get()).isTrue();
}
@@ -146,7 +161,7 @@ public class MessageListenerAdapterTests {
@Test
public void testMappedListenerMethod() throws Exception {
- Map
+ * This class reuses the {@link MessagingMessageListenerAdapter} as much as possible just to avoid duplication.
+ * The {@link Channel} abstraction from AMQP Client 0.9.1 is out use and present here just for API compatibility
+ * and to follow DRY principle.
+ * Can be reworked eventually, when this AMQP 1.0 client won't be based on {@code spring-rabbit} dependency.
*
* @author Artem Bilan
*
@@ -54,6 +67,8 @@ public class RabbitAmqpMessageListenerAdapter extends MessagingMessageListenerAd
private @Nullable Collection