INT-1132: Rescheduling Support for DelayHandler
- Add support for DelayHandler to reschedule persisted Messages on startup via 'implements ApplicationListener<ContextRefreshedEvent>' as late as possible - Change DelayHandler dependency to MessageGroupStore - Add required DelayHandler.messageGroupId property - Make registered by SI-namespace TaskScheduler as default for DelayHandler - Remove 'required' from delayer xml-attribute 'default-delay' as redundant - Additional refactoring & polishing around <delayer> - Polishing delayer's Tests - Tests for 'rescheduling' - Integration test for 'rescheduling' with JdbcMS INT-1132: add 'initializingLatch' to DelayHandler INT-1132: additional polishing INT-1132: add LIFO JMS polling test-case INT-1132: Changes according to PR comments INT-1132: Changes according to PR comments 2 Polishing Remove blocking calls from tests.
This commit is contained in:
committed by
Gary Russell
parent
84f66e059b
commit
d1d4faa151
@@ -2,11 +2,9 @@
|
||||
<beans:beans xmlns="http://www.springframework.org/schema/integration"
|
||||
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||
xmlns:beans="http://www.springframework.org/schema/beans"
|
||||
xmlns:task="http://www.springframework.org/schema/task"
|
||||
xmlns:p="http://www.springframework.org/schema/p"
|
||||
xsi:schemaLocation="http://www.springframework.org/schema/beans
|
||||
http://www.springframework.org/schema/beans/spring-beans.xsd
|
||||
http://www.springframework.org/schema/task
|
||||
http://www.springframework.org/schema/task/spring-task.xsd
|
||||
http://www.springframework.org/schema/integration
|
||||
http://www.springframework.org/schema/integration/spring-integration.xsd">
|
||||
|
||||
@@ -22,8 +20,7 @@
|
||||
default-delay="1234"
|
||||
delay-header-name="foo"
|
||||
order="99"
|
||||
send-timeout="987"
|
||||
wait-for-tasks-to-complete-on-shutdown="true"/>
|
||||
send-timeout="987"/>
|
||||
|
||||
<delayer id="delayerWithCustomScheduler"
|
||||
input-channel="input"
|
||||
@@ -37,7 +34,9 @@
|
||||
default-delay="0"
|
||||
message-store="testMessageStore"/>
|
||||
|
||||
<task:scheduler id="testScheduler" pool-size="7"/>
|
||||
<beans:bean id="testScheduler" class="org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler"
|
||||
p:poolSize="7"
|
||||
p:waitForTasksToCompleteOnShutdown="true"/>
|
||||
|
||||
<beans:bean id="testMessageStore" class="org.springframework.integration.store.SimpleMessageStore"/>
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2010 the original author or authors.
|
||||
* Copyright 2002-2012 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.
|
||||
@@ -17,6 +17,8 @@
|
||||
package org.springframework.integration.config.xml;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertNull;
|
||||
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
@@ -33,6 +35,7 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
|
||||
|
||||
/**
|
||||
* @author Mark Fisher
|
||||
* @author Artem Bilan
|
||||
* @since 1.0.3
|
||||
*/
|
||||
@RunWith(SpringJUnit4ClassRunner.class)
|
||||
@@ -57,8 +60,7 @@ public class DelayerParserTests {
|
||||
assertEquals("foo", accessor.getPropertyValue("delayHeaderName"));
|
||||
assertEquals(new Long(987), new DirectFieldAccessor(
|
||||
accessor.getPropertyValue("messagingTemplate")).getPropertyValue("sendTimeout"));
|
||||
assertEquals(Boolean.TRUE, new DirectFieldAccessor(
|
||||
accessor.getPropertyValue("taskScheduler")).getPropertyValue("waitForTasksToCompleteOnShutdown"));
|
||||
assertNull(accessor.getPropertyValue("taskScheduler"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -73,6 +75,9 @@ public class DelayerParserTests {
|
||||
assertEquals(context.getBean("output"), accessor.getPropertyValue("outputChannel"));
|
||||
assertEquals(new Long(0), accessor.getPropertyValue("defaultDelay"));
|
||||
assertEquals(context.getBean("testScheduler"), accessor.getPropertyValue("taskScheduler"));
|
||||
assertNotNull(accessor.getPropertyValue("taskScheduler"));
|
||||
assertEquals(Boolean.TRUE, new DirectFieldAccessor(
|
||||
accessor.getPropertyValue("taskScheduler")).getPropertyValue("waitForTasksToCompleteOnShutdown"));
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -2,11 +2,9 @@
|
||||
<beans:beans xmlns="http://www.springframework.org/schema/integration"
|
||||
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||
xmlns:beans="http://www.springframework.org/schema/beans"
|
||||
xmlns:task="http://www.springframework.org/schema/task"
|
||||
xmlns:p="http://www.springframework.org/schema/p"
|
||||
xsi:schemaLocation="http://www.springframework.org/schema/beans
|
||||
http://www.springframework.org/schema/beans/spring-beans.xsd
|
||||
http://www.springframework.org/schema/task
|
||||
http://www.springframework.org/schema/task/spring-task.xsd
|
||||
http://www.springframework.org/schema/integration
|
||||
http://www.springframework.org/schema/integration/spring-integration.xsd">
|
||||
|
||||
@@ -22,26 +20,33 @@
|
||||
default-delay="1000"
|
||||
delay-header-name="foo"
|
||||
order="99"
|
||||
send-timeout="1000"
|
||||
wait-for-tasks-to-complete-on-shutdown="true"/>
|
||||
send-timeout="1000"/>
|
||||
|
||||
<channel id="inputB"/>
|
||||
<channel id="outputB"/>
|
||||
<channel id="outputB1">
|
||||
<queue />
|
||||
</channel>
|
||||
|
||||
|
||||
<delayer id="delayerWithCustomScheduler"
|
||||
input-channel="inputB"
|
||||
output-channel="outputB"
|
||||
default-delay="1000"
|
||||
send-timeout="20000"
|
||||
scheduler="multiThreadScheduler"/>
|
||||
scheduler="multiThreadScheduler"/>
|
||||
|
||||
<chain input-channel="delayerInsideChain" output-channel="outputA">
|
||||
<transformer expression="payload.toUpperCase()"/>
|
||||
<delayer id="delayerInsideChain" default-delay="1000"/>
|
||||
<transformer expression="payload.toLowerCase()"/>
|
||||
</chain>
|
||||
|
||||
<beans:bean id="multiThreadScheduler" class="org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler"
|
||||
p:poolSize="5"
|
||||
p:waitForTasksToCompleteOnShutdown="true"/>
|
||||
|
||||
<task:scheduler id="multiThreadScheduler" pool-size="5"/>
|
||||
|
||||
<service-activator input-channel="outputB" output-channel="outputB1" method="processMessage" ref="sampleHandler"/>
|
||||
|
||||
|
||||
<beans:bean id="sampleHandler" class="org.springframework.integration.config.xml.DelayerUsageTests$SampleService"/>
|
||||
|
||||
|
||||
</beans:beans>
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2010 the original author or authors.
|
||||
* Copyright 2002-2012 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.
|
||||
@@ -16,13 +16,15 @@
|
||||
|
||||
package org.springframework.integration.config.xml;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.beans.factory.annotation.Qualifier;
|
||||
import org.springframework.integration.Message;
|
||||
import org.springframework.integration.MessageChannel;
|
||||
import org.springframework.integration.core.PollableChannel;
|
||||
import org.springframework.integration.message.GenericMessage;
|
||||
@@ -32,64 +34,81 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
|
||||
|
||||
/**
|
||||
* @author Oleg Zhurakousky
|
||||
* @author Artem Bilan
|
||||
* @since 1.0.3
|
||||
*/
|
||||
@RunWith(SpringJUnit4ClassRunner.class)
|
||||
@ContextConfiguration
|
||||
public class DelayerUsageTests {
|
||||
|
||||
|
||||
@Autowired @Qualifier("inputA")
|
||||
private MessageChannel inputA;
|
||||
|
||||
@Autowired
|
||||
private MessageChannel delayerInsideChain;
|
||||
|
||||
@Autowired @Qualifier("outputA")
|
||||
private PollableChannel outputA;
|
||||
|
||||
|
||||
@Autowired @Qualifier("inputB")
|
||||
private MessageChannel inputB;
|
||||
|
||||
@Autowired @Qualifier("outputB1")
|
||||
private PollableChannel outputB1;
|
||||
|
||||
|
||||
@Test
|
||||
public void testDelayWithDefaultScheduler(){
|
||||
long start = System.currentTimeMillis();
|
||||
inputA.send(new GenericMessage<String>("Hello"));
|
||||
outputA.receive();
|
||||
assertNotNull(outputA.receive(10000));
|
||||
assertTrue((System.currentTimeMillis() - start) >= 1000);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testDelayWithDefaultSchedulerCustomDelayHeader(){
|
||||
MessageBuilder<String> builder = MessageBuilder.withPayload("Hello");
|
||||
// set custom delay header
|
||||
builder.setHeader("foo", 2000);
|
||||
long start = System.currentTimeMillis();
|
||||
inputA.send(builder.build());
|
||||
outputA.receive();
|
||||
inputA.send(builder.build());
|
||||
assertNotNull(outputA.receive(10000));
|
||||
assertTrue((System.currentTimeMillis() - start) >= 2000);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testDelayWithCustomScheduler(){
|
||||
public void testDelayWithCustomScheduler(){
|
||||
long start = System.currentTimeMillis();
|
||||
inputB.send(new GenericMessage<String>("1"));
|
||||
inputB.send(new GenericMessage<String>("1"));
|
||||
inputB.send(new GenericMessage<String>("2"));
|
||||
inputB.send(new GenericMessage<String>("3"));
|
||||
inputB.send(new GenericMessage<String>("4"));
|
||||
inputB.send(new GenericMessage<String>("5"));
|
||||
inputB.send(new GenericMessage<String>("6"));
|
||||
inputB.send(new GenericMessage<String>("7"));
|
||||
outputB1.receive();
|
||||
outputB1.receive();
|
||||
outputB1.receive();
|
||||
outputB1.receive();
|
||||
outputB1.receive();
|
||||
outputB1.receive();
|
||||
outputB1.receive();
|
||||
|
||||
assertNotNull(outputB1.receive(10000));
|
||||
assertNotNull(outputB1.receive(10000));
|
||||
assertNotNull(outputB1.receive(10000));
|
||||
assertNotNull(outputB1.receive(10000));
|
||||
assertNotNull(outputB1.receive(10000));
|
||||
assertNotNull(outputB1.receive(10000));
|
||||
assertNotNull(outputB1.receive(10000));
|
||||
|
||||
// must execute under 3 seconds, since threadPool is set too 5.
|
||||
// first batch is 5 concurrent invocations on SA, then 2 more
|
||||
// elapsed time for the whole execution should be a bit over 2 seconds depending on the hardware
|
||||
assertTrue(((System.currentTimeMillis() - start) >= 1000) && ((System.currentTimeMillis() - start) < 3000));
|
||||
}
|
||||
|
||||
|
||||
|
||||
@Test //INT-1132
|
||||
public void testDelayerInsideChain(){
|
||||
long start = System.currentTimeMillis();
|
||||
delayerInsideChain.send(new GenericMessage<String>("Hello"));
|
||||
Message<?> message = outputA.receive(10000);
|
||||
assertNotNull(message);
|
||||
assertTrue((System.currentTimeMillis() - start) >= 1000);
|
||||
assertEquals("hello", message.getPayload());
|
||||
}
|
||||
|
||||
public static class SampleService{
|
||||
public String processMessage(String message) throws Exception {
|
||||
Thread.sleep(500);
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2010 the original author or authors.
|
||||
* Copyright 2002-2012 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.
|
||||
@@ -19,221 +19,204 @@ package org.springframework.integration.handler;
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotSame;
|
||||
import static org.junit.Assert.assertSame;
|
||||
import static org.junit.Assert.assertThat;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
import java.util.Date;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
import org.hamcrest.Matchers;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
||||
import org.springframework.beans.DirectFieldAccessor;
|
||||
import org.mockito.Mockito;
|
||||
import org.mockito.invocation.InvocationOnMock;
|
||||
import org.mockito.stubbing.Answer;
|
||||
import org.springframework.context.event.ContextRefreshedEvent;
|
||||
import org.springframework.context.support.StaticApplicationContext;
|
||||
import org.springframework.integration.Message;
|
||||
import org.springframework.integration.MessageDeliveryException;
|
||||
import org.springframework.integration.MessageHandlingException;
|
||||
import org.springframework.integration.channel.DirectChannel;
|
||||
import org.springframework.integration.channel.MessagePublishingErrorHandler;
|
||||
import org.springframework.integration.context.IntegrationContextUtils;
|
||||
import org.springframework.integration.core.MessageHandler;
|
||||
import org.springframework.integration.message.GenericMessage;
|
||||
import org.springframework.integration.store.MessageGroup;
|
||||
import org.springframework.integration.store.MessageGroupStore;
|
||||
import org.springframework.integration.store.SimpleMessageStore;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
import org.springframework.integration.test.util.TestUtils;
|
||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
|
||||
|
||||
/**
|
||||
* @author Mark Fisher
|
||||
* @author Artem Bilan
|
||||
* @since 1.0.3
|
||||
*/
|
||||
public class DelayHandlerTests {
|
||||
|
||||
private static final String DELAYER_MESSAGE_GROUP_ID = "testDelayer.messageGroupId";
|
||||
|
||||
private final DirectChannel input = new DirectChannel();
|
||||
|
||||
private final DirectChannel output = new DirectChannel();
|
||||
|
||||
private final CountDownLatch latch = new CountDownLatch(1);
|
||||
|
||||
private ThreadPoolTaskScheduler taskScheduler;
|
||||
|
||||
private DelayHandler delayHandler;
|
||||
|
||||
private ResultHandler resultHandler = new ResultHandler();
|
||||
|
||||
@Before
|
||||
public void setChannelNames() {
|
||||
public void setup() {
|
||||
input.setBeanName("input");
|
||||
output.setBeanName("output");
|
||||
taskScheduler = new ThreadPoolTaskScheduler();
|
||||
taskScheduler.afterPropertiesSet();
|
||||
delayHandler = new DelayHandler(DELAYER_MESSAGE_GROUP_ID, taskScheduler);
|
||||
delayHandler.setOutputChannel(output);
|
||||
input.subscribe(delayHandler);
|
||||
output.subscribe(resultHandler);
|
||||
}
|
||||
|
||||
private void startDelayerHandler() {
|
||||
delayHandler.afterPropertiesSet();
|
||||
delayHandler.onApplicationEvent(new ContextRefreshedEvent(TestUtils.createTestApplicationContext()));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void noDelayHeaderAndDefaultDelayIsZero() throws Exception {
|
||||
DelayHandler delayHandler = new DelayHandler(0);
|
||||
ResultHandler resultHandler = new ResultHandler();
|
||||
delayHandler.setOutputChannel(output);
|
||||
delayHandler.afterPropertiesSet();
|
||||
input.subscribe(delayHandler);
|
||||
output.subscribe(resultHandler);
|
||||
this.startDelayerHandler();
|
||||
Message<?> message = MessageBuilder.withPayload("test").build();
|
||||
input.send(message);
|
||||
assertSame(message, resultHandler.lastMessage);
|
||||
assertSame(message.getPayload(), resultHandler.lastMessage.getPayload());
|
||||
assertSame(Thread.currentThread(), resultHandler.lastThread);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void noDelayHeaderAndDefaultDelayIsPositive() throws Exception {
|
||||
DelayHandler delayHandler = new DelayHandler(10);
|
||||
ResultHandler resultHandler = new ResultHandler();
|
||||
delayHandler.setOutputChannel(output);
|
||||
delayHandler.afterPropertiesSet();
|
||||
input.subscribe(delayHandler);
|
||||
output.subscribe(resultHandler);
|
||||
delayHandler.setDefaultDelay(10);
|
||||
this.startDelayerHandler();
|
||||
Message<?> message = MessageBuilder.withPayload("test").build();
|
||||
input.send(message);
|
||||
this.waitForLatch(1000);
|
||||
assertSame(message, resultHandler.lastMessage);
|
||||
assertSame(message.getPayload(), resultHandler.lastMessage.getPayload());
|
||||
assertNotSame(Thread.currentThread(), resultHandler.lastThread);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void delayHeaderAndDefaultDelayWouldTimeout() throws Exception {
|
||||
DelayHandler delayHandler = new DelayHandler(5000);
|
||||
delayHandler.setDefaultDelay(5000);
|
||||
delayHandler.setDelayHeaderName("delay");
|
||||
ResultHandler resultHandler = new ResultHandler();
|
||||
delayHandler.setOutputChannel(output);
|
||||
delayHandler.afterPropertiesSet();
|
||||
input.subscribe(delayHandler);
|
||||
output.subscribe(resultHandler);
|
||||
this.startDelayerHandler();
|
||||
Message<?> message = MessageBuilder.withPayload("test")
|
||||
.setHeader("delay", 100).build();
|
||||
input.send(message);
|
||||
this.waitForLatch(1000);
|
||||
assertSame(message, resultHandler.lastMessage);
|
||||
assertSame(message.getPayload(), resultHandler.lastMessage.getPayload());
|
||||
assertNotSame(Thread.currentThread(), resultHandler.lastThread);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void delayHeaderIsNegativeAndDefaultDelayWouldTimeout() throws Exception {
|
||||
DelayHandler delayHandler = new DelayHandler(5000);
|
||||
delayHandler.setDefaultDelay(5000);
|
||||
delayHandler.setDelayHeaderName("delay");
|
||||
ResultHandler resultHandler = new ResultHandler();
|
||||
delayHandler.setOutputChannel(output);
|
||||
delayHandler.afterPropertiesSet();
|
||||
input.subscribe(delayHandler);
|
||||
output.subscribe(resultHandler);
|
||||
this.startDelayerHandler();
|
||||
Message<?> message = MessageBuilder.withPayload("test")
|
||||
.setHeader("delay", -7000).build();
|
||||
input.send(message);
|
||||
this.waitForLatch(1000);
|
||||
assertSame(message, resultHandler.lastMessage);
|
||||
assertSame(message.getPayload(), resultHandler.lastMessage.getPayload());
|
||||
assertSame(Thread.currentThread(), resultHandler.lastThread);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void delayHeaderIsInvalidFallsBackToDefaultDelay() throws Exception {
|
||||
DelayHandler delayHandler = new DelayHandler(5);
|
||||
delayHandler.setDefaultDelay(5);
|
||||
delayHandler.setDelayHeaderName("delay");
|
||||
ResultHandler resultHandler = new ResultHandler();
|
||||
delayHandler.setOutputChannel(output);
|
||||
delayHandler.afterPropertiesSet();
|
||||
input.subscribe(delayHandler);
|
||||
output.subscribe(resultHandler);
|
||||
this.startDelayerHandler();
|
||||
Message<?> message = MessageBuilder.withPayload("test")
|
||||
.setHeader("delay", "not a number").build();
|
||||
input.send(message);
|
||||
this.waitForLatch(1000);
|
||||
assertSame(message, resultHandler.lastMessage);
|
||||
assertSame(message.getPayload(), resultHandler.lastMessage.getPayload());
|
||||
assertNotSame(Thread.currentThread(), resultHandler.lastThread);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void delayHeaderIsDateInTheFutureAndDefaultDelayWouldTimeout() throws Exception {
|
||||
DelayHandler delayHandler = new DelayHandler(5000);
|
||||
delayHandler.setDefaultDelay(5000);
|
||||
delayHandler.setDelayHeaderName("delay");
|
||||
ResultHandler resultHandler = new ResultHandler();
|
||||
delayHandler.setOutputChannel(output);
|
||||
delayHandler.afterPropertiesSet();
|
||||
input.subscribe(delayHandler);
|
||||
output.subscribe(resultHandler);
|
||||
this.startDelayerHandler();
|
||||
Message<?> message = MessageBuilder.withPayload("test")
|
||||
.setHeader("delay", new Date(new Date().getTime() + 150)).build();
|
||||
input.send(message);
|
||||
this.waitForLatch(3000);
|
||||
assertSame(message, resultHandler.lastMessage);
|
||||
assertSame(message.getPayload(), resultHandler.lastMessage.getPayload());
|
||||
assertNotSame(Thread.currentThread(), resultHandler.lastThread);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void delayHeaderIsDateInThePastAndDefaultDelayWouldTimeout() throws Exception {
|
||||
DelayHandler delayHandler = new DelayHandler(5000);
|
||||
delayHandler.setDefaultDelay(5000);
|
||||
delayHandler.setDelayHeaderName("delay");
|
||||
ResultHandler resultHandler = new ResultHandler();
|
||||
delayHandler.setOutputChannel(output);
|
||||
delayHandler.afterPropertiesSet();
|
||||
input.subscribe(delayHandler);
|
||||
output.subscribe(resultHandler);
|
||||
this.startDelayerHandler();
|
||||
Message<?> message = MessageBuilder.withPayload("test")
|
||||
.setHeader("delay", new Date(new Date().getTime() - 60 * 1000)).build();
|
||||
input.send(message);
|
||||
this.waitForLatch(3000);
|
||||
assertSame(message, resultHandler.lastMessage);
|
||||
assertSame(message.getPayload(), resultHandler.lastMessage.getPayload());
|
||||
assertSame(Thread.currentThread(), resultHandler.lastThread);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void delayHeaderIsNullDateAndDefaultDelayIsZero() throws Exception {
|
||||
DelayHandler delayHandler = new DelayHandler(0);
|
||||
delayHandler.setDelayHeaderName("delay");
|
||||
ResultHandler resultHandler = new ResultHandler();
|
||||
delayHandler.setOutputChannel(output);
|
||||
delayHandler.afterPropertiesSet();
|
||||
input.subscribe(delayHandler);
|
||||
output.subscribe(resultHandler);
|
||||
this.startDelayerHandler();
|
||||
Date nullDate = null;
|
||||
Message<?> message = MessageBuilder.withPayload("test")
|
||||
.setHeader("delay", nullDate).build();
|
||||
input.send(message);
|
||||
this.waitForLatch(3000);
|
||||
assertSame(message, resultHandler.lastMessage);
|
||||
assertSame(message.getPayload(), resultHandler.lastMessage.getPayload());
|
||||
assertSame(Thread.currentThread(), resultHandler.lastThread);
|
||||
}
|
||||
|
||||
@Test(expected = TestTimedOutException.class)
|
||||
public void delayHeaderIsFutureDateAndTimesOut() throws Exception {
|
||||
DelayHandler delayHandler = new DelayHandler(0);
|
||||
delayHandler.setDelayHeaderName("delay");
|
||||
ResultHandler resultHandler = new ResultHandler();
|
||||
delayHandler.setOutputChannel(output);
|
||||
delayHandler.afterPropertiesSet();
|
||||
input.subscribe(delayHandler);
|
||||
output.subscribe(resultHandler);
|
||||
this.startDelayerHandler();
|
||||
Date future = new Date(new Date().getTime() + 60 * 1000);
|
||||
Message<?> message = MessageBuilder.withPayload("test")
|
||||
.setHeader("delay", future).build();
|
||||
input.send(message);
|
||||
this.waitForLatch(50);
|
||||
assertSame(message, resultHandler.lastMessage);
|
||||
assertSame(message.getPayload(), resultHandler.lastMessage.getPayload());
|
||||
assertSame(Thread.currentThread(), resultHandler.lastThread);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void delayHeaderIsValidStringAndDefaultDelayWouldTimeout() throws Exception {
|
||||
DelayHandler delayHandler = new DelayHandler(5000);
|
||||
delayHandler.setDefaultDelay(5000);
|
||||
delayHandler.setDelayHeaderName("delay");
|
||||
ResultHandler resultHandler = new ResultHandler();
|
||||
delayHandler.setOutputChannel(output);
|
||||
delayHandler.afterPropertiesSet();
|
||||
input.subscribe(delayHandler);
|
||||
output.subscribe(resultHandler);
|
||||
this.startDelayerHandler();
|
||||
Message<?> message = MessageBuilder.withPayload("test")
|
||||
.setHeader("delay", "20").build();
|
||||
input.send(message);
|
||||
this.waitForLatch(1000);
|
||||
assertSame(message, resultHandler.lastMessage);
|
||||
assertSame(message.getPayload(), resultHandler.lastMessage.getPayload());
|
||||
assertNotSame(Thread.currentThread(), resultHandler.lastThread);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void verifyShutdownWithoutWaitingByDefault() throws Exception {
|
||||
DelayHandler delayHandler = new DelayHandler(5000);
|
||||
delayHandler.afterPropertiesSet();
|
||||
delayHandler.setDefaultDelay(5000);
|
||||
this.startDelayerHandler();
|
||||
delayHandler.handleMessage(new GenericMessage<String>("foo"));
|
||||
delayHandler.destroy();
|
||||
final ThreadPoolTaskScheduler taskScheduler = (ThreadPoolTaskScheduler)
|
||||
new DirectFieldAccessor(delayHandler).getPropertyValue("taskScheduler");
|
||||
taskScheduler.destroy();
|
||||
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
new Thread(new Runnable() {
|
||||
public void run() {
|
||||
@@ -252,13 +235,12 @@ public class DelayHandlerTests {
|
||||
|
||||
@Test
|
||||
public void verifyShutdownWithWait() throws Exception {
|
||||
DelayHandler delayHandler = new DelayHandler(5000);
|
||||
delayHandler.setWaitForTasksToCompleteOnShutdown(true);
|
||||
delayHandler.afterPropertiesSet();
|
||||
delayHandler.setDefaultDelay(5000);
|
||||
taskScheduler.setWaitForTasksToCompleteOnShutdown(true);
|
||||
this.startDelayerHandler();
|
||||
delayHandler.handleMessage(new GenericMessage<String>("foo"));
|
||||
delayHandler.destroy();
|
||||
final ThreadPoolTaskScheduler taskScheduler = (ThreadPoolTaskScheduler)
|
||||
new DirectFieldAccessor(delayHandler).getPropertyValue("taskScheduler");
|
||||
taskScheduler.destroy();
|
||||
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
new Thread(new Runnable() {
|
||||
public void run() {
|
||||
@@ -277,10 +259,8 @@ public class DelayHandlerTests {
|
||||
|
||||
@Test(expected = MessageDeliveryException.class)
|
||||
public void handlerThrowsExceptionWithNoDelay() throws Exception {
|
||||
DelayHandler delayHandler = new DelayHandler(0);
|
||||
delayHandler.setOutputChannel(output);
|
||||
delayHandler.afterPropertiesSet();
|
||||
input.subscribe(delayHandler);
|
||||
this.startDelayerHandler();
|
||||
output.unsubscribe(resultHandler);
|
||||
output.subscribe(new MessageHandler() {
|
||||
public void handleMessage(Message<?> message) {
|
||||
throw new UnsupportedOperationException("intentional test failure");
|
||||
@@ -292,14 +272,14 @@ public class DelayHandlerTests {
|
||||
|
||||
@Test
|
||||
public void errorChannelHeaderAndHandlerThrowsExceptionWithDelay() throws Exception {
|
||||
DelayHandler delayHandler = new DelayHandler(0);
|
||||
delayHandler.setDelayHeaderName("delay");
|
||||
delayHandler.setOutputChannel(output);
|
||||
delayHandler.afterPropertiesSet();
|
||||
DirectChannel errorChannel = new DirectChannel();
|
||||
ResultHandler resultHandler = new ResultHandler();
|
||||
MessagePublishingErrorHandler errorHandler = new MessagePublishingErrorHandler();
|
||||
errorHandler.setDefaultErrorChannel(errorChannel);
|
||||
taskScheduler.setErrorHandler(errorHandler);
|
||||
delayHandler.setDelayHeaderName("delay");
|
||||
this.startDelayerHandler();
|
||||
output.unsubscribe(resultHandler);
|
||||
errorChannel.subscribe(resultHandler);
|
||||
input.subscribe(delayHandler);
|
||||
output.subscribe(new MessageHandler() {
|
||||
public void handleMessage(Message<?> message) {
|
||||
throw new UnsupportedOperationException("intentional test failure");
|
||||
@@ -311,13 +291,10 @@ public class DelayHandlerTests {
|
||||
input.send(message);
|
||||
this.waitForLatch(1000);
|
||||
Message<?> errorMessage = resultHandler.lastMessage;
|
||||
assertEquals(MessageHandlingException.class, errorMessage.getPayload().getClass());
|
||||
MessageHandlingException exceptionPayload = (MessageHandlingException) errorMessage.getPayload();
|
||||
assertSame(message, exceptionPayload.getFailedMessage());
|
||||
assertEquals(MessageDeliveryException.class, exceptionPayload.getCause().getClass());
|
||||
MessageDeliveryException nestedException = (MessageDeliveryException) exceptionPayload.getCause();
|
||||
assertEquals(UnsupportedOperationException.class, nestedException.getCause().getClass());
|
||||
assertSame(message, nestedException.getFailedMessage());
|
||||
assertEquals(MessageDeliveryException.class, errorMessage.getPayload().getClass());
|
||||
MessageDeliveryException exceptionPayload = (MessageDeliveryException) errorMessage.getPayload();
|
||||
assertSame(message.getPayload(), exceptionPayload.getFailedMessage().getPayload());
|
||||
assertEquals(UnsupportedOperationException.class, exceptionPayload.getCause().getClass());
|
||||
assertNotSame(Thread.currentThread(), resultHandler.lastThread);
|
||||
}
|
||||
|
||||
@@ -326,16 +303,16 @@ public class DelayHandlerTests {
|
||||
String errorChannelName = "customErrorChannel";
|
||||
StaticApplicationContext context = new StaticApplicationContext();
|
||||
context.registerSingleton(errorChannelName, DirectChannel.class);
|
||||
context.registerSingleton(IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME, DirectChannel.class);
|
||||
context.refresh();
|
||||
DirectChannel customErrorChannel = (DirectChannel) context.getBean(errorChannelName);
|
||||
DelayHandler delayHandler = new DelayHandler(0);
|
||||
delayHandler.setBeanFactory(context);
|
||||
MessagePublishingErrorHandler errorHandler = new MessagePublishingErrorHandler();
|
||||
errorHandler.setBeanFactory(context);
|
||||
taskScheduler.setErrorHandler(errorHandler);
|
||||
delayHandler.setDelayHeaderName("delay");
|
||||
delayHandler.setOutputChannel(output);
|
||||
delayHandler.afterPropertiesSet();
|
||||
ResultHandler resultHandler = new ResultHandler();
|
||||
this.startDelayerHandler();
|
||||
output.unsubscribe(resultHandler);
|
||||
customErrorChannel.subscribe(resultHandler);
|
||||
input.subscribe(delayHandler);
|
||||
output.subscribe(new MessageHandler() {
|
||||
public void handleMessage(Message<?> message) {
|
||||
throw new UnsupportedOperationException("intentional test failure");
|
||||
@@ -347,32 +324,26 @@ public class DelayHandlerTests {
|
||||
input.send(message);
|
||||
this.waitForLatch(1000);
|
||||
Message<?> errorMessage = resultHandler.lastMessage;
|
||||
assertEquals(MessageHandlingException.class, errorMessage.getPayload().getClass());
|
||||
MessageHandlingException exceptionPayload = (MessageHandlingException) errorMessage.getPayload();
|
||||
assertSame(message, exceptionPayload.getFailedMessage());
|
||||
assertEquals(MessageDeliveryException.class, exceptionPayload.getCause().getClass());
|
||||
MessageDeliveryException nestedException = (MessageDeliveryException) exceptionPayload.getCause();
|
||||
assertEquals(UnsupportedOperationException.class, nestedException.getCause().getClass());
|
||||
assertSame(message, nestedException.getFailedMessage());
|
||||
assertEquals(MessageDeliveryException.class, errorMessage.getPayload().getClass());
|
||||
MessageDeliveryException exceptionPayload = (MessageDeliveryException) errorMessage.getPayload();
|
||||
assertSame(message.getPayload(), exceptionPayload.getFailedMessage().getPayload());
|
||||
assertEquals(UnsupportedOperationException.class, exceptionPayload.getCause().getClass());
|
||||
assertNotSame(Thread.currentThread(), resultHandler.lastThread);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void defaultErrorChannelAndHandlerThrowsExceptionWithDelay() throws Exception {
|
||||
StaticApplicationContext context = new StaticApplicationContext();
|
||||
context.registerSingleton(
|
||||
IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME, DirectChannel.class);
|
||||
context.registerSingleton(IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME, DirectChannel.class);
|
||||
context.refresh();
|
||||
DirectChannel defaultErrorChannel = (DirectChannel) context.getBean(
|
||||
IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME);
|
||||
DelayHandler delayHandler = new DelayHandler(0);
|
||||
delayHandler.setBeanFactory(context);
|
||||
DirectChannel defaultErrorChannel = (DirectChannel) context.getBean(IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME);
|
||||
MessagePublishingErrorHandler errorHandler = new MessagePublishingErrorHandler();
|
||||
errorHandler.setBeanFactory(context);
|
||||
taskScheduler.setErrorHandler(errorHandler);
|
||||
delayHandler.setDelayHeaderName("delay");
|
||||
delayHandler.setOutputChannel(output);
|
||||
delayHandler.afterPropertiesSet();
|
||||
ResultHandler resultHandler = new ResultHandler();
|
||||
this.startDelayerHandler();
|
||||
output.unsubscribe(resultHandler);
|
||||
defaultErrorChannel.subscribe(resultHandler);
|
||||
input.subscribe(delayHandler);
|
||||
output.subscribe(new MessageHandler() {
|
||||
public void handleMessage(Message<?> message) {
|
||||
throw new UnsupportedOperationException("intentional test failure");
|
||||
@@ -383,17 +354,72 @@ public class DelayHandlerTests {
|
||||
input.send(message);
|
||||
this.waitForLatch(1000);
|
||||
Message<?> errorMessage = resultHandler.lastMessage;
|
||||
assertEquals(MessageHandlingException.class, errorMessage.getPayload().getClass());
|
||||
MessageHandlingException exceptionPayload = (MessageHandlingException) errorMessage.getPayload();
|
||||
assertSame(message, exceptionPayload.getFailedMessage());
|
||||
assertEquals(MessageDeliveryException.class, exceptionPayload.getCause().getClass());
|
||||
MessageDeliveryException nestedException = (MessageDeliveryException) exceptionPayload.getCause();
|
||||
assertEquals(UnsupportedOperationException.class, nestedException.getCause().getClass());
|
||||
assertSame(message, nestedException.getFailedMessage());
|
||||
assertEquals(MessageDeliveryException.class, errorMessage.getPayload().getClass());
|
||||
MessageDeliveryException exceptionPayload = (MessageDeliveryException) errorMessage.getPayload();
|
||||
assertSame(message.getPayload(), exceptionPayload.getFailedMessage().getPayload());
|
||||
assertEquals(UnsupportedOperationException.class, exceptionPayload.getCause().getClass());
|
||||
assertNotSame(Thread.currentThread(), resultHandler.lastThread);
|
||||
}
|
||||
|
||||
|
||||
@Test //INT-1132
|
||||
public void testReschedulePersistedMessagesOnStartup() throws Exception {
|
||||
MessageGroupStore messageGroupStore = new SimpleMessageStore();
|
||||
this.delayHandler.setDefaultDelay(200);
|
||||
this.delayHandler.setMessageStore(messageGroupStore);
|
||||
this.startDelayerHandler();
|
||||
Message<?> message = MessageBuilder.withPayload("test").build();
|
||||
this.input.send(message);
|
||||
|
||||
Thread.sleep(100);
|
||||
|
||||
// emulate restart
|
||||
this.taskScheduler.destroy();
|
||||
|
||||
assertEquals(1, messageGroupStore.getMessageGroupCount());
|
||||
assertEquals(DELAYER_MESSAGE_GROUP_ID, messageGroupStore.iterator().next().getGroupId());
|
||||
assertEquals(1, messageGroupStore.messageGroupSize(DELAYER_MESSAGE_GROUP_ID));
|
||||
assertEquals(1, messageGroupStore.getMessageCountForAllMessageGroups());
|
||||
MessageGroup messageGroup = messageGroupStore.getMessageGroup(DELAYER_MESSAGE_GROUP_ID);
|
||||
Message<?> messageInStore = messageGroup.getMessages().iterator().next();
|
||||
Object payload = messageInStore.getPayload();
|
||||
assertEquals("DelayedMessageWrapper", payload.getClass().getSimpleName());
|
||||
assertEquals(message.getPayload(), TestUtils.getPropertyValue(payload, "original.payload"));
|
||||
|
||||
this.taskScheduler.afterPropertiesSet();
|
||||
this.delayHandler = new DelayHandler(DELAYER_MESSAGE_GROUP_ID, this.taskScheduler);
|
||||
this.delayHandler.setOutputChannel(output);
|
||||
this.delayHandler.setDefaultDelay(200);
|
||||
this.delayHandler.setMessageStore(messageGroupStore);
|
||||
this.startDelayerHandler();
|
||||
|
||||
long timeBeforeReceive = System.currentTimeMillis();
|
||||
assertTrue(this.latch.await(10, TimeUnit.SECONDS));
|
||||
long timeAfterReceive = System.currentTimeMillis();
|
||||
assertThat(timeAfterReceive - timeBeforeReceive, Matchers.lessThanOrEqualTo(100L));
|
||||
|
||||
assertSame(message.getPayload(), this.resultHandler.lastMessage.getPayload());
|
||||
assertNotSame(Thread.currentThread(), this.resultHandler.lastThread);
|
||||
assertEquals(1, messageGroupStore.getMessageGroupCount());
|
||||
assertEquals(0, messageGroupStore.messageGroupSize(DELAYER_MESSAGE_GROUP_ID));
|
||||
}
|
||||
|
||||
@Test //INT-1132
|
||||
// Can happen in the parent-child context e.g. Spring-MVC applications
|
||||
public void testDoubleOnApplicationEvent() throws Exception {
|
||||
this.delayHandler = Mockito.spy(this.delayHandler);
|
||||
Mockito.doAnswer(new Answer() {
|
||||
public Object answer(InvocationOnMock invocation) throws Throwable {
|
||||
return null;
|
||||
}
|
||||
}).when(this.delayHandler).reschedulePersistedMessages();
|
||||
|
||||
ContextRefreshedEvent contextRefreshedEvent = new ContextRefreshedEvent(TestUtils.createTestApplicationContext());
|
||||
this.delayHandler.onApplicationEvent(contextRefreshedEvent);
|
||||
this.delayHandler.onApplicationEvent(contextRefreshedEvent);
|
||||
Mockito.verify(this.delayHandler, Mockito.times(1)).reschedulePersistedMessages();
|
||||
}
|
||||
|
||||
private void waitForLatch(long timeout) {
|
||||
try {
|
||||
this.latch.await(timeout, TimeUnit.MILLISECONDS);
|
||||
|
||||
Reference in New Issue
Block a user