Fix after merge.
This commit is contained in:
@@ -22,6 +22,7 @@ import java.util.List;
|
||||
import reactor.core.publisher.Flux;
|
||||
|
||||
import org.springframework.cloud.client.ServiceInstance;
|
||||
import org.springframework.cloud.client.loadbalancer.Request;
|
||||
|
||||
/**
|
||||
* A no-op implementation of {@link ServiceInstanceListSupplier}.
|
||||
@@ -40,4 +41,9 @@ public class NoopServiceInstanceListSupplier implements ServiceInstanceListSuppl
|
||||
return Flux.defer(() -> Flux.just(Collections.emptyList()));
|
||||
}
|
||||
|
||||
@Override
|
||||
public Flux<List<ServiceInstance>> get(Request request) {
|
||||
return Flux.defer(() -> Flux.just(Collections.emptyList()));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -25,10 +25,10 @@ import reactor.core.publisher.Mono;
|
||||
|
||||
import org.springframework.beans.factory.ObjectProvider;
|
||||
import org.springframework.cloud.client.ServiceInstance;
|
||||
import org.springframework.cloud.client.loadbalancer.reactive.DefaultResponse;
|
||||
import org.springframework.cloud.client.loadbalancer.reactive.EmptyResponse;
|
||||
import org.springframework.cloud.client.loadbalancer.reactive.Request;
|
||||
import org.springframework.cloud.client.loadbalancer.reactive.Response;
|
||||
import org.springframework.cloud.client.loadbalancer.DefaultResponse;
|
||||
import org.springframework.cloud.client.loadbalancer.EmptyResponse;
|
||||
import org.springframework.cloud.client.loadbalancer.Request;
|
||||
import org.springframework.cloud.client.loadbalancer.Response;
|
||||
|
||||
/**
|
||||
* A random-based implementation of {@link ReactorServiceInstanceLoadBalancer}.
|
||||
@@ -42,9 +42,6 @@ public class RandomLoadBalancer implements ReactorServiceInstanceLoadBalancer {
|
||||
|
||||
private final String serviceId;
|
||||
|
||||
@Deprecated
|
||||
private ObjectProvider<ServiceInstanceSupplier> serviceInstanceSupplier;
|
||||
|
||||
private ObjectProvider<ServiceInstanceListSupplier> serviceInstanceListSupplierProvider;
|
||||
|
||||
/**
|
||||
@@ -52,59 +49,31 @@ public class RandomLoadBalancer implements ReactorServiceInstanceLoadBalancer {
|
||||
* {@link ServiceInstanceListSupplier} that will be used to get available instances
|
||||
* @param serviceId id of the service for which to choose an instance
|
||||
*/
|
||||
public RandomLoadBalancer(
|
||||
ObjectProvider<ServiceInstanceListSupplier> serviceInstanceListSupplierProvider,
|
||||
public RandomLoadBalancer(ObjectProvider<ServiceInstanceListSupplier> serviceInstanceListSupplierProvider,
|
||||
String serviceId) {
|
||||
this.serviceId = serviceId;
|
||||
this.serviceInstanceListSupplierProvider = serviceInstanceListSupplierProvider;
|
||||
}
|
||||
|
||||
/**
|
||||
* @param serviceId id of the service for which to choose an instance
|
||||
* @param serviceInstanceSupplier a provider of {@link ServiceInstanceSupplier} that
|
||||
* will be used to get available instances
|
||||
* @deprecated Use {@link #RandomLoadBalancer(ObjectProvider, String)}} instead.
|
||||
*/
|
||||
@Deprecated
|
||||
public RandomLoadBalancer(String serviceId,
|
||||
ObjectProvider<ServiceInstanceSupplier> serviceInstanceSupplier) {
|
||||
this.serviceId = serviceId;
|
||||
this.serviceInstanceSupplier = serviceInstanceSupplier;
|
||||
}
|
||||
|
||||
@SuppressWarnings("rawtypes")
|
||||
@Override
|
||||
public Mono<org.springframework.cloud.client.loadbalancer.reactive.Response<ServiceInstance>> choose(
|
||||
Request request) {
|
||||
// TODO: move supplier to Request?
|
||||
// Temporary conditional logic till deprecated members are removed.
|
||||
if (serviceInstanceListSupplierProvider != null) {
|
||||
ServiceInstanceListSupplier supplier = serviceInstanceListSupplierProvider
|
||||
.getIfAvailable(NoopServiceInstanceListSupplier::new);
|
||||
return supplier.get().next()
|
||||
.map(serviceInstances -> processInstanceResponse(supplier,
|
||||
serviceInstances));
|
||||
}
|
||||
ServiceInstanceSupplier supplier = this.serviceInstanceSupplier
|
||||
.getIfAvailable(NoopServiceInstanceSupplier::new);
|
||||
return supplier.get().collectList().map(this::getInstanceResponse);
|
||||
public Mono<Response<ServiceInstance>> choose(Request request) {
|
||||
ServiceInstanceListSupplier supplier = serviceInstanceListSupplierProvider
|
||||
.getIfAvailable(NoopServiceInstanceListSupplier::new);
|
||||
return supplier.get(request).next()
|
||||
.map(serviceInstances -> processInstanceResponse(supplier, serviceInstances));
|
||||
}
|
||||
|
||||
private org.springframework.cloud.client.loadbalancer.reactive.Response<ServiceInstance> processInstanceResponse(
|
||||
ServiceInstanceListSupplier supplier,
|
||||
private Response<ServiceInstance> processInstanceResponse(ServiceInstanceListSupplier supplier,
|
||||
List<ServiceInstance> serviceInstances) {
|
||||
org.springframework.cloud.client.loadbalancer.reactive.Response<ServiceInstance> serviceInstanceResponse = getInstanceResponse(
|
||||
serviceInstances);
|
||||
if (supplier instanceof SelectedInstanceCallback
|
||||
&& serviceInstanceResponse.hasServer()) {
|
||||
((SelectedInstanceCallback) supplier)
|
||||
.selectedServiceInstance(serviceInstanceResponse.getServer());
|
||||
Response<ServiceInstance> serviceInstanceResponse = getInstanceResponse(serviceInstances);
|
||||
if (supplier instanceof SelectedInstanceCallback && serviceInstanceResponse.hasServer()) {
|
||||
((SelectedInstanceCallback) supplier).selectedServiceInstance(serviceInstanceResponse.getServer());
|
||||
}
|
||||
return serviceInstanceResponse;
|
||||
}
|
||||
|
||||
private Response<ServiceInstance> getInstanceResponse(
|
||||
List<ServiceInstance> instances) {
|
||||
private Response<ServiceInstance> getInstanceResponse(List<ServiceInstance> instances) {
|
||||
if (instances.isEmpty()) {
|
||||
if (log.isWarnEnabled()) {
|
||||
log.warn("No servers available for service: " + serviceId);
|
||||
|
||||
@@ -91,7 +91,7 @@ public class RoundRobinLoadBalancer implements ReactorServiceInstanceLoadBalance
|
||||
return serviceInstanceResponse;
|
||||
}
|
||||
|
||||
Response<ServiceInstance> getInstanceResponse(List<ServiceInstance> instances) {
|
||||
private Response<ServiceInstance> getInstanceResponse(List<ServiceInstance> instances) {
|
||||
if (instances.isEmpty()) {
|
||||
if (log.isWarnEnabled()) {
|
||||
log.warn("No servers available for service: " + serviceId);
|
||||
|
||||
@@ -28,7 +28,9 @@ import org.springframework.cloud.client.loadbalancer.Response;
|
||||
import org.springframework.cloud.loadbalancer.support.SimpleObjectProvider;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.verify;
|
||||
import static org.mockito.Mockito.when;
|
||||
|
||||
/**
|
||||
@@ -44,12 +46,9 @@ class RandomLoadBalancerTests {
|
||||
|
||||
@Test
|
||||
void shouldReturnOneServiceInstance() {
|
||||
DiscoveryClientServiceInstanceListSupplier supplier = mock(
|
||||
DiscoveryClientServiceInstanceListSupplier.class);
|
||||
when(supplier.get()).thenReturn(
|
||||
Flux.just(Arrays.asList(serviceInstance, new DefaultServiceInstance())));
|
||||
loadBalancer = new RandomLoadBalancer(new SimpleObjectProvider<>(supplier),
|
||||
"test");
|
||||
DiscoveryClientServiceInstanceListSupplier supplier = mock(DiscoveryClientServiceInstanceListSupplier.class);
|
||||
when(supplier.get(any())).thenReturn(Flux.just(Arrays.asList(serviceInstance, new DefaultServiceInstance())));
|
||||
loadBalancer = new RandomLoadBalancer(new SimpleObjectProvider<>(supplier), "test");
|
||||
|
||||
Response<ServiceInstance> response = loadBalancer.choose().block();
|
||||
|
||||
@@ -67,15 +66,27 @@ class RandomLoadBalancerTests {
|
||||
|
||||
@Test
|
||||
void shouldReturnEmptyResponseWhenNoInstancesAvailable() {
|
||||
DiscoveryClientServiceInstanceListSupplier supplier = mock(
|
||||
DiscoveryClientServiceInstanceListSupplier.class);
|
||||
when(supplier.get()).thenReturn(Flux.just(Collections.emptyList()));
|
||||
loadBalancer = new RandomLoadBalancer(new SimpleObjectProvider<>(supplier),
|
||||
"test");
|
||||
DiscoveryClientServiceInstanceListSupplier supplier = mock(DiscoveryClientServiceInstanceListSupplier.class);
|
||||
when(supplier.get(any())).thenReturn(Flux.just(Collections.emptyList()));
|
||||
loadBalancer = new RandomLoadBalancer(new SimpleObjectProvider<>(supplier), "test");
|
||||
|
||||
Response<ServiceInstance> response = loadBalancer.choose().block();
|
||||
|
||||
assertThat(response.hasServer()).isFalse();
|
||||
}
|
||||
|
||||
@Test
|
||||
void shouldTriggerSelectedInstanceCallback() {
|
||||
SameInstancePreferenceServiceInstanceListSupplier supplier = mock(
|
||||
SameInstancePreferenceServiceInstanceListSupplier.class);
|
||||
when(supplier.get(any())).thenReturn(Flux.just(Collections.singletonList(serviceInstance)));
|
||||
loadBalancer = new RandomLoadBalancer(new SimpleObjectProvider<>(supplier), "test");
|
||||
|
||||
Response<ServiceInstance> response = loadBalancer.choose().block();
|
||||
|
||||
assertThat(response.hasServer()).isTrue();
|
||||
assertThat(response.getServer()).isEqualTo(serviceInstance);
|
||||
verify((SelectedInstanceCallback) supplier).selectedServiceInstance(serviceInstance);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user