Please work

This commit is contained in:
Marcin Grzejszczak
2018-10-11 13:05:16 +02:00
parent c6d96be4fd
commit a34b82f47f
11 changed files with 124 additions and 40 deletions

View File

@@ -343,7 +343,7 @@ The following example shows setting baggage on a span:
----
Span initialSpan = this.tracer.nextSpan().name("span").start();
ExtraFieldPropagation.set(initialSpan.context(), "foo", "bar");
ExtraFieldPropagation.set(initialSpan.context(),"UPPER_CASE", "someValue");
ExtraFieldPropagation.set(initialSpan.context(), "UPPER_CASE", "someValue");
}
----
@@ -380,9 +380,9 @@ spring.sleuth:
[source,java]
----
initialSpan.tag("foo",
ExtraFieldPropagation.get(initialSpan.context(), "foo"));
ExtraFieldPropagation.get(initialSpan.context(), "foo"));
initialSpan.tag("UPPER_CASE",
ExtraFieldPropagation.get(initialSpan.context(), "UPPER_CASE"));
ExtraFieldPropagation.get(initialSpan.context(), "UPPER_CASE"));
----
[[sleuth-adding-project]]

View File

@@ -36,7 +36,7 @@ public interface SpanAdjuster {
* this interface can be used to alter then name. Example:
*
* {@code span -> span.toBuilder().name(scrub(span.getName())).build();}
* @param - span to adjust
* @param span to adjust
* @return - adjusted span
*/
Span adjust(Span span);

View File

@@ -43,6 +43,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.config.BeanDefinition;
import org.springframework.beans.factory.config.BeanPostProcessor;
import org.springframework.boot.autoconfigure.AutoConfigureAfter;
import org.springframework.boot.autoconfigure.condition.ConditionalOnBean;
@@ -53,6 +54,7 @@ import org.springframework.boot.context.properties.EnableConfigurationProperties
import org.springframework.cloud.sleuth.autoconfig.TraceAutoConfiguration;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Role;
import org.springframework.jms.annotation.JmsListenerConfigurer;
import org.springframework.kafka.core.ProducerFactory;
import org.springframework.kafka.listener.AbstractMessageListenerContainer;
@@ -126,6 +128,7 @@ public class TraceMessagingAutoConfiguration {
@Configuration
@ConditionalOnProperty(value = "spring.sleuth.messaging.jms.enabled", matchIfMissing = true)
@ConditionalOnClass(JmsListenerConfigurer.class)
@Role(BeanDefinition.ROLE_INFRASTRUCTURE)
protected static class SleuthJmsConfiguration {
@Bean

View File

@@ -25,7 +25,7 @@ import reactor.util.context.Context;
/**
* A lazy representation of the {@link SpanSubscription}.
*
* @param - type of what subscription returns
* @param <T> of what subscription returns
* @author Marcin Grzejszczak
* @since 2.0.0
*/

View File

@@ -83,7 +83,7 @@ public abstract class ReactorSleuth {
if (beanFactory.isActive()) {
if (log.isTraceEnabled()) {
log.trace(
"Spring Context already refreshed. Creating a scope "
"Spring Context [" + beanFactory + "] already refreshed. Creating a scope "
+ "passing span subscriber with Reactor Context "
+ "[" + sub.currentContext()
+ "] and name [" + scannable.name()

View File

@@ -30,7 +30,7 @@ import reactor.util.context.Context;
/**
* A trace representation of the {@link Subscriber} that always continues a span.
*
* @param - span subscription type
* @param <T> subscription type
* @author Marcin Grzejszczak
* @since 2.0.0
*/

View File

@@ -22,6 +22,8 @@ import java.util.function.Supplier;
import javax.annotation.PreDestroy;
import brave.Tracing;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import reactor.core.publisher.Hooks;
import reactor.core.publisher.Mono;
import reactor.core.scheduler.Schedulers;
@@ -60,6 +62,8 @@ public class TraceReactorAutoConfiguration {
@ConditionalOnBean(Tracing.class)
static class TraceReactorConfiguration {
private static final Log log = LogFactory.getLog(TraceReactorConfiguration.class);
static final String SLEUTH_TRACE_REACTOR_KEY = TraceReactorConfiguration.class
.getName();
@@ -73,6 +77,9 @@ public class TraceReactorAutoConfiguration {
@PreDestroy
public void cleanupHooks() {
if (log.isTraceEnabled()) {
log.trace("Cleaning up hooks");
}
Hooks.resetOnEachOperator(SLEUTH_TRACE_REACTOR_KEY);
Schedulers.resetFactory();
}
@@ -114,9 +121,7 @@ class HookRegisteringBeanDefinitionRegistryPostProcessor
@Override
public ScheduledExecutorService decorateExecutorService(String schedulerType,
Supplier<? extends ScheduledExecutorService> actual) {
return new TraceableScheduledExecutorService(
HookRegisteringBeanDefinitionRegistryPostProcessor.this.context,
actual.get());
return new TraceableScheduledExecutorService(beanFactory, actual.get());
}
};
}

View File

@@ -17,7 +17,6 @@
package org.springframework.cloud.sleuth.annotation;
import java.util.List;
import java.util.concurrent.ConcurrentLinkedQueue;
import brave.Tracer;
import brave.sampler.Sampler;
@@ -26,12 +25,14 @@ import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.assertj.core.api.BDDAssertions;
import org.awaitility.Awaitility;
import org.junit.AfterClass;
import org.junit.Before;
import org.junit.Ignore;
import org.junit.BeforeClass;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import reactor.core.publisher.Hooks;
import reactor.core.publisher.Mono;
import zipkin2.Span;
import zipkin2.reporter.Reporter;
@@ -40,43 +41,51 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.actuate.trace.http.HttpTrace;
import org.springframework.boot.actuate.trace.http.HttpTraceRepository;
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.boot.autoconfigure.ImportAutoConfiguration;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.boot.web.server.LocalServerPort;
import org.springframework.cloud.sleuth.DisableWebFluxSecurity;
import org.springframework.cloud.sleuth.instrument.reactor.TraceReactorAutoConfigurationAccessorConfiguration;
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.SpringJUnit4ClassRunner;
import org.springframework.test.annotation.DirtiesContext;
import org.springframework.test.context.junit4.SpringRunner;
import org.springframework.test.web.reactive.server.WebTestClient;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.reactive.function.client.WebClient;
import static org.assertj.core.api.BDDAssertions.then;
@RunWith(SpringJUnit4ClassRunner.class)
@RunWith(SpringRunner.class)
@SpringBootTest(properties = { "spring.main.web-application-type=reactive" }, classes = {
SleuthSpanCreatorAspectWebFluxTests.TestEndpoint.class,
SleuthSpanCreatorAspectWebFluxTests.TestConfiguration.class }, webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT)
SleuthSpanCreatorAspectWebFluxTests.TestConfiguration.class },
webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT)
@DirtiesContext
public class SleuthSpanCreatorAspectWebFluxTests {
private static final Log log = LogFactory
.getLog(SleuthSpanCreatorAspectWebFluxTests.class);
private final WebClient webClient = WebClient.create();
@Autowired
Tracer tracer;
@Autowired
AccessLoggingHttpTraceRepository repository;
@Autowired
ArrayListSpanReporter reporter;
private WebTestClient webClient;
@LocalServerPort
private int port;
@AfterClass
@BeforeClass
public static void cleanup() {
System.out.println("DUPA2");
Hooks.resetOnLastOperator();
TraceReactorAutoConfigurationAccessorConfiguration.close();
}
private static String toHexString(Long value) {
BDDAssertions.then(value).isNotNull();
return StringUtils.leftPad(Long.toHexString(value), 16, '0');
@@ -86,12 +95,15 @@ public class SleuthSpanCreatorAspectWebFluxTests {
public void setup() {
this.reporter.clear();
this.repository.clear();
log.info("Running app on port [" + this.port + "]");
this.webClient = WebTestClient.bindToServer().baseUrl("http://localhost:" + port)
.build();
}
@Test
public void shouldReturnSpanFromWebFluxTraceContext() {
Mono<Long> mono = webClient.get().uri("http://localhost:" + port + "/test/ping")
.retrieve().bodyToMono(Long.class);
Mono<Long> mono = webClient.get().uri("/test/ping").exchange()
.returnResult(Long.class).getResponseBody().single();
then(this.reporter.getSpans()).isEmpty();
@@ -116,9 +128,8 @@ public class SleuthSpanCreatorAspectWebFluxTests {
@Test
public void shouldReturnSpanFromWebFluxSubscriptionContext() {
Mono<Long> mono = webClient.get()
.uri("http://localhost:" + port + "/test/pingFromContext").retrieve()
.bodyToMono(Long.class);
Mono<Long> mono = webClient.get().uri("/test/pingFromContext").exchange()
.returnResult(Long.class).getResponseBody().single();
then(this.reporter.getSpans()).isEmpty();
@@ -137,9 +148,8 @@ public class SleuthSpanCreatorAspectWebFluxTests {
@Test
public void shouldContinueSpanInWebFlux() {
Mono<Long> mono = webClient.get()
.uri("http://localhost:" + port + "/test/continueSpan").retrieve()
.bodyToMono(Long.class);
Mono<Long> mono = webClient.get().uri("/test/continueSpan").exchange()
.returnResult(Long.class).getResponseBody().single();
then(this.reporter.getSpans()).isEmpty();
@@ -157,9 +167,8 @@ public class SleuthSpanCreatorAspectWebFluxTests {
@Test
public void shouldCreateNewSpanInWebFlux() {
Mono<Long> mono = webClient.get()
.uri("http://localhost:" + port + "/test/newSpan1").retrieve()
.bodyToMono(Long.class);
Mono<Long> mono = webClient.get().uri("/test/newSpan1").exchange()
.returnResult(Long.class).getResponseBody().single();
then(this.reporter.getSpans()).isEmpty();
@@ -178,9 +187,8 @@ public class SleuthSpanCreatorAspectWebFluxTests {
@Test
public void shouldCreateNewSpanInWebFluxInSubscriberContext() {
Mono<Long> mono = webClient.get()
.uri("http://localhost:" + port + "/test/newSpan2").retrieve()
.bodyToMono(Long.class);
Mono<Long> mono = webClient.get().uri("/test/newSpan2").exchange()
.returnResult(Long.class).getResponseBody().single();
then(this.reporter.getSpans()).isEmpty();
@@ -201,8 +209,8 @@ public class SleuthSpanCreatorAspectWebFluxTests {
public void shouldSetupCorrectSpanInHttpTrace() {
repository.clear();
Mono<Long> mono = webClient.get().uri("http://localhost:" + port + "/test/ping")
.retrieve().bodyToMono(Long.class);
Mono<Long> mono = webClient.get().uri("/test/ping").exchange()
.returnResult(Long.class).getResponseBody().single();
then(this.reporter.getSpans()).isEmpty();
@@ -223,6 +231,7 @@ public class SleuthSpanCreatorAspectWebFluxTests {
@Configuration
@EnableAutoConfiguration
@DisableWebFluxSecurity
@ImportAutoConfiguration(TraceReactorAutoConfigurationAccessorConfiguration.class)
protected static class TestConfiguration {
@Bean
@@ -293,11 +302,13 @@ public class SleuthSpanCreatorAspectWebFluxTests {
@GetMapping("/ping")
Mono<Long> ping() {
log.info("ping");
return Mono.just(tracer.currentSpan().context().spanId());
}
@GetMapping("/pingFromContext")
Mono<Long> pingFromContext() {
log.info("pingFromContext");
return Mono.subscriberContext()
.doOnSuccess(context -> log.info("Ping from context"))
.flatMap(context -> Mono

View File

@@ -0,0 +1,17 @@
package org.springframework.cloud.sleuth.instrument.reactor;
import org.junit.Ignore;
import org.junit.runner.RunWith;
import org.junit.runners.Suite;
import org.springframework.cloud.sleuth.annotation.SleuthSpanCreatorAspectWebFluxTests;
import org.springframework.cloud.sleuth.instrument.web.TraceWebFluxTests;
@RunWith(Suite.class)
@Suite.SuiteClasses({ //
TraceWebFluxTests.class, SleuthSpanCreatorAspectWebFluxTests.class //
})
@Ignore
public class AdhocTestSuite {
}

View File

@@ -0,0 +1,43 @@
package org.springframework.cloud.sleuth.instrument.reactor;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import reactor.core.publisher.Hooks;
import reactor.core.scheduler.Schedulers;
import org.springframework.boot.autoconfigure.AutoConfigureBefore;
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
* @author Marcin Grzejszczak
*/
@Configuration
@AutoConfigureBefore(TraceReactorAutoConfiguration.class)
public class TraceReactorAutoConfigurationAccessorConfiguration {
private static final Log log = LogFactory
.getLog(TraceReactorAutoConfigurationAccessorConfiguration.class);
public static void close() {
if (log.isTraceEnabled()) {
log.trace("Cleaning up hooks");
}
new TraceReactorAutoConfiguration.TraceReactorConfiguration().cleanupHooks();
Hooks.resetOnEachOperator();
Hooks.resetOnLastOperator();
Schedulers.resetFactory();
}
@Bean
static HookRegisteringBeanDefinitionRegistryPostProcessor testTraceHookRegisteringBeanDefinitionRegistryPostProcessor(
ConfigurableApplicationContext context) {
log.info("Running clean up and creating the post processor");
close();
return TraceReactorAutoConfiguration.TraceReactorConfiguration
.traceHookRegisteringBeanDefinitionRegistryPostProcessor(context);
}
}

View File

@@ -20,18 +20,22 @@ import brave.Span;
import brave.Tracer;
import brave.sampler.Sampler;
import org.awaitility.Awaitility;
import org.junit.AfterClass;
import org.junit.BeforeClass;
import org.junit.Test;
import org.slf4j.MDC;
import org.springframework.boot.WebApplicationType;
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.boot.autoconfigure.ImportAutoConfiguration;
import org.springframework.boot.builder.SpringApplicationBuilder;
import org.springframework.cloud.sleuth.DisableWebFluxSecurity;
import org.springframework.cloud.sleuth.instrument.reactor.TraceReactorAutoConfigurationAccessorConfiguration;
import org.springframework.cloud.sleuth.instrument.web.client.TraceWebClientAutoConfiguration;
import org.springframework.cloud.sleuth.util.ArrayListSpanReporter;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Import;
import org.springframework.core.env.Environment;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
@@ -45,7 +49,6 @@ import org.springframework.web.reactive.function.server.ServerResponse;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Hooks;
import reactor.core.publisher.Mono;
import reactor.core.scheduler.Schedulers;
import static org.assertj.core.api.BDDAssertions.then;
@@ -54,9 +57,10 @@ public class TraceWebFluxTests {
public static final String EXPECTED_TRACE_ID = "b919095138aa4c6e";
@BeforeClass
@AfterClass
public static void setup() {
Hooks.resetOnLastOperator();
Schedulers.resetFactory();
TraceReactorAutoConfigurationAccessorConfiguration.close();
}
@Test
@@ -169,6 +173,7 @@ public class TraceWebFluxTests {
@Configuration
@EnableAutoConfiguration(exclude = { TraceWebClientAutoConfiguration.class })
@DisableWebFluxSecurity
@ImportAutoConfiguration(TraceReactorAutoConfigurationAccessorConfiguration.class)
static class Config {
@Bean