From cb85fe0f82abb5a832e34aadcfedc31ea3277405 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Thu, 13 Feb 2025 12:13:52 -0500 Subject: [PATCH] GH-9825: DelayerEndpointSpec: Set TaskScheduler to the handler as well Fixes: #9825 Issue link: https://github.com/spring-projects/spring-integration/issues/9825 The `DelayerEndpointSpec` extends `ConsumerEndpointSpec` which has a `taskScheduler()` option. However this is set only to the endpoint for this `MessageHandler`. * Override `taskScheduler()` method on the `DelayerEndpointSpec` to set the provided `TaskScheduler` to the `DelayHandler` as well (cherry picked from commit 12fee0a9fb617dc747018afec4d6d7f9919f0d5f) --- .../integration/dsl/DelayerEndpointSpec.java | 13 +++++++++++++ .../dsl/flows/IntegrationFlowTests.java | 17 +++++++++++++---- 2 files changed, 26 insertions(+), 4 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/DelayerEndpointSpec.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/DelayerEndpointSpec.java index 8f89a508c3..f3c3dfa614 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/DelayerEndpointSpec.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/DelayerEndpointSpec.java @@ -30,6 +30,7 @@ import org.springframework.integration.store.MessageGroupStore; import org.springframework.integration.transaction.TransactionInterceptorBuilder; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; +import org.springframework.scheduling.TaskScheduler; import org.springframework.transaction.TransactionManager; import org.springframework.transaction.interceptor.TransactionInterceptor; import org.springframework.util.Assert; @@ -243,4 +244,16 @@ public class DelayerEndpointSpec extends ConsumerEndpointSpec c.autoStartup(false).id("bridge")) .fixedSubscriberChannel() @@ -820,7 +827,9 @@ public class IntegrationFlowTests { .messageGroupId("delayer") .delayExpression("200") .advice(this.delayedAdvice) - .messageStore(this.messageStore())) + .messageStore(messageStore()) + .taskScheduler(customScheduler) + .id("delayer")) .channel(MessageChannels.queue("bridgeFlow2Output")) .get(); } @@ -833,8 +842,8 @@ public class IntegrationFlowTests { @Bean public IntegrationFlow claimCheckFlow() { return IntegrationFlow.from("claimCheckInput") - .claimCheckIn(this.messageStore()) - .claimCheckOut(this.messageStore()) + .claimCheckIn(messageStore()) + .claimCheckOut(messageStore()) .get(); }