Add experimental WebClient Routing and Response filters
This commit is contained in:
@@ -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<? super HttpClientOptions.Builder> 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<ServerResponse> test() {
|
||||
RouterFunction<ServerResponse> route = RouterFunctions.route(
|
||||
RequestPredicates.path("/testfun"),
|
||||
request -> ServerResponse.ok().body(BodyInserters.fromObject("hello")));
|
||||
return route;
|
||||
}*/
|
||||
|
||||
@Configuration
|
||||
@ConditionalOnClass(Endpoint.class)
|
||||
|
||||
@@ -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));
|
||||
}
|
||||
|
||||
@@ -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();
|
||||
@@ -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<Void> 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;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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<Void> 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");
|
||||
}));
|
||||
}
|
||||
|
||||
}
|
||||
Reference in New Issue
Block a user