diff --git a/spring-cloud-netflix-eureka-client/src/main/java/org/springframework/cloud/netflix/eureka/http/WebClientEurekaHttpClient.java b/spring-cloud-netflix-eureka-client/src/main/java/org/springframework/cloud/netflix/eureka/http/WebClientEurekaHttpClient.java index eb81623c4..6e130c487 100644 --- a/spring-cloud-netflix-eureka-client/src/main/java/org/springframework/cloud/netflix/eureka/http/WebClientEurekaHttpClient.java +++ b/spring-cloud-netflix-eureka-client/src/main/java/org/springframework/cloud/netflix/eureka/http/WebClientEurekaHttpClient.java @@ -16,8 +16,6 @@ package org.springframework.cloud.netflix.eureka.http; -import java.util.Collections; -import java.util.HashMap; import java.util.Map; import com.netflix.appinfo.InstanceInfo; @@ -28,12 +26,12 @@ import com.netflix.discovery.shared.transport.EurekaHttpClient; import com.netflix.discovery.shared.transport.EurekaHttpResponse; import com.netflix.discovery.shared.transport.EurekaHttpResponse.EurekaHttpResponseBuilder; import com.netflix.discovery.util.StringUtil; -import org.apache.commons.logging.Log; -import org.apache.commons.logging.LogFactory; +import reactor.core.publisher.Mono; import org.springframework.http.HttpHeaders; import org.springframework.http.HttpStatus; import org.springframework.http.MediaType; +import org.springframework.http.ResponseEntity; import org.springframework.web.reactive.function.BodyInserters; import org.springframework.web.reactive.function.client.ClientResponse; import org.springframework.web.reactive.function.client.WebClient; @@ -46,8 +44,6 @@ import static com.netflix.discovery.shared.transport.EurekaHttpResponse.anEureka */ public class WebClientEurekaHttpClient implements EurekaHttpClient { - protected final Log logger = LogFactory.getLog(getClass()); - private WebClient webClient; public WebClientEurekaHttpClient(WebClient webClient) { @@ -56,16 +52,18 @@ public class WebClientEurekaHttpClient implements EurekaHttpClient { @Override public EurekaHttpResponse register(InstanceInfo info) { - return webClient.post().uri("apps/" + info.getAppName(), Void.class).body(BodyInserters.fromValue(info)) + return webClient.post().uri("apps/" + info.getAppName()).body(BodyInserters.fromValue(info)) .header(HttpHeaders.ACCEPT_ENCODING, "gzip") - .header(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE).exchange() - .map(response -> eurekaHttpResponse(response)).block(); + .header(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE).retrieve() + .onStatus(HttpStatus::isError, this::ignoreError).toBodilessEntity().map(this::eurekaHttpResponse) + .block(); } @Override public EurekaHttpResponse cancel(String appName, String id) { - return webClient.delete().uri("apps/" + appName + '/' + id, Void.class).exchange() - .map(response -> eurekaHttpResponse(response)).block(); + return webClient.delete().uri("apps/" + appName + '/' + id).retrieve() + .onStatus(HttpStatus::isError, this::ignoreError).toBodilessEntity().map(this::eurekaHttpResponse) + .block(); } @Override @@ -75,14 +73,15 @@ public class WebClientEurekaHttpClient implements EurekaHttpClient { + "&lastDirtyTimestamp=" + info.getLastDirtyTimestamp().toString() + (overriddenStatus != null ? "&overriddenstatus=" + overriddenStatus.name() : ""); - ClientResponse response = webClient.put().uri(urlPath, InstanceInfo.class) + ResponseEntity response = webClient.put().uri(urlPath) .header(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE) - .header(HttpHeaders.ACCEPT, MediaType.APPLICATION_JSON_VALUE).exchange().block(); + .header(HttpHeaders.ACCEPT, MediaType.APPLICATION_JSON_VALUE).retrieve() + .onStatus(HttpStatus::isError, this::ignoreError).toEntity(InstanceInfo.class).block(); EurekaHttpResponseBuilder builder = anEurekaHttpResponse(statusCodeValueOf(response), InstanceInfo.class).headers(headersOf(response)); - InstanceInfo entity = response.toEntity(InstanceInfo.class).block().getBody(); + InstanceInfo entity = response.getBody(); if (entity != null) { builder.entity(entity); @@ -98,9 +97,9 @@ public class WebClientEurekaHttpClient implements EurekaHttpClient { String urlPath = "apps/" + appName + '/' + id + "/status?value=" + newStatus.name() + "&lastDirtyTimestamp=" + info.getLastDirtyTimestamp().toString(); - return webClient.put().uri(urlPath, Void.class) - .header(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE).exchange() - .map(response -> eurekaHttpResponse(response)).block(); + return webClient.put().uri(urlPath).header(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE) + .retrieve().onStatus(HttpStatus::isError, this::ignoreError).toBodilessEntity() + .map(this::eurekaHttpResponse).block(); } @Override @@ -108,9 +107,9 @@ public class WebClientEurekaHttpClient implements EurekaHttpClient { String urlPath = "apps/" + appName + '/' + id + "/status?lastDirtyTimestamp=" + info.getLastDirtyTimestamp().toString(); - return webClient.delete().uri(urlPath, Void.class) - .header(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE).exchange() - .map(response -> eurekaHttpResponse(response)).block(); + return webClient.delete().uri(urlPath).header(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE) + .retrieve().onStatus(HttpStatus::isError, this::ignoreError).toBodilessEntity() + .map(this::eurekaHttpResponse).block(); } @Override @@ -125,13 +124,14 @@ public class WebClientEurekaHttpClient implements EurekaHttpClient { url = url + (urlPath.contains("?") ? "&" : "?") + "regions=" + StringUtil.join(regions); } - ClientResponse response = webClient.get().uri(url, Applications.class) + ResponseEntity response = webClient.get().uri(url) .header(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE) - .header(HttpHeaders.ACCEPT, MediaType.APPLICATION_JSON_VALUE).exchange().block(); + .header(HttpHeaders.ACCEPT, MediaType.APPLICATION_JSON_VALUE).retrieve() + .onStatus(HttpStatus::isError, this::ignoreError).toEntity(Applications.class).block(); int statusCode = statusCodeValueOf(response); - Applications body = response.toEntity(Applications.class).block().getBody(); + Applications body = response.getBody(); return anEurekaHttpResponse(statusCode, statusCode == HttpStatus.OK.value() && body != null ? body : null) .headers(headersOf(response)).build(); @@ -155,11 +155,12 @@ public class WebClientEurekaHttpClient implements EurekaHttpClient { @Override public EurekaHttpResponse getApplication(String appName) { - ClientResponse response = webClient.get().uri("apps/" + appName, Application.class) - .header(HttpHeaders.ACCEPT, MediaType.APPLICATION_JSON_VALUE).exchange().block(); + ResponseEntity response = webClient.get().uri("apps/" + appName) + .header(HttpHeaders.ACCEPT, MediaType.APPLICATION_JSON_VALUE).retrieve() + .onStatus(HttpStatus::isError, this::ignoreError).toEntity(Application.class).block(); int statusCode = statusCodeValueOf(response); - Application body = response.toEntity(Application.class).block().getBody(); + Application body = response.getBody(); Application application = statusCode == HttpStatus.OK.value() && body != null ? body : null; @@ -177,11 +178,12 @@ public class WebClientEurekaHttpClient implements EurekaHttpClient { } private EurekaHttpResponse getInstanceInternal(String urlPath) { - ClientResponse response = webClient.get().uri(urlPath, InstanceInfo.class) - .header(HttpHeaders.ACCEPT, MediaType.APPLICATION_JSON_VALUE).exchange().block(); + ResponseEntity response = webClient.get().uri(urlPath) + .header(HttpHeaders.ACCEPT, MediaType.APPLICATION_JSON_VALUE).retrieve() + .onStatus(HttpStatus::isError, this::ignoreError).toEntity(InstanceInfo.class).block(); int statusCode = statusCodeValueOf(response); - InstanceInfo body = response.toEntity(InstanceInfo.class).block().getBody(); + InstanceInfo body = response.getBody(); return anEurekaHttpResponse(statusCode, statusCode == HttpStatus.OK.value() && body != null ? body : null) .headers(headersOf(response)).build(); @@ -196,26 +198,19 @@ public class WebClientEurekaHttpClient implements EurekaHttpClient { return this.webClient; } - private static Map headersOf(ClientResponse response) { - ClientResponse.Headers httpHeaders = response.headers(); - if (httpHeaders == null) { - return Collections.emptyMap(); - } - HttpHeaders asHeaders = httpHeaders.asHttpHeaders(); - if (asHeaders == null) { - return Collections.emptyMap(); - } - Map headers = new HashMap<>(); - asHeaders.entrySet().stream() - .forEach(entry -> entry.getValue().stream().forEach(v -> headers.put(entry.getKey(), v))); - return headers; + private Mono ignoreError(ClientResponse response) { + return Mono.empty(); } - private int statusCodeValueOf(ClientResponse response) { - return response.statusCode().value(); + private static Map headersOf(ResponseEntity response) { + return response.getHeaders().toSingleValueMap(); } - private EurekaHttpResponse eurekaHttpResponse(ClientResponse response) { + private int statusCodeValueOf(ResponseEntity response) { + return response.getStatusCode().value(); + } + + private EurekaHttpResponse eurekaHttpResponse(ResponseEntity response) { return anEurekaHttpResponse(statusCodeValueOf(response)).headers(headersOf(response)).build(); }