Merge remote-tracking branch 'origin/2.2.x'

This commit is contained in:
Olga Maciaszek-Sharma
2020-03-17 18:44:53 +01:00
2 changed files with 96 additions and 73 deletions

View File

@@ -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

View File

@@ -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<Boolean> 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<List<ServiceInstance>> instances = Flux.just(Lists.list(si1))
.concatWith(Flux.just(Lists.list(si1, si2))
Flux<List<ServiceInstance>> 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);