From bb4064d883bb5d5e19ae6cd22b548862b6663735 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Mon, 29 Nov 2021 19:14:35 +0100 Subject: [PATCH] GH-2090 Fix wild card contentType processing in StreamBridge This fix is dependent on https://github.com/spring-cloud/spring-cloud-function/issues/773 Resolves #2090 --- .../cloud/stream/function/StreamBridge.java | 4 +- .../ImplicitFunctionBindingTests.java | 3 +- .../stream/function/StreamBridgeTests.java | 84 +++++++++++++++++++ 3 files changed, 88 insertions(+), 3 deletions(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java index bab3aa32d..6552e7dde 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java @@ -139,7 +139,9 @@ public final class StreamBridge implements SmartInitializingSingleton { * @return true if data was sent successfully, otherwise false or throws an exception. */ public boolean send(String bindingName, Object data) { - return this.send(bindingName, data, MimeTypeUtils.APPLICATION_JSON); + BindingProperties bindingProperties = this.bindingServiceProperties.getBindingProperties(bindingName); + MimeType contentType = StringUtils.hasText(bindingProperties.getContentType()) ? MimeType.valueOf(bindingProperties.getContentType()) : MimeTypeUtils.APPLICATION_JSON; + return this.send(bindingName, data, contentType); } /** 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 635937b0e..523634436 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 @@ -436,8 +436,7 @@ public class ImplicitFunctionBindingTests { TestChannelBinderConfiguration.getCompleteConfiguration(SingleConsumerConfiguration.class)) .web(WebApplicationType.NONE).run("--spring.cloud.function.definition=consumer", "--spring.jmx.enabled=false", - "--spring.cloud.stream.bindings.input.content-type=text/plain", - "--spring.cloud.stream.bindings.input.consumer.use-native-decoding=true")) { + "--spring.cloud.stream.bindings.consumer-in-0.content-type=text/plain")) { InputDestination source = context.getBean(InputDestination.class); source.send(new GenericMessage("John Doe".getBytes())); diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java index 5a9d19fd9..366e486ca 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java @@ -46,11 +46,16 @@ import org.springframework.integration.config.GlobalChannelInterceptor; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.IntegrationFlows; import org.springframework.integration.handler.LoggingHandler; +import org.springframework.lang.Nullable; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandler; +import org.springframework.messaging.MessageHeaders; +import org.springframework.messaging.converter.AbstractMessageConverter; +import org.springframework.messaging.converter.MessageConverter; import org.springframework.messaging.support.ChannelInterceptor; import org.springframework.messaging.support.MessageBuilder; +import org.springframework.util.MimeType; import org.springframework.util.MimeTypeUtils; import org.springframework.util.ReflectionUtils; @@ -71,6 +76,28 @@ public class StreamBridgeTests { System.clearProperty("spring.cloud.function.definition"); } + @Test + public void testWithOutputContentTypeWildCardBindings() throws Exception { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder(TestChannelBinderConfiguration + .getCompleteConfiguration(ConsumerConfiguration.class, EmptyConfigurationWithCustomConverters.class)) + .web(WebApplicationType.NONE).run( + "--spring.cloud.stream.bindings.foo.content-type=application/*+foo ", + "--spring.cloud.stream.bindings.bar.content-type=application/*+non-registered-foo", + "--spring.jmx.enabled=false")) { + StreamBridge bridge = context.getBean(StreamBridge.class); + bridge.send("foo", "hello foo"); + bridge.send("bar", "hello bar"); + + OutputDestination output = context.getBean(OutputDestination.class); + + assertThat(output.receive(1000, "foo").getHeaders().get(MessageHeaders.CONTENT_TYPE)) + .isEqualTo(MimeType.valueOf("application/json+foo")); + assertThat(output.receive(1000, "bar").getHeaders().get(MessageHeaders.CONTENT_TYPE)) + .isEqualTo(MimeType.valueOf("application/blahblah+non-registered-foo")); + + } + } + @SuppressWarnings("unchecked") @Test public void testNoCachingOfStreamBridgeFunction() throws Exception { @@ -381,6 +408,63 @@ public class StreamBridgeTests { } + @EnableAutoConfiguration + public static class EmptyConfigurationWithCustomConverters { + + @Bean + public MessageConverter fooConverter() { + return new AbstractMessageConverter(MimeType.valueOf("application/json+foo"), MimeType.valueOf("application/json+blah")) { + @Override + protected boolean supports(Class clazz) { + return true; + } + + @Override + @Nullable + protected Object convertFromInternal(Message message, Class targetClass, @Nullable Object conversionHint) { + return message.getPayload(); + } + + @Override + @Nullable + protected Object convertToInternal(Object payload, @Nullable MessageHeaders headers, @Nullable Object conversionHint) { + if (headers.containsKey(MessageHeaders.CONTENT_TYPE) && + (headers.get(MessageHeaders.CONTENT_TYPE).toString().endsWith("+foo") || + headers.get(MessageHeaders.CONTENT_TYPE).toString().endsWith("+blah"))) { + return payload; + } + return null; + } + }; + } + + @Bean + public MessageConverter barConverter() { + return new AbstractMessageConverter(MimeType.valueOf("application/blahblah+non-registered-foo")) { + @Override + protected boolean supports(Class clazz) { + return true; + } + + @Override + @Nullable + protected Object convertFromInternal(Message message, Class targetClass, @Nullable Object conversionHint) { + return message.getPayload(); + } + + @Override + @Nullable + protected Object convertToInternal(Object payload, @Nullable MessageHeaders headers, @Nullable Object conversionHint) { + if (headers.containsKey(MessageHeaders.CONTENT_TYPE) && + (headers.get(MessageHeaders.CONTENT_TYPE).toString().endsWith("+non-registered-foo"))) { + return payload; + } + return null; + } + }; + } + } + @EnableAutoConfiguration public static class ConsumerConfiguration { @Bean