diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java index fef670081d..9fa2be544c 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java @@ -273,18 +273,27 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement } private long determineDelayForMessage(Message message) { + DelayedMessageWrapper delayedMessageWrapper = null; + if (message.getPayload() instanceof DelayedMessageWrapper) { + delayedMessageWrapper = (DelayedMessageWrapper) message.getPayload(); + } + long delay = this.defaultDelay; if (this.delayExpression != null) { Exception delayValueException = null; Object delayValue = null; try { - delayValue = this.delayExpression.getValue(this.evaluationContext, message); + delayValue = this.delayExpression.getValue(this.evaluationContext, + delayedMessageWrapper != null ? delayedMessageWrapper.getOriginal() : message); } catch (EvaluationException e) { delayValueException = e; } if (delayValue instanceof Date) { - delay = ((Date) delayValue).getTime() - new Date().getTime(); + long current = delayedMessageWrapper != null + ? delayedMessageWrapper.getRequestDate() + : System.currentTimeMillis(); + delay = ((Date) delayValue).getTime() - current; } else if (delayValue != null) { try { diff --git a/spring-integration-core/src/test/java/org/springframework/integration/handler/DelayHandlerTests.java b/spring-integration-core/src/test/java/org/springframework/integration/handler/DelayHandlerTests.java index d177fe459d..9915823057 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/handler/DelayHandlerTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/handler/DelayHandlerTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2013 the original author or authors. + * Copyright 2002-2015 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. @@ -24,16 +24,20 @@ import static org.junit.Assert.assertSame; import static org.junit.Assert.assertTrue; import static org.mockito.Mockito.mock; +import java.util.Calendar; import java.util.Date; +import java.util.Queue; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; +import org.junit.After; import org.junit.Before; import org.junit.Test; import org.mockito.Mockito; import org.mockito.invocation.InvocationOnMock; import org.mockito.stubbing.Answer; +import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.BeanFactory; import org.springframework.context.event.ContextRefreshedEvent; import org.springframework.context.support.StaticApplicationContext; @@ -91,6 +95,11 @@ public class DelayHandlerTests { output.subscribe(resultHandler); } + @After + public void tearDown() { + taskScheduler.destroy(); + } + private void setDelayExpression() { Expression expression = new SpelExpressionParser().parseExpression("headers.delay"); this.delayHandler.setDelayExpression(expression); @@ -466,6 +475,37 @@ public class DelayHandlerTests { assertNull(message); } + @Test + public void testRescheduleForTheDateDelay() throws Exception { + this.delayHandler.setDelayExpression(new SpelExpressionParser().parseExpression("payload")); + this.delayHandler.setOutputChannel(new DirectChannel()); + this.delayHandler.setIgnoreExpressionFailures(false); + startDelayerHandler(); + Calendar releaseDate = Calendar.getInstance(); + releaseDate.add(Calendar.HOUR, 1); + this.delayHandler.handleMessage(new GenericMessage(releaseDate.getTime())); + + // emulate restart + this.taskScheduler.destroy(); + MessageGroupStore messageStore = TestUtils.getPropertyValue(this.delayHandler, "messageStore", + MessageGroupStore.class); + MessageGroup messageGroup = messageStore.getMessageGroup(DELAYER_MESSAGE_GROUP_ID); + Message messageInStore = messageGroup.getMessages().iterator().next(); + Object payload = messageInStore.getPayload(); + DirectFieldAccessor dfa = new DirectFieldAccessor(payload); + long requestTime = (Long) dfa.getPropertyValue("requestDate"); + Calendar requestDate = Calendar.getInstance(); + requestDate.setTimeInMillis(requestTime); + requestDate.add(Calendar.HOUR, -2); + dfa.setPropertyValue("requestDate", requestDate.getTimeInMillis()); + this.taskScheduler.afterPropertiesSet(); + this.delayHandler.reschedulePersistedMessages(); + Thread.sleep(10); + Queue works = TestUtils.getPropertyValue(this.taskScheduler, "scheduledExecutor.workQueue", Queue.class); + assertEquals(1, works.size()); + } + + private void waitForLatch(long timeout) { try { this.latch.await(timeout, TimeUnit.MILLISECONDS);