Fix bug of service instance refetching when repeating health check is activated (#904)

Fixes gh-899

(cherry picked from commit 3d257b43003a02b64c4f7fd97680e7981d0b52f4)
This commit is contained in:
Роман Чигвинцев
2021-02-17 17:35:39 +03:00
committed by GitHub
parent 3aff1bea0d
commit 95a22796eb
2 changed files with 37 additions and 2 deletions

View File

@@ -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<List<ServiceInstance>> 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);

View File

@@ -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<Boolean> 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);