diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/GatewayParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/GatewayParser.java index 6047757720..bdfd72e9e4 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/GatewayParser.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/GatewayParser.java @@ -41,7 +41,7 @@ import org.springframework.util.xml.DomUtils; public class GatewayParser extends AbstractSimpleBeanDefinitionParser { private static String[] referenceAttributes = new String[] { - "default-request-channel", "default-reply-channel", "error-channel", "message-mapper" + "default-request-channel", "default-reply-channel", "error-channel", "message-mapper", "async-executor" }; private static String[] innerAttributes = new String[] { @@ -81,6 +81,7 @@ public class GatewayParser extends AbstractSimpleBeanDefinitionParser { IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "request-timeout", "defaultRequestTimeout"); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "reply-timeout", "defaultReplyTimeout"); IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "error-channel"); + IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "async-executor"); } private void postProcessGateway(BeanDefinitionBuilder builder, Element element) { diff --git a/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.1.xsd b/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.1.xsd index 94a59e8a1c..376db49549 100644 --- a/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.1.xsd +++ b/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.1.xsd @@ -594,6 +594,18 @@ + + + + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/GatewayParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/GatewayParserTests.java index 59d14a2b16..ed5b8f39ce 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/GatewayParserTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/GatewayParserTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2010 the original author or authors. + * Copyright 2002-2011 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. @@ -18,18 +18,23 @@ package org.springframework.integration.config.xml; import static org.junit.Assert.assertEquals; +import java.util.concurrent.Callable; import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; import org.junit.Test; - +import org.springframework.beans.factory.BeanNameAware; import org.springframework.context.ApplicationContext; import org.springframework.context.support.ClassPathXmlApplicationContext; +import org.springframework.core.task.SimpleAsyncTaskExecutor; import org.springframework.integration.Message; import org.springframework.integration.MessageChannel; import org.springframework.integration.core.PollableChannel; import org.springframework.integration.gateway.TestService; import org.springframework.integration.message.GenericMessage; import org.springframework.integration.support.MessageBuilder; +import org.springframework.scheduling.annotation.AsyncResult; /** * @author Mark Fisher @@ -67,6 +72,19 @@ public class GatewayParserTests { assertEquals("foo", result); } + @Test + public void testAsyncGateway() throws Exception { + ApplicationContext context = new ClassPathXmlApplicationContext("gatewayParserTests.xml", this.getClass()); + PollableChannel requestChannel = (PollableChannel) context.getBean("requestChannel"); + MessageChannel replyChannel = (MessageChannel) context.getBean("replyChannel"); + this.startResponder(requestChannel, replyChannel); + TestService service = context.getBean("async", TestService.class); + Future> result = service.async("foo"); + Message reply = result.get(1, TimeUnit.SECONDS); + assertEquals("foo", reply.getPayload()); + assertEquals("testExecutor", reply.getHeaders().get("executor")); + } + private void startResponder(final PollableChannel requestChannel, final MessageChannel replyChannel) { Executors.newSingleThreadExecutor().execute(new Runnable() { @@ -79,4 +97,32 @@ public class GatewayParserTests { }); } + + @SuppressWarnings("unused") + private static class TestExecutor extends SimpleAsyncTaskExecutor implements BeanNameAware { + + private static final long serialVersionUID = 1L; + + private volatile String beanName; + + public void setBeanName(String beanName) { + this.beanName = beanName; + } + + @Override + @SuppressWarnings({"rawtypes", "unchecked"}) + public Future submit(Callable task) { + try { + Future result = super.submit(task); + Message message = (Message) result.get(1, TimeUnit.SECONDS); + Message modifiedMessage = MessageBuilder.fromMessage(message) + .setHeader("executor", this.beanName).build(); + return new AsyncResult(modifiedMessage); + } + catch (Exception e) { + throw new IllegalStateException("unexpected exception in testExecutor", e); + } + } + } + } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/gatewayParserTests.xml b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/gatewayParserTests.xml index 3da01a05f4..4a0a9d97b9 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/gatewayParserTests.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/gatewayParserTests.xml @@ -1,11 +1,11 @@ + http://www.springframework.org/schema/integration/spring-integration.xsd"> @@ -29,8 +29,16 @@ default-request-channel="requestChannel" default-reply-channel="replyChannel" default-reply-timeout="5000"/> - + + + + + 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 8ab2ac0acc..7287c51a38 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 @@ -16,6 +16,8 @@ package org.springframework.integration.gateway; +import java.util.concurrent.Future; + import org.springframework.integration.Message; import org.springframework.integration.annotation.Payload; @@ -42,4 +44,6 @@ public interface TestService { @Payload("#method + #args.length") String requestReplyWithPayloadAnnotation(); + Future> async(String s); + }