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 0ee44d46bd..56bc88cfae 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 @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * Copyright 2016-2017 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. @@ -46,7 +46,9 @@ public class PublishSubscribeSpec extends PublishSubscribeChannelSpec subscribersOrderedCall; + + @Test + public void executeFirstFlow() { + this.inputChannel.send(new GenericMessage<>("Test")); + assertThat(this.subscribersOrderedCall, contains(0, 1, 2, 3, 4, 5)); + } + + @Configuration + @EnableIntegration + static class PubSubBugTestContext { + + @Bean + public List subscribersOrderedCall() { + return new LinkedList<>(); + } + + @Bean + public Consumer subscriberConsumerBean() { + return s -> subscribersOrderedCall().add(2); + } + + @Bean + public MessageHandler subscriberMessageHandlerBean() { + return s -> subscribersOrderedCall().add(3); + } + + @Bean + public IntegrationFlow pubSubFlow() { + return f -> f + .publishSubscribeChannel(c -> c + .subscribe(sf -> sf + .handle(m -> subscribersOrderedCall().add(0))) + .subscribe(sf -> sf + .handle((p, h) -> { + subscribersOrderedCall().add(1); + return null; + })) + .subscribe(sf -> sf + .handle(subscriberConsumerBean())) + .subscribe(sf -> sf + .handle(subscriberMessageHandlerBean())) + .subscribe(sf -> sf + .channel("secondInlineSubscriberChannel") + .handle(m -> subscribersOrderedCall().add(4))) + ) + .handle((p, h) -> { + subscribersOrderedCall().add(5); + return null; + }); + } + + } + +}