From b35165ed481681674532a5d70c8fab916440db7f Mon Sep 17 00:00:00 2001 From: Olga Maciaszek-Sharma Date: Fri, 18 Dec 2020 11:15:12 -0600 Subject: [PATCH] Add LB stats support. (#2086) * Add LB stats support. * Adjust to changes in commons. --- .../ReactiveLoadBalancerClientFilter.java | 22 +++++++++------- ...ReactiveLoadBalancerClientFilterTests.java | 25 +++++++++++++------ 2 files changed, 30 insertions(+), 17 deletions(-) diff --git a/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/ReactiveLoadBalancerClientFilter.java b/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/ReactiveLoadBalancerClientFilter.java index 8b2e6c54..725859f9 100644 --- a/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/ReactiveLoadBalancerClientFilter.java +++ b/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/ReactiveLoadBalancerClientFilter.java @@ -31,6 +31,7 @@ import org.springframework.cloud.client.loadbalancer.LoadBalancerLifecycle; import org.springframework.cloud.client.loadbalancer.LoadBalancerLifecycleValidator; import org.springframework.cloud.client.loadbalancer.LoadBalancerProperties; import org.springframework.cloud.client.loadbalancer.LoadBalancerUriTools; +import org.springframework.cloud.client.loadbalancer.Request; import org.springframework.cloud.client.loadbalancer.RequestData; import org.springframework.cloud.client.loadbalancer.RequestDataContext; import org.springframework.cloud.client.loadbalancer.Response; @@ -104,11 +105,13 @@ public class ReactiveLoadBalancerClientFilter implements GlobalFilter, Ordered { Set supportedLifecycleProcessors = LoadBalancerLifecycleValidator .getSupportedLifecycleProcessors(clientFactory.getInstances(serviceId, LoadBalancerLifecycle.class), RequestDataContext.class, ResponseData.class, ServiceInstance.class); - return choose(exchange, serviceId, supportedLifecycleProcessors).doOnNext(response -> { + DefaultRequest lbRequest = new DefaultRequest<>(new RequestDataContext( + new RequestData(exchange.getRequest()), getHint(serviceId, loadBalancerProperties.getHint()))); + return choose(lbRequest, serviceId, supportedLifecycleProcessors).doOnNext(response -> { if (!response.hasServer()) { supportedLifecycleProcessors.forEach(lifecycle -> lifecycle - .onComplete(new CompletionContext<>(CompletionContext.Status.DISCARD, response))); + .onComplete(new CompletionContext<>(CompletionContext.Status.DISCARD, lbRequest, response))); throw NotFoundException.create(properties.isUse404(), "Unable to find instance for " + url.getHost()); } @@ -133,12 +136,15 @@ public class ReactiveLoadBalancerClientFilter implements GlobalFilter, Ordered { } exchange.getAttributes().put(GATEWAY_REQUEST_URL_ATTR, requestUrl); exchange.getAttributes().put(GATEWAY_LOADBALANCER_RESPONSE_ATTR, response); + supportedLifecycleProcessors.forEach(lifecycle -> lifecycle.onStartRequest(lbRequest, response)); }).then(chain.filter(exchange)) - .doOnError(throwable -> supportedLifecycleProcessors.forEach(lifecycle -> lifecycle.onComplete( - new CompletionContext(CompletionContext.Status.FAILED, throwable, + .doOnError(throwable -> supportedLifecycleProcessors.forEach(lifecycle -> lifecycle + .onComplete(new CompletionContext( + CompletionContext.Status.FAILED, throwable, lbRequest, exchange.getAttribute(GATEWAY_LOADBALANCER_RESPONSE_ATTR))))) - .doOnSuccess(aVoid -> supportedLifecycleProcessors.forEach( - lifecycle -> lifecycle.onComplete(new CompletionContext<>(CompletionContext.Status.SUCCESS, + .doOnSuccess(aVoid -> supportedLifecycleProcessors.forEach(lifecycle -> lifecycle + .onComplete(new CompletionContext( + CompletionContext.Status.SUCCESS, lbRequest, exchange.getAttribute(GATEWAY_LOADBALANCER_RESPONSE_ATTR), new ResponseData(exchange.getResponse(), new RequestData(exchange.getRequest())))))); } @@ -147,15 +153,13 @@ public class ReactiveLoadBalancerClientFilter implements GlobalFilter, Ordered { return LoadBalancerUriTools.reconstructURI(serviceInstance, original); } - private Mono> choose(ServerWebExchange exchange, String serviceId, + private Mono> choose(Request lbRequest, String serviceId, Set supportedLifecycleProcessors) { ReactorLoadBalancer loadBalancer = this.clientFactory.getInstance(serviceId, ReactorServiceInstanceLoadBalancer.class); if (loadBalancer == null) { throw new NotFoundException("No loadbalancer available for " + serviceId); } - DefaultRequest lbRequest = new DefaultRequest<>(new RequestDataContext( - new RequestData(exchange.getRequest()), getHint(serviceId, loadBalancerProperties.getHint()))); supportedLifecycleProcessors.forEach(lifecycle -> lifecycle.onStart(lbRequest)); return loadBalancer.choose(lbRequest); } diff --git a/spring-cloud-gateway-server/src/test/java/org/springframework/cloud/gateway/filter/ReactiveLoadBalancerClientFilterTests.java b/spring-cloud-gateway-server/src/test/java/org/springframework/cloud/gateway/filter/ReactiveLoadBalancerClientFilterTests.java index 7aa81317..7d31eb71 100644 --- a/spring-cloud-gateway-server/src/test/java/org/springframework/cloud/gateway/filter/ReactiveLoadBalancerClientFilterTests.java +++ b/spring-cloud-gateway-server/src/test/java/org/springframework/cloud/gateway/filter/ReactiveLoadBalancerClientFilterTests.java @@ -331,9 +331,12 @@ class ReactiveLoadBalancerClientFilterTests { filter.filter(serverWebExchange, chain).subscribe(); verify(lifecycleProcessor).onStart(any(Request.class)); - verify(lifecycleProcessor).onComplete( - argThat(completionContext -> CompletionContext.Status.SUCCESS.equals(completionContext.status()) - && completionContext.getLoadBalancerResponse().getServer().equals(serviceInstance))); + verify(lifecycleProcessor).onStartRequest(any(Request.class), any(Response.class)); + verify(lifecycleProcessor).onComplete(argThat(completionContext -> CompletionContext.Status.SUCCESS + .equals(completionContext.status()) + && completionContext.getLoadBalancerResponse().getServer().equals(serviceInstance) + && HttpMethod.GET.equals( + ((RequestDataContext) completionContext.getLoadBalancerRequest().getContext()).method()))); } @SuppressWarnings({ "unchecked", "rawtypes" }) @@ -346,22 +349,28 @@ class ReactiveLoadBalancerClientFilterTests { filter.filter(serverWebExchange, chain).subscribe(); verify(lifecycleProcessor).onStart(any(Request.class)); - verify(lifecycleProcessor).onComplete( - argThat(completionContext -> CompletionContext.Status.DISCARD.equals(completionContext.status()))); + verify(lifecycleProcessor).onComplete(argThat(completionContext -> CompletionContext.Status.DISCARD + .equals(completionContext.status()) + && HttpMethod.GET.equals( + ((RequestDataContext) completionContext.getLoadBalancerRequest().getContext()).method()))); } @SuppressWarnings({ "unchecked", "rawtypes" }) @Test void loadBalancerLifecycleCallbacksExecutedForFailed() { LoadBalancerLifecycle lifecycleProcessor = mock(LoadBalancerLifecycle.class); - ServiceInstance serviceInstance = null; + ServiceInstance serviceInstance = new DefaultServiceInstance("myservice1", "myservice", "localhost", 8080, + false); ServerWebExchange serverWebExchange = mockExchange(serviceInstance, lifecycleProcessor, true); filter.filter(serverWebExchange, chain).subscribe(); verify(lifecycleProcessor).onStart(any(Request.class)); - verify(lifecycleProcessor).onComplete( - argThat(completionContext -> CompletionContext.Status.FAILED.equals(completionContext.status()))); + verify(lifecycleProcessor).onStartRequest(any(Request.class), any(Response.class)); + verify(lifecycleProcessor).onComplete(argThat(completionContext -> CompletionContext.Status.FAILED + .equals(completionContext.status()) + && HttpMethod.GET.equals( + ((RequestDataContext) completionContext.getLoadBalancerRequest().getContext()).method()))); } @SuppressWarnings({ "rawtypes", "unchecked" })