diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/PublishSubscribeSpec.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/PublishSubscribeSpec.java index f5ee5db071..9a654f26cd 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/PublishSubscribeSpec.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/PublishSubscribeSpec.java @@ -25,6 +25,7 @@ import org.springframework.util.Assert; /** * @author Artem Bilan + * @author Gary Russell * * @since 5.0 */ @@ -32,6 +33,8 @@ public class PublishSubscribeSpec extends PublishSubscribeChannelSpec subscriberFlows = new LinkedHashMap<>(); + private int order; + PublishSubscribeSpec() { super(); } @@ -50,7 +53,7 @@ public class PublishSubscribeSpec extends PublishSubscribeChannelSpec consumer.order(this.order++)); MessageChannel subFlowInput = subFlow.getInputChannel(); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dsl/publishsubscribe/PublishSubscribeTests.java b/spring-integration-core/src/test/java/org/springframework/integration/dsl/publishsubscribe/PublishSubscribeTests.java index a78458b862..1224912fc5 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/dsl/publishsubscribe/PublishSubscribeTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/dsl/publishsubscribe/PublishSubscribeTests.java @@ -31,6 +31,8 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.integration.config.EnableIntegration; import org.springframework.integration.dsl.IntegrationFlow; +import org.springframework.integration.dsl.context.IntegrationFlowContext; +import org.springframework.integration.dsl.context.IntegrationFlowContext.IntegrationFlowRegistration; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandler; import org.springframework.messaging.support.GenericMessage; @@ -52,12 +54,28 @@ public class PublishSubscribeTests { @Autowired private List subscribersOrderedCall; + @Autowired + private PubSubBugTestContext config; + + @Autowired + private IntegrationFlowContext context; + @Test public void executeFirstFlow() { + this.subscribersOrderedCall.clear(); this.inputChannel.send(new GenericMessage<>("Test")); assertThat(this.subscribersOrderedCall).containsExactly(0, 1, 2, 3, 4, 5); } + @Test + public void dynamicFlow() { + this.subscribersOrderedCall.clear(); + IntegrationFlowRegistration reg = this.context.registration(this.config.flow()).register(); + reg.getInputChannel().send(new GenericMessage<>("Test")); + assertThat(this.subscribersOrderedCall).containsExactly(0, 1, 2, 3, 4, 5); + this.context.remove(reg.getId()); + } + @Configuration @EnableIntegration static class PubSubBugTestContext { @@ -79,6 +97,10 @@ public class PublishSubscribeTests { @Bean public IntegrationFlow pubSubFlow() { + return flow(); + } + + IntegrationFlow flow() { return f -> f .publishSubscribeChannel(c -> c .subscribe(sf -> sf