Removes message propagating interceptor
Now that the span data is stored in a header it is safe to remove the slightly clunky propagation implentation that used a subclass. Also removed the Stomp* features because they were diverging from the mainstream integration support and no-one seems to understand why they are needed. If the original author of #96 can explain why they were needed we can ask for a new PR to re-instate a version that works with the new model.
This commit is contained in:
@@ -1,120 +0,0 @@
|
||||
/*
|
||||
* Copyright 2015 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
* You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.cloud.sleuth.instrument.integration;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.TreeMap;
|
||||
|
||||
import org.springframework.cloud.sleuth.Span;
|
||||
import org.springframework.cloud.sleuth.trace.SpanContextHolder;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.simp.SimpMessageHeaderAccessor;
|
||||
import org.springframework.messaging.simp.SimpMessageType;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
/**
|
||||
* Builder class to create STOMP message
|
||||
*
|
||||
* @author Gaurav Rai Mazra
|
||||
*
|
||||
*/
|
||||
public class StompMessageBuilder {
|
||||
|
||||
public static StompMessageBuilder fromMessage(Message<?> message) {
|
||||
return new StompMessageBuilder(message);
|
||||
}
|
||||
|
||||
private Map<String, Object> headers = new TreeMap<String, Object>();
|
||||
private Message<?> message;
|
||||
|
||||
public StompMessageBuilder(final Message<?> message) {
|
||||
this.message = message;
|
||||
this.headers.putAll(message.getHeaders());
|
||||
}
|
||||
|
||||
public StompMessageBuilder setHeader(String key, Object value) {
|
||||
this.headers.put(key, value);
|
||||
return this;
|
||||
}
|
||||
|
||||
public StompMessageBuilder setHeaderIfAbsent(String key, Object value) {
|
||||
if (this.headers.get(key) == null)
|
||||
this.headers.put(key, value);
|
||||
|
||||
return this;
|
||||
}
|
||||
|
||||
public StompMessageBuilder setHeadersFromSpan(final Span span) {
|
||||
if (span != null) {
|
||||
setHeaderIfAbsent(Span.SPAN_ID_NAME, Span.toHex(span.getSpanId()));
|
||||
setHeaderIfAbsent(Span.TRACE_ID_NAME, Span.toHex(span.getTraceId()));
|
||||
setHeaderIfAbsent(Span.SPAN_NAME_NAME, span.getName());
|
||||
Long parentId = getParentId(SpanContextHolder.getCurrentSpan());
|
||||
if (parentId != null)
|
||||
setHeaderIfAbsent(Span.PARENT_ID_NAME, Span.toHex(parentId));
|
||||
|
||||
String processId = span.getProcessId();
|
||||
if (StringUtils.hasText(processId))
|
||||
setHeaderIfAbsent(Span.PROCESS_ID_NAME, processId);
|
||||
}
|
||||
return this;
|
||||
}
|
||||
|
||||
public Message<?> build() {
|
||||
SimpMessageHeaderAccessor headerAccessor = SimpMessageHeaderAccessor.create(SimpMessageType.MESSAGE);
|
||||
for (Map.Entry<String, Object> entry : this.headers.entrySet()) {
|
||||
String key = entry.getKey();
|
||||
if (key != null) {
|
||||
Object value = entry.getValue();
|
||||
pushHeaders(headerAccessor, key, value);
|
||||
}
|
||||
}
|
||||
return org.springframework.messaging.support.MessageBuilder.createMessage(this.message.getPayload(),
|
||||
headerAccessor.getMessageHeaders());
|
||||
}
|
||||
|
||||
private void pushHeaders(final SimpMessageHeaderAccessor accessor, final String key, final Object value) {
|
||||
switch (key) {
|
||||
case SimpMessageHeaderAccessor.DESTINATION_HEADER:
|
||||
case SimpMessageHeaderAccessor.MESSAGE_TYPE_HEADER:
|
||||
case SimpMessageHeaderAccessor.SESSION_ID_HEADER:
|
||||
case SimpMessageHeaderAccessor.SESSION_ATTRIBUTES:
|
||||
case SimpMessageHeaderAccessor.SUBSCRIPTION_ID_HEADER:
|
||||
case SimpMessageHeaderAccessor.USER_HEADER:
|
||||
case SimpMessageHeaderAccessor.CONNECT_MESSAGE_HEADER:
|
||||
case SimpMessageHeaderAccessor.HEART_BEAT_HEADER:
|
||||
case SimpMessageHeaderAccessor.ORIGINAL_DESTINATION:
|
||||
case SimpMessageHeaderAccessor.IGNORE_ERROR:
|
||||
case Span.NOT_SAMPLED_NAME:
|
||||
case Span.PARENT_ID_NAME:
|
||||
case Span.PROCESS_ID_NAME:
|
||||
case Span.SPAN_ID_NAME:
|
||||
case Span.SPAN_NAME_NAME:
|
||||
case Span.TRACE_ID_NAME:
|
||||
accessor.setHeader(key, value);
|
||||
break;
|
||||
default:
|
||||
accessor.setNativeHeader(key, value == null ? null : value.toString());
|
||||
}
|
||||
}
|
||||
|
||||
private Long getParentId(final Span currentSpan) {
|
||||
List<Long> parents = currentSpan.getParents();
|
||||
return parents.isEmpty() ? null : parents.get(0);
|
||||
}
|
||||
}
|
||||
@@ -1,170 +0,0 @@
|
||||
/*
|
||||
* Copyright 2013-2015 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
* You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.cloud.sleuth.instrument.integration;
|
||||
|
||||
import org.springframework.aop.support.AopUtils;
|
||||
import org.springframework.cloud.sleuth.Span;
|
||||
import org.springframework.cloud.sleuth.Tracer;
|
||||
import org.springframework.integration.channel.DirectChannel;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
import org.springframework.messaging.MessageHandler;
|
||||
import org.springframework.messaging.MessageHeaders;
|
||||
import org.springframework.messaging.support.ChannelInterceptorAdapter;
|
||||
import org.springframework.messaging.support.ExecutorChannelInterceptor;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
|
||||
/**
|
||||
* The {@link ExecutorChannelInterceptor} implementation responsible for the {@link Span}
|
||||
* propagation from one message flow's thread to another through the
|
||||
* {@link MessageChannel}s involved in the flow.
|
||||
* <p>
|
||||
* In addition this interceptor cleans up (restores) the {@link Span} in the containers
|
||||
* Threads for channels like
|
||||
* {@link org.springframework.integration.channel.ExecutorChannel} and
|
||||
* {@link org.springframework.integration.channel.QueueChannel}.
|
||||
* @author Spencer Gibb
|
||||
* @since 1.0
|
||||
*/
|
||||
public class TraceContextPropagationChannelInterceptor extends ChannelInterceptorAdapter
|
||||
implements ExecutorChannelInterceptor {
|
||||
|
||||
private final Tracer tracer;
|
||||
|
||||
private final static ThreadLocal<Span> ORIGINAL_CONTEXT = new ThreadLocal<>();
|
||||
|
||||
public TraceContextPropagationChannelInterceptor(Tracer tracer) {
|
||||
this.tracer = tracer;
|
||||
}
|
||||
|
||||
@Override
|
||||
public final Message<?> preSend(Message<?> message, MessageChannel channel) {
|
||||
if (DirectChannel.class.isAssignableFrom(AopUtils.getTargetClass(channel))) {
|
||||
return message;
|
||||
}
|
||||
Span span = this.tracer.getCurrentSpan();
|
||||
if (span != null) {
|
||||
return new MessageWithSpan(message, span);
|
||||
}
|
||||
else {
|
||||
return message;
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public final Message<?> postReceive(Message<?> message, MessageChannel channel) {
|
||||
if (message instanceof MessageWithSpan) {
|
||||
MessageWithSpan messageWithSpan = (MessageWithSpan) message;
|
||||
Message<?> messageToHandle = messageWithSpan.message;
|
||||
populatePropagatedContext(messageWithSpan.span, messageToHandle, channel);
|
||||
return message;
|
||||
}
|
||||
return message;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void afterMessageHandled(Message<?> message, MessageChannel channel,
|
||||
MessageHandler handler, Exception ex) {
|
||||
resetPropagatedContext();
|
||||
}
|
||||
|
||||
@Override
|
||||
public final Message<?> beforeHandle(Message<?> message, MessageChannel channel,
|
||||
MessageHandler handler) {
|
||||
return postReceive(message, channel);
|
||||
}
|
||||
|
||||
private Long getParentId(Span span) {
|
||||
return !span.getParents().isEmpty()
|
||||
? span.getParents().get(0) : null;
|
||||
}
|
||||
|
||||
protected void populatePropagatedContext(Span span, Message<?> message,
|
||||
MessageChannel channel) {
|
||||
if (span != null) {
|
||||
ORIGINAL_CONTEXT.set(this.tracer.continueSpan(span).getSavedSpan());
|
||||
}
|
||||
}
|
||||
|
||||
protected void resetPropagatedContext() {
|
||||
Span originalContext = ORIGINAL_CONTEXT.get();
|
||||
this.tracer.detach(originalContext);
|
||||
ORIGINAL_CONTEXT.remove();
|
||||
}
|
||||
|
||||
private class MessageWithSpan implements Message<Object> {
|
||||
|
||||
private final Message<?> message;
|
||||
|
||||
private final Span span;
|
||||
|
||||
private final MessageHeaders messageHeaders;
|
||||
|
||||
public MessageWithSpan(Message<?> message, Span span) {
|
||||
Assert.notNull(message, "message can not be null");
|
||||
Assert.notNull(span, "span can not be null");
|
||||
this.message = message;
|
||||
this.span = span;
|
||||
|
||||
Map<String, Object> headers = new HashMap<>();
|
||||
headers.putAll(message.getHeaders());
|
||||
|
||||
setHeader(headers, Span.SPAN_ID_NAME, this.span.getSpanId());
|
||||
setHeader(headers, Span.TRACE_ID_NAME, this.span.getTraceId());
|
||||
setHeader(headers, Span.SPAN_NAME_NAME, this.span.getName());
|
||||
Long parentId = getParentId(span);
|
||||
if (parentId != null) {
|
||||
setHeader(headers, Span.PARENT_ID_NAME, parentId);
|
||||
}
|
||||
String processId = span.getProcessId();
|
||||
if (StringUtils.hasText(processId)) {
|
||||
setHeader(headers, Span.PROCESS_ID_NAME, processId);
|
||||
}
|
||||
this.messageHeaders = new MessageHeaders(headers);
|
||||
}
|
||||
|
||||
public void setHeader(Map<String, Object> headers, String name, String value) {
|
||||
if (!headers.containsKey(name)) {
|
||||
headers.put(name, value);
|
||||
}
|
||||
}
|
||||
public void setHeader(Map<String, Object> headers, String name, long value) {
|
||||
setHeader(headers, name, Span.toHex(value));
|
||||
}
|
||||
|
||||
@Override
|
||||
public Object getPayload() {
|
||||
return this.message.getPayload();
|
||||
}
|
||||
|
||||
@Override
|
||||
public MessageHeaders getHeaders() {
|
||||
return this.messageHeaders;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return "MessageWithSpan{" + "message=" + this.message + ", span=" + this.span
|
||||
+ ", messageHeaders=" + this.messageHeaders + '}';
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
@@ -41,13 +41,6 @@ import org.springframework.integration.config.GlobalChannelInterceptor;
|
||||
@EnableConfigurationProperties(TraceKeys.class)
|
||||
public class TraceSpringIntegrationAutoConfiguration {
|
||||
|
||||
@Bean
|
||||
@GlobalChannelInterceptor
|
||||
public TraceContextPropagationChannelInterceptor traceContextPropagationChannelInterceptor(
|
||||
Tracer tracer) {
|
||||
return new TraceContextPropagationChannelInterceptor(tracer);
|
||||
}
|
||||
|
||||
@Bean
|
||||
@GlobalChannelInterceptor
|
||||
public TraceChannelInterceptor traceChannelInterceptor(Tracer tracer,
|
||||
@@ -55,16 +48,4 @@ public class TraceSpringIntegrationAutoConfiguration {
|
||||
return new TraceChannelInterceptor(tracer, traceKeys, random);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public TraceStompMessageChannelInterceptor traceStompMessageChannelInterceptor(
|
||||
Tracer tracer, TraceKeys traceKeys, Random random) {
|
||||
return new TraceStompMessageChannelInterceptor(tracer, traceKeys, random);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public TraceStompMessageContextPropagationChannelInterceptor traceStompMessageContextPropagationChannelInteceptor(
|
||||
Tracer tracer, TraceKeys traceKeys) {
|
||||
return new TraceStompMessageContextPropagationChannelInterceptor(tracer, traceKeys);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -1,66 +0,0 @@
|
||||
/*
|
||||
* Copyright 2015 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
* You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
package org.springframework.cloud.sleuth.instrument.integration;
|
||||
|
||||
import org.springframework.cloud.sleuth.Span;
|
||||
import org.springframework.cloud.sleuth.Tracer;
|
||||
import org.springframework.cloud.sleuth.instrument.TraceKeys;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
import org.springframework.messaging.support.ChannelInterceptor;
|
||||
|
||||
import java.util.Random;
|
||||
|
||||
/**
|
||||
* Interceptor for Stomp Messages sent over websocket
|
||||
*
|
||||
* @author Gaurav Rai Mazra
|
||||
* @author Marcin Grzejszczak
|
||||
*
|
||||
*/
|
||||
public class TraceStompMessageChannelInterceptor extends AbstractTraceChannelInterceptor implements ChannelInterceptor {
|
||||
private ThreadLocal<Span> traceScopeHolder = new ThreadLocal<>();
|
||||
|
||||
public TraceStompMessageChannelInterceptor(Tracer tracer, TraceKeys traceKeys, Random random) {
|
||||
super(tracer, traceKeys, random);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Message<?> preSend(Message<?> message, MessageChannel channel) {
|
||||
if (getTracer().isTracing() || message.getHeaders().containsKey(Span.NOT_SAMPLED_NAME)) {
|
||||
return StompMessageBuilder.fromMessage(message).setHeadersFromSpan(getTracer().getCurrentSpan()).build();
|
||||
}
|
||||
String name = getMessageChannelName(channel);
|
||||
Span span = startSpan(buildSpan(message), name);
|
||||
this.traceScopeHolder.set(span);
|
||||
return StompMessageBuilder.fromMessage(message).setHeadersFromSpan(span).build();
|
||||
}
|
||||
|
||||
private Span startSpan(Span span, String name) {
|
||||
if (span != null) {
|
||||
return getTracer().joinTrace(name, span);
|
||||
}
|
||||
return getTracer().startTrace(name);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void postSend(Message<?> message, MessageChannel channel, boolean sent) {
|
||||
final ThreadLocal<Span> traceScopeHolder = this.traceScopeHolder;
|
||||
Span traceInScope = traceScopeHolder.get();
|
||||
getTracer().close(traceInScope);
|
||||
traceScopeHolder.remove();
|
||||
}
|
||||
}
|
||||
@@ -1,135 +0,0 @@
|
||||
/*
|
||||
* Copyright 2015 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
* You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.cloud.sleuth.instrument.integration;
|
||||
|
||||
import org.springframework.aop.support.AopUtils;
|
||||
import org.springframework.cloud.sleuth.Span;
|
||||
import org.springframework.cloud.sleuth.Tracer;
|
||||
import org.springframework.cloud.sleuth.instrument.TraceKeys;
|
||||
import org.springframework.integration.channel.DirectChannel;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
import org.springframework.messaging.MessageHandler;
|
||||
import org.springframework.messaging.MessageHeaders;
|
||||
import org.springframework.messaging.support.ChannelInterceptorAdapter;
|
||||
import org.springframework.messaging.support.ExecutorChannelInterceptor;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
*
|
||||
* @author Gaurav Rai Mazra
|
||||
* @author Marcin Grzejszczak
|
||||
*
|
||||
*/
|
||||
public class TraceStompMessageContextPropagationChannelInterceptor extends ChannelInterceptorAdapter
|
||||
implements ExecutorChannelInterceptor {
|
||||
|
||||
private final Tracer tracer;
|
||||
private final static ThreadLocal<Span> ORIGINAL_CONTEXT = new ThreadLocal<>();
|
||||
private TraceKeys traceKeys;
|
||||
|
||||
public TraceStompMessageContextPropagationChannelInterceptor(Tracer tracer, TraceKeys traceKeys) {
|
||||
this.tracer = tracer;
|
||||
this.traceKeys = traceKeys;
|
||||
}
|
||||
|
||||
@Override
|
||||
public final Message<?> preSend(Message<?> message, MessageChannel channel) {
|
||||
if (DirectChannel.class.isAssignableFrom(AopUtils.getTargetClass(channel))) {
|
||||
return message;
|
||||
}
|
||||
Span span = this.tracer.getCurrentSpan();
|
||||
if (span != null) {
|
||||
return createMessageWithSpan(message, span);
|
||||
} else {
|
||||
return message;
|
||||
}
|
||||
}
|
||||
|
||||
private Message<?> createMessageWithSpan(Message<?> message, Span span) {
|
||||
MessageWithSpan output = new MessageWithSpan(message, span);
|
||||
addAnnotationsToSpanFromMessage(output, span);
|
||||
return output;
|
||||
}
|
||||
|
||||
private void addAnnotationsToSpanFromMessage(Message<?> message, Span span) {
|
||||
SpanMessageHeaders.addAnnotations(this.traceKeys, message, span);
|
||||
SpanMessageHeaders.addPayloadAnnotations(this.traceKeys, message.getPayload(), span);
|
||||
}
|
||||
|
||||
@Override
|
||||
public final Message<?> postReceive(Message<?> message, MessageChannel channel) {
|
||||
if (message instanceof MessageWithSpan) {
|
||||
MessageWithSpan messageWithSpan = (MessageWithSpan) message;
|
||||
populatePropagatedContext(messageWithSpan.span);
|
||||
return message;
|
||||
}
|
||||
return message;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void afterMessageHandled(Message<?> message, MessageChannel channel, MessageHandler handler, Exception ex) {
|
||||
resetPropagatedContext();
|
||||
}
|
||||
|
||||
@Override
|
||||
public final Message<?> beforeHandle(Message<?> message, MessageChannel channel, MessageHandler handler) {
|
||||
return postReceive(message, channel);
|
||||
}
|
||||
|
||||
protected void populatePropagatedContext(Span span) {
|
||||
if (span != null) {
|
||||
ORIGINAL_CONTEXT.set(this.tracer.continueSpan(span).getSavedSpan());
|
||||
}
|
||||
}
|
||||
|
||||
protected void resetPropagatedContext() {
|
||||
Span originalContext = ORIGINAL_CONTEXT.get();
|
||||
this.tracer.detach(originalContext);
|
||||
ORIGINAL_CONTEXT.remove();
|
||||
}
|
||||
|
||||
private class MessageWithSpan implements Message<Object> {
|
||||
|
||||
private final Message<?> message;
|
||||
private final Span span;
|
||||
|
||||
public MessageWithSpan(Message<?> message, Span span) {
|
||||
Assert.notNull(message, "message can not be null");
|
||||
Assert.notNull(span, "span can not be null");
|
||||
this.span = span;
|
||||
this.message = StompMessageBuilder.fromMessage(message).setHeadersFromSpan(this.span).build();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Object getPayload() {
|
||||
return this.message.getPayload();
|
||||
}
|
||||
|
||||
@Override
|
||||
public MessageHeaders getHeaders() {
|
||||
return this.message.getHeaders();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return "MessageWithSpan{" + "message=" + this.message + ", span=" + this.span + "}";
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
@@ -1,69 +0,0 @@
|
||||
package org.springframework.cloud.sleuth.instrument.integration;
|
||||
|
||||
import static org.assertj.core.api.BDDAssertions.then;
|
||||
|
||||
import org.junit.After;
|
||||
import org.junit.Before;
|
||||
import org.junit.runner.RunWith;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.beans.factory.annotation.Qualifier;
|
||||
import org.springframework.cloud.sleuth.Span;
|
||||
import org.springframework.cloud.sleuth.Tracer;
|
||||
import org.springframework.cloud.sleuth.sampler.AlwaysSampler;
|
||||
import org.springframework.cloud.sleuth.trace.SpanContextHolder;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.support.ExecutorSubscribableChannel;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
|
||||
|
||||
@RunWith(SpringJUnit4ClassRunner.class)
|
||||
public abstract class AbstractTraceStompIntegrationTests {
|
||||
|
||||
@Autowired
|
||||
@Qualifier("executorSubscribableChannel")
|
||||
ExecutorSubscribableChannel channel;
|
||||
@Autowired
|
||||
Tracer tracer;
|
||||
@Autowired StompMessageHandler stompMessageHandler;
|
||||
@Autowired AlwaysSampler sampler;
|
||||
|
||||
@Before
|
||||
public void init() {
|
||||
this.channel.subscribe(this.stompMessageHandler);
|
||||
}
|
||||
|
||||
@After
|
||||
public void close() {
|
||||
SpanContextHolder.removeCurrentSpan();
|
||||
this.channel.unsubscribe(this.stompMessageHandler);
|
||||
}
|
||||
|
||||
Span givenALocallyStartedSpan() {
|
||||
return this.tracer.startTrace("testSendMessage", this.sampler);
|
||||
}
|
||||
|
||||
Message<?> givenMessageToBeSampled() {
|
||||
return StompMessageBuilder.fromMessage(new GenericMessage<>("Message2")).build();
|
||||
}
|
||||
|
||||
void whenTheMessageWasSent(Message<?> message) {
|
||||
this.channel.send(message);
|
||||
then(this.stompMessageHandler.message).isNotNull();
|
||||
}
|
||||
|
||||
Long thenSpanIdFromHeadersIsNotEmpty() {
|
||||
Long header = Span.fromHex(getValueFromHeaders(Span.SPAN_ID_NAME, String.class));
|
||||
then(header).as("Span id should not be empty").isNotNull();
|
||||
return header;
|
||||
}
|
||||
|
||||
Long thenTraceIdFromHeadersIsNotEmpty() {
|
||||
Long header = Span.fromHex(getValueFromHeaders(Span.TRACE_ID_NAME, String.class));
|
||||
then(header).as("Trace id should not be empty").isNotNull();
|
||||
return header;
|
||||
}
|
||||
|
||||
<T> T getValueFromHeaders(String headerName, Class<T> type) {
|
||||
return this.stompMessageHandler.message.getHeaders().get(headerName, type);
|
||||
}
|
||||
}
|
||||
@@ -1,15 +0,0 @@
|
||||
package org.springframework.cloud.sleuth.instrument.integration;
|
||||
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessageHandler;
|
||||
import org.springframework.messaging.MessagingException;
|
||||
|
||||
class StompMessageHandler implements MessageHandler {
|
||||
|
||||
Message<?> message;
|
||||
|
||||
@Override
|
||||
public void handleMessage(Message<?> message) throws MessagingException {
|
||||
this.message = message;
|
||||
}
|
||||
}
|
||||
@@ -1,97 +0,0 @@
|
||||
package org.springframework.cloud.sleuth.instrument.integration;
|
||||
|
||||
import static org.assertj.core.api.BDDAssertions.then;
|
||||
|
||||
import org.junit.Test;
|
||||
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
|
||||
import org.springframework.boot.test.IntegrationTest;
|
||||
import org.springframework.boot.test.SpringApplicationConfiguration;
|
||||
import org.springframework.cloud.sleuth.Span;
|
||||
import org.springframework.cloud.sleuth.instrument.integration.TraceStompMessageChannelInterceptorTests.TestApplication;
|
||||
import org.springframework.cloud.sleuth.sampler.AlwaysSampler;
|
||||
import org.springframework.cloud.sleuth.trace.SpanContextHolder;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.support.ExecutorSubscribableChannel;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
|
||||
/**
|
||||
*
|
||||
* @author Gaurav Rai Mazra
|
||||
*
|
||||
*/
|
||||
|
||||
@SpringApplicationConfiguration(classes = TestApplication.class)
|
||||
@IntegrationTest
|
||||
public class TraceStompMessageChannelInterceptorTests extends AbstractTraceStompIntegrationTests {
|
||||
|
||||
@Test
|
||||
public void should_not_create_span_if_message_contains_not_sampled_header() {
|
||||
Message<?> message = givenMessageNotToBeSampled();
|
||||
|
||||
whenTheMessageWasSent(message);
|
||||
|
||||
thenSpanIdFromHeadersIsEmpty();
|
||||
thenReceivedMessageIsEqualToTheSentOne(message);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void should_create_span_when_headers_dont_contain_not_sampled() {
|
||||
Message<?> message = givenMessageToBeSampled();
|
||||
|
||||
whenTheMessageWasSent(message);
|
||||
|
||||
thenSpanIdFromHeadersIsNotEmpty();
|
||||
thenTraceIdFromHeadersIsNotEmpty();
|
||||
then(SpanContextHolder.getCurrentSpan()).isNull();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void should_propagate_headers_when_message_was_sent_during_local_span_starting() {
|
||||
Span span = givenALocallyStartedSpan();
|
||||
Message<?> message = givenMessageToBeSampled();
|
||||
|
||||
whenTheMessageWasSent(message);
|
||||
this.tracer.close(span);
|
||||
|
||||
Long spanId = thenSpanIdFromHeadersIsNotEmpty();
|
||||
long traceId = thenTraceIdFromHeadersIsNotEmpty();
|
||||
then(traceId).isEqualTo(span.getTraceId());
|
||||
then(spanId).isEqualTo(span.getSpanId());
|
||||
then(SpanContextHolder.getCurrentSpan()).isNull();
|
||||
}
|
||||
|
||||
private Message<?> givenMessageNotToBeSampled() {
|
||||
return StompMessageBuilder.fromMessage(new GenericMessage<>("Message2")).setHeader(Span.NOT_SAMPLED_NAME, "").build();
|
||||
}
|
||||
|
||||
private String thenSpanIdFromHeadersIsEmpty() {
|
||||
String header = getValueFromHeaders(Span.SPAN_ID_NAME, String.class);
|
||||
then(header).as("Span id should be empty").isNullOrEmpty();
|
||||
return header;
|
||||
}
|
||||
|
||||
private void thenReceivedMessageIsEqualToTheSentOne(Message<?> message) {
|
||||
then(message.getPayload()).isEqualTo(this.stompMessageHandler.message.getPayload());
|
||||
}
|
||||
|
||||
@Configuration
|
||||
@EnableAutoConfiguration
|
||||
static class TestApplication {
|
||||
|
||||
@Bean ExecutorSubscribableChannel executorSubscribableChannel(TraceStompMessageChannelInterceptor stompChannelInterceptor) {
|
||||
ExecutorSubscribableChannel channel = new ExecutorSubscribableChannel();
|
||||
channel.addInterceptor(stompChannelInterceptor);
|
||||
return channel;
|
||||
}
|
||||
|
||||
@Bean StompMessageHandler stompMessageHandler() {
|
||||
return new StompMessageHandler();
|
||||
}
|
||||
|
||||
@Bean AlwaysSampler alwaysSampler() {
|
||||
return new AlwaysSampler();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,70 +0,0 @@
|
||||
package org.springframework.cloud.sleuth.instrument.integration;
|
||||
|
||||
import static org.assertj.core.api.BDDAssertions.then;
|
||||
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
|
||||
import org.springframework.boot.test.IntegrationTest;
|
||||
import org.springframework.boot.test.SpringApplicationConfiguration;
|
||||
import org.springframework.cloud.sleuth.Span;
|
||||
import org.springframework.cloud.sleuth.instrument.integration.TraceStompMessageContextPropagationChannelInterceptorTests.TestApplication;
|
||||
import org.springframework.cloud.sleuth.sampler.AlwaysSampler;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.support.ExecutorSubscribableChannel;
|
||||
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
|
||||
|
||||
/**
|
||||
*
|
||||
* @author Gaurav Rai Mazra
|
||||
*
|
||||
*/
|
||||
@RunWith(SpringJUnit4ClassRunner.class)
|
||||
@SpringApplicationConfiguration(classes=TestApplication.class)
|
||||
@IntegrationTest
|
||||
public class TraceStompMessageContextPropagationChannelInterceptorTests extends AbstractTraceStompIntegrationTests {
|
||||
|
||||
@Test
|
||||
public void should_propagate_span_information() {
|
||||
Span span = givenALocallyStartedSpan();
|
||||
Message<?> m = givenMessageToBeSampled();
|
||||
|
||||
whenTheMessageWasSent(m);
|
||||
Long expectedTraceId = span.getTraceId();
|
||||
this.tracer.close(span);
|
||||
|
||||
thenReceivedMessageIsNotNull();
|
||||
long traceId = thenTraceIdFromHeadersIsNotEmpty();
|
||||
then(traceId).isEqualTo(expectedTraceId);
|
||||
thenSpanIdFromHeadersIsNotEmpty();
|
||||
}
|
||||
|
||||
private void thenReceivedMessageIsNotNull() {
|
||||
Message<?> message = this.stompMessageHandler.message;
|
||||
then(message).isNotNull();
|
||||
}
|
||||
|
||||
@Configuration
|
||||
@EnableAutoConfiguration
|
||||
static class TestApplication {
|
||||
|
||||
@Bean ExecutorSubscribableChannel executorSubscribableChannel(
|
||||
TraceStompMessageChannelInterceptor stompChannelInterceptor,
|
||||
TraceStompMessageContextPropagationChannelInterceptor stompMessageContextChannelInterceptor) {
|
||||
ExecutorSubscribableChannel channel = new ExecutorSubscribableChannel();
|
||||
channel.addInterceptor(stompChannelInterceptor);
|
||||
channel.addInterceptor(stompMessageContextChannelInterceptor);
|
||||
return channel;
|
||||
}
|
||||
|
||||
@Bean AlwaysSampler alwaysSampler() {
|
||||
return new AlwaysSampler();
|
||||
}
|
||||
|
||||
@Bean StompMessageHandler stompMessageHandler() {
|
||||
return new StompMessageHandler();
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user