diff --git a/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationProperties.java b/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationProperties.java index 2c021a9021..dfdad45153 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationProperties.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationProperties.java @@ -71,6 +71,11 @@ public final class IntegrationProperties { */ public static final String REQUIRE_COMPONENT_ANNOTATION = INTEGRATION_PROPERTIES_PREFIX + "messagingAnnotations.require.componentAnnotation"; + /** + * Specifies the value of {@link org.springframework.integration.config.annotation.MessagingAnnotationPostProcessor#requireComponentAnnotation}. + */ + public static final String GATEWAY_CONVERT_RECEIVE_MESSAGE = INTEGRATION_PROPERTIES_PREFIX + "messagingGateway.convertReceiveMessage"; + private static Properties defaults; static { diff --git a/spring-integration-core/src/main/java/org/springframework/integration/gateway/GatewayProxyFactoryBean.java b/spring-integration-core/src/main/java/org/springframework/integration/gateway/GatewayProxyFactoryBean.java index f37a286b7b..b7e63f47a5 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/gateway/GatewayProxyFactoryBean.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/gateway/GatewayProxyFactoryBean.java @@ -48,6 +48,7 @@ import org.springframework.expression.common.LiteralExpression; import org.springframework.expression.spel.standard.SpelExpressionParser; import org.springframework.integration.annotation.Gateway; import org.springframework.integration.annotation.GatewayHeader; +import org.springframework.integration.context.IntegrationProperties; import org.springframework.integration.endpoint.AbstractEndpoint; import org.springframework.integration.support.channel.BeanFactoryChannelResolver; import org.springframework.integration.support.management.TrackableComponent; @@ -135,6 +136,8 @@ public class GatewayProxyFactoryBean extends AbstractEndpoint private volatile MethodArgsMessageMapper argsMapper; + private volatile boolean convertReceiveMessage; + /** * Create a Factory whose service interface type can be configured by setter injection. * If none is set, it will fall back to the default service interface type, @@ -358,6 +361,9 @@ public class GatewayProxyFactoryBean extends AbstractEndpoint this.asyncSubmitListenableType = submitType.getClass(); } } + + this.convertReceiveMessage = + getIntegrationProperty(IntegrationProperties.GATEWAY_CONVERT_RECEIVE_MESSAGE, Boolean.class); this.initialized = true; } } @@ -452,7 +458,12 @@ public class GatewayProxyFactoryBean extends AbstractEndpoint if (paramCount == 0 && !hasPayloadExpression) { if (shouldReply) { if (shouldReturnMessage) { - return gateway.receive(); + if (this.convertReceiveMessage) { + return gateway.receive(); + } + else { + return gateway.receiveMessage(); + } } response = gateway.receive(); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/gateway/MessagingGatewaySupport.java b/spring-integration-core/src/main/java/org/springframework/integration/gateway/MessagingGatewaySupport.java index bdc2e8d212..31dac5d3a3 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/gateway/MessagingGatewaySupport.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/gateway/MessagingGatewaySupport.java @@ -59,7 +59,7 @@ public abstract class MessagingGatewaySupport extends AbstractEndpoint private final SimpleMessageConverter messageConverter = new SimpleMessageConverter(); - private final MessagingTemplate messagingTemplate; + protected final MessagingTemplate messagingTemplate; private final HistoryWritingMessagePostProcessor historyWritingPostProcessor = new HistoryWritingMessagePostProcessor(); @@ -393,6 +393,14 @@ public abstract class MessagingGatewaySupport extends AbstractEndpoint return this.messagingTemplate.receiveAndConvert(replyChannel, null); } + protected Message receiveMessage() { + initializeIfNecessary(); + MessageChannel replyChannel = getReplyChannel(); + Assert.state(replyChannel instanceof PollableChannel, + "receive is not supported, because no pollable reply channel has been configured"); + return this.messagingTemplate.receive(replyChannel); + } + protected Object sendAndReceive(Object object) { return this.doSendAndReceive(object, true); } diff --git a/spring-integration-core/src/main/resources/META-INF/spring.integration.default.properties b/spring-integration-core/src/main/resources/META-INF/spring.integration.default.properties index e0be1b5889..b433a27493 100644 --- a/spring-integration-core/src/main/resources/META-INF/spring.integration.default.properties +++ b/spring-integration-core/src/main/resources/META-INF/spring.integration.default.properties @@ -4,3 +4,4 @@ spring.integration.channels.maxBroadcastSubscribers=0x7fffffff spring.integration.taskScheduler.poolSize=10 spring.integration.messagingTemplate.throwExceptionOnLateReply=false spring.integration.messagingAnnotations.require.componentAnnotation=false +spring.integration.messagingGateway.convertReceiveMessage=false diff --git a/spring-integration-core/src/test/java/org/springframework/integration/gateway/GatewayProxyFactoryBeanTests.java b/spring-integration-core/src/test/java/org/springframework/integration/gateway/GatewayProxyFactoryBeanTests.java index 4d92868915..4c1e203e18 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/gateway/GatewayProxyFactoryBeanTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/gateway/GatewayProxyFactoryBeanTests.java @@ -16,14 +16,20 @@ package org.springframework.integration.gateway; +import static org.hamcrest.Matchers.containsString; import static org.hamcrest.Matchers.equalTo; +import static org.hamcrest.Matchers.instanceOf; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertThat; +import static org.junit.Assert.fail; +import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.willAnswer; import static org.mockito.Mockito.mock; import java.lang.reflect.Method; import java.util.Collections; +import java.util.Properties; import java.util.Random; import java.util.concurrent.CountDownLatch; import java.util.concurrent.Executor; @@ -44,6 +50,8 @@ import org.springframework.expression.Expression; import org.springframework.expression.common.LiteralExpression; import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.channel.QueueChannel; +import org.springframework.integration.context.IntegrationContextUtils; +import org.springframework.integration.context.IntegrationProperties; import org.springframework.integration.endpoint.EventDrivenConsumer; import org.springframework.integration.support.utils.IntegrationUtils; import org.springframework.messaging.Message; @@ -137,6 +145,57 @@ public class GatewayProxyFactoryBeanTests { assertEquals("foo", result); } + @Test + public void testReceiveMessage() throws Exception { + QueueChannel replyChannel = new QueueChannel(); + replyChannel.send(new GenericMessage<>("foo")); + GatewayProxyFactoryBean proxyFactory = new GatewayProxyFactoryBean(); + proxyFactory.setServiceInterface(TestService.class); + proxyFactory.setDefaultReplyChannel(replyChannel); + + proxyFactory.setBeanFactory(mock(BeanFactory.class)); + proxyFactory.afterPropertiesSet(); + TestService service = (TestService) proxyFactory.getObject(); + Message message = service.getMessage(); + assertNotNull(message); + assertEquals("foo", message.getPayload()); + } + + @Test + public void testReceiveMessageConvert() throws Exception { + QueueChannel replyChannel = new QueueChannel(); + replyChannel.send(new GenericMessage<>("foo")); + GatewayProxyFactoryBean proxyFactory = new GatewayProxyFactoryBean(); + proxyFactory.setServiceInterface(TestService.class); + proxyFactory.setDefaultReplyChannel(replyChannel); + + BeanFactory beanFactory = mock(BeanFactory.class); + + given(beanFactory.containsBean(IntegrationContextUtils.INTEGRATION_GLOBAL_PROPERTIES_BEAN_NAME)) + .willReturn(true); + + willAnswer(invocation -> { + Properties properties = new Properties(); + properties.setProperty(IntegrationProperties.GATEWAY_CONVERT_RECEIVE_MESSAGE, "true"); + return properties; + }) + .given(beanFactory) + .getBean(IntegrationContextUtils.INTEGRATION_GLOBAL_PROPERTIES_BEAN_NAME, Properties.class); + + proxyFactory.setBeanFactory(beanFactory); + proxyFactory.afterPropertiesSet(); + TestService service = (TestService) proxyFactory.getObject(); + try { + service.getMessage(); + fail("ClassCastException expected"); + } + catch (Exception e) { + assertThat(e, instanceOf(ClassCastException.class)); + assertThat(e.getMessage(), + containsString("java.lang.String cannot be cast to org.springframework.messaging.Message")); + } + } + @Test public void testRequestReplyWithTypeConversion() throws Exception { final QueueChannel requestChannel = new QueueChannel(); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/gateway/TestService.java b/spring-integration-core/src/test/java/org/springframework/integration/gateway/TestService.java index 017effb432..2147aa6e63 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/gateway/TestService.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/gateway/TestService.java @@ -39,6 +39,8 @@ public interface TestService { String solicitResponse(); + Message getMessage(); + Integer requestReplyWithIntegers(Integer input); String requestReplyWithMessageParameter(Message message);