diff --git a/supplier/http-supplier/README.adoc b/supplier/http-supplier/README.adoc index 9e506407..9725e100 100644 --- a/supplier/http-supplier/README.adoc +++ b/supplier/http-supplier/README.adoc @@ -24,6 +24,8 @@ All configuration properties are prefixed with `http`. For more information on the various options available, please see link:src/main/java/org/springframework/cloud/fn/supplier/http/HttpSupplierProperties.java[HttpSupplierProperties]. +The `HeaderMapper` bean can be provided in the target configuration to override a default one in the `HttpSupplierConfiguration`. + ## Tests See this link:src/test/java/org/springframework/cloud/fn/supplier/http/HttpSupplierApplicationTests.java[test suite] for the various ways, this supplier is used. diff --git a/supplier/http-supplier/src/main/java/org/springframework/cloud/fn/supplier/http/HttpSupplierConfiguration.java b/supplier/http-supplier/src/main/java/org/springframework/cloud/fn/supplier/http/HttpSupplierConfiguration.java index a7ccdc95..43585f05 100644 --- a/supplier/http-supplier/src/main/java/org/springframework/cloud/fn/supplier/http/HttpSupplierConfiguration.java +++ b/supplier/http-supplier/src/main/java/org/springframework/cloud/fn/supplier/http/HttpSupplierConfiguration.java @@ -21,6 +21,7 @@ import java.util.function.Supplier; import org.reactivestreams.Publisher; import reactor.core.publisher.Flux; +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -46,37 +47,41 @@ import org.springframework.messaging.MessageHeaders; @Configuration public class HttpSupplierConfiguration { + @Bean + @ConditionalOnMissingBean + public HeaderMapper httpHeaderMapper(HttpSupplierProperties httpSupplierProperties) { + DefaultHttpHeaderMapper defaultHttpHeaderMapper = DefaultHttpHeaderMapper.inboundMapper(); + defaultHttpHeaderMapper.setInboundHeaderNames(httpSupplierProperties.getMappedRequestHeaders()); + return defaultHttpHeaderMapper; + } + @Bean public Publisher> httpSupplierFlow(HttpSupplierProperties httpSupplierProperties, + HeaderMapper httpHeaderMapper, ServerCodecConfigurer serverCodecConfigurer) { return IntegrationFlows.from( - WebFlux.inboundChannelAdapter(httpSupplierProperties.getPathPattern()) - .requestPayloadType(byte[].class) - .statusCodeExpression(new ValueExpression<>(HttpStatus.ACCEPTED)) - .mappedRequestHeaders(httpSupplierProperties.getMappedRequestHeaders()) - .codecConfigurer(serverCodecConfigurer) - .crossOrigin(crossOrigin -> - crossOrigin.origin(httpSupplierProperties.getCors().getAllowedOrigins()) - .allowedHeaders(httpSupplierProperties.getCors().getAllowedHeaders()) - .allowCredentials(httpSupplierProperties.getCors().getAllowCredentials())) - .autoStartup(false)) + WebFlux.inboundChannelAdapter(httpSupplierProperties.getPathPattern()) + .requestPayloadType(byte[].class) + .statusCodeExpression(new ValueExpression<>(HttpStatus.ACCEPTED)) + .headerMapper(httpHeaderMapper) + .codecConfigurer(serverCodecConfigurer) + .crossOrigin(crossOrigin -> + crossOrigin.origin(httpSupplierProperties.getCors().getAllowedOrigins()) + .allowedHeaders(httpSupplierProperties.getCors().getAllowedHeaders()) + .allowCredentials(httpSupplierProperties.getCors().getAllowCredentials())) + .autoStartup(false)) .enrichHeaders((headers) -> headers.headerFunction(MessageHeaders.CONTENT_TYPE, (message) -> - (MediaType.APPLICATION_FORM_URLENCODED_VALUE.equals( - message.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString())) + (MediaType.APPLICATION_FORM_URLENCODED.equals( + message.getHeaders().get(MessageHeaders.CONTENT_TYPE, MediaType.class))) ? MediaType.APPLICATION_JSON : null, true)) .toReactivePublisher(); } - @Bean - public HeaderMapper httpHeaderMapper() { - return DefaultHttpHeaderMapper.inboundMapper(); - } - @Bean public Supplier>> httpSupplier( Publisher> httpRequestPublisher, diff --git a/supplier/http-supplier/src/test/java/org/springframework/cloud/fn/supplier/http/HttpSupplierApplicationTests.java b/supplier/http-supplier/src/test/java/org/springframework/cloud/fn/supplier/http/HttpSupplierApplicationTests.java index 4d749888..b93d80f0 100644 --- a/supplier/http-supplier/src/test/java/org/springframework/cloud/fn/supplier/http/HttpSupplierApplicationTests.java +++ b/supplier/http-supplier/src/test/java/org/springframework/cloud/fn/supplier/http/HttpSupplierApplicationTests.java @@ -20,9 +20,13 @@ import java.nio.charset.StandardCharsets; import java.time.Duration; import java.util.function.Supplier; +import javax.net.ssl.SSLException; + +import io.netty.handler.ssl.SslContext; import io.netty.handler.ssl.SslContextBuilder; import io.netty.handler.ssl.SslProvider; import io.netty.handler.ssl.util.InsecureTrustManagerFactory; +import org.assertj.core.api.InstanceOfAssertFactories; import org.junit.jupiter.api.Test; import reactor.core.publisher.Flux; import reactor.netty.http.client.HttpClient; @@ -41,6 +45,7 @@ import org.springframework.integration.test.util.TestUtils; import org.springframework.integration.webflux.inbound.WebFluxInboundEndpoint; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHeaders; +import org.springframework.test.annotation.DirtiesContext; import org.springframework.web.reactive.function.client.WebClient; import static org.assertj.core.api.Assertions.assertThat; @@ -59,6 +64,7 @@ import static org.assertj.core.api.Assertions.assertThat; "server.ssl.client-auth=want", "spring.codec.max-in-memory-size=10MB" }) +@DirtiesContext public class HttpSupplierApplicationTests { @Autowired @@ -71,7 +77,7 @@ public class HttpSupplierApplicationTests { private int port; @Test - public void testHttpSupplier() { + public void testHttpSupplier() throws SSLException { ServerCodecConfigurer codecConfigurer = TestUtils.getPropertyValue(this.webFluxInboundEndpoint, "codecConfigurer", ServerCodecConfigurer.class); @@ -109,15 +115,26 @@ public class HttpSupplierApplicationTests { assertThat(message) .extracting(Message::getPayload) .isEqualTo("{\"name\":\"test3\"}".getBytes())) + .assertNext((message) -> + assertThat(message) + .satisfies((msg) -> assertThat(msg) + .extracting(Message::getPayload) + .asInstanceOf(InstanceOfAssertFactories.MAP) + .isEmpty()) + .satisfies((msg) -> assertThat(msg.getHeaders()) + .doesNotContainKey(MessageHeaders.CONTENT_TYPE))) .thenCancel() .verifyLater(); + SslContext sslContext = + SslContextBuilder.forClient() + .sslProvider(SslProvider.JDK) + .trustManager(InsecureTrustManagerFactory.INSTANCE) + .build(); + HttpClient httpClient = HttpClient.create() - .secure(sslSpec -> - sslSpec.sslContext(SslContextBuilder.forClient() - .sslProvider(SslProvider.JDK) - .trustManager(InsecureTrustManagerFactory.INSTANCE))); + .secure(sslSpec -> sslSpec.sslContext(sslContext)); WebClient webClient = WebClient.builder() @@ -125,10 +142,31 @@ public class HttpSupplierApplicationTests { .baseUrl("https://localhost:" + port) .build(); - WebClient.RequestBodySpec requestBodySpec = webClient.post().uri("/"); - requestBodySpec.bodyValue("test1").exchange().block(Duration.ofSeconds(10)); - requestBodySpec.bodyValue(new TestPojo("test2")).exchange().block(Duration.ofSeconds(10)); - requestBodySpec.bodyValue(new TestPojo("test3")).exchange().block(Duration.ofSeconds(10)); + webClient.post() + .uri("/") + .bodyValue("test1") + .retrieve() + .toBodilessEntity() + .block(Duration.ofSeconds(10)); + + webClient.post() + .uri("/") + .bodyValue(new TestPojo("test2")) + .retrieve() + .toBodilessEntity() + .block(Duration.ofSeconds(10)); + + webClient.post() + .uri("/") + .bodyValue(new TestPojo("test3")) + .retrieve() + .toBodilessEntity() + .block(Duration.ofSeconds(10)); + + webClient.post().uri("/") + .retrieve() + .toBodilessEntity() + .block(Duration.ofSeconds(10)); stepVerifier.verify(); } @@ -156,6 +194,7 @@ public class HttpSupplierApplicationTests { @SpringBootApplication static class HttpSupplierTestApplication { + } }