diff --git a/docs/src/main/asciidoc/images/SCSt-overview.png b/docs/src/main/asciidoc/images/SCSt-overview.png index 46dd0c809..6cfd685a3 100644 Binary files a/docs/src/main/asciidoc/images/SCSt-overview.png and b/docs/src/main/asciidoc/images/SCSt-overview.png differ diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/ApplicationJsonMessageMarshallingConverter.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/ApplicationJsonMessageMarshallingConverter.java index f94b344cd..e5f50595f 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/ApplicationJsonMessageMarshallingConverter.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/ApplicationJsonMessageMarshallingConverter.java @@ -31,11 +31,13 @@ import com.fasterxml.jackson.databind.ObjectMapper; import org.springframework.cloud.function.context.catalog.FunctionTypeUtils; import org.springframework.core.MethodParameter; import org.springframework.core.ParameterizedTypeReference; +import org.springframework.integration.support.MutableMessageHeaders; import org.springframework.lang.Nullable; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.converter.MappingJackson2MessageConverter; import org.springframework.messaging.converter.MessageConversionException; +import org.springframework.util.MimeType; /** * Variation of {@link MappingJackson2MessageConverter} to support marshalling and @@ -167,4 +169,16 @@ class ApplicationJsonMessageMarshallingConverter extends MappingJackson2MessageC } } + @Override + @Nullable + protected MimeType getMimeType(@Nullable MessageHeaders headers) { + Object contentType = headers.get(MessageHeaders.CONTENT_TYPE); + if (contentType instanceof byte[]) { + contentType = new String((byte[]) contentType, StandardCharsets.UTF_8); + contentType = ((String) contentType).replace("\"", ""); + headers = new MutableMessageHeaders(headers); + headers.put(MessageHeaders.CONTENT_TYPE, contentType); + } + return super.getMimeType(headers); + } } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java index decd81623..8741d01b0 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java @@ -17,6 +17,7 @@ package org.springframework.cloud.stream.function; import java.io.Serializable; +import java.nio.charset.StandardCharsets; import java.time.Duration; import java.util.ArrayList; import java.util.Collection; @@ -51,6 +52,7 @@ import org.springframework.integration.handler.LoggingHandler; import org.springframework.integration.scheduling.PollerMetadata; import org.springframework.integration.support.MessageBuilder; import org.springframework.messaging.Message; +import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.support.GenericMessage; import org.springframework.scheduling.support.PeriodicTrigger; @@ -575,6 +577,27 @@ public class ImplicitFunctionBindingTests { } } + @Test + public void contentTypeAsByteArrayTest() { + System.clearProperty("spring.cloud.function.definition"); + try (ConfigurableApplicationContext context = new SpringApplicationBuilder(TestChannelBinderConfiguration + .getCompleteConfiguration(PojoFunctionConfiguration.class)) + .web(WebApplicationType.NONE).run("--spring.cloud.function.definition=echoPerson", + "--spring.jmx.enabled=false")) { + + InputDestination inputDestination = context.getBean(InputDestination.class); + OutputDestination outputDestination = context.getBean(OutputDestination.class); + + Message inputMessage = MessageBuilder.withPayload("{\"name\":\"Jim Lahey\",\"id\":420}".getBytes()) + .setHeader(MessageHeaders.CONTENT_TYPE, "application/json".getBytes(StandardCharsets.UTF_8)) + .build(); + + inputDestination.send(inputMessage, "echoPerson-in-0"); + + assertThat(outputDestination.receive(100, "echoPerson-out-0").getPayload()).isEqualTo("{\"name\":\"Jim Lahey\",\"id\":420}".getBytes()); + } + } + @EnableAutoConfiguration public static class NoEnableBindingConfiguration { @@ -763,6 +786,11 @@ public class ImplicitFunctionBindingTests { @EnableAutoConfiguration public static class PojoFunctionConfiguration { + @Bean + public Function echoPerson() { + return x -> x; + } + @Bean public Function func() { return x -> {