From d72076c6ed3808cf6756d18981d7b3c46503c7a7 Mon Sep 17 00:00:00 2001 From: alzimmermsft <48699787+alzimmermsft@users.noreply.github.com> Date: Mon, 25 Nov 2019 12:36:03 -0800 Subject: [PATCH 1/3] Updated PagedIterable and IterableStream to prefetch 1 instead of the default --- .../azure/core/http/rest/PagedFluxBase.java | 4 +- .../core/http/rest/PagedIterableBase.java | 8 +- .../com/azure/core/util/IterableStream.java | 4 +- .../core/http/rest/PagedIterableTest.java | 280 +++++++++++------- 4 files changed, 181 insertions(+), 115 deletions(-) diff --git a/sdk/core/azure-core/src/main/java/com/azure/core/http/rest/PagedFluxBase.java b/sdk/core/azure-core/src/main/java/com/azure/core/http/rest/PagedFluxBase.java index 84d30c776feb..2e49329b52b6 100644 --- a/sdk/core/azure-core/src/main/java/com/azure/core/http/rest/PagedFluxBase.java +++ b/sdk/core/azure-core/src/main/java/com/azure/core/http/rest/PagedFluxBase.java @@ -149,7 +149,7 @@ private Publisher extractAndFetchT(PagedResponse page) { if (nextPageLink == null) { return Flux.fromIterable(page.getItems()); } - return Flux.fromIterable(page.getItems()).concatWith(byT(nextPageLink)); + return Flux.fromIterable(page.getItems()).concatWith(Flux.defer(() -> byT(nextPageLink))); } /** @@ -163,6 +163,6 @@ private Publisher extractAndFetchPage(P page) { if (nextPageLink == null) { return Flux.just(page); } - return Flux.just(page).concatWith(byPage(page.getContinuationToken())); + return Flux.just(page).concatWith(Flux.defer(() -> byPage(page.getContinuationToken()))); } } diff --git a/sdk/core/azure-core/src/main/java/com/azure/core/http/rest/PagedIterableBase.java b/sdk/core/azure-core/src/main/java/com/azure/core/http/rest/PagedIterableBase.java index c9eaaeba30f5..8139a4251567 100644 --- a/sdk/core/azure-core/src/main/java/com/azure/core/http/rest/PagedIterableBase.java +++ b/sdk/core/azure-core/src/main/java/com/azure/core/http/rest/PagedIterableBase.java @@ -46,7 +46,7 @@ public PagedIterableBase(PagedFluxBase pagedFluxBase) { * @return {@link Stream} of a Response that extends {@link PagedResponse} */ public Stream

streamByPage() { - return pagedFluxBase.byPage().toStream(); + return pagedFluxBase.byPage().toStream(1); } /** @@ -58,7 +58,7 @@ public Stream

streamByPage() { * with the continuation token */ public Stream

streamByPage(String continuationToken) { - return pagedFluxBase.byPage(continuationToken).toStream(); + return pagedFluxBase.byPage(continuationToken).toStream(1); } /** @@ -68,7 +68,7 @@ public Stream

streamByPage(String continuationToken) { * @return {@link Iterable} interface */ public Iterable

iterableByPage() { - return pagedFluxBase.byPage().toIterable(); + return pagedFluxBase.byPage().toIterable(1); } /** @@ -80,6 +80,6 @@ public Iterable

iterableByPage() { * @return {@link Iterable} interface */ public Iterable

iterableByPage(String continuationToken) { - return pagedFluxBase.byPage(continuationToken).toIterable(); + return pagedFluxBase.byPage(continuationToken).toIterable(1); } } diff --git a/sdk/core/azure-core/src/main/java/com/azure/core/util/IterableStream.java b/sdk/core/azure-core/src/main/java/com/azure/core/util/IterableStream.java index e87b868199a4..9f1e840fd7a3 100644 --- a/sdk/core/azure-core/src/main/java/com/azure/core/util/IterableStream.java +++ b/sdk/core/azure-core/src/main/java/com/azure/core/util/IterableStream.java @@ -46,7 +46,7 @@ public IterableStream(Flux flux) { * @return {@link Stream} of value {@code T}. */ public Stream stream() { - return flux.toStream(); + return flux.toStream(1); } /** @@ -57,7 +57,7 @@ public Stream stream() { */ @Override public Iterator iterator() { - return flux.toIterable().iterator(); + return flux.toIterable(1).iterator(); } } diff --git a/sdk/core/azure-core/src/test/java/com/azure/core/http/rest/PagedIterableTest.java b/sdk/core/azure-core/src/test/java/com/azure/core/http/rest/PagedIterableTest.java index d2fe97586cc6..056023444c94 100644 --- a/sdk/core/azure-core/src/test/java/com/azure/core/http/rest/PagedIterableTest.java +++ b/sdk/core/azure-core/src/test/java/com/azure/core/http/rest/PagedIterableTest.java @@ -3,162 +3,204 @@ package com.azure.core.http.rest; -import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.assertTrue; - import com.azure.core.http.HttpHeaders; import com.azure.core.http.HttpMethod; import com.azure.core.http.HttpRequest; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; -import java.net.MalformedURLException; -import java.net.URL; -import java.util.Iterator; +import java.util.ArrayList; import java.util.List; +import java.util.function.Function; +import java.util.function.Supplier; import java.util.stream.Collectors; import java.util.stream.IntStream; import java.util.stream.Stream; -import org.junit.jupiter.api.Test; -import reactor.core.publisher.Mono; + +import static org.junit.jupiter.api.Assertions.assertEquals; /** * Unit tests for {@link PagedIterable}. */ public class PagedIterableTest { - private List> pagedResponses; private List> pagedStringResponses; - @Test - public void testEmptyResults() { - PagedFlux pagedFlux = getIntegerPagedFlux(0); + private HttpHeaders httpHeaders = new HttpHeaders().put("header1", "value1").put("header2", "value2"); + private HttpRequest httpRequest = new HttpRequest(HttpMethod.GET, "http://localhost"); + private String deserializedHeaders = "header1,value1,header2,value2"; + + @ParameterizedTest + @ValueSource(ints = {0, 5}) + public void streamByPage(int numberOfPages) { + PagedFlux pagedFlux = getIntegerPagedFlux(numberOfPages); PagedIterable pagedIterable = new PagedIterable<>(pagedFlux); - assertEquals(0, pagedIterable.streamByPage().count()); + List> pages = pagedIterable.streamByPage().collect(Collectors.toList()); + + assertEquals(numberOfPages, pages.size()); + assertEquals(pagedResponses, pages); } - @Test - public void testPageStream() { - PagedFlux pagedFlux = getIntegerPagedFlux(5); + @ParameterizedTest + @ValueSource(ints = {0, 5}) + public void iterateByPage(int numberOfPages) { + PagedFlux pagedFlux = getIntegerPagedFlux(numberOfPages); PagedIterable pagedIterable = new PagedIterable<>(pagedFlux); - assertEquals(5, pagedIterable.streamByPage().count()); - assertEquals(pagedResponses, pagedIterable.streamByPage().collect(Collectors.toList())); + List> pages = new ArrayList<>(); + pagedIterable.iterableByPage().iterator().forEachRemaining(pages::add); + + assertEquals(numberOfPages, pages.size()); + assertEquals(pagedResponses, pages); } - @Test - public void testPageIterable() { - PagedFlux pagedFlux = getIntegerPagedFlux(5); + @ParameterizedTest + @ValueSource(ints = {0, 5}) + public void streamByT(int numberOfPages) { + PagedFlux pagedFlux = getIntegerPagedFlux(numberOfPages); PagedIterable pagedIterable = new PagedIterable<>(pagedFlux); - Iterator> iter = pagedIterable.iterableByPage().iterator(); + List values = pagedIterable.stream().collect(Collectors.toList()); - int index = 0; - while (iter.hasNext()) { - PagedResponse pagedResponse = iter.next(); - assertEquals(pagedResponses.get(index++), pagedResponse); - } + assertEquals(numberOfPages * 3, values.size()); + assertEquals(Stream.iterate(0, i -> i + 1).limit(numberOfPages * 3).collect(Collectors.toList()), values); } - @Test - public void testStream() { - PagedFlux pagedFlux = getIntegerPagedFlux(5); + @ParameterizedTest + @ValueSource(ints = {0, 5}) + public void iterateByT(int numberOfPages) { + PagedFlux pagedFlux = getIntegerPagedFlux(numberOfPages); PagedIterable pagedIterable = new PagedIterable<>(pagedFlux); + List values = new ArrayList<>(); + pagedIterable.iterator().forEachRemaining(values::add); - assertEquals(15, pagedIterable.stream().count()); - List ints = Stream.iterate(0, i -> i + 1).limit(15).collect(Collectors.toList()); - assertEquals(ints, pagedIterable.stream().collect(Collectors.toList())); + assertEquals(numberOfPages * 3, values.size()); + assertEquals(Stream.iterate(0, i -> i + 1).limit(numberOfPages * 3).collect(Collectors.toList()), values); } - @Test - public void testIterable() { - PagedFlux pagedFlux = getIntegerPagedFlux(5); + @ParameterizedTest + @ValueSource(ints = {0, 5}) + public void streamByPageMap(int numberOfPages) { + PagedFlux pagedFlux = getIntegerPagedFlux(numberOfPages); PagedIterable pagedIterable = new PagedIterable<>(pagedFlux); - Iterator iter = pagedIterable.iterator(); + List> pages = pagedIterable.mapPage(String::valueOf).streamByPage() + .collect(Collectors.toList()); - int index = 0; - while (iter.hasNext()) { - int val = iter.next(); - assertEquals(index++, val); + assertEquals(numberOfPages, pages.size()); + for (int i = 0; i < numberOfPages; i++) { + assertEquals(pagedStringResponses.get(i).getValue(), pages.get(i).getValue()); } } + @ParameterizedTest + @ValueSource(ints = {0, 5}) + public void iterateByPageMap(int numberOfPages) { + PagedFlux pagedFlux = getIntegerPagedFlux(numberOfPages); + PagedIterable pagedIterable = new PagedIterable<>(pagedFlux); + List> pages = new ArrayList<>(); + pagedIterable.mapPage(String::valueOf).iterableByPage().iterator().forEachRemaining(pages::add); + + assertEquals(numberOfPages, pages.size()); + for (int i = 0; i < numberOfPages; i++) { + assertEquals(pagedStringResponses.get(i).getValue(), pages.get(i).getValue()); + } + } + + @ParameterizedTest + @ValueSource(ints = {0, 5}) + public void streamByTMap(int numberOfPages) { + PagedFlux pagedFlux = getIntegerPagedFlux(numberOfPages); + PagedIterable pagedIterable = new PagedIterable<>(pagedFlux); + List values = pagedIterable.mapPage(String::valueOf).stream().collect(Collectors.toList()); + + assertEquals(numberOfPages * 3, values.size()); + assertEquals(Stream.iterate(0, i -> i + 1).limit(numberOfPages * 3).map(String::valueOf) + .collect(Collectors.toList()), values); + } + + @ParameterizedTest + @ValueSource(ints = {0, 5}) + public void iterateByTMap(int numberOfPages) { + PagedFlux pagedFlux = getIntegerPagedFlux(numberOfPages); + PagedIterable pagedIterable = new PagedIterable<>(pagedFlux); + List values = new ArrayList<>(); + pagedIterable.mapPage(String::valueOf).iterator().forEachRemaining(values::add); + + assertEquals(numberOfPages * 3, values.size()); + assertEquals(Stream.iterate(0, i -> i + 1).limit(numberOfPages * 3).map(String::valueOf) + .collect(Collectors.toList()), values); + } + @Test - public void testMap() { - PagedFlux pagedFlux = getIntegerPagedFlux(5); + public void streamFirstPage() { + TestPagedFlux pagedFlux = getTestPagedFlux(5); PagedIterable pagedIterable = new PagedIterable<>(pagedFlux); - List intStrs = - Stream.iterate(0, i -> i + 1).map(String::valueOf).limit(15).collect(Collectors.toList()); - assertEquals(intStrs, pagedIterable.mapPage(String::valueOf).stream().collect(Collectors.toList())); + + assertEquals(pagedResponses.get(0), pagedIterable.streamByPage().limit(1).collect(Collectors.toList()).get(0)); + assertEquals(0, pagedFlux.getNextPageRetrievals()); } @Test - public void testPageMap() { - PagedFlux pagedFlux = getIntegerPagedFlux(5); + public void iterateFirstPage() { + TestPagedFlux pagedFlux = getTestPagedFlux(5); PagedIterable pagedIterable = new PagedIterable<>(pagedFlux); - int[] index = new int[1]; - assertTrue(pagedIterable.mapPage(String::valueOf).streamByPage().allMatch(pagedResponse -> - pagedStringResponses.get(index[0]++).getValue().equals(pagedResponse.getValue()))); + + assertEquals(pagedResponses.get(0), pagedIterable.iterableByPage().iterator().next()); + assertEquals(0, pagedFlux.getNextPageRetrievals()); } - private PagedFlux getIntegerPagedFlux(int noOfPages) { - try { - HttpHeaders httpHeaders = new HttpHeaders().put("header1", "value1") - .put("header2", "value2"); + @Test + public void streamFirstValue() { + TestPagedFlux pagedFlux = getTestPagedFlux(5); + PagedIterable pagedIterable = new PagedIterable<>(pagedFlux); - HttpRequest httpRequest = new HttpRequest(HttpMethod.GET, new URL("http://localhost")); + Integer firstValue = pagedResponses.get(0).getValue().get(0); + assertEquals(firstValue, pagedIterable.stream().limit(1).collect(Collectors.toList()).get(0)); + } - String deserializedHeaders = "header1,value1,header2,value2"; - pagedResponses = IntStream.range(0, noOfPages) - .boxed() - .map(i -> createPagedResponse(httpRequest, httpHeaders, deserializedHeaders, i, noOfPages)) - .collect(Collectors.toList()); + @Test + public void iterateFirstValue() { + TestPagedFlux pagedFlux = getTestPagedFlux(5); + PagedIterable pagedIterable = new PagedIterable<>(pagedFlux); - pagedStringResponses = IntStream.range(0, noOfPages) - .boxed() - .map(i -> createPagedResponseWithString(httpRequest, httpHeaders, deserializedHeaders, i, noOfPages)) - .collect(Collectors.toList()); + Integer firstValue = pagedResponses.get(0).getValue().get(0); + assertEquals(firstValue, pagedIterable.iterator().next()); + assertEquals(0, pagedFlux.getNextPageRetrievals()); + } - return new PagedFlux<>(() -> pagedResponses.isEmpty() ? Mono.empty() : Mono.just(pagedResponses.get(0)), - continuationToken -> getNextPage(continuationToken, pagedResponses)); - } catch (MalformedURLException e) { - return null; - } + private PagedFlux getIntegerPagedFlux(int numberOfPages) { + createPagedResponse(numberOfPages); + + return new PagedFlux<>(() -> pagedResponses.isEmpty() ? Mono.empty() : Mono.just(pagedResponses.get(0)), + continuationToken -> getNextPage(continuationToken, pagedResponses)); } - private PagedFlux getIntegerPagedFluxSinglePage() { - try { - HttpHeaders httpHeaders = new HttpHeaders().put("header1", "value1") - .put("header2", "value2"); - HttpRequest httpRequest = new HttpRequest(HttpMethod.GET, new URL("http://localhost")); - - String deserializedHeaders = "header1,value1,header2,value2"; - pagedResponses = IntStream.range(0, 1) - .boxed() - .map(i -> createPagedResponse(httpRequest, httpHeaders, deserializedHeaders, i, 1)) - .collect(Collectors.toList()); - - pagedStringResponses = IntStream.range(0, 1) - .boxed() - .map(i -> createPagedResponseWithString(httpRequest, httpHeaders, deserializedHeaders, i, 1)) - .collect(Collectors.toList()); - return new PagedFlux<>(() -> pagedResponses.isEmpty() ? Mono.empty() : Mono.just(pagedResponses.get(0))); - } catch (MalformedURLException e) { - return null; - } + private TestPagedFlux getTestPagedFlux(int numberOfPages) { + createPagedResponse(numberOfPages); + + return new TestPagedFlux<>(() -> pagedResponses.isEmpty() ? Mono.empty() : Mono.just(pagedResponses.get(0)), + continuationToken -> getNextPage(continuationToken, pagedResponses)); } - private PagedResponseBase createPagedResponse(HttpRequest httpRequest, - HttpHeaders httpHeaders, String deserializedHeaders, int i, int noOfPages) { - return new PagedResponseBase<>(httpRequest, 200, - httpHeaders, - getItems(i), - i < noOfPages - 1 ? String.valueOf(i + 1) : null, - deserializedHeaders); + private void createPagedResponse(int numberOfPages) { + pagedResponses = IntStream.range(0, numberOfPages) + .boxed() + .map(i -> + createPagedResponse(httpRequest, httpHeaders, deserializedHeaders, numberOfPages, this::getItems, i)) + .collect(Collectors.toList()); + + pagedStringResponses = IntStream.range(0, numberOfPages) + .boxed() + .map(i -> createPagedResponse(httpRequest, httpHeaders, deserializedHeaders, numberOfPages, + this::getStringItems, i)) + .collect(Collectors.toList()); } - private PagedResponseBase createPagedResponseWithString(HttpRequest httpRequest, - HttpHeaders httpHeaders, String deserializedHeaders, int i, int noOfPages) { - return new PagedResponseBase<>(httpRequest, 200, - httpHeaders, - getStringItems(i), - i < noOfPages - 1 ? String.valueOf(i + 1) : null, + private PagedResponseBase createPagedResponse(HttpRequest httpRequest, HttpHeaders headers, + String deserializedHeaders, int numberOfPages, Function> valueSupplier, int i) { + return new PagedResponseBase<>(httpRequest, 200, headers, valueSupplier.apply(i), + (i < numberOfPages - 1) ? String.valueOf(i + 1) : null, deserializedHeaders); } @@ -169,15 +211,39 @@ private Mono> getNextPage(String continuationToken, return Mono.empty(); } - return Mono.just(pagedResponses.get(Integer.valueOf(continuationToken))); + return Mono.just(pagedResponses.get(Integer.parseInt(continuationToken))); } - private List getItems(Integer i) { + private List getItems(int i) { return IntStream.range(i * 3, i * 3 + 3).boxed().collect(Collectors.toList()); } - private List getStringItems(Integer i) { - return IntStream.range(i * 3, i * 3 + 3).boxed().map(val -> String.valueOf(val)).collect(Collectors.toList()); + private List getStringItems(int i) { + return IntStream.range(i * 3, i * 3 + 3).boxed().map(String::valueOf).collect(Collectors.toList()); } + /* + * Test class used to verify that paged iterable will lazily request next pages. + */ + private static class TestPagedFlux extends PagedFlux { + private int nextPageRetrievals = 0; + + public TestPagedFlux(Supplier>> firstPageRetriever, + Function>> nextPageRetriever) { + super(firstPageRetriever, nextPageRetriever); + } + + @Override + public Flux> byPage(String continuationToken) { + nextPageRetrievals++; + return super.byPage(continuationToken); + } + + /* + * Returns the number of times another page has been retrieved. + */ + int getNextPageRetrievals() { + return nextPageRetrievals; + } + } } From 207c434520c4fd879d80a9231da20d1e82319f95 Mon Sep 17 00:00:00 2001 From: alzimmermsft <48699787+alzimmermsft@users.noreply.github.com> Date: Mon, 25 Nov 2019 13:18:10 -0800 Subject: [PATCH 2/3] Updated to use constant and commented tests --- .../core/http/rest/PagedIterableBase.java | 14 +++++++---- .../com/azure/core/util/IterableStream.java | 10 ++++++-- .../core/http/rest/PagedIterableTest.java | 24 +++++++++++++++++-- 3 files changed, 40 insertions(+), 8 deletions(-) diff --git a/sdk/core/azure-core/src/main/java/com/azure/core/http/rest/PagedIterableBase.java b/sdk/core/azure-core/src/main/java/com/azure/core/http/rest/PagedIterableBase.java index 8139a4251567..cd974741386d 100644 --- a/sdk/core/azure-core/src/main/java/com/azure/core/http/rest/PagedIterableBase.java +++ b/sdk/core/azure-core/src/main/java/com/azure/core/http/rest/PagedIterableBase.java @@ -28,6 +28,12 @@ * @see IterableStream */ public class PagedIterableBase> extends IterableStream { + /* + * This is the default batch size that will be requested when using stream or iterable by page, this will indicate + * to Reactor how many elements should be prefetched before another batch is requested. + */ + private static final int DEFAULT_BATCH_SIZE = 1; + private final PagedFluxBase pagedFluxBase; /** @@ -46,7 +52,7 @@ public PagedIterableBase(PagedFluxBase pagedFluxBase) { * @return {@link Stream} of a Response that extends {@link PagedResponse} */ public Stream

streamByPage() { - return pagedFluxBase.byPage().toStream(1); + return pagedFluxBase.byPage().toStream(DEFAULT_BATCH_SIZE); } /** @@ -58,7 +64,7 @@ public Stream

streamByPage() { * with the continuation token */ public Stream

streamByPage(String continuationToken) { - return pagedFluxBase.byPage(continuationToken).toStream(1); + return pagedFluxBase.byPage(continuationToken).toStream(DEFAULT_BATCH_SIZE); } /** @@ -68,7 +74,7 @@ public Stream

streamByPage(String continuationToken) { * @return {@link Iterable} interface */ public Iterable

iterableByPage() { - return pagedFluxBase.byPage().toIterable(1); + return pagedFluxBase.byPage().toIterable(DEFAULT_BATCH_SIZE); } /** @@ -80,6 +86,6 @@ public Iterable

iterableByPage() { * @return {@link Iterable} interface */ public Iterable

iterableByPage(String continuationToken) { - return pagedFluxBase.byPage(continuationToken).toIterable(1); + return pagedFluxBase.byPage(continuationToken).toIterable(DEFAULT_BATCH_SIZE); } } diff --git a/sdk/core/azure-core/src/main/java/com/azure/core/util/IterableStream.java b/sdk/core/azure-core/src/main/java/com/azure/core/util/IterableStream.java index 9f1e840fd7a3..4ef65b37aa63 100644 --- a/sdk/core/azure-core/src/main/java/com/azure/core/util/IterableStream.java +++ b/sdk/core/azure-core/src/main/java/com/azure/core/util/IterableStream.java @@ -28,6 +28,12 @@ * @see Iterable */ public class IterableStream implements Iterable { + /* + * This is the default batch size that will be requested when using stream or iterable by page, this will indicate + * to Reactor how many elements should be prefetched before another batch is requested. + */ + private static final int DEFAULT_BATCH_SIZE = 1; + private final Flux flux; /** @@ -46,7 +52,7 @@ public IterableStream(Flux flux) { * @return {@link Stream} of value {@code T}. */ public Stream stream() { - return flux.toStream(1); + return flux.toStream(DEFAULT_BATCH_SIZE); } /** @@ -57,7 +63,7 @@ public Stream stream() { */ @Override public Iterator iterator() { - return flux.toIterable(1).iterator(); + return flux.toIterable(DEFAULT_BATCH_SIZE).iterator(); } } diff --git a/sdk/core/azure-core/src/test/java/com/azure/core/http/rest/PagedIterableTest.java b/sdk/core/azure-core/src/test/java/com/azure/core/http/rest/PagedIterableTest.java index 056023444c94..fdd5ffe287ce 100644 --- a/sdk/core/azure-core/src/test/java/com/azure/core/http/rest/PagedIterableTest.java +++ b/sdk/core/azure-core/src/test/java/com/azure/core/http/rest/PagedIterableTest.java @@ -138,7 +138,17 @@ public void streamFirstPage() { PagedIterable pagedIterable = new PagedIterable<>(pagedFlux); assertEquals(pagedResponses.get(0), pagedIterable.streamByPage().limit(1).collect(Collectors.toList()).get(0)); - assertEquals(0, pagedFlux.getNextPageRetrievals()); + + /* + * The goal for this test would be that 0 next page retrieval calls are made but due to how Flux.concatWith + * works it needs to begin the next publisher to determine whether onNext or onComplete should trigger. This + * results in 2 next page retrieval calls for the following reason: + * + * - Makes the initial get first page call, then needs to validate that get next page emits. 1 call made. + * - Retrieving the first page in verification moves the stream iterator to the initial next page, Reactor then + * needs to verify that the page after it emits. 2 calls made. + */ + assertEquals(2, pagedFlux.getNextPageRetrievals()); } @Test @@ -147,7 +157,17 @@ public void iterateFirstPage() { PagedIterable pagedIterable = new PagedIterable<>(pagedFlux); assertEquals(pagedResponses.get(0), pagedIterable.iterableByPage().iterator().next()); - assertEquals(0, pagedFlux.getNextPageRetrievals()); + + /* + * The goal for this test would be that 0 next page retrieval calls are made but due to how Flux.concatWith + * works it needs to begin the next publisher to determine whether onNext or onComplete should trigger. This + * results in 2 next page retrieval calls for the following reason: + * + * - Makes the initial get first page call, then needs to validate that get next page emits. 1 call made. + * - Retrieving the first page in verification moves the stream iterator to the initial next page, Reactor then + * needs to verify that the page after it emits. 2 calls made. + */ + assertEquals(2, pagedFlux.getNextPageRetrievals()); } @Test From 9c866f3fabed5ed6194df44fbf1fb1bb12cbafe9 Mon Sep 17 00:00:00 2001 From: alzimmermsft <48699787+alzimmermsft@users.noreply.github.com> Date: Mon, 25 Nov 2019 13:31:38 -0800 Subject: [PATCH 3/3] Fixed linting issue --- .../test/java/com/azure/core/http/rest/PagedIterableTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/core/azure-core/src/test/java/com/azure/core/http/rest/PagedIterableTest.java b/sdk/core/azure-core/src/test/java/com/azure/core/http/rest/PagedIterableTest.java index fdd5ffe287ce..136abaad61a8 100644 --- a/sdk/core/azure-core/src/test/java/com/azure/core/http/rest/PagedIterableTest.java +++ b/sdk/core/azure-core/src/test/java/com/azure/core/http/rest/PagedIterableTest.java @@ -248,7 +248,7 @@ private List getStringItems(int i) { private static class TestPagedFlux extends PagedFlux { private int nextPageRetrievals = 0; - public TestPagedFlux(Supplier>> firstPageRetriever, + TestPagedFlux(Supplier>> firstPageRetriever, Function>> nextPageRetriever) { super(firstPageRetriever, nextPageRetriever); }