@@ -96,7 +96,15 @@ public class CloudEventsFunctionInvocationHelper implements FunctionInvocationHe
|
|||||||
if (this.messageConverter != null && CLOUD_EVENT_CLASS != null && CLOUD_EVENT_CLASS.isAssignableFrom(result.getClass())) {
|
if (this.messageConverter != null && CLOUD_EVENT_CLASS != null && CLOUD_EVENT_CLASS.isAssignableFrom(result.getClass())) {
|
||||||
convertedResult = this.messageConverter.toMessage(result, input.getHeaders());
|
convertedResult = this.messageConverter.toMessage(result, input.getHeaders());
|
||||||
}
|
}
|
||||||
String targetPrefix = CloudEventMessageUtils.determinePrefixToUse(input.getHeaders(), true);
|
|
||||||
|
String targetPrefix = CloudEventMessageUtils.DEFAULT_ATTR_PREFIX;
|
||||||
|
if (input != null) {
|
||||||
|
targetPrefix = CloudEventMessageUtils.determinePrefixToUse(input.getHeaders(), true);
|
||||||
|
}
|
||||||
|
else if (result instanceof Message) {
|
||||||
|
targetPrefix = CloudEventMessageUtils.determinePrefixToUse(((Message<?>) result).getHeaders(), true);
|
||||||
|
}
|
||||||
|
|
||||||
Assert.hasText(targetPrefix, "Unable to determine prefix for Cloud Event atttributes, "
|
Assert.hasText(targetPrefix, "Unable to determine prefix for Cloud Event atttributes, "
|
||||||
+ "which they must have according to protocol specification. Consider adding 'target-protocol' "
|
+ "which they must have according to protocol specification. Consider adding 'target-protocol' "
|
||||||
+ "header with values of one of the supported protocols - [kafka, amqp, http]");
|
+ "header with values of one of the supported protocols - [kafka, amqp, http]");
|
||||||
|
|||||||
Reference in New Issue
Block a user