Merge branch '1.2.x'

This commit is contained in:
Marcin Grzejszczak
2017-08-03 17:56:47 +02:00
2 changed files with 33 additions and 2 deletions

View File

@@ -22,6 +22,7 @@ import org.springframework.cloud.sleuth.Span;
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;
@@ -66,7 +67,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);
@@ -80,7 +82,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 {