diff --git a/src/main/kotlin/io/emeraldpay/dshackle/reader/CompoundReader.kt b/src/main/kotlin/io/emeraldpay/dshackle/reader/CompoundReader.kt index 26a425af..7c068961 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/reader/CompoundReader.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/reader/CompoundReader.kt @@ -23,7 +23,8 @@ import reactor.core.publisher.Mono import java.time.Duration /** - * Composition of multiple readers. Reader returns first value returned by any of the source readers. + * Composition of multiple readers. + * Reader returns first value returned by any of the source readers by checking one by one until one of them returns a non-empty result. */ class CompoundReader( private vararg val readers: Reader @@ -38,12 +39,13 @@ class CompoundReader( return Mono.empty() } return Flux.fromIterable(readers.asIterable()) - .flatMap { rdr -> + .flatMap({ rdr -> rdr.read(key) .timeout(Defaults.timeoutInternal, Mono.empty()) .doOnError { t -> log.warn("Failed to read from $rdr", t) } .onErrorResume { Mono.empty() } - }.next() + }, 1) + .next() } } \ No newline at end of file diff --git a/src/test/groovy/io/emeraldpay/dshackle/reader/CompoundReaderSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/reader/CompoundReaderSpec.groovy index 76098b17..31c3cae3 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/reader/CompoundReaderSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/reader/CompoundReaderSpec.groovy @@ -72,21 +72,20 @@ class CompoundReaderSpec extends Specification { .verify(Duration.ofSeconds(1)) } - def "Return second"() { + def "Doesn't call others after getting first"() { setup: - def reader = new CompoundReader(reader3, reader2) - when: - def act = reader.read("test") - then: - StepVerifier.create(act) - .expectNext("test-2") - .expectComplete() - .verify(Duration.ofSeconds(1)) - } - - def "Return third"() { - setup: - def reader = new CompoundReader(reader3, reader2, reader1) + def call2 = false + def reader2 = new Reader() { + @Override + Mono read(String key) { + call2 = true + return Mono.just("test-2").delaySubscription(Duration.ofMillis(200)) + } + } + def reader = new CompoundReader( + reader1, + reader2 + ) when: def act = reader.read("test") then: @@ -94,16 +93,29 @@ class CompoundReaderSpec extends Specification { .expectNext("test-1") .expectComplete() .verify(Duration.ofSeconds(1)) + !call2 } - def "Ignore empty"() { + def "Return first even if it's slow"() { setup: - def reader = new CompoundReader(reader3, reader1Empty, reader2, reader1Empty) + def reader = new CompoundReader(reader3, reader2) when: def act = reader.read("test") then: StepVerifier.create(act) - .expectNext("test-2") + .expectNext("test-3") + .expectComplete() + .verify(Duration.ofSeconds(1)) + } + + def "Ignore empty"() { + setup: + def reader = new CompoundReader(reader1Empty, reader3, reader2, reader1Empty) + when: + def act = reader.read("test") + then: + StepVerifier.create(act) + .expectNext("test-3") .expectComplete() .verify(Duration.ofSeconds(1)) }