diff --git a/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/NettyRoutingFilter.java b/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/NettyRoutingFilter.java index d41c1706..780c318b 100644 --- a/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/NettyRoutingFilter.java +++ b/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/NettyRoutingFilter.java @@ -130,61 +130,62 @@ public class NettyRoutingFilter implements GlobalFilter, Ordered { boolean preserveHost = exchange.getAttributeOrDefault(PRESERVE_HOST_HEADER_ATTRIBUTE, false); Route route = exchange.getAttribute(GATEWAY_ROUTE_ATTR); - Flux responseFlux = httpClient(route, exchange).flatMapMany(hc -> hc.headers(headers -> { - headers.add(httpHeaders); - // Will either be set below, or later by Netty - headers.remove(HttpHeaders.HOST); - if (preserveHost) { - String host = request.getHeaders().getFirst(HttpHeaders.HOST); - headers.add(HttpHeaders.HOST, host); - } - }).request(method).uri(url).send((req, nettyOutbound) -> { - if (log.isTraceEnabled()) { - nettyOutbound.withConnection(connection -> log.trace("outbound route: " - + connection.channel().id().asShortText() + ", inbound: " + exchange.getLogPrefix())); - } - return nettyOutbound.send(request.getBody().map(this::getByteBuf)); - }).responseConnection((res, connection) -> { + Flux responseFlux = getHttpClientMono(route, exchange) + .flatMapMany(httpClient -> httpClient.headers(headers -> { + headers.add(httpHeaders); + // Will either be set below, or later by Netty + headers.remove(HttpHeaders.HOST); + if (preserveHost) { + String host = request.getHeaders().getFirst(HttpHeaders.HOST); + headers.add(HttpHeaders.HOST, host); + } + }).request(method).uri(url).send((req, nettyOutbound) -> { + if (log.isTraceEnabled()) { + nettyOutbound.withConnection(connection -> log.trace("outbound route: " + + connection.channel().id().asShortText() + ", inbound: " + exchange.getLogPrefix())); + } + return nettyOutbound.send(request.getBody().map(this::getByteBuf)); + }).responseConnection((res, connection) -> { - // 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); - exchange.getAttributes().put(CLIENT_RESPONSE_CONN_ATTR, connection); + // 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); + exchange.getAttributes().put(CLIENT_RESPONSE_CONN_ATTR, connection); - ServerHttpResponse response = exchange.getResponse(); - // put headers and status so filters can modify the response - HttpHeaders headers = new HttpHeaders(); + ServerHttpResponse response = exchange.getResponse(); + // put headers and status so filters can modify the response + HttpHeaders headers = new HttpHeaders(); - res.responseHeaders().forEach(entry -> headers.add(entry.getKey(), entry.getValue())); + res.responseHeaders().forEach(entry -> headers.add(entry.getKey(), entry.getValue())); - String contentTypeValue = headers.getFirst(HttpHeaders.CONTENT_TYPE); - if (StringUtils.hasLength(contentTypeValue)) { - exchange.getAttributes().put(ORIGINAL_RESPONSE_CONTENT_TYPE_ATTR, contentTypeValue); - } + String contentTypeValue = headers.getFirst(HttpHeaders.CONTENT_TYPE); + if (StringUtils.hasLength(contentTypeValue)) { + exchange.getAttributes().put(ORIGINAL_RESPONSE_CONTENT_TYPE_ATTR, contentTypeValue); + } - setResponseStatus(res, response); + setResponseStatus(res, response); - // make sure headers filters run after setting status so it is - // available in response - HttpHeaders filteredResponseHeaders = HttpHeadersFilter.filter(getHeadersFilters(), headers, exchange, - Type.RESPONSE); + // make sure headers filters run after setting status so it is + // available in response + HttpHeaders filteredResponseHeaders = HttpHeadersFilter.filter(getHeadersFilters(), headers, + exchange, Type.RESPONSE); - if (!filteredResponseHeaders.containsKey(HttpHeaders.TRANSFER_ENCODING) - && filteredResponseHeaders.containsKey(HttpHeaders.CONTENT_LENGTH)) { - // It is not valid to have both the transfer-encoding header and - // the content-length header. - // Remove the transfer-encoding header in the response if the - // content-length header is present. - response.getHeaders().remove(HttpHeaders.TRANSFER_ENCODING); - } + if (!filteredResponseHeaders.containsKey(HttpHeaders.TRANSFER_ENCODING) + && filteredResponseHeaders.containsKey(HttpHeaders.CONTENT_LENGTH)) { + // It is not valid to have both the transfer-encoding header and + // the content-length header. + // Remove the transfer-encoding header in the response if the + // content-length header is present. + response.getHeaders().remove(HttpHeaders.TRANSFER_ENCODING); + } - exchange.getAttributes().put(CLIENT_RESPONSE_HEADER_NAMES, filteredResponseHeaders.keySet()); + exchange.getAttributes().put(CLIENT_RESPONSE_HEADER_NAMES, filteredResponseHeaders.keySet()); - response.getHeaders().addAll(filteredResponseHeaders); + response.getHeaders().addAll(filteredResponseHeaders); - return Mono.just(res); - })); + return Mono.just(res); + })); Duration responseTimeout = getResponseTimeout(route); if (responseTimeout != null) { @@ -239,7 +240,7 @@ public class NettyRoutingFilter implements GlobalFilter, Ordered { * @param exchange the current ServerWebExchange. * @return the configured HttpClient. */ - protected Mono httpClient(Route route, ServerWebExchange exchange) { + protected Mono getHttpClientMono(Route route, ServerWebExchange exchange) { return Mono.just(getHttpClient(route, exchange)); }