diff --git a/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/config/GatewayAutoConfiguration.java b/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/config/GatewayAutoConfiguration.java index 08d56749..71c6138d 100644 --- a/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/config/GatewayAutoConfiguration.java +++ b/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/config/GatewayAutoConfiguration.java @@ -32,9 +32,11 @@ import org.springframework.boot.context.properties.EnableConfigurationProperties import org.springframework.cloud.gateway.actuate.GatewayEndpoint; import org.springframework.cloud.gateway.filter.GlobalFilter; import org.springframework.cloud.gateway.filter.NettyRoutingFilter; +import org.springframework.cloud.gateway.filter.NettyWriteResponseFilter; import org.springframework.cloud.gateway.filter.RouteToRequestUrlFilter; +import org.springframework.cloud.gateway.filter.WebClientHttpRoutingFilter; +import org.springframework.cloud.gateway.filter.WebClientWriteResponseFilter; import org.springframework.cloud.gateway.filter.WebsocketRoutingFilter; -import org.springframework.cloud.gateway.filter.WriteResponseFilter; import org.springframework.cloud.gateway.filter.factory.AddRequestHeaderWebFilterFactory; import org.springframework.cloud.gateway.filter.factory.AddRequestParameterWebFilterFactory; import org.springframework.cloud.gateway.filter.factory.AddResponseHeaderWebFilterFactory; @@ -89,6 +91,7 @@ import org.springframework.scripting.support.ResourceScriptSource; import com.netflix.hystrix.HystrixObservableCommand; +import org.springframework.web.reactive.function.client.WebClient; import org.springframework.web.reactive.socket.client.ReactorNettyWebSocketClient; import org.springframework.web.reactive.socket.client.WebSocketClient; import org.springframework.web.reactive.socket.server.WebSocketService; @@ -131,6 +134,11 @@ public class GatewayAutoConfiguration { return new NettyRoutingFilter(httpClient); } + @Bean + public NettyWriteResponseFilter nettyWriteResponseFilter() { + return new NettyWriteResponseFilter(); + } + @Bean public ReactorNettyWebSocketClient reactorNettyWebSocketClient(@Qualifier("nettyClientOptions") Consumer options) { return new ReactorNettyWebSocketClient(options); @@ -208,11 +216,18 @@ public class GatewayAutoConfiguration { return new WebsocketRoutingFilter(webSocketClient, webSocketService); } - @Bean - public WriteResponseFilter writeResponseFilter() { - return new WriteResponseFilter(); + /*@Bean + //TODO: default over netty? configurable + public WebClientHttpRoutingFilter webClientHttpRoutingFilter() { + //TODO: WebClient bean + return new WebClientHttpRoutingFilter(WebClient.builder().build()); } + @Bean + public WebClientWriteResponseFilter webClientWriteResponseFilter() { + return new WebClientWriteResponseFilter(); + }*/ + // Predicate Factory beans @Bean @@ -365,13 +380,6 @@ public class GatewayAutoConfiguration { } } - /*@Bean - public RouterFunction test() { - RouterFunction route = RouterFunctions.route( - RequestPredicates.path("/testfun"), - request -> ServerResponse.ok().body(BodyInserters.fromObject("hello"))); - return route; - }*/ @Configuration @ConditionalOnClass(Endpoint.class) diff --git a/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/NettyRoutingFilter.java b/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/NettyRoutingFilter.java index af751dc0..2f6affd2 100644 --- a/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/NettyRoutingFilter.java +++ b/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/NettyRoutingFilter.java @@ -108,7 +108,7 @@ public class NettyRoutingFilter implements GlobalFilter, Ordered { response.setStatusCode(HttpStatus.valueOf(res.status().code())); // Defer committing the response until all route filters have run - // Put client response as ServerWebExchange attribute and write response later WriteResponseFilter + // Put client response as ServerWebExchange attribute and write response later NettyWriteResponseFilter exchange.getAttributes().put(CLIENT_RESPONSE_ATTR, res); }).then(chain.filter(exchange)); } diff --git a/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/WriteResponseFilter.java b/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/NettyWriteResponseFilter.java similarity index 91% rename from spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/WriteResponseFilter.java rename to spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/NettyWriteResponseFilter.java index 6641b5c9..4bcd76c9 100644 --- a/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/WriteResponseFilter.java +++ b/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/NettyWriteResponseFilter.java @@ -17,8 +17,6 @@ package org.springframework.cloud.gateway.filter; -import java.util.Optional; - import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.springframework.core.Ordered; @@ -37,9 +35,9 @@ import reactor.ipc.netty.http.client.HttpClientResponse; /** * @author Spencer Gibb */ -public class WriteResponseFilter implements GlobalFilter, Ordered { +public class NettyWriteResponseFilter implements GlobalFilter, Ordered { - private static final Log log = LogFactory.getLog(WriteResponseFilter.class); + private static final Log log = LogFactory.getLog(NettyWriteResponseFilter.class); public static final int WRITE_RESPONSE_FILTER_ORDER = -1; @@ -58,7 +56,7 @@ public class WriteResponseFilter implements GlobalFilter, Ordered { if (clientResponse == null) { return Mono.empty(); } - log.trace("WriteResponseFilter start"); + log.trace("NettyWriteResponseFilter start"); ServerHttpResponse response = exchange.getResponse(); NettyDataBufferFactory factory = (NettyDataBufferFactory) response.bufferFactory(); diff --git a/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/WebClientHttpRoutingFilter.java b/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/WebClientHttpRoutingFilter.java new file mode 100644 index 00000000..55702964 --- /dev/null +++ b/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/WebClientHttpRoutingFilter.java @@ -0,0 +1,107 @@ +/* + * Copyright 2013-2017 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.springframework.cloud.gateway.filter; + +import java.net.URI; + +import org.springframework.core.Ordered; +import org.springframework.http.HttpHeaders; +import org.springframework.http.HttpMethod; +import org.springframework.http.server.reactive.ServerHttpRequest; +import org.springframework.http.server.reactive.ServerHttpResponse; +import org.springframework.web.reactive.function.BodyInserters; +import org.springframework.web.reactive.function.client.WebClient; +import org.springframework.web.reactive.function.client.WebClient.RequestBodySpec; +import org.springframework.web.reactive.function.client.WebClient.RequestHeadersSpec; +import org.springframework.web.server.ServerWebExchange; +import org.springframework.web.server.WebFilterChain; + +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 WebClientHttpRoutingFilter implements GlobalFilter, Ordered { + + private final WebClient webClient; + + public WebClientHttpRoutingFilter(WebClient webClient) { + this.webClient = webClient; + } + + @Override + public int getOrder() { + return 2000000; + } + + @Override + public Mono filter(ServerWebExchange exchange, WebFilterChain chain) { + URI requestUrl = exchange.getRequiredAttribute(GATEWAY_REQUEST_URL_ATTR); + + String scheme = requestUrl.getScheme(); + if (!scheme.equals("http") && !scheme.equals("https)")) { + return chain.filter(exchange); + } + + ServerHttpRequest request = exchange.getRequest(); + + //TODO: support forms + + HttpMethod method = request.getMethod(); + + RequestBodySpec bodySpec = this.webClient.method(method) + .uri(requestUrl) + .headers(httpHeaders -> { + httpHeaders.addAll(request.getHeaders()); + httpHeaders.remove(HttpHeaders.HOST); + }); + + RequestHeadersSpec headersSpec; + if (requiresBody(method)) { + headersSpec = bodySpec.body(BodyInserters.fromDataBuffers(request.getBody())); + } else { + headersSpec = bodySpec; + } + + return headersSpec.exchange() + .log("webClient route") + .flatMap(res -> { + ServerHttpResponse response = exchange.getResponse(); + response.getHeaders().putAll(res.headers().asHttpHeaders()); + response.setStatusCode(res.statusCode()); + // Defer committing the response until all route filters have run + // Put client response as ServerWebExchange attribute and write response later NettyWriteResponseFilter + exchange.getAttributes().put(CLIENT_RESPONSE_ATTR, res); + return chain.filter(exchange); + }); + } + + private boolean requiresBody(HttpMethod method) { + switch (method) { + case PUT: + case POST: + case PATCH: + return true; + default: + return false; + } + } +} diff --git a/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/WebClientWriteResponseFilter.java b/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/WebClientWriteResponseFilter.java new file mode 100644 index 00000000..377d37b4 --- /dev/null +++ b/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/WebClientWriteResponseFilter.java @@ -0,0 +1,63 @@ +/* + * Copyright 2013-2017 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +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.http.server.reactive.ServerHttpResponse; +import org.springframework.web.reactive.function.BodyExtractors; +import org.springframework.web.reactive.function.client.ClientResponse; +import org.springframework.web.server.ServerWebExchange; +import org.springframework.web.server.WebFilterChain; + +import static org.springframework.cloud.gateway.support.ServerWebExchangeUtils.CLIENT_RESPONSE_ATTR; + +import reactor.core.publisher.Mono; + +/** + * @author Spencer Gibb + */ +public class WebClientWriteResponseFilter implements GlobalFilter, Ordered { + + private static final Log log = LogFactory.getLog(WebClientWriteResponseFilter.class); + + public static final int WRITE_RESPONSE_FILTER_ORDER = -1; + + @Override + public int getOrder() { + return WRITE_RESPONSE_FILTER_ORDER; + } + + @Override + public Mono filter(ServerWebExchange exchange, WebFilterChain chain) { + // NOTICE: nothing in "pre" filter stage as CLIENT_RESPONSE_ATTR is not added + // until the WebHandler is run + return chain.filter(exchange).then(Mono.defer(() -> { + ClientResponse clientResponse = exchange.getAttribute(CLIENT_RESPONSE_ATTR); + if (clientResponse == null) { + return Mono.empty(); + } + log.trace("WebClientWriteResponseFilter start"); + ServerHttpResponse response = exchange.getResponse(); + + return response.writeWith(clientResponse.body(BodyExtractors.toDataBuffers())).log("webClient response"); + })); + } + +}