From ba68a74dfaed08b2e8c50ef2e2a03fef8fb57305 Mon Sep 17 00:00:00 2001 From: Spencer Gibb Date: Mon, 16 Jan 2017 22:05:23 -0700 Subject: [PATCH] Move to reactor netty httpclient from WebClient --- .../config/GatewayAutoConfiguration.java | 28 ++++---- .../filter/RouteToRequestUrlFilter.java | 6 +- .../gateway/filter/WriteResponseFilter.java | 15 ++-- .../gateway/handler/RoutingWebHandler.java | 69 +++++++++++++++++++ .../handler/WebClientRoutingWebHandler.java | 51 -------------- 5 files changed, 100 insertions(+), 69 deletions(-) create mode 100644 src/main/java/org/springframework/cloud/gateway/handler/RoutingWebHandler.java delete mode 100644 src/main/java/org/springframework/cloud/gateway/handler/WebClientRoutingWebHandler.java diff --git a/src/main/java/org/springframework/cloud/gateway/config/GatewayAutoConfiguration.java b/src/main/java/org/springframework/cloud/gateway/config/GatewayAutoConfiguration.java index f6bef8aa..593fa02e 100644 --- a/src/main/java/org/springframework/cloud/gateway/config/GatewayAutoConfiguration.java +++ b/src/main/java/org/springframework/cloud/gateway/config/GatewayAutoConfiguration.java @@ -25,7 +25,7 @@ import org.springframework.cloud.gateway.filter.route.SetResponseHeaderRouteFilt import org.springframework.cloud.gateway.filter.route.SetStatusRouteFilter; import org.springframework.cloud.gateway.handler.FilteringWebHandler; import org.springframework.cloud.gateway.handler.RoutePredicateHandlerMapping; -import org.springframework.cloud.gateway.handler.WebClientRoutingWebHandler; +import org.springframework.cloud.gateway.handler.RoutingWebHandler; import org.springframework.cloud.gateway.handler.predicate.CookieRoutePredicate; import org.springframework.cloud.gateway.handler.predicate.HeaderRoutePredicate; import org.springframework.cloud.gateway.handler.predicate.HostRoutePredicate; @@ -35,8 +35,8 @@ import org.springframework.cloud.gateway.handler.predicate.RoutePredicate; import org.springframework.cloud.gateway.handler.predicate.UrlRoutePredicate; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; -import org.springframework.http.client.reactive.ReactorClientHttpConnector; -import org.springframework.web.reactive.function.client.WebClient; + +import reactor.ipc.netty.http.client.HttpClient; /** * @author Spencer Gibb @@ -47,8 +47,11 @@ public class GatewayAutoConfiguration { @Bean @ConditionalOnMissingBean - public WebClient webClient() { - return WebClient.builder(new ReactorClientHttpConnector(opts -> opts.disablePool())).build(); + public HttpClient httpClient() { + return HttpClient.create(opts -> { + //opts.poolResources(PoolResources.elastic("proxy")); + opts.disablePool(); + }); } @Bean @@ -63,21 +66,21 @@ public class GatewayAutoConfiguration { } @Bean - public WebClientRoutingWebHandler webClientRoutingWebHandler(WebClient webClient) { - return new WebClientRoutingWebHandler(webClient); + public RoutingWebHandler routingWebHandler(HttpClient httpClient) { + return new RoutingWebHandler(httpClient); } @Bean - public FilteringWebHandler filteringWebHandler(WebClientRoutingWebHandler webHandler, - List globalFilters, - Map routeFilters) { + public FilteringWebHandler filteringWebHandler(RoutingWebHandler webHandler, + List globalFilters, + Map routeFilters) { return new FilteringWebHandler(webHandler, globalFilters, routeFilters); } @Bean public RoutePredicateHandlerMapping routePredicateHandlerMapping(FilteringWebHandler webHandler, - Map predicates, - RouteReader routeReader) { + Map predicates, + RouteReader routeReader) { return new RoutePredicateHandlerMapping(webHandler, predicates, routeReader); } @@ -188,3 +191,4 @@ public class GatewayAutoConfiguration { } } + diff --git a/src/main/java/org/springframework/cloud/gateway/filter/RouteToRequestUrlFilter.java b/src/main/java/org/springframework/cloud/gateway/filter/RouteToRequestUrlFilter.java index 10f641fe..815cc992 100644 --- a/src/main/java/org/springframework/cloud/gateway/filter/RouteToRequestUrlFilter.java +++ b/src/main/java/org/springframework/cloud/gateway/filter/RouteToRequestUrlFilter.java @@ -11,6 +11,8 @@ import org.springframework.web.server.ServerWebExchange; import org.springframework.web.server.WebFilterChain; import org.springframework.web.util.UriComponentsBuilder; +import static org.springframework.cloud.gateway.support.ServerWebExchangeUtils.GATEWAY_REQUEST_URL_ATTR; +import static org.springframework.cloud.gateway.support.ServerWebExchangeUtils.GATEWAY_ROUTE_ATTR; import static org.springframework.cloud.gateway.support.ServerWebExchangeUtils.getAttribute; import reactor.core.publisher.Mono; @@ -30,7 +32,7 @@ public class RouteToRequestUrlFilter implements GlobalFilter, Ordered { @Override public Mono filter(ServerWebExchange exchange, WebFilterChain chain) { - Route route = getAttribute(exchange, ServerWebExchangeUtils.GATEWAY_ROUTE_ATTR, Route.class); + Route route = getAttribute(exchange, GATEWAY_ROUTE_ATTR, Route.class); if (route == null) { return chain.filter(exchange); } @@ -39,7 +41,7 @@ public class RouteToRequestUrlFilter implements GlobalFilter, Ordered { .uri(route.getUri()) .build(true) .toUri(); - exchange.getAttributes().put(ServerWebExchangeUtils.GATEWAY_REQUEST_URL_ATTR, requestUrl); + exchange.getAttributes().put(GATEWAY_REQUEST_URL_ATTR, requestUrl); return chain.filter(exchange); } diff --git a/src/main/java/org/springframework/cloud/gateway/filter/WriteResponseFilter.java b/src/main/java/org/springframework/cloud/gateway/filter/WriteResponseFilter.java index 3882e5cc..39ab7937 100644 --- a/src/main/java/org/springframework/cloud/gateway/filter/WriteResponseFilter.java +++ b/src/main/java/org/springframework/cloud/gateway/filter/WriteResponseFilter.java @@ -3,9 +3,9 @@ package org.springframework.cloud.gateway.filter; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.springframework.core.Ordered; -import org.springframework.core.io.buffer.DataBuffer; +import org.springframework.core.io.buffer.NettyDataBuffer; +import org.springframework.core.io.buffer.NettyDataBufferFactory; import org.springframework.http.server.reactive.ServerHttpResponse; -import org.springframework.web.reactive.function.client.ClientResponse; import org.springframework.web.server.ServerWebExchange; import org.springframework.web.server.WebFilterChain; @@ -15,6 +15,7 @@ import static org.springframework.cloud.gateway.support.ServerWebExchangeUtils.g import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; +import reactor.ipc.netty.http.client.HttpClientResponse; /** * @author Spencer Gibb @@ -35,7 +36,7 @@ public class WriteResponseFilter implements GlobalFilter, Ordered { // NOTICE: nothing in "pre" filter stage as CLIENT_RESPONSE_ATTR is not added // until the WebHandler is run return chain.filter(exchange).then(() -> { - ClientResponse clientResponse = getAttribute(exchange, CLIENT_RESPONSE_ATTR, ClientResponse.class); + HttpClientResponse clientResponse = getAttribute(exchange, CLIENT_RESPONSE_ATTR, HttpClientResponse.class); if (clientResponse == null) { return Mono.empty(); } @@ -48,7 +49,13 @@ public class WriteResponseFilter implements GlobalFilter, Ordered { return Mono.empty(); }); - Flux body = clientResponse.body((inputMessage, context) -> inputMessage.getBody()); + NettyDataBufferFactory factory = (NettyDataBufferFactory) response.bufferFactory(); + //TODO: what if it's not netty + + final Flux body = clientResponse.receive() + .retain() //TODO: needed? + .map(factory::wrap); + return response.writeWith(body); }); } diff --git a/src/main/java/org/springframework/cloud/gateway/handler/RoutingWebHandler.java b/src/main/java/org/springframework/cloud/gateway/handler/RoutingWebHandler.java new file mode 100644 index 00000000..2d02743b --- /dev/null +++ b/src/main/java/org/springframework/cloud/gateway/handler/RoutingWebHandler.java @@ -0,0 +1,69 @@ +package org.springframework.cloud.gateway.handler; + +import java.net.URI; +import java.util.Optional; + +import org.springframework.core.io.buffer.DataBuffer; +import org.springframework.http.HttpHeaders; +import org.springframework.http.HttpStatus; +import org.springframework.http.server.reactive.ServerHttpRequest; +import org.springframework.http.server.reactive.ServerHttpResponse; +import org.springframework.web.server.ServerWebExchange; +import org.springframework.web.server.WebHandler; + +import static org.springframework.cloud.gateway.support.ServerWebExchangeUtils.CLIENT_RESPONSE_ATTR; +import static org.springframework.cloud.gateway.support.ServerWebExchangeUtils.GATEWAY_REQUEST_URL_ATTR; + +import io.netty.buffer.Unpooled; +import io.netty.handler.codec.http.DefaultHttpHeaders; +import io.netty.handler.codec.http.HttpMethod; +import reactor.core.publisher.Mono; +import reactor.ipc.netty.NettyPipeline; +import reactor.ipc.netty.http.client.HttpClient; + +/** + * @author Spencer Gibb + */ +public class RoutingWebHandler implements WebHandler { + + private final HttpClient httpClient; + + public RoutingWebHandler(HttpClient httpClient) { + this.httpClient = httpClient; + } + + @Override + public Mono handle(ServerWebExchange exchange) { + Optional requestUrl = exchange.getAttribute(GATEWAY_REQUEST_URL_ATTR); + ServerHttpRequest request = exchange.getRequest(); + + final HttpMethod method = HttpMethod.valueOf(request.getMethod().toString()); + final String url = requestUrl.get().toString(); + + final DefaultHttpHeaders httpHeaders = new DefaultHttpHeaders(); + request.getHeaders().forEach(httpHeaders::set); + + return this.httpClient.request(method, url, req -> + req.options(NettyPipeline.SendOptions::flushOnEach) + .headers(httpHeaders) + //.sendHeaders() + .send(request.getBody() + .map(DataBuffer::asByteBuffer) + .map(Unpooled::wrappedBuffer))) + .then(res -> { + // Defer committing the response until all route filters have run + // Put client response as ServerWebExchange attribute and write response later WriteResponseFilter + exchange.getAttributes().put(CLIENT_RESPONSE_ATTR, res); + + ServerHttpResponse response = exchange.getResponse(); + // put headers and status so filters can modify the response + final HttpHeaders headers = new HttpHeaders(); + res.responseHeaders().forEach(entry -> headers.add(entry.getKey(), entry.getValue())); + + response.getHeaders().putAll(headers); + response.setStatusCode(HttpStatus.valueOf(res.status().code())); + + return Mono.empty(); + }); + } +} diff --git a/src/main/java/org/springframework/cloud/gateway/handler/WebClientRoutingWebHandler.java b/src/main/java/org/springframework/cloud/gateway/handler/WebClientRoutingWebHandler.java deleted file mode 100644 index 163ea8dd..00000000 --- a/src/main/java/org/springframework/cloud/gateway/handler/WebClientRoutingWebHandler.java +++ /dev/null @@ -1,51 +0,0 @@ -package org.springframework.cloud.gateway.handler; - -import java.net.URI; -import java.util.Optional; - -import org.springframework.http.server.reactive.ServerHttpRequest; -import org.springframework.http.server.reactive.ServerHttpResponse; -import org.springframework.web.reactive.function.client.ClientRequest; -import org.springframework.web.reactive.function.client.WebClient; -import org.springframework.web.server.ServerWebExchange; -import org.springframework.web.server.WebHandler; - -import static org.springframework.cloud.gateway.support.ServerWebExchangeUtils.CLIENT_RESPONSE_ATTR; -import static org.springframework.cloud.gateway.support.ServerWebExchangeUtils.GATEWAY_REQUEST_URL_ATTR; - -import reactor.core.publisher.Mono; - -/** - * @author Spencer Gibb - */ -public class WebClientRoutingWebHandler implements WebHandler { - - private final WebClient webClient; - - public WebClientRoutingWebHandler(WebClient webClient) { - this.webClient = webClient; - } - - @Override - public Mono handle(ServerWebExchange exchange) { - Optional requestUrl = exchange.getAttribute(GATEWAY_REQUEST_URL_ATTR); - ServerHttpRequest request = exchange.getRequest(); - ClientRequest clientRequest = ClientRequest - .method(request.getMethod(), requestUrl.get()) - .headers(request.getHeaders()) - .body((r, context) -> r.writeWith(request.getBody())); - - return this.webClient.exchange(clientRequest).then(clientResponse -> { - // Defer committing the response until all route filters have run - // Put client response as ServerWebExchange attribute and write response later WriteResponseFilter - - exchange.getAttributes().put(CLIENT_RESPONSE_ATTR, clientResponse); - - ServerHttpResponse response = exchange.getResponse(); - // put headers and status so filters can modify the response - response.getHeaders().putAll(clientResponse.headers().asHttpHeaders()); - response.setStatusCode(clientResponse.statusCode()); - return Mono.empty(); - }); - } -}