From 494c290614b61cfea409fce94a21d7b2b67a5a3c Mon Sep 17 00:00:00 2001 From: arimu1 <19286898+arimu1@users.noreply.github.com> Date: Sun, 13 Sep 2026 07:04:37 +0700 Subject: [PATCH] fix(cloudevent): honor target-protocol for Kafka ce_ prefix Restore reading the target-protocol header when choosing Cloud Event attribute prefixes so cross-binder flows (e.g. AMQP to Kafka) can use ce_ without requiring kafka_* headers on the message. Fixes gh-1461 Signed-off-by: arimu1 <19286898+arimu1@users.noreply.github.com> --- .../cloudevent/CloudEventMessageUtils.java | 21 +++++++++++++++++++ .../CloudEventsFunctionInvocationHelper.java | 18 +++++++++++++++- .../context/message/MessageUtils.java | 7 ++++++- ...CloudEventMessageUtilsAndBuilderTests.java | 15 +++++++++++++ 4 files changed, 59 insertions(+), 2 deletions(-) diff --git a/spring-cloud-function-context/src/main/java/org/springframework/cloud/function/cloudevent/CloudEventMessageUtils.java b/spring-cloud-function-context/src/main/java/org/springframework/cloud/function/cloudevent/CloudEventMessageUtils.java index 867b135f9..57141dee1 100644 --- a/spring-cloud-function-context/src/main/java/org/springframework/cloud/function/cloudevent/CloudEventMessageUtils.java +++ b/spring-cloud-function-context/src/main/java/org/springframework/cloud/function/cloudevent/CloudEventMessageUtils.java @@ -22,6 +22,7 @@ import java.util.Collections; import java.util.HashMap; import java.util.Iterator; +import java.util.Locale; import java.util.Map; import java.util.Objects; import java.util.stream.Collectors; @@ -370,6 +371,10 @@ static String determinePrefixToUse(Map messageHeaders) { } static String extractTargetProtocol(Map messageHeaders) { + String targetProtocol = resolveTargetProtocolHeader(messageHeaders); + if (StringUtils.hasText(targetProtocol)) { + return targetProtocol; + } Iterator keyIterator = messageHeaders.keySet().iterator(); for (; keyIterator.hasNext();) { String key = keyIterator.next(); @@ -383,6 +388,22 @@ else if (key.startsWith("amqp")) { return null; } + private static String resolveTargetProtocolHeader(Map messageHeaders) { + Object targetProtocol = messageHeaders.get(MessageUtils.TARGET_PROTOCOL); + if (targetProtocol == null) { + for (Map.Entry entry : messageHeaders.entrySet()) { + if (MessageUtils.TARGET_PROTOCOL.equalsIgnoreCase(entry.getKey())) { + targetProtocol = entry.getValue(); + break; + } + } + } + if (targetProtocol instanceof String protocol && StringUtils.hasText(protocol)) { + return protocol.toLowerCase(Locale.ROOT); + } + return null; + } + static String determinePrefixToUse(Map messageHeaders, boolean strict) { String targetProtocol = extractTargetProtocol(messageHeaders); String prefix = determinePrefixToUse(targetProtocol); diff --git a/spring-cloud-function-context/src/main/java/org/springframework/cloud/function/cloudevent/CloudEventsFunctionInvocationHelper.java b/spring-cloud-function-context/src/main/java/org/springframework/cloud/function/cloudevent/CloudEventsFunctionInvocationHelper.java index f552590e2..f4ee68074 100644 --- a/spring-cloud-function-context/src/main/java/org/springframework/cloud/function/cloudevent/CloudEventsFunctionInvocationHelper.java +++ b/spring-cloud-function-context/src/main/java/org/springframework/cloud/function/cloudevent/CloudEventsFunctionInvocationHelper.java @@ -17,6 +17,8 @@ package org.springframework.cloud.function.cloudevent; import java.net.URI; +import java.util.HashMap; +import java.util.Map; import java.util.UUID; import org.apache.commons.logging.Log; @@ -99,7 +101,8 @@ public Message postProcessResult(Object result, Message input) { String targetPrefix = CloudEventMessageUtils.DEFAULT_ATTR_PREFIX; if (input != null) { - targetPrefix = CloudEventMessageUtils.determinePrefixToUse(input.getHeaders(), true); + targetPrefix = CloudEventMessageUtils.determinePrefixToUse( + headersForTargetPrefix(input, result), true); } else if (result instanceof Message resultMessage) { targetPrefix = CloudEventMessageUtils.determinePrefixToUse(resultMessage.getHeaders(), true); @@ -119,6 +122,19 @@ public void setApplicationContext(ApplicationContext applicationContext) throws this.applicationContext = (ConfigurableApplicationContext) applicationContext; } + + private static Map headersForTargetPrefix(Message input, Object result) { + Map headers = new HashMap<>(input.getHeaders()); + if (result instanceof Message resultMessage) { + resultMessage.getHeaders().forEach((key, value) -> { + if (value != null) { + headers.putIfAbsent(key, value); + } + }); + } + return headers; + } + private Message doPostProcessResult(Object result, String targetPrefix) { Message resultMessage = null; //result instanceof Message ? (Message) result : null; CloudEventMessageBuilder messageBuilder; diff --git a/spring-cloud-function-context/src/main/java/org/springframework/cloud/function/context/message/MessageUtils.java b/spring-cloud-function-context/src/main/java/org/springframework/cloud/function/context/message/MessageUtils.java index 0d073b4d0..2d4e29bfc 100644 --- a/spring-cloud-function-context/src/main/java/org/springframework/cloud/function/context/message/MessageUtils.java +++ b/spring-cloud-function-context/src/main/java/org/springframework/cloud/function/context/message/MessageUtils.java @@ -32,7 +32,12 @@ public abstract class MessageUtils { */ public static String MESSAGE_TYPE = "message-type"; /** - * Value for 'target-protocol' typically use as header key. + * Header key for the target messaging protocol of an outgoing Cloud Event (e.g. kafka, amqp, http). + */ + public static String TARGET_PROTOCOL = "target-protocol"; + + /** + * Value for 'source-type' typically use as header key. */ public static String SOURCE_TYPE = "source-type"; diff --git a/spring-cloud-function-context/src/test/java/org/springframework/cloud/function/cloudevent/CloudEventMessageUtilsAndBuilderTests.java b/spring-cloud-function-context/src/test/java/org/springframework/cloud/function/cloudevent/CloudEventMessageUtilsAndBuilderTests.java index 71e9cf8af..0a481390f 100644 --- a/spring-cloud-function-context/src/test/java/org/springframework/cloud/function/cloudevent/CloudEventMessageUtilsAndBuilderTests.java +++ b/spring-cloud-function-context/src/test/java/org/springframework/cloud/function/cloudevent/CloudEventMessageUtilsAndBuilderTests.java @@ -21,6 +21,7 @@ import org.junit.jupiter.api.Test; +import org.springframework.cloud.function.context.message.MessageUtils; import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; @@ -98,6 +99,20 @@ public void testAttributeRecognitionAndCanonicalization() { assertThat(CloudEventMessageUtils.getSpecVersion(httpMessage)).isEqualTo("1.0"); } + + @Test + void determinePrefixToUsePrefersTargetProtocolHeaderOverSourceProtocolHeaders() { + Message ceMessage = CloudEventMessageBuilder.withData("foo") + .build(CloudEventMessageUtils.AMQP_ATTR_PREFIX); + ceMessage = MessageBuilder.fromMessage(ceMessage) + .setHeader(MessageUtils.TARGET_PROTOCOL, "kafka") + .setHeader("amqp_correlationId", "123") + .build(); + + assertThat(CloudEventMessageUtils.determinePrefixToUse(ceMessage.getHeaders(), true)) + .isEqualTo(CloudEventMessageUtils.KAFKA_ATTR_PREFIX); + } + @Test void buildWithKafkaPrefixUsesContentTypeHeaderForDataContentType() { Message kafkaMessage = CloudEventMessageBuilder.withData("hello")