diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptor.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptor.java index 5569ce561..ce2d3027e 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptor.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptor.java @@ -324,7 +324,8 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter Object payload = message.getPayload(); if (payload instanceof MessagingException) { MessagingException e = (MessagingException) payload; - return e.getFailedMessage(); + Message failedMessage = e.getFailedMessage(); + return failedMessage != null ? failedMessage : message; } return message; } diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptorTest.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptorTest.java index 0eee574a8..534e1224d 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptorTest.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptorTest.java @@ -263,6 +263,28 @@ public class TracingChannelInterceptorTest { assertThat(this.message.getHeaders().getErrorChannel()).isSameAs(errorsReplyChannel); } + @Test + public void errorMessageHeadersWithNullPayloadRetained() { + this.channel.addInterceptor(this.interceptor); + Map errorChannelHeaders = new HashMap<>(); + errorChannelHeaders.put(TraceMessageHeaders.TRACE_ID_NAME, "000000000000000a"); + errorChannelHeaders.put(TraceMessageHeaders.SPAN_ID_NAME, "000000000000000a"); + this.channel.send(new ErrorMessage(new MessagingException("exception"), + errorChannelHeaders)); + + this.message = this.channel.receive(); + + assertThat(this.message).isNotNull(); + String spanId = this.message.getHeaders().get(TraceMessageHeaders.SPAN_ID_NAME, + String.class); + assertThat(spanId).isNotNull(); + String traceId = this.message.getHeaders().get(TraceMessageHeaders.TRACE_ID_NAME, + String.class); + assertThat(traceId).isEqualTo("000000000000000a"); + assertThat(spanId).isNotEqualTo("000000000000000a"); + assertThat(this.spans).hasSize(2); + } + ChannelInterceptor producerSideOnly(ChannelInterceptor delegate) { return new ChannelInterceptorAdapter() { @Override