From 1e5b41cfeec6b75fed8955945b952c3b270c46a7 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Tue, 6 Jun 2023 09:54:28 +0200 Subject: [PATCH] GH-SCF-1045 Fix type discovery in DefaultPollableMessageSource --- .../stream/binder/PollableConsumerTests.java | 5 ++--- .../binder/DefaultPollableMessageSource.java | 22 ++++++++++--------- 2 files changed, 14 insertions(+), 13 deletions(-) diff --git a/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/binder/PollableConsumerTests.java b/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/binder/PollableConsumerTests.java index 4fb0d81f5..9c2ecb2c0 100644 --- a/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/binder/PollableConsumerTests.java +++ b/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/binder/PollableConsumerTests.java @@ -24,7 +24,6 @@ import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; import org.junit.jupiter.api.BeforeAll; -import org.junit.jupiter.api.Disabled; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.TestInstance; @@ -210,7 +209,7 @@ public class PollableConsumerTests { } @Test - @Disabled +// @Disabled void testConvertList() { TestChannelBinder binder = createBinder(); MessageConverterConfigurer configurer = this.context @@ -242,7 +241,7 @@ public class PollableConsumerTests { } @Test - @Disabled +// @Disabled void testConvertMap() { TestChannelBinder binder = createBinder(); MessageConverterConfigurer configurer = this.context diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultPollableMessageSource.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultPollableMessageSource.java index 4788d99b4..b6895bc08 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultPollableMessageSource.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultPollableMessageSource.java @@ -27,6 +27,7 @@ import org.apache.commons.logging.LogFactory; import org.springframework.aop.framework.ProxyFactory; import org.springframework.aop.support.NameMatchMethodPointcutAdvisor; +import org.springframework.cloud.function.context.catalog.FunctionTypeUtils; import org.springframework.context.Lifecycle; import org.springframework.core.AttributeAccessor; import org.springframework.core.ParameterizedTypeReference; @@ -116,7 +117,7 @@ public class DefaultPollableMessageSource @Override public Object invoke(MethodInvocation invocation) throws Throwable { Object result = invocation.proceed(); - if (result instanceof Message received) { + if (result instanceof Message received) { for (ChannelInterceptor interceptor : this.interceptors) { received = interceptor.preSend(received, dummyChannel); if (received == null) { @@ -308,16 +309,17 @@ public class DefaultPollableMessageSource private Message receive(ParameterizedTypeReference type) { Message message = this.source.receive(); if (message != null && type != null && this.messageConverter != null) { - Class targetType = type == null ? Object.class - : type.getType() instanceof Class clazz ? clazz - : Object.class; - Object payload = this.messageConverter.fromMessage(message, targetType, type); - if (payload == null) { - throw new MessageConversionException(message, - "No converter could convert Message"); + Object payload; + if (FunctionTypeUtils.isTypeCollection(type.getType()) || FunctionTypeUtils.isTypeMap(type.getType())) { + payload = this.messageConverter.fromMessage(message, FunctionTypeUtils.getRawType(type.getType()), type.getType()); } - message = MessageBuilder.withPayload(payload) - .copyHeaders(message.getHeaders()).build(); + else { + payload = this.messageConverter.fromMessage(message, (Class) type.getType()); + } + if (payload == null) { + throw new MessageConversionException(message, "No converter could convert Message"); + } + message = MessageBuilder.withPayload(payload).copyHeaders(message.getHeaders()).build(); } return message; }