From cedacc8b17e6121342f4945b8c9b3593dbaa550e Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Mon, 21 Dec 2020 13:03:20 +0100 Subject: [PATCH] Fix FunctionConfiguration to account for Cloud Events recognition --- pom.xml | 2 +- ...cationJsonMessageMarshallingConverter.java | 7 +++- .../function/FunctionConfiguration.java | 41 ++++++++++++------- 3 files changed, 34 insertions(+), 16 deletions(-) diff --git a/pom.xml b/pom.xml index 8d9568019..25ad5c116 100644 --- a/pom.xml +++ b/pom.xml @@ -25,7 +25,7 @@ 1.8 2020.0.0-RC2 2.1 - 3.1.0-RC1 + 3.1.0-SNAPSHOT true true true 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 18e560afd..3bd8949ab 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 @@ -110,7 +110,12 @@ class ApplicationJsonMessageMarshallingConverter extends MappingJackson2MessageC StandardCharsets.UTF_8); } else { - result = super.convertFromInternal(message, targetClass, conversionHint); + try { + result = super.convertFromInternal(message, targetClass, conversionHint); + } + catch (Exception e) { + // ignore + } } } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java index 495407913..d71d953be 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java @@ -53,6 +53,7 @@ import org.springframework.boot.autoconfigure.AutoConfigureAfter; import org.springframework.boot.autoconfigure.AutoConfigureBefore; import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.cloud.function.cloudevent.CloudEventMessageUtils; import org.springframework.cloud.function.context.FunctionCatalog; import org.springframework.cloud.function.context.FunctionProperties; import org.springframework.cloud.function.context.FunctionRegistry; @@ -90,16 +91,18 @@ import org.springframework.core.env.Environment; import org.springframework.core.type.MethodMetadata; import org.springframework.integration.channel.AbstractMessageChannel; import org.springframework.integration.channel.AbstractSubscribableChannel; +import org.springframework.integration.core.MessagingTemplate; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.IntegrationFlowBuilder; import org.springframework.integration.dsl.IntegrationFlows; -import org.springframework.integration.handler.ServiceActivatingHandler; +import org.springframework.integration.handler.AbstractMessageHandler; import org.springframework.integration.support.MessageBuilder; import org.springframework.integration.util.IntegrationReactiveUtils; import org.springframework.lang.Nullable; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHeaders; +import org.springframework.messaging.MessagingException; import org.springframework.messaging.SubscribableChannel; import org.springframework.scheduling.TaskScheduler; import org.springframework.util.Assert; @@ -445,6 +448,9 @@ public class FunctionConfiguration { Map headersMap = (Map) ReflectionUtils .getField(headersField, ((Message) message).getHeaders()); headersMap.putIfAbsent(MessageUtils.TARGET_PROTOCOL, targetProtocol); + if (CloudEventMessageUtils.isCloudEvent((message))) { + headersMap.putIfAbsent(MessageUtils.MESSAGE_TYPE, CloudEventMessageUtils.CLOUDEVENT_VALUE); + } return message; }); } @@ -527,16 +533,13 @@ public class FunctionConfiguration { String inputDestinationName = inputBindingNames.iterator().next(); Object inputDestination = this.applicationContext.getBean(inputDestinationName); if (inputDestination != null && inputDestination instanceof SubscribableChannel) { - ServiceActivatingHandler handler = createFunctionHandler(function, inputDestinationName, outputDestinationName); - if (StringUtils.hasText(outputDestinationName)) { // consumer implicit or function<.., mono> - handler.setOutputChannelName(outputDestinationName); - } + AbstractMessageHandler handler = createFunctionHandler(function, inputDestinationName, outputDestinationName); ((SubscribableChannel) inputDestination).subscribe(handler); } } } - private ServiceActivatingHandler createFunctionHandler(FunctionInvocationWrapper function, + private AbstractMessageHandler createFunctionHandler(FunctionInvocationWrapper function, String inputChannelName, String outputChannelName) { ConsumerProperties consumerProperties = StringUtils.hasText(inputChannelName) ? this.serviceProperties.getBindingProperties(inputChannelName).getConsumer() @@ -545,22 +548,27 @@ public class FunctionConfiguration { ? this.serviceProperties.getBindingProperties(outputChannelName).getProducer() : null; - ServiceActivatingHandler handler = new ServiceActivatingHandler(new FunctionWrapper(function, consumerProperties, - producerProperties, applicationContext, this.determineTargetProtocol(outputChannelName))) { + FunctionWrapper functionInvocationWrapper = (new FunctionWrapper(function, consumerProperties, + producerProperties, applicationContext, this.determineTargetProtocol(outputChannelName))); + + MessagingTemplate template = new MessagingTemplate(); + template.setBeanFactory(applicationContext.getBeanFactory()); + AbstractMessageHandler handler = new AbstractMessageHandler() { @Override - protected void sendOutputs(Object result, Message requestMessage) { + public void handleMessageInternal(Message message) throws MessagingException { + Object result = functionInvocationWrapper.apply((Message) message); if (result instanceof Iterable) { for (Object resultElement : (Iterable) result) { - this.doSendMessage(resultElement, requestMessage); + this.doSendMessage(resultElement, message); } } else if (ObjectUtils.isArray(result) && !(result instanceof byte[])) { for (int i = 0; i < ((Object[]) result).length; i++) { - this.doSendMessage(((Object[]) result)[i], requestMessage); + this.doSendMessage(((Object[]) result)[i], message); } } else { - this.doSendMessage(result, requestMessage); + this.doSendMessage(result, message); } } @@ -573,11 +581,13 @@ public class FunctionConfiguration { } outputChannel.send(((Message) result)); } - else { - super.sendOutputs(result, requestMessage); + else if (StringUtils.hasText(outputChannelName)) { + template.send(outputChannelName, (Message) result); } } + }; + handler.setBeanFactory(this.applicationContext); handler.afterPropertiesSet(); return handler; @@ -694,6 +704,9 @@ public class FunctionConfiguration { if (StringUtils.hasText(targetProtocol)) { headersMap.putIfAbsent(MessageUtils.TARGET_PROTOCOL, targetProtocol); } + if (CloudEventMessageUtils.isCloudEvent(message)) { + headersMap.putIfAbsent(MessageUtils.MESSAGE_TYPE, CloudEventMessageUtils.CLOUDEVENT_VALUE); + } if (message != null && consumerProperties != null) { headersMap.put(FunctionProperties.SKIP_CONVERSION_HEADER, consumerProperties.isUseNativeDecoding()); }