diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceChannelInterceptor.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceChannelInterceptor.java index 98358be53..a290b5a27 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceChannelInterceptor.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceChannelInterceptor.java @@ -26,7 +26,9 @@ import org.springframework.cloud.sleuth.sampler.NeverSampler; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandler; +import org.springframework.messaging.support.GenericMessage; import org.springframework.messaging.support.MessageBuilder; +import org.springframework.messaging.support.MessageHeaderAccessor; /** * A channel interceptor that automatically starts / continues / closes and detaches @@ -84,7 +86,9 @@ public class TraceChannelInterceptor extends AbstractTraceChannelInterceptor { messageBuilder.setHeader(TraceMessageHeaders.MESSAGE_SENT_FROM_CLIENT, true); } getSpanInjector().inject(span, messageBuilder); - return messageBuilder.build(); + MessageHeaderAccessor headers = MessageHeaderAccessor.getMutableAccessor(message); + headers.copyHeaders(messageBuilder.build().getHeaders()); + return new GenericMessage(message.getPayload(), headers.getMessageHeaders()); } private Span startSpan(Span span, String name, Message message) { diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceChannelInterceptorTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceChannelInterceptorTests.java index 2ab5c63af..567f90e39 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceChannelInterceptorTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceChannelInterceptorTests.java @@ -44,6 +44,7 @@ import org.springframework.integration.support.MessageBuilder; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHandler; import org.springframework.messaging.MessagingException; +import org.springframework.messaging.support.MessageHeaderAccessor; import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @@ -117,6 +118,15 @@ public class TraceChannelInterceptorTests implements MessageHandler { then(this.span.isExportable()).isFalse(); } + @Test + public void messageHeadersStillMutable() { + this.tracedChannel.send(MessageBuilder.withPayload("hi") + .setHeader(Span.SAMPLED_NAME, Span.SPAN_NOT_SAMPLED).build()); + assertNotNull("message was null", this.message); + MessageHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(this.message, MessageHeaderAccessor.class); + assertNotNull("Message header accessor should be still available", accessor); + } + @Test public void parentSpanIncluded() { this.tracedChannel.send(MessageBuilder.withPayload("hi")