diff --git a/supplier/file-supplier/src/main/java/org/springframework/cloud/fn/supplier/file/FileSupplierConfiguration.java b/supplier/file-supplier/src/main/java/org/springframework/cloud/fn/supplier/file/FileSupplierConfiguration.java index 5761713f..50b275e5 100644 --- a/supplier/file-supplier/src/main/java/org/springframework/cloud/fn/supplier/file/FileSupplierConfiguration.java +++ b/supplier/file-supplier/src/main/java/org/springframework/cloud/fn/supplier/file/FileSupplierConfiguration.java @@ -85,7 +85,8 @@ public class FileSupplierConfiguration { monoSink.success(this.fileMessageSource.receive()))) .subscribeOn(Schedulers.boundedElastic()) .repeatWhenEmpty(it -> it.delayElements(this.fileSupplierProperties.getDelayWhenEmpty())) - .repeat(); + .repeat() + .doOnSubscribe(s -> this.fileMessageSource.start()); } @Bean diff --git a/supplier/file-supplier/src/test/java/org/springframework/cloud/fn/supplier/file/DefaultFileSupplierTests.java b/supplier/file-supplier/src/test/java/org/springframework/cloud/fn/supplier/file/DefaultFileSupplierTests.java index d908a211..234e5fd8 100644 --- a/supplier/file-supplier/src/test/java/org/springframework/cloud/fn/supplier/file/DefaultFileSupplierTests.java +++ b/supplier/file-supplier/src/test/java/org/springframework/cloud/fn/supplier/file/DefaultFileSupplierTests.java @@ -24,7 +24,10 @@ import org.junit.jupiter.api.Test; import reactor.core.publisher.Flux; import reactor.test.StepVerifier; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.integration.file.FileHeaders; +import org.springframework.integration.file.FileReadingMessageSource; import org.springframework.messaging.Message; import static org.assertj.core.api.Assertions.assertThat; @@ -36,6 +39,10 @@ import static org.assertj.core.api.Assertions.assertThat; */ public class DefaultFileSupplierTests extends AbstractFileSupplierTests { + @Autowired + @Qualifier("fileMessageSource") + private FileReadingMessageSource fileMessageSource; + @Test public void testBasicFlow() throws IOException { @@ -74,5 +81,8 @@ public class DefaultFileSupplierTests extends AbstractFileSupplierTests { .verifyLater(); Files.write(tempFile, "testing".getBytes()); stepVerifier.verify(); + + assertThat(this.fileMessageSource.isRunning()).isTrue(); } + }