Added support for failed messages

without this change when messages arrive at a error channel a new traceid is set
with this change we continue the trace

fixes #667
This commit is contained in:
Marcin Grzejszczak
2017-08-03 17:50:47 +02:00
parent e8d30be75e
commit aadf0a8334
2 changed files with 33 additions and 2 deletions

View File

@@ -24,6 +24,7 @@ import org.springframework.cloud.sleuth.Tracer;
import org.springframework.cloud.sleuth.sampler.NeverSampler;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageDeliveryException;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.support.GenericMessage;
import org.springframework.messaging.support.MessageBuilder;
@@ -75,7 +76,8 @@ public class TraceChannelInterceptor extends AbstractTraceChannelInterceptor {
@Override
public Message<?> preSend(Message<?> message, MessageChannel channel) {
MessageBuilder<?> messageBuilder = MessageBuilder.fromMessage(message);
Message<?> retrievedMessage = getMessage(message);
MessageBuilder<?> messageBuilder = MessageBuilder.fromMessage(retrievedMessage);
Span parentSpan = getTracer().isTracing() ? getTracer().getCurrentSpan()
: buildSpan(new MessagingTextMap(messageBuilder));
String name = getMessageChannelName(channel);
@@ -89,7 +91,16 @@ public class TraceChannelInterceptor extends AbstractTraceChannelInterceptor {
getSpanInjector().inject(span, new MessagingTextMap(messageBuilder));
MessageHeaderAccessor headers = MessageHeaderAccessor.getMutableAccessor(message);
headers.copyHeaders(messageBuilder.build().getHeaders());
return new GenericMessage<Object>(message.getPayload(), headers.getMessageHeaders());
return new GenericMessage<>(retrievedMessage.getPayload(), headers.getMessageHeaders());
}
private Message getMessage(Message<?> message) {
Object payload = message.getPayload();
if (payload instanceof MessageDeliveryException) {
MessageDeliveryException e = (MessageDeliveryException) payload;
return e.getFailedMessage();
}
return message;
}
private Span startSpan(Span span, String name, Message<?> message) {

View File

@@ -44,6 +44,7 @@ import org.springframework.integration.core.MessagingTemplate;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageDeliveryException;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.MessagingException;
import org.springframework.messaging.support.ChannelInterceptor;
@@ -321,6 +322,25 @@ public class TraceChannelInterceptorTests implements MessageHandler {
this.tracedChannel.removeInterceptor(immutableMessageInterceptor);
}
@Test
public void workWithMessageDeliveryException() throws Exception {
Message<?> message = new GenericMessage<>(new MessageDeliveryException(
MessageBuilder.withPayload("hi")
.setHeader(TraceMessageHeaders.TRACE_ID_NAME, Span.idToHex(10L))
.setHeader(TraceMessageHeaders.SPAN_ID_NAME, Span.idToHex(20L)).build()
));
this.tracedChannel.send(message);
String spanId = this.message.getHeaders().get(TraceMessageHeaders.SPAN_ID_NAME, String.class);
then(spanId).isNotNull();
long traceId = Span
.hexToId(this.message.getHeaders().get(TraceMessageHeaders.TRACE_ID_NAME, String.class));
then(traceId).isEqualTo(10L);
then(spanId).isNotEqualTo(20L);
then(this.accumulator.getSpans()).hasSize(1);
}
@Configuration
@EnableAutoConfiguration
static class App {