From 1bc1ae9edf20f8405eed27e960a9708ff8a36b38 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 31 Aug 2023 11:27:13 -0700 Subject: [PATCH] GH-1210: Add Kotlin suspend functions support * GH-1210: Add Kotlin suspend functions support Fixes https://github.com/spring-projects/spring-amqp/issues/1210 Kotlin Coroutines are essentially `Future` wrapping. Therefore, it is natural to have `suspend` support on `@RabbitListener` methods as we do now for `CompletableFuture` and `Mono` * Introduce some utilities since we cannot reuse existing from Spring Messaging: they are there about Kotlin Coroutines only for reactive handlers * Some code clean up in the `RabbitListenerAnnotationBeanPostProcessor` for the latest Java * Add optional dep for `kotlinx-coroutines-reactor` and document the feature * Fix since javadocs to reflect 2.4.x. Co-authored-by: Artem Bilan --- build.gradle | 2 + ...itListenerAnnotationBeanPostProcessor.java | 105 ++----------- .../AbstractAdaptableMessageListener.java | 8 +- .../AmqpMessageHandlerMethodFactory.java | 145 ++++++++++++++++++ ...inuationHandlerMethodArgumentResolver.java | 50 ++++++ .../KotlinAwareInvocableHandlerMethod.java | 49 ++++++ .../annotation/EnableRabbitKotlinTests.kt | 26 ++-- src/reference/asciidoc/amqp.adoc | 6 +- src/reference/asciidoc/whats-new.adoc | 2 +- 9 files changed, 284 insertions(+), 109 deletions(-) create mode 100644 spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/AmqpMessageHandlerMethodFactory.java create mode 100644 spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/ContinuationHandlerMethodArgumentResolver.java create mode 100644 spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/KotlinAwareInvocableHandlerMethod.java diff --git a/build.gradle b/build.gradle index 0e2eb16b..b18697aa 100644 --- a/build.gradle +++ b/build.gradle @@ -52,6 +52,7 @@ ext { jaywayJsonPathVersion = '2.4.0' junit4Version = '4.13.2' junitJupiterVersion = '5.8.2' + kotlinCoroutinesVersion = '1.6.4' log4jVersion = '2.17.2' logbackVersion = '1.2.3' lz4Version = '1.8.0' @@ -379,6 +380,7 @@ project('spring-rabbit') { } optionalApi "com.jayway.jsonpath:json-path:$jaywayJsonPathVersion" optionalApi "org.apache.commons:commons-pool2:$commonsPoolVersion" + optionalApi "org.jetbrains.kotlinx:kotlinx-coroutines-reactor:$kotlinCoroutinesVersion" testApi project(':spring-rabbit-junit') testImplementation("com.willowtreeapps.assertk:assertk-jvm:$assertkVersion") diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/annotation/RabbitListenerAnnotationBeanPostProcessor.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/annotation/RabbitListenerAnnotationBeanPostProcessor.java index 2d03aecf..0256c414 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/annotation/RabbitListenerAnnotationBeanPostProcessor.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/annotation/RabbitListenerAnnotationBeanPostProcessor.java @@ -18,8 +18,6 @@ package org.springframework.amqp.rabbit.annotation; import java.lang.reflect.AnnotatedElement; import java.lang.reflect.Method; -import java.lang.reflect.ParameterizedType; -import java.lang.reflect.Type; import java.nio.charset.Charset; import java.nio.charset.StandardCharsets; import java.util.ArrayList; @@ -30,7 +28,6 @@ import java.util.HashMap; import java.util.HashSet; import java.util.List; import java.util.Map; -import java.util.Optional; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; @@ -56,6 +53,7 @@ import org.springframework.amqp.rabbit.listener.MultiMethodRabbitListenerEndpoin import org.springframework.amqp.rabbit.listener.RabbitListenerContainerFactory; import org.springframework.amqp.rabbit.listener.RabbitListenerEndpointRegistrar; import org.springframework.amqp.rabbit.listener.RabbitListenerEndpointRegistry; +import org.springframework.amqp.rabbit.listener.adapter.AmqpMessageHandlerMethodFactory; import org.springframework.amqp.rabbit.listener.adapter.ReplyPostProcessor; import org.springframework.amqp.rabbit.listener.api.RabbitListenerErrorHandler; import org.springframework.amqp.support.converter.MessageConverter; @@ -76,7 +74,6 @@ import org.springframework.beans.factory.config.ConfigurableBeanFactory; import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; import org.springframework.context.EnvironmentAware; import org.springframework.context.expression.StandardBeanExpressionResolver; -import org.springframework.core.MethodParameter; import org.springframework.core.Ordered; import org.springframework.core.annotation.AnnotationUtils; import org.springframework.core.annotation.MergedAnnotations; @@ -88,12 +85,9 @@ import org.springframework.core.env.Environment; import org.springframework.core.task.TaskExecutor; import org.springframework.format.support.DefaultFormattingConversionService; import org.springframework.lang.Nullable; -import org.springframework.messaging.Message; import org.springframework.messaging.converter.GenericMessageConverter; import org.springframework.messaging.handler.annotation.support.DefaultMessageHandlerMethodFactory; import org.springframework.messaging.handler.annotation.support.MessageHandlerMethodFactory; -import org.springframework.messaging.handler.annotation.support.MethodArgumentNotValidException; -import org.springframework.messaging.handler.annotation.support.PayloadMethodArgumentResolver; import org.springframework.messaging.handler.invocation.HandlerMethodArgumentResolver; import org.springframework.messaging.handler.invocation.InvocableHandlerMethod; import org.springframework.util.Assert; @@ -101,8 +95,6 @@ import org.springframework.util.ClassUtils; import org.springframework.util.CollectionUtils; import org.springframework.util.ReflectionUtils; import org.springframework.util.StringUtils; -import org.springframework.validation.BindingResult; -import org.springframework.validation.ObjectError; import org.springframework.validation.Validator; /** @@ -440,14 +432,10 @@ public class RabbitListenerAnnotationBeanPostProcessor List resolvedQueues = resolveQueues(rabbitListener, declarables); if (!resolvedQueues.isEmpty()) { if (resolvedQueues.get(0) instanceof String) { - endpoint.setQueueNames(resolvedQueues.stream() - .map(o -> (String) o) - .collect(Collectors.toList()).toArray(new String[0])); + endpoint.setQueueNames(resolvedQueues.stream().map(o -> (String) o).toArray(String[]::new)); } else { - endpoint.setQueues(resolvedQueues.stream() - .map(o -> (Queue) o) - .collect(Collectors.toList()).toArray(new Queue[0])); + endpoint.setQueues(resolvedQueues.stream().map(o -> (Queue) o).toArray(Queue[]::new)); } } endpoint.setConcurrency(resolveExpressionAsStringOrInteger(rabbitListener.concurrency(), "concurrency")); @@ -664,12 +652,10 @@ public class RabbitListenerAnnotationBeanPostProcessor String[] queues = rabbitListener.queues(); QueueBinding[] bindings = rabbitListener.bindings(); org.springframework.amqp.rabbit.annotation.Queue[] queuesToDeclare = rabbitListener.queuesToDeclare(); - List queueNames = new ArrayList(); - List queueBeans = new ArrayList(); - if (queues.length > 0) { - for (int i = 0; i < queues.length; i++) { - resolveQueues(queues[i], queueNames, queueBeans); - } + List queueNames = new ArrayList<>(); + List queueBeans = new ArrayList<>(); + for (String queue : queues) { + resolveQueues(queue, queueNames, queueBeans); } if (!queueNames.isEmpty()) { // revert to the previous behavior of just using the name when there is mixture of String and Queue @@ -681,8 +667,8 @@ public class RabbitListenerAnnotationBeanPostProcessor throw new BeanInitializationException( "@RabbitListener can have only one of 'queues', 'queuesToDeclare', or 'bindings'"); } - for (int i = 0; i < queuesToDeclare.length; i++) { - queueNames.add(declareQueue(queuesToDeclare[i], declarables)); + for (org.springframework.amqp.rabbit.annotation.Queue queue : queuesToDeclare) { + queueNames.add(declareQueue(queue, declarables)); } } if (bindings.length > 0) { @@ -752,7 +738,7 @@ public class RabbitListenerAnnotationBeanPostProcessor declareExchangeAndBinding(binding, queueName, declarables); } } - return queues.toArray(new String[queues.size()]); + return queues.toArray(new String[0]); } private String declareQueue(org.springframework.amqp.rabbit.annotation.Queue bindingQueue, @@ -859,7 +845,7 @@ public class RabbitListenerAnnotationBeanPostProcessor } private Map resolveArguments(Argument[] arguments) { - Map map = new HashMap(); + Map map = new HashMap<>(); for (Argument arg : arguments) { String key = resolveExpressionAsString(arg.name(), "@Argument.name"); if (StringUtils.hasText(key)) { @@ -1025,7 +1011,7 @@ public class RabbitListenerAnnotationBeanPostProcessor } private MessageHandlerMethodFactory createDefaultMessageHandlerMethodFactory() { - DefaultMessageHandlerMethodFactory defaultFactory = new DefaultMessageHandlerMethodFactory(); + DefaultMessageHandlerMethodFactory defaultFactory = new AmqpMessageHandlerMethodFactory(); Validator validator = RabbitListenerAnnotationBeanPostProcessor.this.registrar.getValidator(); if (validator != null) { defaultFactory.setValidator(validator); @@ -1038,74 +1024,14 @@ public class RabbitListenerAnnotationBeanPostProcessor List customArgumentsResolver = new ArrayList<>( RabbitListenerAnnotationBeanPostProcessor.this.registrar.getCustomMethodArgumentResolvers()); defaultFactory.setCustomArgumentResolvers(customArgumentsResolver); - GenericMessageConverter messageConverter = new GenericMessageConverter( - this.defaultFormattingConversionService); - defaultFactory.setMessageConverter(messageConverter); - // Has to be at the end - look at PayloadMethodArgumentResolver documentation - customArgumentsResolver.add(new OptionalEmptyAwarePayloadArgumentResolver(messageConverter, validator)); + defaultFactory.setMessageConverter(new GenericMessageConverter(this.defaultFormattingConversionService)); + defaultFactory.afterPropertiesSet(); return defaultFactory; } } - private static class OptionalEmptyAwarePayloadArgumentResolver extends PayloadMethodArgumentResolver { - - OptionalEmptyAwarePayloadArgumentResolver( - org.springframework.messaging.converter.MessageConverter messageConverter, - @Nullable Validator validator) { - - super(messageConverter, validator); - } - - @Override - public Object resolveArgument(MethodParameter parameter, Message message) throws Exception { // NOSONAR - Object resolved = null; - try { - resolved = super.resolveArgument(parameter, message); - } - catch (MethodArgumentNotValidException ex) { - Type type = parameter.getGenericParameterType(); - if (isOptional(message, type)) { - BindingResult bindingResult = ex.getBindingResult(); - if (bindingResult != null) { - List allErrors = bindingResult.getAllErrors(); - if (allErrors.size() == 1) { - String defaultMessage = allErrors.get(0).getDefaultMessage(); - if ("Payload value must not be empty".equals(defaultMessage)) { - return Optional.empty(); - } - } - } - } - throw ex; - } - /* - * Replace Optional.empty() list elements with null. - */ - if (resolved instanceof List) { - List list = ((List) resolved); - for (int i = 0; i < list.size(); i++) { - if (list.get(i).equals(Optional.empty())) { - list.set(i, null); - } - } - } - return resolved; - } - - private boolean isOptional(Message message, Type type) { - return (Optional.class.equals(type) || (type instanceof ParameterizedType - && Optional.class.equals(((ParameterizedType) type).getRawType()))) - && message.getPayload().equals(Optional.empty()); - } - - @Override - protected boolean isEmptyPayload(Object payload) { - return payload == null || payload.equals(Optional.empty()); - } - - } /** * The metadata holder of the class with {@link RabbitListener} * and {@link RabbitHandler} annotations. @@ -1145,6 +1071,9 @@ public class RabbitListenerAnnotationBeanPostProcessor /** * A method annotated with {@link RabbitListener}, together with the annotations. + * + * @param method the method with annotations + * @param annotations on the method */ private static class ListenerMethod { diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/AbstractAdaptableMessageListener.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/AbstractAdaptableMessageListener.java index 5c0d74bd..44cccd1f 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/AbstractAdaptableMessageListener.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/AbstractAdaptableMessageListener.java @@ -365,7 +365,7 @@ public abstract class AbstractAdaptableMessageListener implements ChannelAwareMe * response message back. * @param resultArg the result object to handle (never null) * @param request the original request message - * @param channel the Rabbit channel to operate on (may be null) + * @param channel the Rabbit channel to operate on (maybe null) * @param source the source data for the method invocation - e.g. * {@code o.s.messaging.Message}; may be null * @see #buildMessage @@ -404,8 +404,8 @@ public abstract class AbstractAdaptableMessageListener implements ChannelAwareMe } else if (monoPresent && MonoHandler.isMono(resultArg.getReturnValue())) { if (!this.isManualAck) { - this.logger.warn("Container AcknowledgeMode must be MANUAL for a Mono return type; " - + "otherwise the container will ack the message immediately"); + this.logger.warn("Container AcknowledgeMode must be MANUAL for a Mono return type" + + "(or Kotlin suspend function); otherwise the container will ack the message immediately"); } MonoHandler.subscribe(resultArg.getReturnValue(), r -> asyncSuccess(resultArg, request, channel, source, r), @@ -461,7 +461,7 @@ public abstract class AbstractAdaptableMessageListener implements ChannelAwareMe } private void asyncFailure(Message request, Channel channel, Throwable t) { - this.logger.error("Future or Mono was completed with an exception for " + request, t); + this.logger.error("Future, Mono, or suspend function was completed with an exception for " + request, t); try { channel.basicNack(request.getMessageProperties().getDeliveryTag(), false, ContainerUtils.shouldRequeue(this.defaultRequeueRejected, t, this.logger)); diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/AmqpMessageHandlerMethodFactory.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/AmqpMessageHandlerMethodFactory.java new file mode 100644 index 00000000..5ef855d9 --- /dev/null +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/AmqpMessageHandlerMethodFactory.java @@ -0,0 +1,145 @@ +/* + * Copyright 2023 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.amqp.rabbit.listener.adapter; + +import java.lang.reflect.Method; +import java.lang.reflect.ParameterizedType; +import java.lang.reflect.Type; +import java.util.List; +import java.util.Optional; + +import org.springframework.core.KotlinDetector; +import org.springframework.core.MethodParameter; +import org.springframework.lang.Nullable; +import org.springframework.messaging.Message; +import org.springframework.messaging.converter.MessageConverter; +import org.springframework.messaging.handler.annotation.support.DefaultMessageHandlerMethodFactory; +import org.springframework.messaging.handler.annotation.support.MethodArgumentNotValidException; +import org.springframework.messaging.handler.annotation.support.PayloadMethodArgumentResolver; +import org.springframework.messaging.handler.invocation.HandlerMethodArgumentResolver; +import org.springframework.messaging.handler.invocation.HandlerMethodArgumentResolverComposite; +import org.springframework.messaging.handler.invocation.InvocableHandlerMethod; +import org.springframework.validation.BindingResult; +import org.springframework.validation.ObjectError; +import org.springframework.validation.Validator; + +/** + * Extension of the {@link DefaultMessageHandlerMethodFactory} for Spring AMQP requirements. + * + * @author Artem Bilan + * + * @since 2.4.16 + */ +public class AmqpMessageHandlerMethodFactory extends DefaultMessageHandlerMethodFactory { + + private final HandlerMethodArgumentResolverComposite argumentResolvers = + new HandlerMethodArgumentResolverComposite(); + + private MessageConverter messageConverter; + + private Validator validator; + + @Override + public void setMessageConverter(MessageConverter messageConverter) { + super.setMessageConverter(messageConverter); + this.messageConverter = messageConverter; + } + + @Override + public void setValidator(Validator validator) { + super.setValidator(validator); + this.validator = validator; + } + + @Override + protected List initArgumentResolvers() { + List resolvers = super.initArgumentResolvers(); + if (KotlinDetector.isKotlinPresent()) { + // Insert before PayloadMethodArgumentResolver + resolvers.add(resolvers.size() - 1, new ContinuationHandlerMethodArgumentResolver()); + } + // Has to be at the end, but before PayloadMethodArgumentResolver + resolvers.add(resolvers.size() - 1, + new OptionalEmptyAwarePayloadArgumentResolver(this.messageConverter, this.validator)); + this.argumentResolvers.addResolvers(resolvers); + return resolvers; + } + + @Override + public InvocableHandlerMethod createInvocableHandlerMethod(Object bean, Method method) { + InvocableHandlerMethod handlerMethod = new KotlinAwareInvocableHandlerMethod(bean, method); + handlerMethod.setMessageMethodArgumentResolvers(this.argumentResolvers); + return handlerMethod; + } + + private static class OptionalEmptyAwarePayloadArgumentResolver extends PayloadMethodArgumentResolver { + + OptionalEmptyAwarePayloadArgumentResolver(MessageConverter messageConverter, @Nullable Validator validator) { + super(messageConverter, validator); + } + + @Override + public Object resolveArgument(MethodParameter parameter, Message message) throws Exception { // NOSONAR + Object resolved; + try { + resolved = super.resolveArgument(parameter, message); + } + catch (MethodArgumentNotValidException ex) { + Type type = parameter.getGenericParameterType(); + if (isOptional(message, type)) { + BindingResult bindingResult = ex.getBindingResult(); + if (bindingResult != null) { + List allErrors = bindingResult.getAllErrors(); + if (allErrors.size() == 1) { + String defaultMessage = allErrors.get(0).getDefaultMessage(); + if ("Payload value must not be empty".equals(defaultMessage)) { + return Optional.empty(); + } + } + } + } + throw ex; + } + /* + * Replace Optional.empty() list elements with null. + */ + if (resolved instanceof List) { + List list = (List) resolved; + for (int i = 0; i < list.size(); i++) { + if (list.get(i).equals(Optional.empty())) { + list.set(i, null); + } + } + } + return resolved; + } + + private boolean isOptional(Message message, Type type) { + return (Optional.class.equals(type) || + (type instanceof ParameterizedType + && Optional.class.equals(((ParameterizedType) type).getRawType()))) + && message.getPayload().equals(Optional.empty()); + } + + @Override + protected boolean isEmptyPayload(Object payload) { + return payload == null || payload.equals(Optional.empty()); + } + + } + +} diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/ContinuationHandlerMethodArgumentResolver.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/ContinuationHandlerMethodArgumentResolver.java new file mode 100644 index 00000000..72a6aa03 --- /dev/null +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/ContinuationHandlerMethodArgumentResolver.java @@ -0,0 +1,50 @@ +/* + * Copyright 2023 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.amqp.rabbit.listener.adapter; + +import org.springframework.core.MethodParameter; +import org.springframework.messaging.Message; +import org.springframework.messaging.handler.invocation.HandlerMethodArgumentResolver; + +import reactor.core.publisher.Mono; + +/** + * No-op resolver for method arguments of type {@link kotlin.coroutines.Continuation}. + *

+ * This class is similar to + * {@link org.springframework.messaging.handler.annotation.reactive.ContinuationHandlerMethodArgumentResolver} + * but for regular {@link HandlerMethodArgumentResolver} contract. + * + * @author Artem Bilan + * + * @since 2.4.16 + * + * @see org.springframework.messaging.handler.annotation.reactive.ContinuationHandlerMethodArgumentResolver + */ +public class ContinuationHandlerMethodArgumentResolver implements HandlerMethodArgumentResolver { + + @Override + public boolean supportsParameter(MethodParameter parameter) { + return "kotlin.coroutines.Continuation".equals(parameter.getParameterType().getName()); + } + + @Override + public Object resolveArgument(MethodParameter parameter, Message message) { + return Mono.empty(); + } + +} diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/KotlinAwareInvocableHandlerMethod.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/KotlinAwareInvocableHandlerMethod.java new file mode 100644 index 00000000..626477ef --- /dev/null +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/KotlinAwareInvocableHandlerMethod.java @@ -0,0 +1,49 @@ +/* + * Copyright 2023 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.amqp.rabbit.listener.adapter; + +import java.lang.reflect.Method; + +import org.springframework.core.CoroutinesUtils; +import org.springframework.core.KotlinDetector; +import org.springframework.messaging.handler.invocation.InvocableHandlerMethod; + +/** + * An {@link InvocableHandlerMethod} extension for supporting Kotlin {@code suspend} function. + * + * @author Artem Bilan + * + * @since 2.4.16 + */ +public class KotlinAwareInvocableHandlerMethod extends InvocableHandlerMethod { + + public KotlinAwareInvocableHandlerMethod(Object bean, Method method) { + super(bean, method); + } + + @Override + protected Object doInvoke(Object... args) throws Exception { + Method method = getBridgedMethod(); + if (KotlinDetector.isSuspendingFunction(method)) { + return CoroutinesUtils.invokeSuspendingFunction(method, getBean(), args); + } + else { + return super.doInvoke(args); + } + } + +} diff --git a/spring-rabbit/src/test/kotlin/org/springframework/amqp/rabbit/annotation/EnableRabbitKotlinTests.kt b/spring-rabbit/src/test/kotlin/org/springframework/amqp/rabbit/annotation/EnableRabbitKotlinTests.kt index 9960d52c..981b2848 100644 --- a/spring-rabbit/src/test/kotlin/org/springframework/amqp/rabbit/annotation/EnableRabbitKotlinTests.kt +++ b/spring-rabbit/src/test/kotlin/org/springframework/amqp/rabbit/annotation/EnableRabbitKotlinTests.kt @@ -20,6 +20,7 @@ import assertk.assertThat import assertk.assertions.isEqualTo import assertk.assertions.isTrue import org.junit.jupiter.api.Test +import org.springframework.amqp.core.AcknowledgeMode import org.springframework.amqp.rabbit.config.SimpleRabbitListenerContainerFactory import org.springframework.amqp.rabbit.connection.CachingConnectionFactory import org.springframework.amqp.rabbit.core.RabbitTemplate @@ -56,14 +57,14 @@ class EnableRabbitKotlinTests { private lateinit var config: Config @Test - fun `send and wait for consume` () { + fun `send and wait for consume`() { val template = RabbitTemplate(this.config.cf()) template.convertAndSend("kotlinQueue", "test") assertThat(this.config.latch.await(10, TimeUnit.SECONDS)).isTrue(); } @Test - fun `send and wait for consume with EH` () { + fun `send and wait for consume with EH`() { val template = RabbitTemplate(this.config.cf()) template.convertAndSend("kotlinQueue1", "test") assertThat(this.config.ehLatch.await(10, TimeUnit.SECONDS)).isTrue(); @@ -78,27 +79,22 @@ class EnableRabbitKotlinTests { val latch = CountDownLatch(1) @RabbitListener(queues = ["kotlinQueue"]) - fun handle(@Suppress("UNUSED_PARAMETER") data: String) { + suspend fun handle(@Suppress("UNUSED_PARAMETER") data: String) { this.latch.countDown() } @Bean - fun rabbitListenerContainerFactory(cf: CachingConnectionFactory): SimpleRabbitListenerContainerFactory { - val factory = SimpleRabbitListenerContainerFactory() - factory.setConnectionFactory(cf) - return factory - } + fun rabbitListenerContainerFactory(cf: CachingConnectionFactory) = + SimpleRabbitListenerContainerFactory().also { + it.setAcknowledgeMode(AcknowledgeMode.MANUAL) + it.setConnectionFactory(cf) + } @Bean - fun cf(): CachingConnectionFactory { - return CachingConnectionFactory( - RabbitAvailableCondition.getBrokerRunning().connectionFactory) - } + fun cf() = CachingConnectionFactory(RabbitAvailableCondition.getBrokerRunning().connectionFactory) @Bean - fun multi(): Multi { - return Multi() - } + fun multi() = Multi() @Bean fun proxyListenerPostProcessor(): BeanPostProcessor? { diff --git a/src/reference/asciidoc/amqp.adoc b/src/reference/asciidoc/amqp.adoc index 60e8e69b..3904eada 100644 --- a/src/reference/asciidoc/amqp.adoc +++ b/src/reference/asciidoc/amqp.adoc @@ -3075,7 +3075,7 @@ Each queue needed a separate property. ====== Reply Management The existing support in `MessageListenerAdapter` already lets your method have a non-void return type. -When that is the case, the result of the invocation is encapsulated in a message sent to the the address specified in the `ReplyToAddress` header of the original message, or to the default address configured on the listener. +When that is the case, the result of the invocation is encapsulated in a message sent to the address specified in the `ReplyToAddress` header of the original message, or to the default address configured on the listener. You can set that default address by using the `@SendTo` annotation of the messaging abstraction. Assuming our `processOrder` method should now return an `OrderStatus`, we can write it as follows to automatically send a reply: @@ -3635,6 +3635,10 @@ If some exception occurs within the listener method that prevents creation of th Starting with versions 2.2.21, 2.3.13, 2.4.1, the `AcknowledgeMode` will be automatically set the `MANUAL` when async return types are detected. In addition, incoming messages with fatal exceptions will be negatively acknowledged individually, previously any prior unacknowledged message were also negatively acknowledged. +Starting with version 2.4.16, the `@RabbitListener` (and `@RabbitHandler`) methods can be marked with Kotlin `suspend` and the whole handling process and reply producing (optional) happens on respective Kotlin coroutine. +The above mentioned rules about `AcknowledgeMode.MANUAL` still apply. +The `org.jetbrains.kotlinx:kotlinx-coroutines-reactor` dependency must be present in classpath to allow `suspend` function invocations. + [[threading]] ===== Threading and Asynchronous Consumers diff --git a/src/reference/asciidoc/whats-new.adoc b/src/reference/asciidoc/whats-new.adoc index 671cd01d..01f1ca72 100644 --- a/src/reference/asciidoc/whats-new.adoc +++ b/src/reference/asciidoc/whats-new.adoc @@ -24,7 +24,7 @@ See <> for more information. ==== Remoting Support -Support remoting using Spring Framework's RMI support is deprecated and will be removed in 3.0. +Support of remoting using Spring Framework's RMI support is deprecated and will be removed in 3.0. See <> for more information. ==== Message Converter Changes