diff --git a/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/HeaderTests.java b/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/HeaderTests.java index 45302209f0..1b6366b5aa 100644 --- a/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/HeaderTests.java +++ b/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/HeaderTests.java @@ -129,6 +129,24 @@ void checkMessageWrappedFunctionalConsumer() { assertThat(headers.get(MessageHeaders.CONTENT_TYPE)).isEqualTo("application/json"); } + @Test + void timestampHeaderIsPresentOnConsumedMessage() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(FunctionUpperCaseConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false", + "--spring.cloud.function.definition=uppercase")) { + + InputDestination input = context.getBean(InputDestination.class); + input.send(new GenericMessage<>("hello".getBytes()), "uppercase-in-0"); + + OutputDestination output = context.getBean(OutputDestination.class); + Message result = output.receive(1000, "uppercase-out-0"); + + assertThat(result.getHeaders().getTimestamp()).isNotNull(); + } + } + @EnableAutoConfiguration public static class EmptyConfiguration { diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java index 2c0e9bec0f..5ff18d010b 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java @@ -113,6 +113,7 @@ import org.springframework.messaging.MessagingException; import org.springframework.messaging.SubscribableChannel; import org.springframework.messaging.support.ChannelInterceptor; +import org.springframework.messaging.support.GenericMessage; import org.springframework.scheduling.TaskScheduler; import org.springframework.scheduling.Trigger; import org.springframework.scheduling.support.CronTrigger; @@ -409,12 +410,19 @@ private Message wrapToMessageIfNecessary(T value) { } private static

Message

sanitize(Message

inputMessage) { - return MessageBuilder + Message

sanitized = MessageBuilder .fromMessage(inputMessage) .removeHeader("spring.cloud.stream.sendto.destination") -// .setHeader(MessageUtils.SOURCE_TYPE, inputMessage.getHeaders().get(MessageUtils.TARGET_PROTOCOL)) -// .removeHeader(MessageUtils.TARGET_PROTOCOL) .build(); + if (sanitized == inputMessage) { + // MessageBuilder.build() returns the same instance when no header was + // actually modified (fast-path in BaseMessageBuilder), which skips the + // GenericMessage constructor and therefore skips stamping a "timestamp" + // header. Force a fresh GenericMessage so the header is always present, + // restoring pre-4.3.3 behavior. See GH-3211. + sanitized = new GenericMessage<>(inputMessage.getPayload(), inputMessage.getHeaders()); + } + return sanitized; } private static class FunctionToDestinationBinder implements InitializingBean, ApplicationContextAware {