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
This commit is contained in:
Oleg Zhurakousky
2021-11-29 19:14:35 +01:00
parent ada1f4c3a5
commit bb4064d883
3 changed files with 88 additions and 3 deletions

View File

@@ -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);
}
/**

View File

@@ -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<byte[]>("John Doe".getBytes()));

View File

@@ -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