diff --git a/spring-integration-core/src/main/java/org/springframework/integration/core/AsyncMessagingOperations.java b/spring-integration-core/src/main/java/org/springframework/integration/core/AsyncMessagingOperations.java index ee8217ff02..7b29161f94 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/core/AsyncMessagingOperations.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/core/AsyncMessagingOperations.java @@ -27,11 +27,11 @@ import org.springframework.integration.MessageChannel; */ public interface AsyncMessagingOperations { -

Future> asyncReceive(); + Future> asyncReceive(); -

Future> asyncReceive(PollableChannel channel); + Future> asyncReceive(PollableChannel channel); -

Future> asyncReceive(String channelName); + Future> asyncReceive(String channelName); Future asyncReceiveAndConvert(); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/core/AsyncMessagingTemplate.java b/spring-integration-core/src/main/java/org/springframework/integration/core/AsyncMessagingTemplate.java index 608024019f..12fc4fb0cc 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/core/AsyncMessagingTemplate.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/core/AsyncMessagingTemplate.java @@ -42,25 +42,25 @@ public class AsyncMessagingTemplate extends MessagingTemplate implements AsyncMe (AsyncTaskExecutor) executor : new TaskExecutorAdapter(executor); } - public

Future> asyncReceive() { - return this.executor.submit(new Callable>() { - public Message

call() throws Exception { + public Future> asyncReceive() { + return this.executor.submit(new Callable>() { + public Message call() throws Exception { return receive(); } }); } - public

Future> asyncReceive(final PollableChannel channel) { - return this.executor.submit(new Callable>() { - public Message

call() throws Exception { + public Future> asyncReceive(final PollableChannel channel) { + return this.executor.submit(new Callable>() { + public Message call() throws Exception { return receive(channel); } }); } - public

Future> asyncReceive(final String channelName) { - return this.executor.submit(new Callable>() { - public Message

call() throws Exception { + public Future> asyncReceive(final String channelName) { + return this.executor.submit(new Callable>() { + public Message call() throws Exception { return receive(channelName); } }); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/core/AsyncMessagingTemplateTests.java b/spring-integration-core/src/test/java/org/springframework/integration/core/AsyncMessagingTemplateTests.java index ed0d3a3f8f..40a9948a8a 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/core/AsyncMessagingTemplateTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/core/AsyncMessagingTemplateTests.java @@ -23,6 +23,7 @@ import static org.junit.Assert.fail; import java.util.concurrent.CancellationException; import java.util.concurrent.ExecutionException; +import java.util.concurrent.Executors; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; @@ -31,9 +32,12 @@ import org.junit.Test; import org.springframework.context.support.StaticApplicationContext; import org.springframework.integration.Message; +import org.springframework.integration.MessageChannel; import org.springframework.integration.MessagingException; import org.springframework.integration.channel.DirectChannel; +import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.handler.AbstractReplyProducingMessageHandler; +import org.springframework.integration.message.GenericMessage; import org.springframework.integration.support.MessageBuilder; import org.springframework.util.Assert; @@ -43,6 +47,57 @@ import org.springframework.util.Assert; */ public class AsyncMessagingTemplateTests { + @Test + public void asyncReceiveWithDefaultChannel() throws Exception { + QueueChannel channel = new QueueChannel(); + AsyncMessagingTemplate template = new AsyncMessagingTemplate(); + template.setDefaultChannel(channel); + Future> result = template.asyncReceive(); + sendMessageAfterDelay(channel, new GenericMessage("test"), 200); + long start = System.currentTimeMillis(); + assertNotNull(result.get(1000, TimeUnit.MILLISECONDS)); + long elapsed = System.currentTimeMillis() - start; + assertEquals("test", result.get().getPayload()); + assertTrue(elapsed >= 200); + } + + @Test + public void asyncReceiveWithExplicitChannel() throws Exception { + QueueChannel channel = new QueueChannel(); + AsyncMessagingTemplate template = new AsyncMessagingTemplate(); + Future> result = template.asyncReceive(channel); + sendMessageAfterDelay(channel, new GenericMessage("test"), 200); + long start = System.currentTimeMillis(); + assertNotNull(result.get(1000, TimeUnit.MILLISECONDS)); + long elapsed = System.currentTimeMillis() - start; + assertEquals("test", result.get().getPayload()); + assertTrue(elapsed >= 200); + } + + @Test + public void asyncReceiveWithResolvedChannel() throws Exception { + StaticApplicationContext context = new StaticApplicationContext(); + context.registerSingleton("testChannel", QueueChannel.class); + context.refresh(); + QueueChannel channel = context.getBean("testChannel", QueueChannel.class); + AsyncMessagingTemplate template = new AsyncMessagingTemplate(); + template.setBeanFactory(context); + Future> result = template.asyncReceive("testChannel"); + sendMessageAfterDelay(channel, new GenericMessage("test"), 200); + long start = System.currentTimeMillis(); + assertNotNull(result.get(1000, TimeUnit.MILLISECONDS)); + long elapsed = System.currentTimeMillis() - start; + assertTrue(elapsed >= 200); + assertEquals("test", result.get().getPayload()); + } + + @Test(expected = TimeoutException.class) + public void asyncReceiveWithTimeoutException() throws Exception { + AsyncMessagingTemplate template = new AsyncMessagingTemplate(); + Future> result = template.asyncReceive(new QueueChannel()); + result.get(100, TimeUnit.MILLISECONDS); + } + @Test public void asyncSendAndReceiveWithDefaultChannel() throws Exception { DirectChannel channel = new DirectChannel(); @@ -177,6 +232,21 @@ public class AsyncMessagingTemplateTests { } + private static void sendMessageAfterDelay(final MessageChannel channel, final GenericMessage message, final int delay) { + Executors.newSingleThreadExecutor().execute(new Runnable() { + public void run() { + try { + Thread.sleep(delay); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return; + } + channel.send(message); + } + }); + } + private static class EchoHandler extends AbstractReplyProducingMessageHandler { private final long delay;