From 7287a3c767b6a24594b7b08d0a9b3f8f27583ab1 Mon Sep 17 00:00:00 2001 From: Dave Syer Date: Tue, 1 Dec 2020 09:31:36 +0000 Subject: [PATCH 1/2] Re-implement /pause endpoint A new interface PauseHandler provides a callback for /pause and /resume. The only implementation provided out of the box is for Spring Integration's Pausable (so it will work with Spring Cloud Stream Kafka for instance). Fixes gh-788 --- .../RefreshEndpointAutoConfiguration.java | 42 ++++++++++++++++++- .../cloud/context/restart/PauseHandler.java | 28 +++++++++++++ .../context/restart/RestartEndpoint.java | 17 +++++--- 3 files changed, 81 insertions(+), 6 deletions(-) create mode 100644 spring-cloud-context/src/main/java/org/springframework/cloud/context/restart/PauseHandler.java diff --git a/spring-cloud-context/src/main/java/org/springframework/cloud/autoconfigure/RefreshEndpointAutoConfiguration.java b/spring-cloud-context/src/main/java/org/springframework/cloud/autoconfigure/RefreshEndpointAutoConfiguration.java index 05c1dedb..2f34cbe7 100644 --- a/spring-cloud-context/src/main/java/org/springframework/cloud/autoconfigure/RefreshEndpointAutoConfiguration.java +++ b/spring-cloud-context/src/main/java/org/springframework/cloud/autoconfigure/RefreshEndpointAutoConfiguration.java @@ -16,6 +16,10 @@ package org.springframework.cloud.autoconfigure; +import java.util.ArrayList; +import java.util.List; +import java.util.stream.Collectors; + import org.springframework.beans.factory.ObjectProvider; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.actuate.autoconfigure.endpoint.EndpointAutoConfiguration; @@ -30,6 +34,7 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingClas import org.springframework.cloud.bootstrap.config.PropertySourceBootstrapConfiguration; import org.springframework.cloud.context.properties.ConfigurationPropertiesRebinder; import org.springframework.cloud.context.refresh.ContextRefresher; +import org.springframework.cloud.context.restart.PauseHandler; import org.springframework.cloud.context.restart.RestartEndpoint; import org.springframework.cloud.context.scope.refresh.RefreshScope; import org.springframework.cloud.endpoint.RefreshEndpoint; @@ -37,6 +42,7 @@ import org.springframework.cloud.health.RefreshScopeHealthIndicator; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; +import org.springframework.integration.core.Pausable; import org.springframework.integration.monitor.IntegrationMBeanExporter; /** @@ -82,9 +88,19 @@ public class RefreshEndpointAutoConfiguration { @ConditionalOnClass(IntegrationMBeanExporter.class) class RestartEndpointWithIntegrationConfiguration { - @Autowired(required = false) private IntegrationMBeanExporter exporter; + RestartEndpointWithIntegrationConfiguration( + @Autowired(required = false) IntegrationMBeanExporter exporter) { + this.exporter = exporter; + } + + @Bean + public PauseHandler integrationPauseHandler(ObjectProvider pausables) { + return new IntegrationPauseHandler( + pausables.orderedStream().collect(Collectors.toList())); + } + @Bean @ConditionalOnAvailableEndpoint @ConditionalOnMissingBean @@ -96,6 +112,30 @@ class RestartEndpointWithIntegrationConfiguration { return endpoint; } + private class IntegrationPauseHandler implements PauseHandler { + + private List pausables = new ArrayList<>(); + + IntegrationPauseHandler(List pausables) { + this.pausables.addAll(pausables); + } + + @Override + public void pause() { + for (Pausable pausable : this.pausables) { + pausable.pause(); + } + } + + @Override + public void resume() { + for (int i = this.pausables.size(); i-- > 0;) { + this.pausables.get(i).resume(); + } + } + + } + } @Configuration(proxyBeanMethods = false) diff --git a/spring-cloud-context/src/main/java/org/springframework/cloud/context/restart/PauseHandler.java b/spring-cloud-context/src/main/java/org/springframework/cloud/context/restart/PauseHandler.java new file mode 100644 index 00000000..57f46473 --- /dev/null +++ b/spring-cloud-context/src/main/java/org/springframework/cloud/context/restart/PauseHandler.java @@ -0,0 +1,28 @@ +/* + * Copyright 2020-2020 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.context.restart; + +/** + * @author Dave Syer + */ +public interface PauseHandler { + + void pause(); + + void resume(); + +} diff --git a/spring-cloud-context/src/main/java/org/springframework/cloud/context/restart/RestartEndpoint.java b/spring-cloud-context/src/main/java/org/springframework/cloud/context/restart/RestartEndpoint.java index 8f049b5c..d6f29090 100644 --- a/spring-cloud-context/src/main/java/org/springframework/cloud/context/restart/RestartEndpoint.java +++ b/spring-cloud-context/src/main/java/org/springframework/cloud/context/restart/RestartEndpoint.java @@ -19,10 +19,12 @@ package org.springframework.cloud.context.restart; import java.io.Closeable; import java.io.IOException; import java.util.Collections; +import java.util.List; +import java.util.stream.Collectors; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; - +// import org.springframework.beans.BeansException; import org.springframework.beans.factory.config.BeanPostProcessor; import org.springframework.boot.SpringApplication; @@ -63,6 +65,8 @@ public class RestartEndpoint implements ApplicationListener pauseHandlers = Collections.emptyList(); + private long timeout; // @ManagedAttribute @@ -88,6 +92,8 @@ public class RestartEndpoint implements ApplicationListener 0;) { + PauseHandler handler = this.pauseHandlers.get(i); + handler.resume(); } } From cbedae2e82fcfe15ecb61d23e992f405d5238257 Mon Sep 17 00:00:00 2001 From: Olga Maciaszek-Sharma Date: Tue, 1 Dec 2020 04:14:14 -0600 Subject: [PATCH 2/2] Allow refetching instances for healthcheck (#855) * Allow refetching instances by HealthCheckServiceInstanceListSupplier. * Add docs and javadocs. * Fix docs after review. --- docs/src/main/asciidoc/_configprops.adoc | 7 +- .../main/asciidoc/spring-cloud-commons.adoc | 12 ++- .../reactive/LoadBalancerProperties.java | 44 +++++++++ ...ealthCheckServiceInstanceListSupplier.java | 17 +++- ...CheckServiceInstanceListSupplierTests.java | 97 +++++++++++++------ 5 files changed, 135 insertions(+), 42 deletions(-) diff --git a/docs/src/main/asciidoc/_configprops.adoc b/docs/src/main/asciidoc/_configprops.adoc index 92d5e393..0f7126ae 100644 --- a/docs/src/main/asciidoc/_configprops.adoc +++ b/docs/src/main/asciidoc/_configprops.adoc @@ -10,10 +10,6 @@ |spring.cloud.discovery.client.health-indicator.enabled | true | |spring.cloud.discovery.client.health-indicator.include-description | false | |spring.cloud.discovery.client.simple.instances | | -|spring.cloud.discovery.client.simple.local.instance-id | | The unique identifier or name for the service instance. -|spring.cloud.discovery.client.simple.local.metadata | | Metadata for the service instance. Can be used by discovery clients to modify their behaviour per instance, e.g. when load balancing. -|spring.cloud.discovery.client.simple.local.service-id | | The identifier or name for the service. Multiple instances might share the same service ID. -|spring.cloud.discovery.client.simple.local.uri | | The URI of the service instance. Will be parsed to extract the scheme, host, and port. |spring.cloud.discovery.client.simple.order | | |spring.cloud.discovery.enabled | true | Enables discovery client health indicators. |spring.cloud.features.enabled | true | Enables the features endpoint. @@ -34,6 +30,9 @@ |spring.cloud.loadbalancer.health-check.initial-delay | 0 | Initial delay value for the HealthCheck scheduler. |spring.cloud.loadbalancer.health-check.interval | 25s | Interval for rerunning the HealthCheck scheduler. |spring.cloud.loadbalancer.health-check.path | | +|spring.cloud.loadbalancer.health-check.refetch-instances | false | Indicates whether the instances should be refetched by the HealthCheckServiceInstanceListSupplier. This can be used if the instances can be updated and the underlying delegate does not provide an ongoing flux. +|spring.cloud.loadbalancer.health-check.refetch-instances-interval | 25s | Interval for refetching available service instances. +|spring.cloud.loadbalancer.health-check.repeat-health-check | true | Indicates whether health checks should keep repeating. It might be useful to set it to false if periodically refetching the instances, as every refetch will also trigger a healthcheck. |spring.cloud.loadbalancer.retry.enabled | true | |spring.cloud.loadbalancer.retry.max-retries-on-next-service-instance | 1 | Number of retries to be executed on the next ServiceInstance. A ServiceInstance is chosen before each retry call. |spring.cloud.loadbalancer.retry.max-retries-on-same-service-instance | 0 | Number of retries to be executed on the same ServiceInstance. diff --git a/docs/src/main/asciidoc/spring-cloud-commons.adoc b/docs/src/main/asciidoc/spring-cloud-commons.adoc index 62a0a29b..02da9030 100644 --- a/docs/src/main/asciidoc/spring-cloud-commons.adoc +++ b/docs/src/main/asciidoc/spring-cloud-commons.adoc @@ -953,16 +953,22 @@ TIP: This mechanism is particularly helpful while using the `SimpleDiscoveryClie clients backed by an actual Service Registry, it's not necessary to use, as we already get healthy instances after querying the external ServiceDiscovery. -TIP:: This supplier is also recommended for setups with a small number of instances per service +TIP: This supplier is also recommended for setups with a small number of instances per service in order to avoid retrying calls on a failing instance. +WARNING: If using any of the Service Discovery-backed suppliers, adding this health-check mechanism is usually not necessary, as we retrieve the health state of the instances directly +from the Service Registry. + +TIP: The `HealthCheckServiceInstanceListSupplier` relies on having updated instances provided by a delegate flux. In the rare cases when you want to use a delegate that does not refresh the instances, even though the list of instances may change (such as the `ReactiveDiscoveryClientServiceInstanceListSupplier` provided by us), you can set `spring.cloud.loadbalancer.health-check.refetch-instances` to `true` to have the instance list refreshed by the `HealthCheckServiceInstanceListSupplier`. You can then also adjust the refretch intervals by modifying the value of `spring.cloud.loadbalancer.health-check.refetch-instances-interval` and opt to disable the additional healthcheck repetitions by setting `spring.cloud.loadbalancer.repeat-health-check` to `fasle` as every instances refetch + will also trigger a healthcheck. + `HealthCheckServiceInstanceListSupplier` uses properties prefixed with `spring.cloud.loadbalancer.health-check`. You can set the `initialDelay` and `interval` for the scheduler. You can set the default path for the healthcheck URL by setting the value of the `spring.cloud.loadbalancer.health-check.path.default`. You can also set a specific value for any given service by setting the value of the `spring.cloud.loadbalancer.health-check.path.[SERVICE_ID]`, substituting the `[SERVICE_ID]` with the correct ID of your service. If the path is not set, `/actuator/health` is used by default. -TIP:: If you rely on the default path (`/actuator/health`), make sure you add `spring-boot-starter-actuator` to your collaborator's dependencies, unless you are planning to add such an endpoint on your own. +TIP: If you rely on the default path (`/actuator/health`), make sure you add `spring-boot-starter-actuator` to your collaborator's dependencies, unless you are planning to add such an endpoint on your own. In order to use the health-check scheduler approach, you will have to instantiate a `HealthCheckServiceInstanceListSupplier` bean in a <>. @@ -987,7 +993,7 @@ public class CustomLoadBalancerConfiguration { } ---- -NOTE:: `HealthCheckServiceInstanceListSupplier` has its own caching mechanism based on Reactor Flux `replay()`, therefore, if it's being used, you may want to skip wrapping that supplier with `CachingServiceInstanceListSupplier`. +WARNING: `HealthCheckServiceInstanceListSupplier` has its own caching mechanism based on Reactor Flux `replay()`. Therefore, if it's being used, you may want to skip wrapping that supplier with `CachingServiceInstanceListSupplier`. [[spring-cloud-loadbalancer-starter]] diff --git a/spring-cloud-commons/src/main/java/org/springframework/cloud/client/loadbalancer/reactive/LoadBalancerProperties.java b/spring-cloud-commons/src/main/java/org/springframework/cloud/client/loadbalancer/reactive/LoadBalancerProperties.java index 6c6ea92b..bbf35bb2 100644 --- a/spring-cloud-commons/src/main/java/org/springframework/cloud/client/loadbalancer/reactive/LoadBalancerProperties.java +++ b/spring-cloud-commons/src/main/java/org/springframework/cloud/client/loadbalancer/reactive/LoadBalancerProperties.java @@ -56,8 +56,44 @@ public class LoadBalancerProperties { */ private Duration interval = Duration.ofSeconds(25); + /** + * Interval for refetching available service instances. + */ + private Duration refetchInstancesInterval = Duration.ofSeconds(25); + private Map path = new LinkedCaseInsensitiveMap<>(); + /** + * Indicates whether the instances should be refetched by the + * HealthCheckServiceInstanceListSupplier. This can be used if the + * instances can be updated and the underlying delegate does not provide an + * ongoing flux. + */ + private boolean refetchInstances = false; + + /** + * Indicates whether health checks should keep repeating. It might be useful to + * set it to false if periodically refetching the instances, as every + * refetch will also trigger a healthcheck. + */ + private boolean repeatHealthCheck = true; + + public boolean getRefetchInstances() { + return refetchInstances; + } + + public void setRefetchInstances(boolean refetchInstances) { + this.refetchInstances = refetchInstances; + } + + public boolean getRepeatHealthCheck() { + return repeatHealthCheck; + } + + public void setRepeatHealthCheck(boolean repeatHealthCheck) { + this.repeatHealthCheck = repeatHealthCheck; + } + public int getInitialDelay() { return initialDelay; } @@ -66,6 +102,14 @@ public class LoadBalancerProperties { this.initialDelay = initialDelay; } + public Duration getRefetchInstancesInterval() { + return refetchInstancesInterval; + } + + public void setRefetchInstancesInterval(Duration refetchInstancesInterval) { + this.refetchInstancesInterval = refetchInstancesInterval; + } + public Map getPath() { return path; } 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 85ad174b..a661c714 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 @@ -26,6 +26,7 @@ import org.apache.commons.logging.LogFactory; import reactor.core.Disposable; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; +import reactor.retry.Repeat; import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.InitializingBean; @@ -64,14 +65,19 @@ public class HealthCheckServiceInstanceListSupplier public HealthCheckServiceInstanceListSupplier(ServiceInstanceListSupplier delegate, LoadBalancerProperties.HealthCheck healthCheck, WebClient webClient) { super(delegate); - this.healthCheck = healthCheck; defaultHealthCheckPath = healthCheck.getPath().getOrDefault("default", "/actuator/health"); this.webClient = webClient; - aliveInstancesReplay = Flux.defer(delegate) - .delaySubscription(Duration.ofMillis(healthCheck.getInitialDelay())) + this.healthCheck = healthCheck; + Repeat aliveInstancesReplayRepeat = Repeat + .onlyIf(repeatContext -> this.healthCheck.getRefetchInstances()) + .fixedBackoff(healthCheck.getRefetchInstancesInterval()); + Flux> aliveInstancesFlux = Flux.defer(delegate) .switchMap(serviceInstances -> healthCheckFlux(serviceInstances).map( alive -> Collections.unmodifiableList(new ArrayList<>(alive)))) + .repeatWhen(aliveInstancesReplayRepeat); + aliveInstancesReplay = aliveInstancesFlux + .delaySubscription(Duration.ofMillis(healthCheck.getInitialDelay())) .replay(1).refCount(1); } @@ -86,6 +92,9 @@ public class HealthCheckServiceInstanceListSupplier protected Flux> healthCheckFlux( List instances) { + Repeat healthCheckFluxRepeat = Repeat + .onlyIf(repeatContext -> healthCheck.getRepeatHealthCheck()) + .fixedBackoff(healthCheck.getInterval()); return Flux.defer(() -> { List> checks = new ArrayList<>(instances.size()); for (ServiceInstance instance : instances) { @@ -117,7 +126,7 @@ public class HealthCheckServiceInstanceListSupplier result.add(alive); return result; }).defaultIfEmpty(result); - }).repeatWhen(restart -> restart.delayElements(healthCheck.getInterval())); + }).repeatWhen(healthCheckFluxRepeat); } @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 b04a6828..2318ba38 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 @@ -17,6 +17,7 @@ package org.springframework.cloud.loadbalancer.core; import java.time.Duration; +import java.util.Collections; import java.util.List; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; @@ -47,6 +48,8 @@ import org.springframework.web.bind.annotation.RestController; import org.springframework.web.reactive.function.client.WebClient; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; /** * Tests for {@link HealthCheckServiceInstanceListSupplier}. @@ -79,7 +82,7 @@ class HealthCheckServiceInstanceListSupplierTests { } @AfterEach - void tearDown() throws Exception { + void tearDown() { if (listSupplier != null) { listSupplier.destroy(); listSupplier = null; @@ -140,14 +143,14 @@ class HealthCheckServiceInstanceListSupplierTests { SERVICE_ID, "127.0.0.2", port, false); StepVerifier.withVirtualTime(() -> { - ServiceInstanceListSupplier delegate = Mockito - .mock(ServiceInstanceListSupplier.class); + ServiceInstanceListSupplier delegate = mock( + ServiceInstanceListSupplier.class); Mockito.when(delegate.getServiceId()).thenReturn(SERVICE_ID); Mockito.when(delegate.get()).thenReturn( Flux.just(Lists.list(serviceInstance1, serviceInstance2))); - HealthCheckServiceInstanceListSupplier mock = Mockito - .mock(HealthCheckServiceInstanceListSupplier.class); + HealthCheckServiceInstanceListSupplier mock = mock( + HealthCheckServiceInstanceListSupplier.class); Mockito.doReturn(Mono.just(true)).when(mock).isAlive(serviceInstance1); Mockito.doReturn(Mono.just(false)).when(mock).isAlive(serviceInstance2); @@ -176,14 +179,14 @@ class HealthCheckServiceInstanceListSupplierTests { SERVICE_ID, "127.0.0.2", port, false); StepVerifier.withVirtualTime(() -> { - ServiceInstanceListSupplier delegate = Mockito - .mock(ServiceInstanceListSupplier.class); + ServiceInstanceListSupplier delegate = mock( + ServiceInstanceListSupplier.class); Mockito.when(delegate.getServiceId()).thenReturn(SERVICE_ID); Mockito.when(delegate.get()).thenReturn( Flux.just(Lists.list(serviceInstance1, serviceInstance2))); - HealthCheckServiceInstanceListSupplier mock = Mockito - .mock(HealthCheckServiceInstanceListSupplier.class); + HealthCheckServiceInstanceListSupplier mock = mock( + HealthCheckServiceInstanceListSupplier.class); Mockito.doReturn(Mono.just(true)).when(mock).isAlive(serviceInstance1); Mockito.doReturn(Mono.just(true)).when(mock).isAlive(serviceInstance2); @@ -213,14 +216,14 @@ class HealthCheckServiceInstanceListSupplierTests { SERVICE_ID, "127.0.0.2", port, false); StepVerifier.withVirtualTime(() -> { - ServiceInstanceListSupplier delegate = Mockito - .mock(ServiceInstanceListSupplier.class); + ServiceInstanceListSupplier delegate = mock( + ServiceInstanceListSupplier.class); Mockito.when(delegate.getServiceId()).thenReturn(SERVICE_ID); Mockito.when(delegate.get()).thenReturn( Flux.just(Lists.list(serviceInstance1, serviceInstance2))); - HealthCheckServiceInstanceListSupplier mock = Mockito - .mock(HealthCheckServiceInstanceListSupplier.class); + HealthCheckServiceInstanceListSupplier mock = mock( + HealthCheckServiceInstanceListSupplier.class); Mockito.doReturn(Mono.just(true)).when(mock).isAlive(serviceInstance1); Mockito.doReturn(Mono.error(new RuntimeException("boom"))).when(mock) .isAlive(serviceInstance2); @@ -250,8 +253,8 @@ class HealthCheckServiceInstanceListSupplierTests { SERVICE_ID, "127.0.0.2", port, false); StepVerifier.withVirtualTime(() -> { - ServiceInstanceListSupplier delegate = Mockito - .mock(ServiceInstanceListSupplier.class); + ServiceInstanceListSupplier delegate = mock( + ServiceInstanceListSupplier.class); Mockito.when(delegate.getServiceId()).thenReturn(SERVICE_ID); Mockito.when(delegate.get()).thenReturn( Flux.just(Lists.list(serviceInstance1, serviceInstance2))); @@ -282,8 +285,8 @@ class HealthCheckServiceInstanceListSupplierTests { SERVICE_ID, "127.0.0.1", port, false); StepVerifier.withVirtualTime(() -> { - ServiceInstanceListSupplier delegate = Mockito - .mock(ServiceInstanceListSupplier.class); + ServiceInstanceListSupplier delegate = mock( + ServiceInstanceListSupplier.class); Mockito.when(delegate.getServiceId()).thenReturn(SERVICE_ID); Mockito.when(delegate.get()) .thenReturn(Flux.just(Lists.list(serviceInstance1))); @@ -314,14 +317,14 @@ class HealthCheckServiceInstanceListSupplierTests { SERVICE_ID, "127.0.0.2", port, false); StepVerifier.withVirtualTime(() -> { - ServiceInstanceListSupplier delegate = Mockito - .mock(ServiceInstanceListSupplier.class); + ServiceInstanceListSupplier delegate = mock( + ServiceInstanceListSupplier.class); Mockito.when(delegate.getServiceId()).thenReturn(SERVICE_ID); Mockito.when(delegate.get()).thenReturn( Flux.just(Lists.list(serviceInstance1, serviceInstance2))); - HealthCheckServiceInstanceListSupplier mock = Mockito - .mock(HealthCheckServiceInstanceListSupplier.class); + HealthCheckServiceInstanceListSupplier mock = mock( + HealthCheckServiceInstanceListSupplier.class); Mockito.doReturn(Mono.just(false), Mono.just(true)).when(mock) .isAlive(serviceInstance1); Mockito.doReturn(Mono.error(new RuntimeException("boom"))).when(mock) @@ -352,14 +355,14 @@ class HealthCheckServiceInstanceListSupplierTests { SERVICE_ID, "127.0.0.1", port, false); StepVerifier.withVirtualTime(() -> { - ServiceInstanceListSupplier delegate = Mockito - .mock(ServiceInstanceListSupplier.class); + ServiceInstanceListSupplier delegate = mock( + ServiceInstanceListSupplier.class); Mockito.when(delegate.getServiceId()).thenReturn(SERVICE_ID); Mockito.when(delegate.get()) .thenReturn(Flux.just(Lists.list(serviceInstance1))); - HealthCheckServiceInstanceListSupplier mock = Mockito - .mock(HealthCheckServiceInstanceListSupplier.class); + HealthCheckServiceInstanceListSupplier mock = mock( + HealthCheckServiceInstanceListSupplier.class); Mockito.when(mock.isAlive(serviceInstance1)).thenReturn(Mono.never(), Mono.just(true)); @@ -391,8 +394,8 @@ class HealthCheckServiceInstanceListSupplierTests { SERVICE_ID, "127.0.0.2", port, false); StepVerifier.withVirtualTime(() -> { - ServiceInstanceListSupplier delegate = Mockito - .mock(ServiceInstanceListSupplier.class); + ServiceInstanceListSupplier delegate = mock( + ServiceInstanceListSupplier.class); Mockito.when(delegate.getServiceId()).thenReturn(SERVICE_ID); Flux> instances = Flux .just(Lists.list(serviceInstance1)) @@ -421,6 +424,39 @@ class HealthCheckServiceInstanceListSupplierTests { .verify(VERIFY_TIMEOUT); } + @Test + void shouldRefetchInstances() { + healthCheck.setInitialDelay(1000); + healthCheck.setRepeatHealthCheck(false); + 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); @@ -430,8 +466,8 @@ class HealthCheckServiceInstanceListSupplierTests { AtomicInteger emitCounter = new AtomicInteger(); StepVerifier.withVirtualTime(() -> { - ServiceInstanceListSupplier delegate = Mockito - .mock(ServiceInstanceListSupplier.class); + ServiceInstanceListSupplier delegate = mock( + ServiceInstanceListSupplier.class); Mockito.when(delegate.getServiceId()).thenReturn(SERVICE_ID); Mockito.when(delegate.get()) .thenReturn(Flux.just(Lists.list(serviceInstance1))); @@ -468,8 +504,7 @@ class HealthCheckServiceInstanceListSupplierTests { final AtomicInteger instancesCanceled = new AtomicInteger(); final AtomicBoolean subscribed = new AtomicBoolean(); - ServiceInstanceListSupplier delegate = Mockito - .mock(ServiceInstanceListSupplier.class); + ServiceInstanceListSupplier delegate = mock(ServiceInstanceListSupplier.class); Mockito.when(delegate.get()) .thenReturn(Flux.>never() .doOnSubscribe(subscription -> subscribed.set(true))