From 8c88cf14d83c9edc85486e36455a0c942445a675 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Mon, 28 Jun 2021 15:47:06 +0200 Subject: [PATCH] GH-2179 provide mechanism to remove short-lived headers Headers such as 'spring.cloud.stream.sendto.destination' are no longer valid once the message has been sent, os it will be removed Resolves #2179 --- .../stream/function/FunctionConfiguration.java | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) 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 61d677e88..1a44e70f6 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 @@ -352,6 +352,11 @@ public class FunctionConfiguration { : MessageBuilder.withPayload(value).build(); } + @SuppressWarnings({ "unchecked", "rawtypes" }) + private static Message sanitize(Message inputMessage) { + return MessageBuilder.fromMessage(inputMessage).removeHeader("spring.cloud.stream.sendto.destination").build(); + } + private static class FunctionToDestinationBinder implements InitializingBean, ApplicationContextAware { protected final Log logger = LogFactory.getLog(getClass()); @@ -466,7 +471,12 @@ public class FunctionConfiguration { + inputBindingName + ".consumer.concurrency=" + consumerProperties.getConcurrency() + "'"); } SubscribableChannel inputChannel = this.applicationContext.getBean(inputBindingName, SubscribableChannel.class); - return IntegrationReactiveUtils.messageChannelToFlux(inputChannel); + return IntegrationReactiveUtils.messageChannelToFlux(inputChannel).map(m -> { + if (m instanceof Message) { + m = sanitize(m); + } + return m; + }); }) .map(publisher -> { if (targetProtocolEnhancer.get() != null) { @@ -713,6 +723,7 @@ public class FunctionConfiguration { @SuppressWarnings("unchecked") @Override public Object apply(Message message) { + message = sanitize(message); Map headersMap = (Map) ReflectionUtils .getField(this.headersField, message.getHeaders()); if (StringUtils.hasText(targetProtocol)) {