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 31bd8c5a..7fb49a10 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 @@ -66,11 +66,11 @@ public class HealthCheckServiceInstanceListSupplier LoadBalancerProperties.HealthCheck healthCheck, WebClient webClient) { this.delegate = delegate; this.healthCheck = healthCheck; - this.defaultHealthCheckPath = healthCheck.getPath().getOrDefault("default", + defaultHealthCheckPath = healthCheck.getPath().getOrDefault("default", "/actuator/health"); this.webClient = webClient; - this.aliveInstancesReplay = Flux.defer(delegate) - .delaySubscription(Duration.ofMillis(this.healthCheck.getInitialDelay())) + aliveInstancesReplay = Flux.defer(delegate) + .delaySubscription(Duration.ofMillis(healthCheck.getInitialDelay())) .switchMap(serviceInstances -> healthCheckFlux(serviceInstances).map( alive -> Collections.unmodifiableList(new ArrayList<>(alive)))) .replay(1).refCount(1); @@ -97,12 +97,12 @@ public class HealthCheckServiceInstanceListSupplier instance.getServiceId(), instance.getUri()), error); } return Mono.empty(); - }).timeout(this.healthCheck.getInterval(), Mono.defer(() -> { + }).timeout(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", instance.getServiceId(), instance.getUri(), - this.healthCheck.getInterval())); + healthCheck.getInterval())); } return Mono.empty(); })).handle((isHealthy, sink) -> { @@ -118,7 +118,7 @@ public class HealthCheckServiceInstanceListSupplier result.add(alive); return result; }).defaultIfEmpty(result); - }).repeatWhen(restart -> restart.delayElements(this.healthCheck.getInterval())); + }).repeatWhen(restart -> restart.delayElements(healthCheck.getInterval())); } @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 bda3279b..accd279d 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 @@ -137,21 +137,22 @@ class HealthCheckServiceInstanceListSupplierTests { void shouldReturnOnlyAliveService() { healthCheck.setInitialDelay(1000); - ServiceInstance si1 = new DefaultServiceInstance("ignored-service-1", SERVICE_ID, - "127.0.0.1", port, false); - ServiceInstance si2 = new DefaultServiceInstance("ignored-service-2", SERVICE_ID, - "127.0.0.2", port, false); + 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 = Mockito .mock(ServiceInstanceListSupplier.class); Mockito.when(delegate.getServiceId()).thenReturn(SERVICE_ID); - Mockito.when(delegate.get()).thenReturn(Flux.just(Lists.list(si1, si2))); + Mockito.when(delegate.get()).thenReturn( + Flux.just(Lists.list(serviceInstance1, serviceInstance2))); HealthCheckServiceInstanceListSupplier mock = Mockito .mock(HealthCheckServiceInstanceListSupplier.class); - Mockito.doReturn(Mono.just(true)).when(mock).isAlive(si1); - Mockito.doReturn(Mono.just(false)).when(mock).isAlive(si2); + Mockito.doReturn(Mono.just(true)).when(mock).isAlive(serviceInstance1); + Mockito.doReturn(Mono.just(false)).when(mock).isAlive(serviceInstance2); listSupplier = new HealthCheckServiceInstanceListSupplier(delegate, healthCheck, webClient) { @@ -164,28 +165,30 @@ class HealthCheckServiceInstanceListSupplierTests { return listSupplier.get(); }).expectSubscription() .expectNoEvent(Duration.ofMillis(healthCheck.getInitialDelay())) - .expectNext(Lists.list(si1)).expectNoEvent(healthCheck.getInterval()) - .thenCancel().verify(VERIFY_TIMEOUT); + .expectNext(Lists.list(serviceInstance1)) + .expectNoEvent(healthCheck.getInterval()).thenCancel() + .verify(VERIFY_TIMEOUT); } @Test void shouldEmitOnEachAliveServiceInBatch() { healthCheck.setInitialDelay(1000); - ServiceInstance si1 = new DefaultServiceInstance("ignored-service-1", SERVICE_ID, - "127.0.0.1", port, false); - ServiceInstance si2 = new DefaultServiceInstance("ignored-service-2", SERVICE_ID, - "127.0.0.2", port, false); + 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 = Mockito .mock(ServiceInstanceListSupplier.class); Mockito.when(delegate.getServiceId()).thenReturn(SERVICE_ID); - Mockito.when(delegate.get()).thenReturn(Flux.just(Lists.list(si1, si2))); + Mockito.when(delegate.get()).thenReturn( + Flux.just(Lists.list(serviceInstance1, serviceInstance2))); HealthCheckServiceInstanceListSupplier mock = Mockito .mock(HealthCheckServiceInstanceListSupplier.class); - Mockito.doReturn(Mono.just(true)).when(mock).isAlive(si1); - Mockito.doReturn(Mono.just(true)).when(mock).isAlive(si2); + Mockito.doReturn(Mono.just(true)).when(mock).isAlive(serviceInstance1); + Mockito.doReturn(Mono.just(true)).when(mock).isAlive(serviceInstance2); listSupplier = new HealthCheckServiceInstanceListSupplier(delegate, healthCheck, webClient) { @@ -198,7 +201,8 @@ class HealthCheckServiceInstanceListSupplierTests { return listSupplier.get(); }).expectSubscription() .expectNoEvent(Duration.ofMillis(healthCheck.getInitialDelay())) - .expectNext(Lists.list(si1)).expectNext(Lists.list(si1, si2)) + .expectNext(Lists.list(serviceInstance1)) + .expectNext(Lists.list(serviceInstance1, serviceInstance2)) .expectNoEvent(healthCheck.getInterval()).thenCancel() .verify(VERIFY_TIMEOUT); } @@ -206,22 +210,23 @@ class HealthCheckServiceInstanceListSupplierTests { @Test void shouldNotFailIfIsAliveReturnsError() { healthCheck.setInitialDelay(1000); - ServiceInstance si1 = new DefaultServiceInstance("ignored-service-1", SERVICE_ID, - "127.0.0.1", port, false); - ServiceInstance si2 = new DefaultServiceInstance("ignored-service-2", SERVICE_ID, - "127.0.0.2", port, false); + 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 = Mockito .mock(ServiceInstanceListSupplier.class); Mockito.when(delegate.getServiceId()).thenReturn(SERVICE_ID); - Mockito.when(delegate.get()).thenReturn(Flux.just(Lists.list(si1, si2))); + Mockito.when(delegate.get()).thenReturn( + Flux.just(Lists.list(serviceInstance1, serviceInstance2))); HealthCheckServiceInstanceListSupplier mock = Mockito .mock(HealthCheckServiceInstanceListSupplier.class); - Mockito.doReturn(Mono.just(true)).when(mock).isAlive(si1); + Mockito.doReturn(Mono.just(true)).when(mock).isAlive(serviceInstance1); Mockito.doReturn(Mono.error(new RuntimeException("boom"))).when(mock) - .isAlive(si2); + .isAlive(serviceInstance2); listSupplier = new HealthCheckServiceInstanceListSupplier(delegate, healthCheck, webClient) { @@ -234,28 +239,30 @@ class HealthCheckServiceInstanceListSupplierTests { return listSupplier.get(); }).expectSubscription() .expectNoEvent(Duration.ofMillis(healthCheck.getInitialDelay())) - .expectNext(Lists.list(si1)).expectNoEvent(healthCheck.getInterval()) - .thenCancel().verify(VERIFY_TIMEOUT); + .expectNext(Lists.list(serviceInstance1)) + .expectNoEvent(healthCheck.getInterval()).thenCancel() + .verify(VERIFY_TIMEOUT); } @Test void shouldEmitAllInstancesIfAllIsAliveChecksFailed() { healthCheck.setInitialDelay(1000); - ServiceInstance si1 = new DefaultServiceInstance("ignored-service-1", SERVICE_ID, - "127.0.0.1", port, false); - ServiceInstance si2 = new DefaultServiceInstance("ignored-service-2", SERVICE_ID, - "127.0.0.2", port, false); + 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 = Mockito .mock(ServiceInstanceListSupplier.class); Mockito.when(delegate.getServiceId()).thenReturn(SERVICE_ID); - Mockito.when(delegate.get()).thenReturn(Flux.just(Lists.list(si1, si2))); + Mockito.when(delegate.get()).thenReturn( + Flux.just(Lists.list(serviceInstance1, serviceInstance2))); listSupplier = new HealthCheckServiceInstanceListSupplier(delegate, healthCheck, webClient) { @Override protected Mono isAlive(ServiceInstance serviceInstance) { - if (serviceInstance == si1) { + if (serviceInstance == serviceInstance1) { return Mono.just(false); } else { @@ -274,14 +281,15 @@ class HealthCheckServiceInstanceListSupplierTests { @Test void shouldMakeInitialDaleyAfterPropertiesSet() { healthCheck.setInitialDelay(1000); - ServiceInstance si1 = new DefaultServiceInstance("ignored-service-1", SERVICE_ID, - "127.0.0.1", port, false); + ServiceInstance serviceInstance1 = new DefaultServiceInstance("ignored-service-1", + SERVICE_ID, "127.0.0.1", port, false); StepVerifier.withVirtualTime(() -> { ServiceInstanceListSupplier delegate = Mockito .mock(ServiceInstanceListSupplier.class); Mockito.when(delegate.getServiceId()).thenReturn(SERVICE_ID); - Mockito.when(delegate.get()).thenReturn(Flux.just(Lists.list(si1))); + Mockito.when(delegate.get()) + .thenReturn(Flux.just(Lists.list(serviceInstance1))); listSupplier = new HealthCheckServiceInstanceListSupplier(delegate, healthCheck, webClient) { @Override @@ -295,29 +303,32 @@ class HealthCheckServiceInstanceListSupplierTests { return listSupplier.get(); }).expectSubscription() .expectNoEvent(Duration.ofMillis(healthCheck.getInitialDelay())) - .expectNext(Lists.list(si1)).expectNoEvent(healthCheck.getInterval()) - .thenCancel().verify(VERIFY_TIMEOUT); + .expectNext(Lists.list(serviceInstance1)) + .expectNoEvent(healthCheck.getInterval()).thenCancel() + .verify(VERIFY_TIMEOUT); } @Test void shouldRepeatIsAliveChecksIndefinitely() { healthCheck.setInitialDelay(1000); - ServiceInstance si1 = new DefaultServiceInstance("ignored-service-1", SERVICE_ID, - "127.0.0.1", port, false); - ServiceInstance si2 = new DefaultServiceInstance("ignored-service-2", SERVICE_ID, - "127.0.0.2", port, false); + 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 = Mockito .mock(ServiceInstanceListSupplier.class); Mockito.when(delegate.getServiceId()).thenReturn(SERVICE_ID); - Mockito.when(delegate.get()).thenReturn(Flux.just(Lists.list(si1, si2))); + Mockito.when(delegate.get()).thenReturn( + Flux.just(Lists.list(serviceInstance1, serviceInstance2))); HealthCheckServiceInstanceListSupplier mock = Mockito .mock(HealthCheckServiceInstanceListSupplier.class); - Mockito.doReturn(Mono.just(false), Mono.just(true)).when(mock).isAlive(si1); + Mockito.doReturn(Mono.just(false), Mono.just(true)).when(mock) + .isAlive(serviceInstance1); Mockito.doReturn(Mono.error(new RuntimeException("boom"))).when(mock) - .isAlive(si2); + .isAlive(serviceInstance2); listSupplier = new HealthCheckServiceInstanceListSupplier(delegate, healthCheck, webClient) { @@ -331,25 +342,29 @@ class HealthCheckServiceInstanceListSupplierTests { }).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(serviceInstance1)) + .expectNoEvent(healthCheck.getInterval()) + .expectNext(Lists.list(serviceInstance1)).thenCancel() + .verify(VERIFY_TIMEOUT); } @Test void shouldTimeoutIsAliveCheck() { healthCheck.setInitialDelay(1000); - ServiceInstance si1 = new DefaultServiceInstance("ignored-service-1", SERVICE_ID, - "127.0.0.1", port, false); + ServiceInstance serviceInstance1 = new DefaultServiceInstance("ignored-service-1", + SERVICE_ID, "127.0.0.1", port, false); StepVerifier.withVirtualTime(() -> { ServiceInstanceListSupplier delegate = Mockito .mock(ServiceInstanceListSupplier.class); Mockito.when(delegate.getServiceId()).thenReturn(SERVICE_ID); - Mockito.when(delegate.get()).thenReturn(Flux.just(Lists.list(si1))); + Mockito.when(delegate.get()) + .thenReturn(Flux.just(Lists.list(serviceInstance1))); HealthCheckServiceInstanceListSupplier mock = Mockito .mock(HealthCheckServiceInstanceListSupplier.class); - Mockito.when(mock.isAlive(si1)).thenReturn(Mono.never(), Mono.just(true)); + Mockito.when(mock.isAlive(serviceInstance1)).thenReturn(Mono.never(), + Mono.just(true)); listSupplier = new HealthCheckServiceInstanceListSupplier(delegate, healthCheck, webClient) { @@ -363,25 +378,28 @@ class HealthCheckServiceInstanceListSupplierTests { }).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)) - .thenCancel().verify(VERIFY_TIMEOUT); + .expectNoEvent(healthCheck.getInterval()) + .expectNext(Lists.list(serviceInstance1)) + .expectNoEvent(healthCheck.getInterval()) + .expectNext(Lists.list(serviceInstance1)).thenCancel() + .verify(VERIFY_TIMEOUT); } @Test void shouldUpdateInstances() { healthCheck.setInitialDelay(1000); - ServiceInstance si1 = new DefaultServiceInstance("ignored-service-1", SERVICE_ID, - "127.0.0.1", port, false); - ServiceInstance si2 = new DefaultServiceInstance("ignored-service-2", SERVICE_ID, - "127.0.0.2", port, false); + 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 = Mockito .mock(ServiceInstanceListSupplier.class); Mockito.when(delegate.getServiceId()).thenReturn(SERVICE_ID); - Flux> instances = Flux.just(Lists.list(si1)) - .concatWith(Flux.just(Lists.list(si1, si2)) + Flux> instances = Flux + .just(Lists.list(serviceInstance1)) + .concatWith(Flux.just(Lists.list(serviceInstance1, serviceInstance2)) .delayElements(healthCheck.getInterval().dividedBy(2))); Mockito.when(delegate.get()).thenReturn(instances); @@ -396,18 +414,21 @@ class HealthCheckServiceInstanceListSupplierTests { return listSupplier.get(); }).expectSubscription() .expectNoEvent(Duration.ofMillis(healthCheck.getInitialDelay())) - .expectNext(Lists.list(si1)) + .expectNext(Lists.list(serviceInstance1)) .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(serviceInstance1)) + .expectNext(Lists.list(serviceInstance1, serviceInstance2)) + .expectNoEvent(healthCheck.getInterval()) + .expectNext(Lists.list(serviceInstance1)) + .expectNext(Lists.list(serviceInstance1, serviceInstance2)).thenCancel() + .verify(VERIFY_TIMEOUT); } @Test void shouldCacheResultIfAfterPropertiesSetInvoked() { healthCheck.setInitialDelay(1000); - ServiceInstance si1 = new DefaultServiceInstance("ignored-service-1", SERVICE_ID, - "127.0.0.1", port, false); + ServiceInstance serviceInstance1 = new DefaultServiceInstance("ignored-service-1", + SERVICE_ID, "127.0.0.1", port, false); AtomicInteger emitCounter = new AtomicInteger(); @@ -415,7 +436,8 @@ class HealthCheckServiceInstanceListSupplierTests { ServiceInstanceListSupplier delegate = Mockito .mock(ServiceInstanceListSupplier.class); Mockito.when(delegate.getServiceId()).thenReturn(SERVICE_ID); - Mockito.when(delegate.get()).thenReturn(Flux.just(Lists.list(si1))); + Mockito.when(delegate.get()) + .thenReturn(Flux.just(Lists.list(serviceInstance1))); listSupplier = new HealthCheckServiceInstanceListSupplier(delegate, healthCheck, webClient) { @@ -437,7 +459,8 @@ class HealthCheckServiceInstanceListSupplierTests { return listSupplier.get().take(1).concatWith(listSupplier.get().take(1)); }).expectSubscription() .expectNoEvent(Duration.ofMillis(healthCheck.getInitialDelay())) - .expectNext(Lists.list(si1)).expectNext(Lists.list(si1)).thenCancel() + .expectNext(Lists.list(serviceInstance1)) + .expectNext(Lists.list(serviceInstance1)).thenCancel() .verify(VERIFY_TIMEOUT); Assertions.assertThat(emitCounter).hasValue(1);