From 95a22796eba288d53445f6d0a5a4667f8b4c5252 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=A0=D0=BE=D0=BC=D0=B0=D0=BD=20=D0=A7=D0=B8=D0=B3=D0=B2?= =?UTF-8?q?=D0=B8=D0=BD=D1=86=D0=B5=D0=B2?= Date: Wed, 17 Feb 2021 17:35:39 +0300 Subject: [PATCH] Fix bug of service instance refetching when repeating health check is activated (#904) Fixes gh-899 (cherry picked from commit 3d257b43003a02b64c4f7fd97680e7981d0b52f4) --- ...ealthCheckServiceInstanceListSupplier.java | 5 +-- ...CheckServiceInstanceListSupplierTests.java | 34 +++++++++++++++++++ 2 files changed, 37 insertions(+), 2 deletions(-) diff --git a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/HealthCheckServiceInstanceListSupplier.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/HealthCheckServiceInstanceListSupplier.java index a661c714..28f73836 100644 --- a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/HealthCheckServiceInstanceListSupplier.java +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/HealthCheckServiceInstanceListSupplier.java @@ -43,6 +43,7 @@ import org.springframework.web.util.UriComponentsBuilder; * * @author Olga Maciaszek-Sharma * @author Roman Matiushchenko + * @author Roman Chigvintsev * @since 2.2.0 */ public class HealthCheckServiceInstanceListSupplier @@ -73,9 +74,9 @@ public class HealthCheckServiceInstanceListSupplier .onlyIf(repeatContext -> this.healthCheck.getRefetchInstances()) .fixedBackoff(healthCheck.getRefetchInstancesInterval()); Flux> aliveInstancesFlux = Flux.defer(delegate) + .repeatWhen(aliveInstancesReplayRepeat) .switchMap(serviceInstances -> healthCheckFlux(serviceInstances).map( - alive -> Collections.unmodifiableList(new ArrayList<>(alive)))) - .repeatWhen(aliveInstancesReplayRepeat); + alive -> Collections.unmodifiableList(new ArrayList<>(alive)))); aliveInstancesReplay = aliveInstancesFlux .delaySubscription(Duration.ofMillis(healthCheck.getInitialDelay())) .replay(1).refCount(1); diff --git a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/HealthCheckServiceInstanceListSupplierTests.java b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/HealthCheckServiceInstanceListSupplierTests.java index 2318ba38..78954c94 100644 --- a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/HealthCheckServiceInstanceListSupplierTests.java +++ b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/HealthCheckServiceInstanceListSupplierTests.java @@ -56,6 +56,7 @@ import static org.mockito.Mockito.when; * * @author Olga Maciaszek-Sharma * @author Roman Matiushchenko + * @author Roman Chigvintsev */ @ExtendWith(SpringExtension.class) @SpringBootTest( @@ -457,6 +458,39 @@ class HealthCheckServiceInstanceListSupplierTests { .verify(VERIFY_TIMEOUT); } + @Test + void shouldRefetchInstancesWithRepeatingHealthCheck() { + healthCheck.setInitialDelay(1000); + healthCheck.setRepeatHealthCheck(true); + healthCheck.setRefetchInstancesInterval(Duration.ofSeconds(1)); + healthCheck.setRefetchInstances(true); + ServiceInstance serviceInstance1 = new DefaultServiceInstance("ignored-service-1", + SERVICE_ID, "127.0.0.1", port, false); + ServiceInstance serviceInstance2 = new DefaultServiceInstance("ignored-service-2", + SERVICE_ID, "127.0.0.2", port, false); + + StepVerifier.withVirtualTime(() -> { + ServiceInstanceListSupplier delegate = mock( + ServiceInstanceListSupplier.class); + when(delegate.get()) + .thenReturn(Flux.just(Collections.singletonList(serviceInstance1))) + .thenReturn(Flux.just(Collections.singletonList(serviceInstance2))); + listSupplier = new HealthCheckServiceInstanceListSupplier(delegate, + healthCheck, webClient) { + @Override + protected Mono isAlive(ServiceInstance serviceInstance) { + return Mono.just(true); + } + }; + return listSupplier.get(); + }).expectSubscription() + .expectNoEvent(Duration.ofMillis(healthCheck.getInitialDelay())) + .expectNext(Lists.list(serviceInstance1)) + .thenAwait(healthCheck.getRefetchInstancesInterval()) + .expectNext(Lists.list(serviceInstance2)).thenCancel() + .verify(VERIFY_TIMEOUT); + } + @Test void shouldCacheResultIfAfterPropertiesSetInvoked() { healthCheck.setInitialDelay(1000);