From 3b4e902c2f6169c2bc2aecd386ba1df960b5d711 Mon Sep 17 00:00:00 2001 From: Mark Paluch Date: Tue, 19 May 2020 16:24:05 +0200 Subject: [PATCH] Introduce ReactiveVaultEndpointProvider. We now provide a reactive variant of VaultEndpointProvider to allow for non-blocking lookup of VaultEndpoint. Closes gh-555. --- .../vault/client/ReactiveVaultClients.java | 121 ++++++++++++++++-- .../client/ReactiveVaultEndpointProvider.java | 36 ++++++ .../vault/client/VaultClients.java | 31 +++-- .../vault/client/WebClientBuilder.java | 26 +++- .../AbstractReactiveVaultConfiguration.java | 43 +++++-- .../ReactiveVaultClientsIntegrationTests.java | 67 ++++++++++ src/main/asciidoc/new-features.adoc | 1 + 7 files changed, 289 insertions(+), 36 deletions(-) create mode 100644 spring-vault-core/src/main/java/org/springframework/vault/client/ReactiveVaultEndpointProvider.java create mode 100644 spring-vault-core/src/test/java/org/springframework/vault/client/ReactiveVaultClientsIntegrationTests.java diff --git a/spring-vault-core/src/main/java/org/springframework/vault/client/ReactiveVaultClients.java b/spring-vault-core/src/main/java/org/springframework/vault/client/ReactiveVaultClients.java index d3dbcf5f..a75caf06 100644 --- a/spring-vault-core/src/main/java/org/springframework/vault/client/ReactiveVaultClients.java +++ b/spring-vault-core/src/main/java/org/springframework/vault/client/ReactiveVaultClients.java @@ -15,7 +15,10 @@ */ package org.springframework.vault.client; +import java.net.URI; + import reactor.core.publisher.Mono; +import reactor.core.scheduler.Schedulers; import org.springframework.core.codec.ByteArrayDecoder; import org.springframework.core.codec.ByteArrayEncoder; @@ -30,6 +33,8 @@ import org.springframework.web.reactive.function.client.ExchangeFilterFunction; import org.springframework.web.reactive.function.client.ExchangeStrategies; import org.springframework.web.reactive.function.client.WebClient; import org.springframework.web.util.UriBuilderFactory; +import org.springframework.web.util.UriComponents; +import org.springframework.web.util.UriComponentsBuilder; /** * Vault Client factory to create {@link WebClient} configured to the needs of accessing @@ -61,7 +66,9 @@ public class ReactiveVaultClients { * {@link ClientHttpConnector}. The client accepts relative URIs without a leading * slash that are expanded to use {@link VaultEndpoint}. *

- * Requires Jackson 2 for Object-to-JSON mapping. + * Requires Jackson 2 for Object-to-JSON mapping. {@link VaultEndpointProvider} is + * called on {@link Schedulers#boundedElastic()} to ensure that I/O threads are never + * blocked. * * @param endpointProvider must not be {@literal null}. * @param connector must not be {@literal null}. @@ -69,10 +76,24 @@ public class ReactiveVaultClients { */ public static WebClient createWebClient(VaultEndpointProvider endpointProvider, ClientHttpConnector connector) { + return createWebClient(wrap(endpointProvider), connector); + } - Assert.notNull(endpointProvider, "VaultEndpointProvider must not be null"); - Assert.notNull(connector, "ClientHttpConnector must not be null"); - + /** + * Create a {@link WebClient} configured with {@link VaultEndpoint} and + * {@link ClientHttpConnector}. The client accepts relative URIs without a leading + * slash that are expanded to use {@link VaultEndpoint}. + *

+ * Requires Jackson 2 for Object-to-JSON mapping. + * + * @param endpointProvider must not be {@literal null}. + * @param connector must not be {@literal null}. + * @return the configured {@link WebClient}. + * @since 2.3 + */ + public static WebClient createWebClient( + ReactiveVaultEndpointProvider endpointProvider, + ClientHttpConnector connector) { return createWebClientBuilder(endpointProvider, connector).build(); } @@ -88,14 +109,13 @@ public class ReactiveVaultClients { * @return the prepared {@link WebClient.Builder}. */ static WebClient.Builder createWebClientBuilder( - VaultEndpointProvider endpointProvider, ClientHttpConnector connector) { + ReactiveVaultEndpointProvider endpointProvider, + ClientHttpConnector connector) { - Assert.notNull(endpointProvider, "VaultEndpointProvider must not be null"); + Assert.notNull(endpointProvider, + "ReactiveVaultEndpointProvider must not be null"); Assert.notNull(connector, "ClientHttpConnector must not be null"); - UriBuilderFactory uriBuilderFactory = VaultClients - .createUriBuilderFactory(endpointProvider); - ExchangeStrategies strategies = ExchangeStrategies.builder() .codecs(configurer -> { @@ -110,8 +130,50 @@ public class ReactiveVaultClients { }).build(); - return WebClient.builder().uriBuilderFactory(uriBuilderFactory) - .exchangeStrategies(strategies).clientConnector(connector); + WebClient.Builder builder = WebClient.builder().exchangeStrategies(strategies) + .clientConnector(connector); + + boolean simpleSource = false; + if (endpointProvider instanceof VaultEndpointProviderAdapter) { + + if (((VaultEndpointProviderAdapter) endpointProvider).source instanceof SimpleVaultEndpointProvider) { + simpleSource = true; + + UriBuilderFactory uriBuilderFactory = VaultClients + .createUriBuilderFactory( + ((VaultEndpointProviderAdapter) endpointProvider).source); + builder.uriBuilderFactory(uriBuilderFactory); + } + } + + if (!simpleSource) { + builder.filter((request, next) -> { + + URI uri = request.url(); + + if (!uri.isAbsolute()) { + + return endpointProvider.getVaultEndpoint().flatMap(endpoint -> { + + UriComponents uriComponents = UriComponentsBuilder + .fromUri(uri).scheme(endpoint.getScheme()) + .host(endpoint.getHost()).port(endpoint.getPort()) + .replacePath(endpoint.getPath()).path(VaultClients + .normalizePath(endpoint.getPath(), uri.getPath())) + .build(); + + ClientRequest requestToSend = ClientRequest.from(request) + .url(uriComponents.toUri()).build(); + + return next.exchange(requestToSend); + }); + + } + return next.exchange(request); + }); + } + + return builder; } /** @@ -140,4 +202,41 @@ public class ReactiveVaultClients { }); }); } + + /** + * Wrap a {@link VaultEndpointProvider} into a {@link ReactiveVaultEndpointProvider} + * to invoke {@link VaultEndpointProvider#getVaultEndpoint()} on a dedicated + * {@link Schedulers#boundedElastic() scheduler}. + * + * @param endpointProvider must not be {@literal null}. + * @return {@link ReactiveVaultEndpointProvider} wrapping + * {@link VaultEndpointProvider}. + * @since 2.3 + */ + public static ReactiveVaultEndpointProvider wrap( + VaultEndpointProvider endpointProvider) { + + Assert.notNull(endpointProvider, "VaultEndpointProvider must not be null"); + + return new VaultEndpointProviderAdapter(endpointProvider); + } + + private static class VaultEndpointProviderAdapter + implements ReactiveVaultEndpointProvider { + + private final VaultEndpointProvider source; + private final Mono mono; + + VaultEndpointProviderAdapter(VaultEndpointProvider provider) { + + this.source = provider; + this.mono = Mono.fromSupplier(provider::getVaultEndpoint) + .subscribeOn(Schedulers.boundedElastic()); + } + + @Override + public Mono getVaultEndpoint() { + return mono; + } + } } diff --git a/spring-vault-core/src/main/java/org/springframework/vault/client/ReactiveVaultEndpointProvider.java b/spring-vault-core/src/main/java/org/springframework/vault/client/ReactiveVaultEndpointProvider.java new file mode 100644 index 00000000..c03385e0 --- /dev/null +++ b/spring-vault-core/src/main/java/org/springframework/vault/client/ReactiveVaultEndpointProvider.java @@ -0,0 +1,36 @@ +/* + * Copyright 2020 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 + * + * https://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.vault.client; + +import reactor.core.publisher.Mono; + +/** + * Component that provides reactively a {@link VaultEndpoint}. Allows to use a different + * {@link VaultEndpoint} for each Vault request. + * + * @author Mark Paluch + * @since 2.3 + */ +@FunctionalInterface +public interface ReactiveVaultEndpointProvider { + + /** + * Provides access to {@link VaultEndpoint}. + * + * @return a mono emitting the {@link VaultEndpoint}. + */ + Mono getVaultEndpoint(); +} diff --git a/spring-vault-core/src/main/java/org/springframework/vault/client/VaultClients.java b/spring-vault-core/src/main/java/org/springframework/vault/client/VaultClients.java index ad222fa5..b3d2b3da 100644 --- a/spring-vault-core/src/main/java/org/springframework/vault/client/VaultClients.java +++ b/spring-vault-core/src/main/java/org/springframework/vault/client/VaultClients.java @@ -249,16 +249,7 @@ public class VaultClients { } if (baseUrl != null) { - - if (uriTemplate.startsWith("/") && baseUrl.endsWith("/")) { - return uriTemplate.substring(1); - } - - if (!uriTemplate.startsWith("/") && !baseUrl.endsWith("/")) { - return "/" + uriTemplate; - } - - return uriTemplate; + return normalizePath(baseUrl, uriTemplate); } try { @@ -277,4 +268,24 @@ public class VaultClients { return uriTemplate; } + + /** + * Normalize the URI {@code path} so that it can be combined with {@code prefix}. + * + * @param prefix + * @param path + * @return + */ + static String normalizePath(String prefix, String path) { + + if (path.startsWith("/") && prefix.endsWith("/")) { + return path.substring(1); + } + + if (!path.startsWith("/") && !prefix.endsWith("/")) { + return "/" + path; + } + + return path; + } } diff --git a/spring-vault-core/src/main/java/org/springframework/vault/client/WebClientBuilder.java b/spring-vault-core/src/main/java/org/springframework/vault/client/WebClientBuilder.java index 5908226d..5d8196c5 100644 --- a/spring-vault-core/src/main/java/org/springframework/vault/client/WebClientBuilder.java +++ b/spring-vault-core/src/main/java/org/springframework/vault/client/WebClientBuilder.java @@ -49,7 +49,7 @@ import org.springframework.web.reactive.function.client.WebClient; */ public class WebClientBuilder { - private @Nullable VaultEndpointProvider endpointProvider; + private @Nullable ReactiveVaultEndpointProvider endpointProvider; private Supplier httpConnector = () -> ClientHttpConnectorFactory .create(new ClientOptions(), SslConfiguration.unconfigured()); @@ -84,14 +84,27 @@ public class WebClientBuilder { /** * Set the {@link VaultEndpointProvider} that should be used with the - * {@link WebClient}. + * {@link WebClient}. {@link VaultEndpointProvider#getVaultEndpoint()} is called on + * {@link reactor.core.scheduler.Schedulers#boundedElastic() a dedicated Thread} to + * ensure that I/O threads are never blocked. * * @param provider the {@link VaultEndpoint} provider. * @return {@code this} {@link WebClientBuilder}. */ public WebClientBuilder endpointProvider(VaultEndpointProvider provider) { + return endpointProvider(ReactiveVaultClients.wrap(provider)); + } - Assert.notNull(provider, "VaultEndpointProvider must not be null"); + /** + * Set the {@link ReactiveVaultEndpointProvider} that should be used with the + * {@link WebClient}. + * + * @param provider the {@link VaultEndpoint} provider. + * @return {@code this} {@link WebClientBuilder}. + */ + public WebClientBuilder endpointProvider(ReactiveVaultEndpointProvider provider) { + + Assert.notNull(provider, "ReactiveVaultEndpointProvider must not be null"); this.endpointProvider = provider; @@ -199,9 +212,6 @@ public class WebClientBuilder { */ public WebClient build() { - Assert.state(this.endpointProvider != null, - "VaultEndpointProvider must not be null"); - WebClient.Builder builder = createWebClientBuilder(); if (!defaultHeaders.isEmpty()) { @@ -215,7 +225,6 @@ public class WebClientBuilder { headers.add(key, value); } })).build()); - }); } @@ -234,6 +243,9 @@ public class WebClientBuilder { */ protected WebClient.Builder createWebClientBuilder() { + Assert.state(this.endpointProvider != null, + "VaultEndpointProvider must not be null"); + ClientHttpConnector connector = this.httpConnector.get(); return ReactiveVaultClients.createWebClientBuilder(endpointProvider, connector); diff --git a/spring-vault-core/src/main/java/org/springframework/vault/config/AbstractReactiveVaultConfiguration.java b/spring-vault-core/src/main/java/org/springframework/vault/config/AbstractReactiveVaultConfiguration.java index 3abe5e81..3481aac7 100644 --- a/spring-vault-core/src/main/java/org/springframework/vault/config/AbstractReactiveVaultConfiguration.java +++ b/spring-vault-core/src/main/java/org/springframework/vault/config/AbstractReactiveVaultConfiguration.java @@ -34,7 +34,8 @@ import org.springframework.vault.authentication.SessionManager; import org.springframework.vault.authentication.TokenAuthentication; import org.springframework.vault.authentication.VaultTokenSupplier; import org.springframework.vault.client.ClientHttpConnectorFactory; -import org.springframework.vault.client.VaultEndpoint; +import org.springframework.vault.client.ReactiveVaultClients; +import org.springframework.vault.client.ReactiveVaultEndpointProvider; import org.springframework.vault.client.VaultEndpointProvider; import org.springframework.vault.client.WebClientBuilder; import org.springframework.vault.client.WebClientCustomizer; @@ -67,17 +68,45 @@ import org.springframework.web.reactive.function.client.WebClient; public abstract class AbstractReactiveVaultConfiguration extends AbstractVaultConfiguration { + /** + * @return a {@link ReactiveVaultEndpointProvider} returning the value of + * {@link #vaultEndpointProvider()}. + * + * @see #vaultEndpoint() + * @see #vaultEndpointProvider() + * @since 2.3 + */ + public ReactiveVaultEndpointProvider reactiveVaultEndpointProvider() { + return ReactiveVaultClients.wrap(vaultEndpointProvider()); + } + /** * Create a {@link WebClientBuilder} initialized with {@link VaultEndpointProvider} * and {@link ClientHttpConnector}. May be overridden by subclasses. * * @return the {@link WebClientBuilder}. - * @see #vaultEndpointProvider() + * @see #reactiveVaultEndpointProvider() * @see #clientHttpConnector() * @since 2.2 */ protected WebClientBuilder webClientBuilder(VaultEndpointProvider endpointProvider, ClientHttpConnector httpConnector) { + return webClientBuilder(ReactiveVaultClients.wrap(endpointProvider), + httpConnector); + } + + /** + * Create a {@link WebClientBuilder} initialized with {@link VaultEndpointProvider} + * and {@link ClientHttpConnector}. May be overridden by subclasses. + * + * @return the {@link WebClientBuilder}. + * @see #reactiveVaultEndpointProvider() + * @see #clientHttpConnector() + * @since 2.3 + */ + protected WebClientBuilder webClientBuilder( + ReactiveVaultEndpointProvider endpointProvider, + ClientHttpConnector httpConnector) { ObjectProvider customizers = getBeanFactory() .getBeanProvider(WebClientCustomizer.class); @@ -103,7 +132,7 @@ public abstract class AbstractReactiveVaultConfiguration ClientHttpConnector httpConnector = clientHttpConnector(); return new DefaultWebClientFactory(httpConnector, clientHttpConnector -> { - return webClientBuilder(vaultEndpointProvider(), clientHttpConnector); + return webClientBuilder(reactiveVaultEndpointProvider(), clientHttpConnector); }); } @@ -112,17 +141,15 @@ public abstract class AbstractReactiveVaultConfiguration * * @return the {@link ReactiveVaultTemplate}. * @see #vaultEndpoint() + * @see #reactiveVaultEndpointProvider() * @see #clientHttpConnector() * @see #reactiveSessionManager() */ @Bean public ReactiveVaultTemplate reactiveVaultTemplate() { - VaultEndpointProvider provider = vaultEndpointProvider(); - VaultEndpoint vaultEndpoint = provider.getVaultEndpoint(); - return new ReactiveVaultTemplate( - webClientBuilder(() -> vaultEndpoint, clientHttpConnector()), + webClientBuilder(reactiveVaultEndpointProvider(), clientHttpConnector()), getReactiveSessionManager()); } @@ -206,7 +233,7 @@ public abstract class AbstractReactiveVaultConfiguration /** * Return the {@link WebClientFactory}. - * + * * @return the {@link WebClientFactory} bean. * @since 2.3 */ diff --git a/spring-vault-core/src/test/java/org/springframework/vault/client/ReactiveVaultClientsIntegrationTests.java b/spring-vault-core/src/test/java/org/springframework/vault/client/ReactiveVaultClientsIntegrationTests.java new file mode 100644 index 00000000..cf8ff0fb --- /dev/null +++ b/spring-vault-core/src/test/java/org/springframework/vault/client/ReactiveVaultClientsIntegrationTests.java @@ -0,0 +1,67 @@ +/* + * Copyright 2020 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 + * + * https://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.vault.client; + +import java.util.concurrent.atomic.AtomicReference; + +import org.junit.jupiter.api.Test; +import reactor.core.publisher.Mono; +import reactor.test.StepVerifier; + +import org.springframework.vault.support.ClientOptions; +import org.springframework.vault.util.IntegrationTestSupport; +import org.springframework.vault.util.Settings; +import org.springframework.vault.util.TestRestTemplateFactory; +import org.springframework.web.reactive.function.client.WebClient; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * Integration tests for {@link ReactiveVaultClients}. + * + * @author Mark Paluch + */ +class ReactiveVaultClientsIntegrationTests extends IntegrationTestSupport { + + @Test + void shouldUseVaultEndpointProvider() { + + AtomicReference resolver = new AtomicReference<>(); + + WebClient client = ReactiveVaultClients.createWebClient(() -> { + + return Mono.fromSupplier(() -> { + resolver.set(Thread.currentThread()); + return TestRestTemplateFactory.TEST_VAULT_ENDPOINT; + }); + }, ClientHttpConnectorFactory.create(new ClientOptions(), + Settings.createSslConfiguration())); + + client.get().uri("/sys/health").exchange() + .flatMap(it -> it.bodyToMono(String.class)).as(StepVerifier::create) + .consumeNextWith(actual -> { + assertThat(actual).contains("initialized").contains("standby"); + }).verifyComplete(); + + client.get().uri("sys/health").exchange() + .flatMap(it -> it.bodyToMono(String.class)).as(StepVerifier::create) + .consumeNextWith(actual -> { + assertThat(actual).contains("initialized").contains("standby"); + }).verifyComplete(); + + assertThat(resolver).hasValue(Thread.currentThread()); + } +} diff --git a/src/main/asciidoc/new-features.adoc b/src/main/asciidoc/new-features.adoc index 6edaff86..63eb480c 100644 --- a/src/main/asciidoc/new-features.adoc +++ b/src/main/asciidoc/new-features.adoc @@ -5,6 +5,7 @@ === What's new in Spring Vault 2.3 * Support for PEM-encoded certificates for keystore and truststore usage. +* `ReactiveVaultEndpointProvider` for non-blocking lookup of `VaultEndpoint`. [[new-features.2-2-0]] === What's new in Spring Vault 2.2