From 900cb865455713d62abf8977eaf7551a0948fca4 Mon Sep 17 00:00:00 2001 From: buildmaster Date: Fri, 18 Jun 2021 17:28:30 +0000 Subject: [PATCH 1/7] Bumping versions --- benchmarks/pom.xml | 2 +- pom.xml | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/benchmarks/pom.xml b/benchmarks/pom.xml index 4e6f0a5c8..f0ee4c318 100644 --- a/benchmarks/pom.xml +++ b/benchmarks/pom.xml @@ -41,7 +41,7 @@ 4.9.0 0.2.0.RELEASE 1.26 - 3.1.3 + 3.1.4-SNAPSHOT diff --git a/pom.xml b/pom.xml index ac04794c1..5393834aa 100644 --- a/pom.xml +++ b/pom.xml @@ -66,7 +66,7 @@ 3.0.4-SNAPSHOT 3.0.4-SNAPSHOT 2.0.3-SNAPSHOT - 3.1.3 + 3.1.4-SNAPSHOT 3.1.4-SNAPSHOT 3.0.4-SNAPSHOT 3.0.4-SNAPSHOT From e17ae7445aa847d04ab178fbb6e057ff144b708d Mon Sep 17 00:00:00 2001 From: buildmaster Date: Tue, 22 Jun 2021 17:34:30 +0000 Subject: [PATCH 2/7] Bumping versions --- benchmarks/pom.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/benchmarks/pom.xml b/benchmarks/pom.xml index f0ee4c318..50c51b9ed 100644 --- a/benchmarks/pom.xml +++ b/benchmarks/pom.xml @@ -28,7 +28,7 @@ org.springframework.boot spring-boot-starter-parent - 2.4.6 + 2.4.7 From b876551a78ad46184c19cf483e2c55aad8af1113 Mon Sep 17 00:00:00 2001 From: Marcin Grzejszczak Date: Wed, 23 Jun 2021 17:20:01 +0200 Subject: [PATCH 3/7] Fixing sc function (#1983) - Removed an unnecessary child span; the parent span was never ended - Added tests --- .../TraceFunctionAutoConfiguration.java | 66 +++++++++- .../sleuth/brave/bridge/BravePropagator.java | 1 - .../FunctionMessageSpanCustomizer.java | 58 +++++++++ .../messaging/MessagingSleuthOperators.java | 18 +++ .../messaging/TraceFunctionAroundWrapper.java | 36 ++++-- .../messaging/TraceMessageHandler.java | 80 ++++++------ .../messaging/TracingChannelInterceptor.java | 1 - .../sleuth/brave/BraveTestSpanHandler.java | 30 +++++ .../TraceFunctionAroundWrapperTests.java | 13 ++ .../cloud/sleuth/test/TestPropagator.java | 54 ++++++++ .../cloud/sleuth/test/TestSpanBuilder.java | 88 +++++++++++++ .../cloud/sleuth/test/TestSpanHandler.java | 3 + .../cloud/sleuth/test/TestTracer.java | 116 ++++++++++++++++++ .../test/TestTracingBeanPostProcessor.java | 44 +++++++ 14 files changed, 557 insertions(+), 51 deletions(-) create mode 100644 spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/FunctionMessageSpanCustomizer.java create mode 100644 tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestPropagator.java create mode 100644 tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestSpanBuilder.java create mode 100644 tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestTracer.java create mode 100644 tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestTracingBeanPostProcessor.java diff --git a/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/messaging/TraceFunctionAutoConfiguration.java b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/messaging/TraceFunctionAutoConfiguration.java index f0ec48468..07997b9ee 100644 --- a/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/messaging/TraceFunctionAutoConfiguration.java +++ b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/messaging/TraceFunctionAutoConfiguration.java @@ -16,19 +16,28 @@ package org.springframework.cloud.sleuth.autoconfig.instrument.messaging; +import java.util.ArrayList; +import java.util.List; + +import org.springframework.beans.factory.ObjectProvider; import org.springframework.boot.autoconfigure.AutoConfigureAfter; import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingClass; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.cloud.context.scope.refresh.RefreshScopeRefreshedEvent; import org.springframework.cloud.function.context.catalog.FunctionAroundWrapper; +import org.springframework.cloud.sleuth.Span; import org.springframework.cloud.sleuth.Tracer; import org.springframework.cloud.sleuth.autoconfig.brave.BraveAutoConfiguration; +import org.springframework.cloud.sleuth.instrument.messaging.FunctionMessageSpanCustomizer; import org.springframework.cloud.sleuth.instrument.messaging.TraceFunctionAroundWrapper; import org.springframework.cloud.sleuth.propagation.Propagator; +import org.springframework.cloud.stream.messaging.DirectWithAttributesChannel; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.core.env.Environment; +import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageHeaderAccessor; /** @@ -48,8 +57,61 @@ public class TraceFunctionAutoConfiguration { @Bean TraceFunctionAroundWrapper traceFunctionAroundWrapper(Environment environment, Tracer tracer, Propagator propagator, - Propagator.Setter injector, Propagator.Getter extractor) { - return new TraceFunctionAroundWrapper(environment, tracer, propagator, injector, extractor); + Propagator.Setter injector, Propagator.Getter extractor, + ObjectProvider> customizers) { + return new TraceFunctionAroundWrapper(environment, tracer, propagator, injector, extractor, + customizers.getIfAvailable(ArrayList::new)); + } + + @Configuration(proxyBeanMethods = false) + @ConditionalOnClass(DirectWithAttributesChannel.class) + static class TraceFunctionStreamConfiguration { + + @Configuration(proxyBeanMethods = false) + @ConditionalOnClass(name = "org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfiguration") + @ConditionalOnMissingClass("org.springframework.cloud.stream.binder.rabbit.properties.RabbitBinderConfigurationProperties") + static class KafkaOnlyStreamConfiguration { + + @Bean + FunctionMessageSpanCustomizer traceKafkaFunctionMessageSpanCustomizer() { + return new FunctionMessageSpanCustomizer() { + @Override + public void customizeInputMessageSpan(Span span, Message message) { + span.remoteServiceName("kafka"); + } + + @Override + public void customizeOutputMessageSpan(Span span, Message message) { + span.remoteServiceName("kafka"); + } + }; + } + + } + + @Configuration(proxyBeanMethods = false) + @ConditionalOnClass( + name = "org.springframework.cloud.stream.binder.rabbit.properties.RabbitBinderConfigurationProperties") + @ConditionalOnMissingClass("org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfiguration") + static class RabbitOnlyStreamConfiguration { + + @Bean + FunctionMessageSpanCustomizer traceRabbitFunctionMessageSpanCustomizer() { + return new FunctionMessageSpanCustomizer() { + @Override + public void customizeInputMessageSpan(Span span, Message message) { + span.remoteServiceName("rabbitmq"); + } + + @Override + public void customizeOutputMessageSpan(Span span, Message message) { + span.remoteServiceName("rabbitmq"); + } + }; + } + + } + } } diff --git a/spring-cloud-sleuth-brave/src/main/java/org/springframework/cloud/sleuth/brave/bridge/BravePropagator.java b/spring-cloud-sleuth-brave/src/main/java/org/springframework/cloud/sleuth/brave/bridge/BravePropagator.java index 64d2afec4..c5e98779f 100644 --- a/spring-cloud-sleuth-brave/src/main/java/org/springframework/cloud/sleuth/brave/bridge/BravePropagator.java +++ b/spring-cloud-sleuth-brave/src/main/java/org/springframework/cloud/sleuth/brave/bridge/BravePropagator.java @@ -54,7 +54,6 @@ public class BravePropagator implements Propagator { public Span.Builder extract(C carrier, Getter getter) { TraceContextOrSamplingFlags extract = this.tracing.propagation().extractor(getter::get).extract(carrier); if (extract.samplingFlags() == SamplingFlags.EMPTY) { - this.tracing.tracer().nextSpan(); return new BraveSpanBuilder(this.tracing.tracer()); } return BraveSpanBuilder.toBuilder(this.tracing.tracer(), extract); diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/FunctionMessageSpanCustomizer.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/FunctionMessageSpanCustomizer.java new file mode 100644 index 000000000..ca3c5a1f9 --- /dev/null +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/FunctionMessageSpanCustomizer.java @@ -0,0 +1,58 @@ +/* + * Copyright 2013-2021 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 + * + * https://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.messaging; + +import org.springframework.cloud.sleuth.Span; +import org.springframework.messaging.Message; + +/** + * Allows customization of messaging spans for Spring Cloud Function instrumentation. + * + * @author Marcin Grzejszczak + * @since 3.0.4 + */ +public interface FunctionMessageSpanCustomizer { + + /** + * Customizes the span created after wrapping the input message in a span + * representation. + * @param span current span to customize + * @param message received or sent message + */ + default void customizeInputMessageSpan(Span span, Message message) { + + } + + /** + * Customizes the span wrapping the function execution. + * @param span current span to customize + * @param message message to be sent + */ + default void customizeFunctionSpan(Span span, Message message) { + + } + + /** + * Customizes the span created for the output message. + * @param span current span to customize + * @param message message to be sent + */ + default void customizeOutputMessageSpan(Span span, Message message) { + + } + +} diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/MessagingSleuthOperators.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/MessagingSleuthOperators.java index 8a548b8c8..08e18e04a 100644 --- a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/MessagingSleuthOperators.java +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/MessagingSleuthOperators.java @@ -188,6 +188,23 @@ public final class MessagingSleuthOperators { * @return instrumented message */ public static Message handleOutputMessage(BeanFactory beanFactory, Message message, Throwable throwable) { + return handleOutputMessage(beanFactory, message, span -> { + }, throwable); + } + + /** + * Creates an output message with tracer headers and reports the corresponding + * producer span. If the message contains a header called {@code destination} it will + * be used to tag the span with destination name. + * @param beanFactory - bean factory + * @param message - message to which tracer headers should be injected + * @param spanCustomizer - customizer of the output span + * @param throwable - exception that took place while processing the message + * @param - message payload + * @return instrumented message + */ + public static Message handleOutputMessage(BeanFactory beanFactory, Message message, + Consumer spanCustomizer, Throwable throwable) { TraceMessageHandler traceMessageHandler = TraceMessageHandler.forNonSpringIntegration(beanFactory); Span span = traceMessageHandler.parentSpan(message); span = span != null ? span : traceMessageHandler.consumerSpan(message); @@ -198,6 +215,7 @@ public final class MessagingSleuthOperators { } MessageAndSpan messageAndSpan = traceMessageHandler.wrapOutputMessage(message, span, String.valueOf(message.getHeaders().getOrDefault("destination", ""))); + spanCustomizer.accept(messageAndSpan.span); traceMessageHandler.afterMessageHandled(messageAndSpan.span, throwable); return messageAndSpan.msg; } diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceFunctionAroundWrapper.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceFunctionAroundWrapper.java index 0ce06cc92..aca731b49 100644 --- a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceFunctionAroundWrapper.java +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceFunctionAroundWrapper.java @@ -16,6 +16,8 @@ package org.springframework.cloud.sleuth.instrument.messaging; +import java.util.Collections; +import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; @@ -58,17 +60,26 @@ public class TraceFunctionAroundWrapper extends FunctionAroundWrapper private final TraceMessageHandler traceMessageHandler; + private final List customizers; + final Map functionToDestinationCache = new ConcurrentHashMap<>(); public TraceFunctionAroundWrapper(Environment environment, Tracer tracer, Propagator propagator, Propagator.Setter injector, Propagator.Getter extractor) { + this(environment, tracer, propagator, injector, extractor, Collections.emptyList()); + } + + public TraceFunctionAroundWrapper(Environment environment, Tracer tracer, Propagator propagator, + Propagator.Setter injector, Propagator.Getter extractor, + List customizers) { this.environment = environment; this.tracer = tracer; this.propagator = propagator; this.injector = injector; this.extractor = extractor; + this.customizers = customizers; this.traceMessageHandler = TraceMessageHandler.forNonSpringIntegration(this.tracer, this.propagator, - this.injector, this.extractor); + this.injector, this.extractor, this.customizers); } @Override @@ -76,23 +87,26 @@ public class TraceFunctionAroundWrapper extends FunctionAroundWrapper MessageAndSpans invocationMessage = null; Span span; if (message == null && targetFunction.isSupplier()) { // Supplier - span = traceMessageHandler.tracer.nextSpan().name(targetFunction.getFunctionDefinition()); + if (log.isDebugEnabled()) { + log.debug("Creating a span for a supplier"); + } + span = this.tracer.nextSpan().name(targetFunction.getFunctionDefinition()); + customizedInputMessageSpan(span, null); } else { if (log.isDebugEnabled()) { log.debug("Will retrieve the tracing headers from the message"); } - invocationMessage = traceMessageHandler.wrapInputMessage(message, + invocationMessage = this.traceMessageHandler.wrapInputMessage(message, inputDestination(targetFunction.getFunctionDefinition())); if (log.isDebugEnabled()) { log.debug("Wrapped input msg " + invocationMessage); } span = invocationMessage.childSpan; } - Object result; Throwable throwable = null; - try (Tracer.SpanInScope ws = tracer.withSpan(span.start())) { + try (Tracer.SpanInScope ws = this.tracer.withSpan(span.start())) { result = invocationMessage == null ? targetFunction.get() : targetFunction.apply(invocationMessage.msg); } catch (Exception e) { @@ -100,7 +114,7 @@ public class TraceFunctionAroundWrapper extends FunctionAroundWrapper throw e; } finally { - traceMessageHandler.afterMessageHandled(span, throwable); + this.traceMessageHandler.afterMessageHandled(span, throwable); } if (result == null) { if (log.isDebugEnabled()) { @@ -109,10 +123,12 @@ public class TraceFunctionAroundWrapper extends FunctionAroundWrapper return null; } Message msgResult = toMessage(result); - MessageAndSpan wrappedOutputMessage; + if (log.isDebugEnabled()) { + log.debug("Will instrument the output message"); + } if (invocationMessage != null) { - wrappedOutputMessage = traceMessageHandler.wrapOutputMessage(msgResult, invocationMessage.parentSpan, + wrappedOutputMessage = this.traceMessageHandler.wrapOutputMessage(msgResult, invocationMessage.parentSpan, outputDestination(targetFunction.getFunctionDefinition())); } else { @@ -129,6 +145,10 @@ public class TraceFunctionAroundWrapper extends FunctionAroundWrapper return traceMessageHandler.wrapOutputMessage(resultMessage, spanFromMessage, outputDestination(name)); } + private void customizedInputMessageSpan(Span spanToCustomize, Message msg) { + this.customizers.forEach(cust -> cust.customizeInputMessageSpan(spanToCustomize, msg)); + } + private Message toMessage(Object result) { if (!(result instanceof Message)) { return MessageBuilder.withPayload(result).build(); diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessageHandler.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessageHandler.java index 9c3a0cc91..6c9d84e84 100644 --- a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessageHandler.java +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessageHandler.java @@ -81,10 +81,12 @@ class TraceMessageHandler { private final Function outputMessageSpanFunction; + private final List customizers; + TraceMessageHandler(Tracer tracer, Propagator propagator, Propagator.Setter injector, Propagator.Getter extractor, Function preSendFunction, TriConsumer preSendMessageManipulator, - Function outputMessageSpanFunction) { + Function outputMessageSpanFunction, List customizers) { this.tracer = tracer; this.propagator = propagator; this.injector = injector; @@ -93,18 +95,20 @@ class TraceMessageHandler { this.preSendFunction = preSendFunction; this.preSendMessageManipulator = preSendMessageManipulator; this.outputMessageSpanFunction = outputMessageSpanFunction; + this.customizers = customizers; } static TraceMessageHandler forNonSpringIntegration(Tracer tracer, Propagator propagator, - Propagator.Setter injector, Propagator.Getter extractor) { - Function preSendFunction = span -> tracer.nextSpan(span).name("handle").start(); + Propagator.Setter injector, Propagator.Getter extractor, + List customizers) { + Function preSendFunction = span -> tracer.nextSpan(span).name("function").start(); TriConsumer preSendMessageManipulator = (headers, parentSpan, childSpan) -> { headers.setHeader("traceHandlerParentSpan", parentSpan); headers.setHeader(Span.class.getName(), childSpan); }; Function postReceiveFunction = span -> tracer.spanBuilder().setParent(span.context()); return new TraceMessageHandler(tracer, propagator, injector, extractor, preSendFunction, - preSendMessageManipulator, postReceiveFunction); + preSendMessageManipulator, postReceiveFunction, customizers); } @SuppressWarnings("unchecked") @@ -112,7 +116,7 @@ class TraceMessageHandler { Propagator.Setter setter = firstBeanOrException(beanFactory, Propagator.Setter.class); Propagator.Getter getter = firstBeanOrException(beanFactory, Propagator.Getter.class); return forNonSpringIntegration(beanFactory.getBean(Tracer.class), beanFactory.getBean(Propagator.class), setter, - getter); + getter, customizers(beanFactory)); } private static T firstBeanOrException(BeanFactory beanFactory, Class clazz) { @@ -125,6 +129,16 @@ class TraceMessageHandler { return object; } + private static List customizers(BeanFactory beanFactory) { + List customizers = new ArrayList<>(); + ObjectProvider provider = beanFactory + .getBeanProvider(FunctionMessageSpanCustomizer.class); + for (FunctionMessageSpanCustomizer functionMessageSpanCustomizer : provider) { + customizers.add(functionMessageSpanCustomizer); + } + return customizers; + } + /** * Wraps the given input message with tracing headers and returns a corresponding * span. @@ -134,41 +148,33 @@ class TraceMessageHandler { */ MessageAndSpans wrapInputMessage(Message message, String destinationName) { MessageHeaderAccessor headers = mutableHeaderAccessor(message); - Span extracted = this.propagator.extract(headers, this.extractor).start(); - // Start and finish a consumer span as we will immediately process it. - Span.Builder consumerSpanBuilder = this.tracer.spanBuilder().setParent(extracted.context()); - Span consumerSpan = consumerSpan(destinationName, extracted, consumerSpanBuilder); - // create and scope a span for the message processor - Span span = this.preSendFunction.apply(consumerSpan); - // remove any trace headers, but don't re-inject as we are synchronously - // processing the - // message and can rely on scoping to access this span later. - clearTracingHeaders(headers); - this.preSendMessageManipulator.accept(headers, consumerSpan, span); + Span.Builder consumerSpanBuilder = this.propagator.extract(headers, this.extractor); + Span consumerSpan = consumerSpan(destinationName, consumerSpanBuilder, message); if (log.isDebugEnabled()) { - log.debug("Created a handle span after retrieving the message " + consumerSpanBuilder); + log.debug("Built a consumer span " + consumerSpan); } + Span childSpan = this.preSendFunction.apply(consumerSpan); + clearTracingHeaders(headers); + this.preSendMessageManipulator.accept(headers, consumerSpan, childSpan); + this.customizers.forEach(customizer -> customizer.customizeFunctionSpan(childSpan, message)); if (message instanceof ErrorMessage) { return new MessageAndSpans(new ErrorMessage((Throwable) message.getPayload(), headers.getMessageHeaders()), - consumerSpan, span); + consumerSpan, childSpan); } headers.setImmutable(); return new MessageAndSpans(new GenericMessage<>(message.getPayload(), headers.getMessageHeaders()), - consumerSpan, span); + consumerSpan, childSpan); } - private Span consumerSpan(String destinationName, Span extracted, Span.Builder consumerSpanBuilder) { - Span consumerSpan; - if (!extracted.isNoop()) { - consumerSpanBuilder.kind(Span.Kind.CONSUMER).start(); - addTags(consumerSpanBuilder, destinationName); - consumerSpanBuilder.remoteServiceName(REMOTE_SERVICE_NAME); - consumerSpan = consumerSpanBuilder.start(); - consumerSpan.end(); - } - else { - consumerSpan = consumerSpanBuilder.start(); - } + private Span consumerSpan(String destinationName, Span.Builder consumerSpanBuilder, Message message) { + consumerSpanBuilder.kind(Span.Kind.CONSUMER).name("handle"); + addTags(consumerSpanBuilder, destinationName); + consumerSpanBuilder.remoteServiceName(REMOTE_SERVICE_NAME); + // this is the consumer part of the producer->consumer mechanism + Span consumerSpan = consumerSpanBuilder.start(); + this.customizers.forEach(customizer -> customizer.customizeInputMessageSpan(consumerSpan, message)); + // we're ending this immediately just to have a properly nested graph + consumerSpan.end(); return consumerSpan; } @@ -191,12 +197,6 @@ class TraceMessageHandler { } } - private void addTags(Span result, String destinationName) { - if (StringUtils.hasText(destinationName)) { - result.tag("channel", SpanNameUtil.shorten(destinationName)); - } - } - /** * Called either when message got received and processed or message got sent. * @param span - span that corresponds to the given operation @@ -233,7 +233,7 @@ class TraceMessageHandler { MessageHeaderAccessor headers = mutableHeaderAccessor(retrievedMessage); Span.Builder span = this.outputMessageSpanFunction.apply(parentSpan); clearTracingHeaders(headers); - Span producerSpan = createProducerSpan(headers, span, destinationName); + Span producerSpan = createProducerSpan(headers, span, destinationName, message); this.propagator.inject(producerSpan.context(), headers, this.injector); if (log.isDebugEnabled()) { log.debug("Created a new span output message " + span); @@ -241,12 +241,14 @@ class TraceMessageHandler { return new MessageAndSpan(outputMessage(message, retrievedMessage, headers), producerSpan); } - private Span createProducerSpan(MessageHeaderAccessor headers, Span.Builder spanBuilder, String destinationName) { + private Span createProducerSpan(MessageHeaderAccessor headers, Span.Builder spanBuilder, String destinationName, + Message message) { spanBuilder.kind(Span.Kind.PRODUCER).name("send").remoteServiceName(toRemoteServiceName(headers)); Span span = spanBuilder.start(); if (!span.isNoop()) { addTags(spanBuilder, destinationName); } + this.customizers.forEach(customizer -> customizer.customizeOutputMessageSpan(span, message)); return span; } diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptor.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptor.java index 25ed6e3b9..21447c6e5 100644 --- a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptor.java +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptor.java @@ -109,7 +109,6 @@ public final class TracingChannelInterceptor implements ExecutorChannelIntercept public TracingChannelInterceptor(Tracer tracer, Propagator propagator, Propagator.Setter setter, Propagator.Getter getter, Function remoteServiceNameMapper, MessageSpanCustomizer messageSpanCustomizer) { - this.tracer = tracer; this.propagator = propagator; this.injector = setter; diff --git a/tests/common/src/main/java/org/springframework/cloud/sleuth/brave/BraveTestSpanHandler.java b/tests/common/src/main/java/org/springframework/cloud/sleuth/brave/BraveTestSpanHandler.java index 796631e74..b2f638f69 100644 --- a/tests/common/src/main/java/org/springframework/cloud/sleuth/brave/BraveTestSpanHandler.java +++ b/tests/common/src/main/java/org/springframework/cloud/sleuth/brave/BraveTestSpanHandler.java @@ -18,6 +18,7 @@ package org.springframework.cloud.sleuth.brave; import java.util.Iterator; import java.util.List; +import java.util.Queue; import java.util.stream.Collectors; import brave.test.IntegrationTestSpanHandler; @@ -27,6 +28,8 @@ import org.springframework.cloud.sleuth.brave.bridge.BraveAccessor; import org.springframework.cloud.sleuth.exporter.FinishedSpan; import org.springframework.cloud.sleuth.test.TestSpanHandler; +import static org.assertj.core.api.BDDAssertions.then; + public class BraveTestSpanHandler implements TestSpanHandler { final brave.test.TestSpanHandler spans; @@ -81,6 +84,33 @@ public class BraveTestSpanHandler implements TestSpanHandler { return BraveAccessor.finishedSpan(this.spans.get(index)); } + @Override + public void assertAllSpansWereFinishedOrAbandoned(Queue createdSpans) { + List finishedSpans = reportedSpans(); + then(finishedSpans).as("There should be that many finished spans as many created ones") + .hasSize(createdSpans.size()); + // finished -> a,b,c ; created -> b,c,d => matchedFinished = b,c + List matchedFinishedSpans = finishedSpans.stream() + .filter(f -> createdSpans.stream().anyMatch(cs -> f.getSpanId().equals(cs.context().spanId()))) + .collect(Collectors.toList()); + // finished -> a,b,c ; created -> b,c,d => matchedCreated = b,c + List matchedCreatedSpans = createdSpans.stream() + .filter(cs -> finishedSpans.stream().anyMatch(f -> cs.context().spanId().equals(f.getSpanId()))) + .collect(Collectors.toList()); + // finished -> a,b,c ; created -> b,c,d => missingFinished = a + List missingFinishedSpans = finishedSpans.stream() + .filter(f -> matchedFinishedSpans.stream().noneMatch(m -> m.getSpanId().equals(f.getSpanId()))) + .collect(Collectors.toList()); + // finished -> a,b,c ; created -> b,c,d => missingCreated = d + List missingCreatedSpans = createdSpans.stream().filter( + f -> matchedCreatedSpans.stream().noneMatch(m -> m.context().spanId().equals(f.context().spanId()))) + .collect(Collectors.toList()); + if (!missingFinishedSpans.isEmpty() || !missingCreatedSpans.isEmpty()) { + throw new AssertionError("There were unmatched created spans " + missingCreatedSpans + + " and/or finished span " + missingFinishedSpans); + } + } + @Override public Iterator iterator() { return reportedSpans().iterator(); diff --git a/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceFunctionAroundWrapperTests.java b/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceFunctionAroundWrapperTests.java index c9b802c07..fa94083a6 100644 --- a/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceFunctionAroundWrapperTests.java +++ b/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceFunctionAroundWrapperTests.java @@ -26,6 +26,8 @@ import org.springframework.boot.builder.SpringApplicationBuilder; import org.springframework.cloud.function.context.FunctionCatalog; import org.springframework.cloud.function.context.catalog.SimpleFunctionRegistry.FunctionInvocationWrapper; import org.springframework.cloud.sleuth.test.TestSpanHandler; +import org.springframework.cloud.sleuth.test.TestTracer; +import org.springframework.cloud.sleuth.test.TestTracingBeanPostProcessor; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.messaging.Message; @@ -49,10 +51,13 @@ public abstract class TraceFunctionAroundWrapperTests { FunctionCatalog catalog = context.getBean(FunctionCatalog.class); FunctionInvocationWrapper function = catalog.lookup("greeter"); function.setSkipOutputConversion(true); + Message result = (Message) function.get(); + assertThat(result.getPayload()).isEqualTo("hello"); assertThat(spanHandler.reportedSpans().size()).isEqualTo(2); assertThat(((String) result.getHeaders().get("b3"))).contains(spanHandler.get(0).getTraceId()); + spanHandler.assertAllSpansWereFinishedOrAbandoned(context.getBean(TestTracer.class).createdSpans()); } } @@ -66,10 +71,13 @@ public abstract class TraceFunctionAroundWrapperTests { FunctionCatalog catalog = context.getBean(FunctionCatalog.class); FunctionInvocationWrapper function = catalog.lookup("uppercase"); function.setSkipOutputConversion(true); + Message result = (Message) function.apply(MessageBuilder.withPayload("hello").build()); + assertThat(result.getPayload()).isEqualTo("HELLO"); assertThat(spanHandler.reportedSpans().size()).isEqualTo(3); assertThat(((String) result.getHeaders().get("b3"))).contains(spanHandler.get(0).getTraceId()); + spanHandler.assertAllSpansWereFinishedOrAbandoned(context.getBean(TestTracer.class).createdSpans()); } } @@ -88,6 +96,11 @@ public abstract class TraceFunctionAroundWrapperTests { return v -> v.toUpperCase(); } + @Bean + static TestTracingBeanPostProcessor testTracerBeanPostProcessor() { + return new TestTracingBeanPostProcessor(); + } + } }; diff --git a/tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestPropagator.java b/tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestPropagator.java new file mode 100644 index 000000000..002646a7e --- /dev/null +++ b/tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestPropagator.java @@ -0,0 +1,54 @@ +/* + * Copyright 2013-2021 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 + * + * https://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.test; + +import java.util.List; + +import org.springframework.cloud.sleuth.Span; +import org.springframework.cloud.sleuth.TraceContext; +import org.springframework.cloud.sleuth.propagation.Propagator; + +/** + * {@link Propagator} that stores information about started spans. + */ +public class TestPropagator implements Propagator { + + private final Propagator delegate; + + private final TestTracer testTracer; + + public TestPropagator(Propagator delegate, TestTracer testTracer) { + this.delegate = delegate; + this.testTracer = testTracer; + } + + @Override + public List fields() { + return this.delegate.fields(); + } + + @Override + public void inject(TraceContext context, C carrier, Setter setter) { + this.delegate.inject(context, carrier, setter); + } + + @Override + public Span.Builder extract(C carrier, Getter getter) { + return new TestSpanBuilder(this.delegate.extract(carrier, getter), this.testTracer); + } + +} diff --git a/tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestSpanBuilder.java b/tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestSpanBuilder.java new file mode 100644 index 000000000..b367931c1 --- /dev/null +++ b/tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestSpanBuilder.java @@ -0,0 +1,88 @@ +/* + * Copyright 2013-2021 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 + * + * https://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.test; + +import org.springframework.cloud.sleuth.Span; +import org.springframework.cloud.sleuth.TraceContext; + +class TestSpanBuilder implements Span.Builder { + + private final Span.Builder delegate; + + private final TestTracer testTracer; + + TestSpanBuilder(Span.Builder delegate, TestTracer testTracer) { + this.delegate = delegate; + this.testTracer = testTracer; + } + + @Override + public Span.Builder setParent(TraceContext context) { + delegate.setParent(context); + return this; + } + + @Override + public Span.Builder setNoParent() { + delegate.setNoParent(); + return this; + } + + @Override + public Span.Builder name(String name) { + delegate.name(name); + return this; + } + + @Override + public Span.Builder event(String value) { + delegate.event(value); + return this; + } + + @Override + public Span.Builder tag(String key, String value) { + delegate.tag(key, value); + return this; + } + + @Override + public Span.Builder error(Throwable throwable) { + delegate.error(throwable); + return this; + } + + @Override + public Span.Builder kind(Span.Kind spanKind) { + delegate.kind(spanKind); + return this; + } + + @Override + public Span.Builder remoteServiceName(String remoteServiceName) { + delegate.remoteServiceName(remoteServiceName); + return this; + } + + @Override + public Span start() { + Span span = delegate.start(); + this.testTracer.createdSpans.add(span); + return span; + } + +} diff --git a/tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestSpanHandler.java b/tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestSpanHandler.java index 17681a5ea..d227a64aa 100644 --- a/tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestSpanHandler.java +++ b/tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestSpanHandler.java @@ -17,6 +17,7 @@ package org.springframework.cloud.sleuth.test; import java.util.List; +import java.util.Queue; import org.springframework.cloud.sleuth.Span; import org.springframework.cloud.sleuth.exporter.FinishedSpan; @@ -35,4 +36,6 @@ public interface TestSpanHandler extends Iterable { FinishedSpan get(int index); + void assertAllSpansWereFinishedOrAbandoned(Queue createdSpans); + } diff --git a/tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestTracer.java b/tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestTracer.java new file mode 100644 index 000000000..d420f750b --- /dev/null +++ b/tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestTracer.java @@ -0,0 +1,116 @@ +/* + * Copyright 2013-2021 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 + * + * https://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.test; + +import java.util.LinkedList; +import java.util.Map; +import java.util.Queue; + +import org.springframework.cloud.sleuth.BaggageInScope; +import org.springframework.cloud.sleuth.ScopedSpan; +import org.springframework.cloud.sleuth.Span; +import org.springframework.cloud.sleuth.SpanCustomizer; +import org.springframework.cloud.sleuth.TraceContext; +import org.springframework.cloud.sleuth.Tracer; +import org.springframework.lang.Nullable; + +public class TestTracer implements Tracer, AutoCloseable { + + private final Tracer delegate; + + final Queue createdSpans = new LinkedList<>(); + + public TestTracer(Tracer delegate) { + this.delegate = delegate; + } + + @Override + public Map getAllBaggage() { + return delegate.getAllBaggage(); + } + + @Override + public BaggageInScope getBaggage(String name) { + return delegate.getBaggage(name); + } + + @Override + public BaggageInScope getBaggage(TraceContext traceContext, String name) { + return delegate.getBaggage(traceContext, name); + } + + @Override + public BaggageInScope createBaggage(String name) { + return delegate.createBaggage(name); + } + + @Override + public BaggageInScope createBaggage(String name, String value) { + return delegate.createBaggage(name, value); + } + + @Override + public Span nextSpan() { + Span span = delegate.nextSpan(); + this.createdSpans.add(span); + return span; + } + + @Override + public Span nextSpan(Span parent) { + Span span = delegate.nextSpan(parent); + this.createdSpans.add(span); + return span; + } + + @Override + public SpanInScope withSpan(Span span) { + return delegate.withSpan(span); + } + + @Override + public ScopedSpan startScopedSpan(String name) { + return delegate.startScopedSpan(name); + } + + @Override + public Span.Builder spanBuilder() { + return new TestSpanBuilder(delegate.spanBuilder(), this); + } + + @Override + @Nullable + public SpanCustomizer currentSpanCustomizer() { + return delegate.currentSpanCustomizer(); + } + + @Override + @Nullable + public Span currentSpan() { + return delegate.currentSpan(); + } + + @Override + public void close() throws Exception { + this.createdSpans.clear(); + } + + public Queue createdSpans() { + return createdSpans; + } + +} diff --git a/tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestTracingBeanPostProcessor.java b/tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestTracingBeanPostProcessor.java new file mode 100644 index 000000000..6e68f1d8b --- /dev/null +++ b/tests/common/src/main/java/org/springframework/cloud/sleuth/test/TestTracingBeanPostProcessor.java @@ -0,0 +1,44 @@ +/* + * Copyright 2013-2021 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 + * + * https://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.test; + +import org.springframework.beans.BeansException; +import org.springframework.beans.factory.config.BeanPostProcessor; +import org.springframework.cloud.sleuth.Tracer; +import org.springframework.cloud.sleuth.propagation.Propagator; + +/** + * Wraps all tracing related components into test representations. That way additional + * assertions can take place. + */ +public class TestTracingBeanPostProcessor implements BeanPostProcessor { + + TestTracer testTracer; + + @Override + public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException { + if (bean instanceof Tracer && !(bean instanceof TestTracer)) { + this.testTracer = new TestTracer((Tracer) bean); + return this.testTracer; + } + else if (bean instanceof Propagator && !(bean instanceof TestPropagator)) { + return new TestPropagator((Propagator) bean, this.testTracer); + } + return bean; + } + +} From 08b4c1368e639b7eeb95b69008d0be4d3be838ae Mon Sep 17 00:00:00 2001 From: Andrew Flower Date: Fri, 25 Jun 2021 00:42:53 +0900 Subject: [PATCH 4/7] Broaden caught exceptions (#1985) Handle BeanCreationException instead of just BeanCurrentlyInCreationException. This handles the case when a cyclic reference (causing a BeanCurrentlyInCreationException) occurs further down the line - when it is not the actuator of that references the in-creation HttpTracing bean, but a dependency of that actuator. In this case the exception that bubbles up here is an UnsatisfiedDependencyException, which is a BeanCreationException --- .../autoconfig/instrument/web/SkipPatternConfiguration.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/web/SkipPatternConfiguration.java b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/web/SkipPatternConfiguration.java index 49dc59cf6..f635072ab 100644 --- a/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/web/SkipPatternConfiguration.java +++ b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/web/SkipPatternConfiguration.java @@ -23,7 +23,7 @@ import java.util.StringJoiner; import java.util.regex.Pattern; import java.util.stream.Collectors; -import org.springframework.beans.factory.BeanCurrentlyInCreationException; +import org.springframework.beans.factory.BeanCreationException; import org.springframework.beans.factory.ObjectProvider; import org.springframework.boot.actuate.autoconfigure.endpoint.web.WebEndpointProperties; import org.springframework.boot.actuate.autoconfigure.web.server.ConditionalOnManagementPort; @@ -83,7 +83,7 @@ class SkipPatternConfiguration { } return () -> result; } - catch (BeanCurrentlyInCreationException e) { + catch (BeanCreationException e) { // Most likely, there is an actuator endpoint that indirectly references an // instrumented HTTP client. return () -> consolidateSkipPatterns(patterns); From c8961588044dab6c9a1b43e64b745b7e9860bfed Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Mon, 28 Jun 2021 19:22:58 +0200 Subject: [PATCH 5/7] GH-1975 Fix test after SCF change Fix sleuth tests related to the change in SCF - https://github.com/spring-cloud/spring-cloud-function/commit/c86890806e856ef90f64f2b2e6142731eef74747 --- .../instrument/messaging/TraceFunctionAroundWrapperTests.java | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceFunctionAroundWrapperTests.java b/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceFunctionAroundWrapperTests.java index fa94083a6..2197a8171 100644 --- a/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceFunctionAroundWrapperTests.java +++ b/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceFunctionAroundWrapperTests.java @@ -50,11 +50,10 @@ public abstract class TraceFunctionAroundWrapperTests { assertThat(spanHandler.reportedSpans()).isEmpty(); FunctionCatalog catalog = context.getBean(FunctionCatalog.class); FunctionInvocationWrapper function = catalog.lookup("greeter"); - function.setSkipOutputConversion(true); Message result = (Message) function.get(); - assertThat(result.getPayload()).isEqualTo("hello"); + assertThat(result.getPayload()).isEqualTo("hello".getBytes()); assertThat(spanHandler.reportedSpans().size()).isEqualTo(2); assertThat(((String) result.getHeaders().get("b3"))).contains(spanHandler.get(0).getTraceId()); spanHandler.assertAllSpansWereFinishedOrAbandoned(context.getBean(TestTracer.class).createdSpans()); @@ -70,7 +69,6 @@ public abstract class TraceFunctionAroundWrapperTests { assertThat(spanHandler.reportedSpans()).isEmpty(); FunctionCatalog catalog = context.getBean(FunctionCatalog.class); FunctionInvocationWrapper function = catalog.lookup("uppercase"); - function.setSkipOutputConversion(true); Message result = (Message) function.apply(MessageBuilder.withPayload("hello").build()); From e9d55a3af0b501f0fc994f5b134bf3ce0ac62872 Mon Sep 17 00:00:00 2001 From: Marcin Grzejszczak Date: Fri, 2 Jul 2021 13:13:42 +0200 Subject: [PATCH 6/7] Added awaitility to benchmarks --- benchmarks/pom.xml | 9 ++++----- .../SleuthBenchmarkingStreamApplication.java | 20 +++++++++++-------- .../jmh/stream/MicroBenchmarkStreamTests.java | 5 +++-- 3 files changed, 19 insertions(+), 15 deletions(-) diff --git a/benchmarks/pom.xml b/benchmarks/pom.xml index 50c51b9ed..0972ea89a 100644 --- a/benchmarks/pom.xml +++ b/benchmarks/pom.xml @@ -41,7 +41,6 @@ 4.9.0 0.2.0.RELEASE 1.26 - 3.1.4-SNAPSHOT @@ -55,8 +54,8 @@ org.springframework.cloud - spring-cloud-stream-dependencies - ${spring-cloud-stream.version} + spring-cloud-sleuth + ${project.version} pom import @@ -65,7 +64,7 @@ - ${project.groupId} + org.springframework.cloud spring-cloud-starter-sleuth @@ -150,7 +149,7 @@ org.awaitility awaitility - test + compile diff --git a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/app/stream/SleuthBenchmarkingStreamApplication.java b/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/app/stream/SleuthBenchmarkingStreamApplication.java index 32f1fa718..6c176269d 100644 --- a/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/app/stream/SleuthBenchmarkingStreamApplication.java +++ b/benchmarks/src/main/java/org/springframework/cloud/sleuth/benchmarks/app/stream/SleuthBenchmarkingStreamApplication.java @@ -23,6 +23,7 @@ import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import java.util.function.Function; +import org.awaitility.Awaitility; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import reactor.core.publisher.Flux; @@ -60,15 +61,16 @@ public class SleuthBenchmarkingStreamApplication { // "DECORATE_ON_LAST"); // System.setProperty("spring.sleuth.reactor.instrumentation-type", "MANUAL"); System.setProperty("spring.sleuth.reactor.instrumentation-type", "DECORATE_QUEUES"); + System.setProperty("spring.sleuth.integration.enabled", "true"); System.setProperty("spring.sleuth.function.type", "DECORATE_QUEUES"); ConfigurableApplicationContext context = SpringApplication.run(SleuthBenchmarkingStreamApplication.class, args); - for (int i = 0; i < 1; i++) { - InputDestination input = context.getBean(InputDestination.class); - input.send(MessageBuilder.withPayload("hello".getBytes()) - .setHeader("b3", "4883117762eb9420-4883117762eb9420-1").build()); - log.info("Retrieving the message for tests"); - OutputDestination output = context.getBean(OutputDestination.class); - Message message = output.receive(200L); + InputDestination input = context.getBean(InputDestination.class); + input.send(MessageBuilder.withPayload("hello".getBytes()) + .setHeader("b3", "4883117762eb9420-4883117762eb9420-1").build()); + log.info("Retrieving the message for tests"); + OutputDestination output = context.getBean(OutputDestination.class); + Awaitility.await().untilAsserted( () -> { + Message message = output.receive(1L); log.info("Got the message from output"); assertThat(message).isNotNull(); log.info("Message is not null"); @@ -77,7 +79,9 @@ public class SleuthBenchmarkingStreamApplication { String b3 = message.getHeaders().get("b3", String.class); log.info("Checking the b3 header [" + b3 + "]"); assertThat(b3).startsWith("4883117762eb9420"); - } + }); + context.close(); + System.exit(0); } @Bean diff --git a/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/stream/MicroBenchmarkStreamTests.java b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/stream/MicroBenchmarkStreamTests.java index 6ce4905df..be5393769 100644 --- a/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/stream/MicroBenchmarkStreamTests.java +++ b/benchmarks/src/test/java/org/springframework/cloud/sleuth/benchmarks/jmh/stream/MicroBenchmarkStreamTests.java @@ -23,6 +23,7 @@ import java.util.concurrent.TimeUnit; import brave.Tracing; import jmh.mbr.junit5.Microbenchmark; +import org.awaitility.Awaitility; import org.junit.platform.commons.annotation.Testable; import org.openjdk.jmh.annotations.Benchmark; import org.openjdk.jmh.annotations.BenchmarkMode; @@ -114,12 +115,12 @@ public class MicroBenchmarkStreamTests { void run() { sendInputMessage(); - assertThatOutputMessageGotReceived(); + Awaitility.await().untilAsserted(this::assertThatOutputMessageGotReceived); } private void assertThatOutputMessageGotReceived() { // System.out.println("Retrieving the message for tests"); - Message message = output.receive(200L); + Message message = output.receive(1L); // System.out.println("Got the message from output"); assertThat(message).isNotNull(); // System.out.println("Message is not null"); From 8b61684168be90e06b861c26205d25f61ab3a0a2 Mon Sep 17 00:00:00 2001 From: Marcin Grzejszczak Date: Fri, 2 Jul 2021 17:14:05 +0200 Subject: [PATCH 7/7] Ensures that there are no issues with context setup when picking custom propagation type (#1989) fixes gh-1987 --- .../brave/BraveBaggageConfiguration.java | 7 +- .../CustomPropagationFactoryTests.java | 102 +++++++++++++++ .../CompositePropagationFactorySupplier.java | 8 +- .../brave/propagation/PropagationType.java | 3 +- ...positePropagationFactorySupplierTests.java | 120 ------------------ 5 files changed, 113 insertions(+), 127 deletions(-) create mode 100644 spring-cloud-sleuth-autoconfigure/src/test/java/org/springframework/cloud/sleuth/autoconfig/brave/baggage/CustomPropagationFactoryTests.java delete mode 100644 spring-cloud-sleuth-brave/src/test/java/org/springframework/cloud/sleuth/brave/bridge/CompositePropagationFactorySupplierTests.java diff --git a/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/brave/BraveBaggageConfiguration.java b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/brave/BraveBaggageConfiguration.java index e307d2b3d..16d9e03df 100644 --- a/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/brave/BraveBaggageConfiguration.java +++ b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/brave/BraveBaggageConfiguration.java @@ -105,7 +105,11 @@ class BraveBaggageConfiguration { // See #1643 @Bean @ConditionalOnMissingBean - PropagationFactorySupplier defaultPropagationFactorySupplier() { + PropagationFactorySupplier defaultPropagationFactorySupplier(SleuthPropagationProperties properties) { + if (properties.getType().contains(PropagationType.CUSTOM)) { + throw new IllegalStateException( + "Please register a bean with the following signature [extends Propagation.Factory implements Propagation] to override the default Sleuth behaviour or [implements PropagationFactorySupplier] to reuse it."); + } return () -> B3Propagation.newFactoryBuilder().injectFormat(B3Propagation.Format.SINGLE_NO_PARENT).build(); } @@ -131,7 +135,6 @@ class BraveBaggageConfiguration { @Qualifier(PROPAGATION_KEYS) List propagationKeys, SleuthBaggageProperties sleuthBaggageProperties, SleuthPropagationProperties sleuthPropagationProperties, PropagationFactorySupplier supplier, @Nullable List baggagePropagationCustomizers) { - Set localFields = redirectOldPropertyToNew(LOCAL_KEYS, localKeys, "spring.sleuth.baggage.local-fields", sleuthBaggageProperties.getLocalFields()); for (String fieldName : localFields) { diff --git a/spring-cloud-sleuth-autoconfigure/src/test/java/org/springframework/cloud/sleuth/autoconfig/brave/baggage/CustomPropagationFactoryTests.java b/spring-cloud-sleuth-autoconfigure/src/test/java/org/springframework/cloud/sleuth/autoconfig/brave/baggage/CustomPropagationFactoryTests.java new file mode 100644 index 000000000..ad08b8385 --- /dev/null +++ b/spring-cloud-sleuth-autoconfigure/src/test/java/org/springframework/cloud/sleuth/autoconfig/brave/baggage/CustomPropagationFactoryTests.java @@ -0,0 +1,102 @@ +/* + * Copyright 2013-2021 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 + * + * https://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.autoconfig.brave.baggage; + +import java.util.Collections; +import java.util.List; + +import brave.internal.propagation.StringPropagationAdapter; +import brave.propagation.Propagation; +import brave.propagation.TraceContext; +import brave.propagation.TraceContextOrSamplingFlags; +import org.assertj.core.api.BDDAssertions; +import org.junit.jupiter.api.Test; + +import org.springframework.boot.actuate.autoconfigure.security.servlet.ManagementWebSecurityAutoConfiguration; +import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.boot.autoconfigure.mongo.MongoAutoConfiguration; +import org.springframework.boot.autoconfigure.quartz.QuartzAutoConfiguration; +import org.springframework.boot.test.context.runner.ApplicationContextRunner; +import org.springframework.cloud.gateway.config.GatewayAutoConfiguration; +import org.springframework.cloud.gateway.config.GatewayClassPathWarningAutoConfiguration; +import org.springframework.cloud.gateway.config.GatewayMetricsAutoConfiguration; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; + +public class CustomPropagationFactoryTests { + + @Test + void should_fail_to_start_the_context_when_propagation_type_custom_and_no_custom_propagation_provided() { + new ApplicationContextRunner().withUserConfiguration(Config.class) + .withPropertyValues("spring.sleuth.propagation.type=custom") + .run(context -> BDDAssertions.then(context).hasFailed()); + } + + @Test + void should_start_the_context_when_propagation_type_custom_and_no_custom_propagation_provided() { + new ApplicationContextRunner().withUserConfiguration(CustomConfig.class) + .withPropertyValues("spring.sleuth.propagation.type=custom").run(context -> BDDAssertions.then(context) + .hasNotFailed().getBean(CustomConfig.CustomPropagation.class)); + } + + @Configuration(proxyBeanMethods = false) + @EnableAutoConfiguration(exclude = { GatewayClassPathWarningAutoConfiguration.class, GatewayAutoConfiguration.class, + GatewayMetricsAutoConfiguration.class, ManagementWebSecurityAutoConfiguration.class, + MongoAutoConfiguration.class, QuartzAutoConfiguration.class }) + static class Config { + + } + + @Configuration(proxyBeanMethods = false) + @EnableAutoConfiguration(exclude = { GatewayClassPathWarningAutoConfiguration.class, GatewayAutoConfiguration.class, + GatewayMetricsAutoConfiguration.class, ManagementWebSecurityAutoConfiguration.class, + MongoAutoConfiguration.class, QuartzAutoConfiguration.class }) + static class CustomConfig { + + @Bean + CustomPropagation customPropagation() { + return new CustomPropagation(); + } + + static class CustomPropagation extends Propagation.Factory implements Propagation { + + @Override + public List keys() { + return Collections.emptyList(); + } + + @Override + public TraceContext.Injector injector(Setter setter) { + return (traceContext, request) -> { + }; + } + + @Override + public TraceContext.Extractor extractor(Getter getter) { + return request -> TraceContextOrSamplingFlags.EMPTY; + } + + @Override + public Propagation create(KeyFactory keyFactory) { + return StringPropagationAdapter.create(this, keyFactory); + } + + } + + } + +} diff --git a/spring-cloud-sleuth-brave/src/main/java/org/springframework/cloud/sleuth/brave/bridge/CompositePropagationFactorySupplier.java b/spring-cloud-sleuth-brave/src/main/java/org/springframework/cloud/sleuth/brave/bridge/CompositePropagationFactorySupplier.java index 269dedef2..3afdd8d3b 100644 --- a/spring-cloud-sleuth-brave/src/main/java/org/springframework/cloud/sleuth/brave/bridge/CompositePropagationFactorySupplier.java +++ b/spring-cloud-sleuth-brave/src/main/java/org/springframework/cloud/sleuth/brave/bridge/CompositePropagationFactorySupplier.java @@ -84,7 +84,7 @@ class CompositePropagationFactory extends Propagation.Factory implements Propaga W3CPropagation w3CPropagation = new W3CPropagation(braveBaggageManager, localFields); this.mapping.put(PropagationType.W3C, new AbstractMap.SimpleEntry<>(w3CPropagation, w3CPropagation.get())); LazyPropagationFactory lazyPropagationFactory = new LazyPropagationFactory( - beanFactory.getBeanProvider(Factory.class)); + beanFactory.getBeanProvider(PropagationFactorySupplier.class)); this.mapping.put(PropagationType.CUSTOM, new AbstractMap.SimpleEntry<>(lazyPropagationFactory, lazyPropagationFactory.get())); } @@ -161,17 +161,17 @@ class CompositePropagationFactory extends Propagation.Factory implements Propaga @SuppressWarnings("unchecked") private static final class LazyPropagationFactory extends Propagation.Factory { - private final ObjectProvider delegate; + private final ObjectProvider delegate; private volatile Propagation.Factory propagationFactory; - private LazyPropagationFactory(ObjectProvider delegate) { + private LazyPropagationFactory(ObjectProvider delegate) { this.delegate = delegate; } private Propagation.Factory propagationFactory() { if (this.propagationFactory == null) { - this.propagationFactory = this.delegate.getIfAvailable(() -> NoOpPropagation.INSTANCE); + this.propagationFactory = this.delegate.getIfAvailable(() -> () -> NoOpPropagation.INSTANCE).get(); } return this.propagationFactory; } diff --git a/spring-cloud-sleuth-brave/src/main/java/org/springframework/cloud/sleuth/brave/propagation/PropagationType.java b/spring-cloud-sleuth-brave/src/main/java/org/springframework/cloud/sleuth/brave/propagation/PropagationType.java index ffc744e8c..02a13ed19 100644 --- a/spring-cloud-sleuth-brave/src/main/java/org/springframework/cloud/sleuth/brave/propagation/PropagationType.java +++ b/spring-cloud-sleuth-brave/src/main/java/org/springframework/cloud/sleuth/brave/propagation/PropagationType.java @@ -40,7 +40,8 @@ public enum PropagationType { W3C, /** - * Custom propagation type. + * Custom propagation type. If picked, requires bean registration overriding the + * default propagation mechanisms. */ CUSTOM diff --git a/spring-cloud-sleuth-brave/src/test/java/org/springframework/cloud/sleuth/brave/bridge/CompositePropagationFactorySupplierTests.java b/spring-cloud-sleuth-brave/src/test/java/org/springframework/cloud/sleuth/brave/bridge/CompositePropagationFactorySupplierTests.java deleted file mode 100644 index 236bf90d7..000000000 --- a/spring-cloud-sleuth-brave/src/test/java/org/springframework/cloud/sleuth/brave/bridge/CompositePropagationFactorySupplierTests.java +++ /dev/null @@ -1,120 +0,0 @@ -/* - * Copyright 2013-2021 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 - * - * https://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.brave.bridge; - -import java.util.Collections; -import java.util.List; -import java.util.Objects; - -import brave.internal.codec.HexCodec; -import brave.internal.propagation.StringPropagationAdapter; -import brave.propagation.Propagation; -import brave.propagation.TraceContext; -import brave.propagation.TraceContextOrSamplingFlags; -import org.assertj.core.api.BDDAssertions; -import org.junit.jupiter.api.Test; -import org.mockito.Mockito; - -import org.springframework.beans.factory.BeanFactory; -import org.springframework.cloud.loadbalancer.support.SimpleObjectProvider; -import org.springframework.cloud.sleuth.brave.propagation.PropagationType; -import org.springframework.util.StringUtils; - -class CompositePropagationFactorySupplierTests { - - @Test - void should_pick_custom_registered_propagation_when_custom_mode_picked() { - BeanFactory beanFactory = Mockito.mock(BeanFactory.class); - Mockito.when(beanFactory.getBeanProvider(BraveBaggageManager.class)) - .thenReturn(new SimpleObjectProvider(new BraveBaggageManager())); - Mockito.when(beanFactory.getBeanProvider(Propagation.Factory.class)) - .thenReturn(new SimpleObjectProvider(new CustomTracePropagation())); - Mockito.when(beanFactory.getBeanProvider(Propagation.class)) - .thenReturn(new SimpleObjectProvider(new CustomTracePropagation())); - - CompositePropagationFactorySupplier supplier = new CompositePropagationFactorySupplier(beanFactory, - Collections.emptyList(), Collections.singletonList(PropagationType.CUSTOM)); - - BDDAssertions.then(supplier.get().get().keys()).containsExactly(CustomTraceExtractor.CUSTOM_TRACE_HEADER); - } - -} - -class CustomTracePropagation extends Propagation.Factory implements Propagation { - - public static final List KEYS = Collections.singletonList(CustomTraceExtractor.CUSTOM_TRACE_HEADER); - - @Override - public List keys() { - return KEYS; - } - - @Override - public TraceContext.Injector injector(Setter setter) { - return (traceContext, request) -> { - String trace = traceContext.traceIdString() + ":" + traceContext.spanIdString(); - setter.put(request, CustomTraceExtractor.CUSTOM_TRACE_HEADER, trace); - }; - } - - @Override - public TraceContext.Extractor extractor(Getter getter) { - Objects.requireNonNull(getter); - return new CustomTraceExtractor<>(getter); - } - - @Override - public Propagation create(KeyFactory keyFactory) { - return StringPropagationAdapter.create(this, keyFactory); - } - -} - -class CustomTraceExtractor implements TraceContext.Extractor { - - static final String CUSTOM_TRACE_HEADER = "x-custom-trace"; - - final Propagation.Getter getter; - - CustomTraceExtractor(Propagation.Getter getter) { - this.getter = getter; - } - - @Override - @SuppressWarnings("ReturnCount") - public TraceContextOrSamplingFlags extract(R request) { - String traceString = getter.get(request, CUSTOM_TRACE_HEADER); - if (!StringUtils.hasText(traceString)) { - return TraceContextOrSamplingFlags.EMPTY; - } - String[] trace = traceString.split(":"); - if (trace.length != 2) { - return TraceContextOrSamplingFlags.EMPTY; - } - - try { - TraceContext traceContext = TraceContext.newBuilder().traceId(HexCodec.lowerHexToUnsignedLong(trace[0])) - .spanId(HexCodec.lowerHexToUnsignedLong(trace[1])).build(); - - return TraceContextOrSamplingFlags.create(traceContext); - } - catch (NumberFormatException ex) { - return TraceContextOrSamplingFlags.EMPTY; - } - } - -}