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 86efcaae..31bd8c5a 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 @@ -71,11 +71,9 @@ public class HealthCheckServiceInstanceListSupplier this.webClient = webClient; this.aliveInstancesReplay = Flux.defer(delegate) .delaySubscription(Duration.ofMillis(this.healthCheck.getInitialDelay())) - .switchMap(serviceInstances -> healthCheckFlux(serviceInstances) - .map(alive -> Collections.unmodifiableList(new ArrayList<>(alive))) - ) - .replay(1) - .refCount(1); + .switchMap(serviceInstances -> healthCheckFlux(serviceInstances).map( + alive -> Collections.unmodifiableList(new ArrayList<>(alive)))) + .replay(1).refCount(1); } @Override @@ -87,7 +85,8 @@ public class HealthCheckServiceInstanceListSupplier this.healthCheckDisposable = aliveInstancesReplay.subscribe(); } - protected Flux> healthCheckFlux(List instances) { + protected Flux> healthCheckFlux( + List instances) { return Flux.defer(() -> { List> checks = new ArrayList<>(instances.size()); for (ServiceInstance instance : instances) { @@ -98,8 +97,7 @@ public class HealthCheckServiceInstanceListSupplier instance.getServiceId(), instance.getUri()), error); } return Mono.empty(); - }) - .timeout(this.healthCheck.getInterval(), Mono.defer(() -> { + }).timeout(this.healthCheck.getInterval(), Mono.defer(() -> { if (LOG.isDebugEnabled()) { LOG.debug(String.format( "The instance for service %s: %s did not respond for %s during health check", @@ -107,8 +105,7 @@ public class HealthCheckServiceInstanceListSupplier this.healthCheck.getInterval())); } return Mono.empty(); - })) - .handle((isHealthy, sink) -> { + })).handle((isHealthy, sink) -> { if (isHealthy) { sink.next(instance); } @@ -120,10 +117,8 @@ public class HealthCheckServiceInstanceListSupplier return Flux.merge(checks).map(alive -> { result.add(alive); return result; - }) - .defaultIfEmpty(result); - }) - .repeatWhen(restart -> restart.delayElements(this.healthCheck.getInterval())); + }).defaultIfEmpty(result); + }).repeatWhen(restart -> restart.delayElements(this.healthCheck.getInterval())); } @Override @@ -145,9 +140,8 @@ public class HealthCheckServiceInstanceListSupplier .uri(UriComponentsBuilder.fromUri(serviceInstance.getUri()) .path(healthCheckPath).build().toUri()) .exchange() - .flatMap(clientResponse -> clientResponse.releaseBody() - .thenReturn(HttpStatus.OK.value() == clientResponse.rawStatusCode()) - ); + .flatMap(clientResponse -> clientResponse.releaseBody().thenReturn( + HttpStatus.OK.value() == clientResponse.rawStatusCode())); } @Override 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 7d85ccb1..82320d59 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 @@ -162,13 +162,10 @@ class HealthCheckServiceInstanceListSupplierTests { }; return listSupplier.get(); - }) - .expectSubscription() + }).expectSubscription() .expectNoEvent(Duration.ofMillis(healthCheck.getInitialDelay())) - .expectNext(Lists.list(si1)) - .expectNoEvent(healthCheck.getInterval()) - .thenCancel() - .verify(VERIFY_TIMEOUT); + .expectNext(Lists.list(si1)).expectNoEvent(healthCheck.getInterval()) + .thenCancel().verify(VERIFY_TIMEOUT); } @Test @@ -199,13 +196,10 @@ class HealthCheckServiceInstanceListSupplierTests { }; return listSupplier.get(); - }) - .expectSubscription() + }).expectSubscription() .expectNoEvent(Duration.ofMillis(healthCheck.getInitialDelay())) - .expectNext(Lists.list(si1)) - .expectNext(Lists.list(si1, si2)) - .expectNoEvent(healthCheck.getInterval()) - .thenCancel() + .expectNext(Lists.list(si1)).expectNext(Lists.list(si1, si2)) + .expectNoEvent(healthCheck.getInterval()).thenCancel() .verify(VERIFY_TIMEOUT); } @@ -238,13 +232,10 @@ class HealthCheckServiceInstanceListSupplierTests { }; return listSupplier.get(); - }) - .expectSubscription() + }).expectSubscription() .expectNoEvent(Duration.ofMillis(healthCheck.getInitialDelay())) - .expectNext(Lists.list(si1)) - .expectNoEvent(healthCheck.getInterval()) - .thenCancel() - .verify(VERIFY_TIMEOUT); + .expectNext(Lists.list(si1)).expectNoEvent(healthCheck.getInterval()) + .thenCancel().verify(VERIFY_TIMEOUT); } @Test @@ -274,13 +265,10 @@ class HealthCheckServiceInstanceListSupplierTests { }; return listSupplier.get(); - }) - .expectSubscription() + }).expectSubscription() .expectNoEvent(Duration.ofMillis(healthCheck.getInitialDelay())) - .expectNext(Lists.list()) - .expectNoEvent(healthCheck.getInterval()) - .thenCancel() - .verify(VERIFY_TIMEOUT); + .expectNext(Lists.list()).expectNoEvent(healthCheck.getInterval()) + .thenCancel().verify(VERIFY_TIMEOUT); } @Test @@ -305,13 +293,10 @@ class HealthCheckServiceInstanceListSupplierTests { listSupplier.afterPropertiesSet(); return listSupplier.get(); - }) - .expectSubscription() + }).expectSubscription() .expectNoEvent(Duration.ofMillis(healthCheck.getInitialDelay())) - .expectNext(Lists.list(si1)) - .expectNoEvent(healthCheck.getInterval()) - .thenCancel() - .verify(VERIFY_TIMEOUT); + .expectNext(Lists.list(si1)).expectNoEvent(healthCheck.getInterval()) + .thenCancel().verify(VERIFY_TIMEOUT); } @Test @@ -343,15 +328,11 @@ class HealthCheckServiceInstanceListSupplierTests { }; return listSupplier.get(); - }) - .expectSubscription() + }).expectSubscription() .expectNoEvent(Duration.ofMillis(healthCheck.getInitialDelay())) - .expectNext(Lists.list()) - .expectNoEvent(healthCheck.getInterval()) - .expectNext(Lists.list(si1)) - .expectNoEvent(healthCheck.getInterval()) - .expectNext(Lists.list(si1)) - .thenCancel().verify(VERIFY_TIMEOUT); + .expectNext(Lists.list()).expectNoEvent(healthCheck.getInterval()) + .expectNext(Lists.list(si1)).expectNoEvent(healthCheck.getInterval()) + .expectNext(Lists.list(si1)).thenCancel().verify(VERIFY_TIMEOUT); } @Test @@ -379,15 +360,11 @@ class HealthCheckServiceInstanceListSupplierTests { }; return listSupplier.get(); - }) - .expectSubscription() + }).expectSubscription() .expectNoEvent(Duration.ofMillis(healthCheck.getInitialDelay())) - .expectNoEvent(healthCheck.getInterval()) - .expectNext(Lists.list()) - .expectNoEvent(healthCheck.getInterval()) - .expectNext(Lists.list(si1)) - .expectNoEvent(healthCheck.getInterval()) - .expectNext(Lists.list(si1)) + .expectNoEvent(healthCheck.getInterval()).expectNext(Lists.list()) + .expectNoEvent(healthCheck.getInterval()).expectNext(Lists.list(si1)) + .expectNoEvent(healthCheck.getInterval()).expectNext(Lists.list(si1)) .thenCancel().verify(VERIFY_TIMEOUT); } @@ -417,18 +394,13 @@ class HealthCheckServiceInstanceListSupplierTests { }; return listSupplier.get(); - }) - .expectSubscription() + }).expectSubscription() .expectNoEvent(Duration.ofMillis(healthCheck.getInitialDelay())) .expectNext(Lists.list(si1)) .thenAwait(healthCheck.getInterval().dividedBy(2)) - .expectNext(Lists.list(si1)) - .expectNext(Lists.list(si1, si2)) - .expectNoEvent(healthCheck.getInterval()) - .expectNext(Lists.list(si1)) - .expectNext(Lists.list(si1, si2)) - .thenCancel() - .verify(VERIFY_TIMEOUT); + .expectNext(Lists.list(si1)).expectNext(Lists.list(si1, si2)) + .expectNoEvent(healthCheck.getInterval()).expectNext(Lists.list(si1)) + .expectNext(Lists.list(si1, si2)).thenCancel().verify(VERIFY_TIMEOUT); } @Test @@ -453,20 +425,19 @@ class HealthCheckServiceInstanceListSupplierTests { } @Override - protected Flux> healthCheckFlux(List instances) { - return super.healthCheckFlux(instances).doOnNext(it -> emitCounter.incrementAndGet()); + protected Flux> healthCheckFlux( + List instances) { + return super.healthCheckFlux(instances) + .doOnNext(it -> emitCounter.incrementAndGet()); } }; listSupplier.afterPropertiesSet(); return listSupplier.get().take(1).concatWith(listSupplier.get().take(1)); - }) - .expectSubscription() + }).expectSubscription() .expectNoEvent(Duration.ofMillis(healthCheck.getInitialDelay())) - .expectNext(Lists.list(si1)) - .expectNext(Lists.list(si1)) - .thenCancel() + .expectNext(Lists.list(si1)).expectNext(Lists.list(si1)).thenCancel() .verify(VERIFY_TIMEOUT); Assertions.assertThat(emitCounter).hasValue(1); @@ -479,20 +450,19 @@ class HealthCheckServiceInstanceListSupplierTests { ServiceInstanceListSupplier delegate = Mockito .mock(ServiceInstanceListSupplier.class); - Mockito.when(delegate.get()) - .thenReturn(Flux.>never() - .log("test") - .doOnCancel(instancesCanceled::incrementAndGet)); + Mockito.when(delegate.get()).thenReturn(Flux.>never() + .log("test").doOnCancel(instancesCanceled::incrementAndGet)); - listSupplier = new HealthCheckServiceInstanceListSupplier(delegate, healthCheck, webClient); + listSupplier = new HealthCheckServiceInstanceListSupplier(delegate, healthCheck, + webClient); listSupplier.afterPropertiesSet(); Assertions.assertThat(instancesCanceled).hasValue(0); listSupplier.destroy(); - Awaitility.await() - .pollDelay(Duration.ofMillis(100)).atMost(VERIFY_TIMEOUT).untilAsserted( + Awaitility.await().pollDelay(Duration.ofMillis(100)).atMost(VERIFY_TIMEOUT) + .untilAsserted( () -> Assertions.assertThat(instancesCanceled).hasValue(1)); }