diff --git a/docs/src/main/asciidoc/spring-cloud-sleuth.adoc b/docs/src/main/asciidoc/spring-cloud-sleuth.adoc index 99c9ed7d0..2b3d99417 100644 --- a/docs/src/main/asciidoc/spring-cloud-sleuth.adoc +++ b/docs/src/main/asciidoc/spring-cloud-sleuth.adoc @@ -876,6 +876,34 @@ In the following example, we register the `TracingFilter` bean, add the `ZIPKIN- include::{project-root}/tests/spring-cloud-sleuth-instrumentation-mvc-tests/src/test/java/org/springframework/cloud/sleuth/instrument/web/TraceFilterIntegrationTests.java[tags=response_headers,indent=0] ---- +=== Messaging + +Sleuth automatically configures the `MessagingTracing` bean which serves as a +foundation for Messaging instrumentation such as Kafka or JMS. + +If a customization of producer / consumer sampling of messaging traces is required, +just register a bean of type `brave.sampler.SamplerFunction` and +name the bean `sleuthProducerSampler` for producer sampler and `sleuthConsumerSampler` +for consumer sampler. + +For your convenience the `@ProducerSampler` and `@ConsumerSampler` +annotations can be used to inject the proper beans or to reference the bean +names via their static String `NAME` fields. + +Ex. Here's a sampler that traces 100 consumer requests per second, except for +the "alerts" channel. Other requests will use a global rate provided by the +`Tracing` component. + +[source,java] +---- +@Configuration +class Config { +include::{project-root}/tests/spring-cloud-sleuth-instrumentation-messaging-tests/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfigurationIntegrationTests.java[tags=custom_messaging_server_sampler,indent=2] +} +---- + +For more, see https://github.com/openzipkin/brave/tree/master/instrumentation/messaging#sampling-policy + === RPC Sleuth automatically configures the `RpcTracing` bean which serves as a diff --git a/spring-cloud-sleuth-core/pom.xml b/spring-cloud-sleuth-core/pom.xml index a89083377..805b11f57 100644 --- a/spring-cloud-sleuth-core/pom.xml +++ b/spring-cloud-sleuth-core/pom.xml @@ -211,6 +211,10 @@ io.zipkin.brave brave-context-log4j2 + + io.zipkin.brave + brave-instrumentation-messaging + io.zipkin.brave brave-instrumentation-rpc diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/ConsumerSampler.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/ConsumerSampler.java new file mode 100644 index 000000000..bebabc715 --- /dev/null +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/ConsumerSampler.java @@ -0,0 +1,50 @@ +/* + * Copyright 2013-2019 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 java.lang.annotation.Documented; +import java.lang.annotation.ElementType; +import java.lang.annotation.Inherited; +import java.lang.annotation.Retention; +import java.lang.annotation.RetentionPolicy; +import java.lang.annotation.Target; + +import brave.sampler.SamplerFunction; + +import org.springframework.beans.factory.annotation.Qualifier; + +/** + * Annotate a consumer {@link SamplerFunction} that should be injected to + * {@link brave.messaging.MessagingTracing.Builder#consumerSampler(SamplerFunction)}. + * + * @since 2.2.0 + * @see Qualifier + */ +@Target({ ElementType.FIELD, ElementType.METHOD, ElementType.PARAMETER, ElementType.TYPE, + ElementType.ANNOTATION_TYPE }) +@Retention(RetentionPolicy.RUNTIME) +@Inherited +@Documented +@Qualifier(ConsumerSampler.NAME) +public @interface ConsumerSampler { + + /** + * Default name for message consumer sampler. + */ + String NAME = "sleuthConsumerSampler"; + +} diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/ProducerSampler.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/ProducerSampler.java new file mode 100644 index 000000000..c289f1acf --- /dev/null +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/ProducerSampler.java @@ -0,0 +1,50 @@ +/* + * Copyright 2013-2019 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 java.lang.annotation.Documented; +import java.lang.annotation.ElementType; +import java.lang.annotation.Inherited; +import java.lang.annotation.Retention; +import java.lang.annotation.RetentionPolicy; +import java.lang.annotation.Target; + +import brave.sampler.SamplerFunction; + +import org.springframework.beans.factory.annotation.Qualifier; + +/** + * Annotate a producer {@link SamplerFunction} that should be injected to + * {@link brave.messaging.MessagingTracing.Builder#producerSampler(SamplerFunction)}. + * + * @since 2.2.0 + * @see Qualifier + */ +@Target({ ElementType.FIELD, ElementType.METHOD, ElementType.PARAMETER, ElementType.TYPE, + ElementType.ANNOTATION_TYPE }) +@Retention(RetentionPolicy.RUNTIME) +@Inherited +@Documented +@Qualifier(ProducerSampler.NAME) +public @interface ProducerSampler { + + /** + * Default name for messaging producer sampler. + */ + String NAME = "sleuthProducerSampler"; + +} diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfiguration.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfiguration.java index 2852e4107..6bddda2ff 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfiguration.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfiguration.java @@ -17,6 +17,7 @@ package org.springframework.cloud.sleuth.instrument.messaging; import java.lang.reflect.Field; +import java.util.ArrayList; import java.util.Collections; import java.util.List; @@ -25,7 +26,11 @@ import brave.Tracer; import brave.Tracing; import brave.jms.JmsTracing; import brave.kafka.clients.KafkaTracing; -import brave.propagation.Propagation; +import brave.messaging.MessagingRequest; +import brave.messaging.MessagingTracing; +import brave.messaging.MessagingTracingCustomizer; +import brave.propagation.Propagation.Getter; +import brave.sampler.SamplerFunction; import brave.spring.rabbit.SpringRabbitTracing; import org.aopalliance.intercept.MethodInterceptor; import org.aopalliance.intercept.MethodInvocation; @@ -44,6 +49,7 @@ import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.aop.framework.ProxyFactoryBean; import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanFactory; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.config.BeanDefinition; import org.springframework.beans.factory.config.BeanPostProcessor; import org.springframework.boot.autoconfigure.AutoConfigureAfter; @@ -67,6 +73,7 @@ import org.springframework.kafka.listener.MessageListener; import org.springframework.kafka.listener.MessageListenerContainer; import org.springframework.kafka.listener.adapter.MessagingMessageListenerAdapter; import org.springframework.kafka.support.DefaultKafkaHeaderMapper; +import org.springframework.lang.Nullable; import org.springframework.messaging.Message; import org.springframework.messaging.MessagingException; import org.springframework.messaging.converter.MessageConverter; @@ -89,6 +96,29 @@ import org.springframework.util.ReflectionUtils; @EnableConfigurationProperties(SleuthMessagingProperties.class) public class TraceMessagingAutoConfiguration { + @Autowired(required = false) + List messagingTracingCustomizers = new ArrayList<>(); + + @Bean + @ConditionalOnMissingBean + // NOTE: stable bean name as might be used outside sleuth + MessagingTracing messagingTracing(Tracing tracing, + @Nullable @ProducerSampler SamplerFunction producerSampler, + @Nullable @ConsumerSampler SamplerFunction consumerSampler) { + + MessagingTracing.Builder builder = MessagingTracing.newBuilder(tracing); + if (producerSampler != null) { + builder.producerSampler(producerSampler); + } + if (consumerSampler != null) { + builder.consumerSampler(consumerSampler); + } + for (MessagingTracingCustomizer customizer : this.messagingTracingCustomizers) { + customizer.customize(builder); + } + return builder.build(); + } + @Configuration(proxyBeanMethods = false) @ConditionalOnProperty(value = "spring.sleuth.messaging.rabbit.enabled", matchIfMissing = true) @@ -105,9 +135,9 @@ public class TraceMessagingAutoConfiguration { @Bean @ConditionalOnMissingBean - SpringRabbitTracing springRabbitTracing(Tracing tracing, + SpringRabbitTracing springRabbitTracing(MessagingTracing messagingTracing, SleuthMessagingProperties properties) { - return SpringRabbitTracing.newBuilder(tracing) + return SpringRabbitTracing.newBuilder(messagingTracing) .remoteServiceName( properties.getMessaging().getRabbit().getRemoteServiceName()) .build(); @@ -123,8 +153,9 @@ public class TraceMessagingAutoConfiguration { @Bean @ConditionalOnMissingBean - KafkaTracing kafkaTracing(Tracing tracing, SleuthMessagingProperties properties) { - return KafkaTracing.newBuilder(tracing) + KafkaTracing kafkaTracing(MessagingTracing messagingTracing, + SleuthMessagingProperties properties) { + return KafkaTracing.newBuilder(messagingTracing) .remoteServiceName( properties.getMessaging().getKafka().getRemoteServiceName()) .build(); @@ -164,8 +195,9 @@ public class TraceMessagingAutoConfiguration { @Bean @ConditionalOnMissingBean - JmsTracing jmsTracing(Tracing tracing, SleuthMessagingProperties properties) { - return JmsTracing.newBuilder(tracing) + JmsTracing jmsTracing(MessagingTracing messagingTracing, + SleuthMessagingProperties properties) { + return JmsTracing.newBuilder(messagingTracing) .remoteServiceName( properties.getMessaging().getJms().getRemoteServiceName()) .build(); @@ -207,9 +239,9 @@ public class TraceMessagingAutoConfiguration { @Bean TracingMethodMessageHandlerAdapter tracingMethodMessageHandlerAdapter( - Tracing tracing, - Propagation.Getter traceMessagePropagationGetter) { - return new TracingMethodMessageHandlerAdapter(tracing, + MessagingTracing messagingTracing, + Getter traceMessagePropagationGetter) { + return new TracingMethodMessageHandlerAdapter(messagingTracing, traceMessagePropagationGetter); } @@ -477,7 +509,7 @@ class SqsQueueMessageHandlerFactory extends QueueMessageHandlerFactory { class SqsQueueMessageHandler extends QueueMessageHandler { // copied from QueueMessageHandler - private static final String LOGICAL_RESOURCE_ID = "LogicalResourceId"; + static final String LOGICAL_RESOURCE_ID = "LogicalResourceId"; private TracingMethodMessageHandlerAdapter handlerAdapter; diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingMethodMessageHandlerAdapter.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingMethodMessageHandlerAdapter.java index e3defb22a..6167050dd 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingMethodMessageHandlerAdapter.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingMethodMessageHandlerAdapter.java @@ -21,8 +21,10 @@ import java.util.function.BiConsumer; import brave.Span; import brave.Tracer; import brave.Tracing; -import brave.propagation.Propagation; -import brave.propagation.TraceContext; +import brave.messaging.ConsumerRequest; +import brave.messaging.MessagingTracing; +import brave.propagation.Propagation.Getter; +import brave.propagation.TraceContext.Extractor; import brave.propagation.TraceContextOrSamplingFlags; import org.springframework.messaging.Message; @@ -30,6 +32,7 @@ import org.springframework.messaging.MessageHandler; import org.springframework.messaging.support.MessageHeaderAccessor; import static brave.Span.Kind.CONSUMER; +import static org.springframework.cloud.sleuth.instrument.messaging.SqsQueueMessageHandler.LOGICAL_RESOURCE_ID; /** * Adds tracing extraction to an instance of @@ -46,22 +49,26 @@ import static brave.Span.Kind.CONSUMER; */ class TracingMethodMessageHandlerAdapter { - private Tracing tracing; + private final Tracing tracing; - private Tracer tracer; + private final Tracer tracer; - private TraceContext.Extractor extractor; + private final Extractor extractor; - TracingMethodMessageHandlerAdapter(Tracing tracing, - Propagation.Getter traceMessagePropagationGetter) { - this.tracing = tracing; + private final Getter getter; + + TracingMethodMessageHandlerAdapter(MessagingTracing messagingTracing, + Getter getter) { + this.tracing = messagingTracing.tracing(); this.tracer = tracing.tracer(); - this.extractor = tracing.propagation().extractor(traceMessagePropagationGetter); + this.extractor = tracing.propagation().extractor(MessageConsumerRequest.GETTER); + this.getter = getter; } void wrapMethodMessageHandler(Message message, MessageHandler messageHandler, BiConsumer> messageSpanTagger) { - TraceContextOrSamplingFlags extracted = extractAndClearHeaders(message); + MessageConsumerRequest request = new MessageConsumerRequest(message, this.getter); + TraceContextOrSamplingFlags extracted = extractAndClearHeaders(request); Span consumerSpan = tracer.nextSpan(extracted); Span listenerSpan = tracer.newChild(consumerSpan.context()); @@ -95,15 +102,77 @@ class TracingMethodMessageHandlerAdapter { } } - private TraceContextOrSamplingFlags extractAndClearHeaders(Message message) { - MessageHeaderAccessor headers = MessageHeaderAccessor.getMutableAccessor(message); - TraceContextOrSamplingFlags extracted = extractor.extract(headers); + private TraceContextOrSamplingFlags extractAndClearHeaders( + MessageConsumerRequest request) { + TraceContextOrSamplingFlags extracted = extractor.extract(request); for (String propagationKey : tracing.propagation().keys()) { - headers.removeHeader(propagationKey); + request.removeHeader(propagationKey); } return extracted; } } + +final class MessageConsumerRequest extends ConsumerRequest { + + static final Getter GETTER = new Getter() { + @Override + public String get(MessageConsumerRequest request, String name) { + return request.getHeader(name); + } + + @Override + public String toString() { + return "MessageConsumerRequest::getHeader"; + } + }; + + final Message delegate; + + final MessageHeaderAccessor mutableHeaders; + + final Getter getter; + + MessageConsumerRequest(Message delegate, + Getter getter) { + this.delegate = delegate; + this.mutableHeaders = MessageHeaderAccessor.getMutableAccessor(delegate); + this.getter = getter; + } + + @Override + public Span.Kind spanKind() { + return Span.Kind.CONSUMER; + } + + @Override + public Object unwrap() { + return this.delegate; + } + + @Override + public String operation() { + return "receive"; + } + + @Override + public String channelKind() { + return "queue"; + } + + @Override + public String channelName() { + return this.delegate.getHeaders().get(LOGICAL_RESOURCE_ID).toString(); + } + + String getHeader(String name) { + return this.getter.get(this.mutableHeaders, name); + } + + void removeHeader(String name) { + this.mutableHeaders.removeHeader(name); + } + +} diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/autoconfig/TraceAutoConfigurationCustomizersTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/autoconfig/TraceAutoConfigurationCustomizersTests.java index 63799ff98..8a0c2689a 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/autoconfig/TraceAutoConfigurationCustomizersTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/autoconfig/TraceAutoConfigurationCustomizersTests.java @@ -18,8 +18,10 @@ package org.springframework.cloud.sleuth.autoconfig; import brave.TracingCustomizer; import brave.http.HttpTracingCustomizer; +import brave.messaging.MessagingTracingCustomizer; import brave.propagation.CurrentTraceContextCustomizer; import brave.propagation.ExtraFieldCustomizer; +import brave.propagation.Propagation; import brave.rpc.RpcTracingCustomizer; import brave.sampler.Sampler; import org.junit.Test; @@ -27,11 +29,13 @@ import org.junit.Test; import org.springframework.boot.autoconfigure.AutoConfigurations; import org.springframework.boot.test.context.assertj.AssertableApplicationContext; import org.springframework.boot.test.context.runner.ApplicationContextRunner; +import org.springframework.cloud.sleuth.instrument.messaging.TraceMessagingAutoConfiguration; import org.springframework.cloud.sleuth.instrument.rpc.TraceRpcAutoConfiguration; import org.springframework.cloud.sleuth.instrument.web.TraceHttpAutoConfiguration; import org.springframework.cloud.sleuth.instrument.web.TraceWebAutoConfiguration; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.messaging.support.MessageHeaderAccessor; import static org.assertj.core.api.BDDAssertions.then; @@ -40,7 +44,9 @@ public class TraceAutoConfigurationCustomizersTests { private final ApplicationContextRunner contextRunner = new ApplicationContextRunner() .withConfiguration(AutoConfigurations.of(TraceAutoConfiguration.class, TraceWebAutoConfiguration.class, TraceHttpAutoConfiguration.class, - TraceRpcAutoConfiguration.class)) + TraceRpcAutoConfiguration.class, + FakeSpringMessagingAutoConfiguration.class, + TraceMessagingAutoConfiguration.class)) .withUserConfiguration(Customizers.class); @Test @@ -75,6 +81,17 @@ public class TraceAutoConfigurationCustomizersTests { then(bean.rpcCustomizerApplied).isTrue(); } + // SQS has a dependency on the getter and this is better than exposing things public + @Configuration + static class FakeSpringMessagingAutoConfiguration { + + @Bean + Propagation.Getter traceMessagePropagationGetter() { + return (headers, key) -> null; + } + + } + @Configuration static class Customizers { @@ -88,6 +105,8 @@ public class TraceAutoConfigurationCustomizersTests { boolean rpcCustomizerApplied; + boolean messagingCustomizerApplied; + @Bean TracingCustomizer sleuthTracingCustomizer() { return builder -> tracingCustomizerApplied = true; @@ -108,6 +127,11 @@ public class TraceAutoConfigurationCustomizersTests { return builder -> httpCustomizerApplied = true; } + @Bean + MessagingTracingCustomizer sleuthMessagingCustomizer() { + return builder -> messagingCustomizerApplied = true; + } + @Bean RpcTracingCustomizer sleuthRpcTracingCustomizer() { return builder -> rpcCustomizerApplied = true; diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfigurationIntegrationTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfigurationIntegrationTests.java new file mode 100644 index 000000000..c90c04060 --- /dev/null +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfigurationIntegrationTests.java @@ -0,0 +1,74 @@ +/* + * Copyright 2013-2019 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 brave.messaging.MessagingRequest; +import brave.messaging.MessagingRuleSampler; +import brave.sampler.Matchers; +import brave.sampler.RateLimitingSampler; +import brave.sampler.Sampler; +import brave.sampler.SamplerFunction; +import org.junit.Test; +import org.junit.runner.RunWith; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.cloud.sleuth.util.ArrayListSpanReporter; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.test.context.junit4.SpringRunner; + +import static brave.messaging.MessagingRequestMatchers.channelNameEquals; +import static org.assertj.core.api.BDDAssertions.then; + +@RunWith(SpringRunner.class) +@SpringBootTest(classes = TraceMessagingAutoConfigurationIntegrationTests.Config.class, + webEnvironment = SpringBootTest.WebEnvironment.NONE) +public class TraceMessagingAutoConfigurationIntegrationTests { + + @Autowired + @ConsumerSampler + SamplerFunction sampler; + + @Test + public void should_inject_messaging_sampler() { + then(this.sampler).isNotNull(); + } + + @EnableAutoConfiguration + @Configuration + public static class Config { + + @Bean + ArrayListSpanReporter reporter() { + return new ArrayListSpanReporter(); + } + + // tag::custom_messaging_consumer_sampler[] + @Bean(name = ConsumerSampler.NAME) + SamplerFunction myMessagingSampler() { + return MessagingRuleSampler.newBuilder() + .putRule(channelNameEquals("alerts"), Sampler.NEVER_SAMPLE) + .putRule(Matchers.alwaysMatch(), RateLimitingSampler.create(100)) + .build(); + } + // end::custom_messaging_consumer_sampler[] + + } + +} diff --git a/spring-cloud-sleuth-dependencies/pom.xml b/spring-cloud-sleuth-dependencies/pom.xml index ee08eec5d..2c643a603 100644 --- a/spring-cloud-sleuth-dependencies/pom.xml +++ b/spring-cloud-sleuth-dependencies/pom.xml @@ -31,8 +31,8 @@ spring-cloud-sleuth-dependencies Spring Cloud Sleuth Dependencies - 5.8.0 - 0.34.1 + 5.9.0 + 0.35.0 3.4.1 diff --git a/tests/spring-cloud-sleuth-instrumentation-messaging-tests/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfigurationTests.java b/tests/spring-cloud-sleuth-instrumentation-messaging-tests/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfigurationTests.java index f52e7b465..96e7d7d9d 100644 --- a/tests/spring-cloud-sleuth-instrumentation-messaging-tests/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfigurationTests.java +++ b/tests/spring-cloud-sleuth-instrumentation-messaging-tests/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfigurationTests.java @@ -18,7 +18,11 @@ package org.springframework.cloud.sleuth.instrument.messaging; import brave.Tracer; import brave.kafka.clients.KafkaTracing; +import brave.messaging.MessagingRequest; +import brave.messaging.MessagingTracing; import brave.sampler.Sampler; +import brave.sampler.SamplerFunction; +import brave.sampler.SamplerFunctions; import brave.spring.rabbit.SpringRabbitTracing; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerRecord; @@ -32,8 +36,11 @@ import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.autoconfigure.AutoConfigurations; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.boot.test.context.runner.ApplicationContextRunner; +import org.springframework.cloud.sleuth.autoconfig.TraceAutoConfiguration; import org.springframework.cloud.sleuth.util.ArrayListSpanReporter; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -102,6 +109,55 @@ public class TraceMessagingAutoConfigurationTests { then(this.testSleuthKafkaHeaderMapperBeanPostProcessor.tracingCalled).isTrue(); } + @Test + public void defaultsToBraveProducerSampler() { + contextRunner().run((context) -> { + SamplerFunction producerSampler = context + .getBean(MessagingTracing.class).producerSampler(); + + then(producerSampler).isSameAs(SamplerFunctions.deferDecision()); + }); + } + + @Test + public void configuresUserProvidedProducerSampler() { + contextRunner().withUserConfiguration(ProducerSamplerConfig.class) + .run((context) -> { + SamplerFunction producerSampler = context + .getBean(MessagingTracing.class).producerSampler(); + + then(producerSampler).isSameAs(ProducerSamplerConfig.INSTANCE); + }); + } + + @Test + public void defaultsToBraveConsumerSampler() { + contextRunner().run((context) -> { + SamplerFunction consumerSampler = context + .getBean(MessagingTracing.class).consumerSampler(); + + then(consumerSampler).isSameAs(SamplerFunctions.deferDecision()); + }); + } + + @Test + public void configuresUserProvidedConsumerSampler() { + contextRunner().withUserConfiguration(ConsumerSamplerConfig.class) + .run((context) -> { + SamplerFunction consumerSampler = context + .getBean(MessagingTracing.class).consumerSampler(); + + then(consumerSampler).isSameAs(ConsumerSamplerConfig.INSTANCE); + }); + } + + private ApplicationContextRunner contextRunner(String... propertyValues) { + return new ApplicationContextRunner().withPropertyValues(propertyValues) + .withConfiguration(AutoConfigurations.of(TraceAutoConfiguration.class, + TraceMessagingAutoConfiguration.class, + TraceMessagingAutoConfiguration.class)); + } + @Configuration @EnableAutoConfiguration protected static class Config { @@ -225,3 +281,27 @@ class TestSleuthKafkaHeaderMapperBeanPostProcessor } } + +@Configuration +class ProducerSamplerConfig { + + static final SamplerFunction INSTANCE = request -> null; + + @Bean(ProducerSampler.NAME) + SamplerFunction sleuthProducerSampler() { + return INSTANCE; + } + +} + +@Configuration +class ConsumerSamplerConfig { + + static final SamplerFunction INSTANCE = request -> null; + + @Bean(ConsumerSampler.NAME) + SamplerFunction sleuthConsumerSampler() { + return INSTANCE; + } + +}