Bringing back reactor & webflux support

This commit is contained in:
Marcin Grzejszczak
2017-10-08 20:38:18 +02:00
parent 525e32a7a2
commit 67e7d07b80
7 changed files with 648 additions and 0 deletions

View File

@@ -0,0 +1,52 @@
package org.springframework.cloud.sleuth.instrument.reactor;
import java.util.function.Function;
import java.util.function.Predicate;
import org.reactivestreams.Publisher;
import org.springframework.cloud.sleuth.Tracer;
import reactor.core.Fuseable;
import reactor.core.Scannable;
import reactor.core.publisher.Operators;
/**
* Reactive Span pointcuts factories
*
* @author Stephane Maldini
* @since 2.0.0
*/
public abstract class ReactorSleuth {
/**
* Return a span operator pointcut given a {@link Tracer}. This can be used in reactor
* via {@link reactor.core.publisher.Flux#transform(Function)}, {@link
* reactor.core.publisher.Mono#transform(Function)}, {@link
* reactor.core.publisher.Hooks#onEachOperator(Function)} or {@link
* reactor.core.publisher.Hooks#onLastOperator(Function)}.
*
* @param tracer the {@link Tracer} instance to use in this span operator
* @param <T> an arbitrary type that is left unchanged by the span operator
*
* @return a new Span operator pointcut
*/
public static <T> Function<? super Publisher<T>, ? extends Publisher<T>> spanOperator(Tracer tracer) {
return Operators.lift(POINTCUT_FILTER, ((scannable, sub) -> {
//do not trace fused flows
if(scannable instanceof Fuseable && sub instanceof Fuseable.QueueSubscription){
return sub;
}
return new SpanSubscriber<>(
sub,
sub.currentContext(),
tracer,
scannable.name());
}));
}
private static final Predicate<Scannable> POINTCUT_FILTER =
s -> !(s instanceof Fuseable.ScalarCallable);
private ReactorSleuth() {
}
}

View File

@@ -0,0 +1,177 @@
package org.springframework.cloud.sleuth.instrument.reactor;
import java.util.concurrent.atomic.AtomicBoolean;
import org.reactivestreams.Subscriber;
import org.reactivestreams.Subscription;
import org.springframework.cloud.sleuth.Span;
import org.springframework.cloud.sleuth.Tracer;
import reactor.core.CoreSubscriber;
import reactor.util.Logger;
import reactor.util.Loggers;
import reactor.util.context.Context;
/**
* A trace representation of the {@link Subscriber}
*
* @author Stephane Maldini
* @author Marcin Grzejszczak
* @since 2.0.0
*/
final class SpanSubscriber<T> extends AtomicBoolean implements Subscription,
CoreSubscriber<T> {
private static final Logger log = Loggers.getLogger(SpanSubscriber.class);
private final Span span;
private final Span rootSpan;
private final Subscriber<? super T> subscriber;
private final Context context;
private final Tracer tracer;
private Subscription s;
SpanSubscriber(Subscriber<? super T> subscriber, Context ctx, Tracer tracer,
String name) {
this.subscriber = subscriber;
this.tracer = tracer;
Span root = ctx.getOrDefault(Span.class, tracer.getCurrentSpan());
if (log.isTraceEnabled()) {
log.trace("Span from context [{}]", root);
}
this.rootSpan = root;
if (log.isTraceEnabled()) {
log.trace("Stored context root span [{}]", this.rootSpan);
}
this.span = tracer.createSpan(name, root);
if (log.isTraceEnabled()) {
log.trace("Created span [{}], with name [{}]", this.span, name);
}
this.context = ctx.put(Span.class, this.span);
}
@Override public void onSubscribe(Subscription subscription) {
if (log.isTraceEnabled()) {
log.trace("On subscribe");
}
this.s = subscription;
this.tracer.continueSpan(this.span);
if (log.isTraceEnabled()) {
log.trace("On subscribe - span continued");
}
this.subscriber.onSubscribe(this);
}
@Override public void request(long n) {
if (log.isTraceEnabled()) {
log.trace("Request");
}
this.tracer.continueSpan(this.span);
if (log.isTraceEnabled()) {
log.trace("Request - continued");
}
this.s.request(n);
// We're in the main thread so we don't want to pollute it with wrong spans
// that's why we need to detach the current one and continue with its parent
Span localRootSpan = this.span;
while (localRootSpan != null) {
if (this.rootSpan != null) {
if (localRootSpan.getSpanId() != this.rootSpan.getSpanId() &&
!isRootParentSpan(localRootSpan)) {
localRootSpan = continueDetachedSpan(localRootSpan);
} else {
localRootSpan = null;
}
} else if (!isRootParentSpan(localRootSpan)) {
localRootSpan = continueDetachedSpan(localRootSpan);
} else {
localRootSpan = null;
}
}
if (log.isTraceEnabled()) {
log.trace("Request after cleaning. Current span [{}]",
this.tracer.getCurrentSpan());
}
}
private boolean isRootParentSpan(Span localRootSpan) {
return localRootSpan.getSpanId() == localRootSpan.getTraceId();
}
private Span continueDetachedSpan(Span localRootSpan) {
if (log.isTraceEnabled()) {
log.trace("Will detach span {}", localRootSpan);
}
Span detachedSpan = this.tracer.detach(localRootSpan);
return this.tracer.continueSpan(detachedSpan);
}
@Override public void cancel() {
try {
if (log.isTraceEnabled()) {
log.trace("Cancel");
}
this.s.cancel();
}
finally {
cleanup();
}
}
@Override public void onNext(T o) {
this.subscriber.onNext(o);
}
@Override public void onError(Throwable throwable) {
try {
this.subscriber.onError(throwable);
}
finally {
cleanup();
}
}
@Override public void onComplete() {
try {
this.subscriber.onComplete();
}
finally {
cleanup();
}
}
void cleanup() {
if (compareAndSet(false, true)) {
if (log.isTraceEnabled()) {
log.trace("Cleaning up");
}
if (this.tracer.getCurrentSpan() != this.span) {
if (log.isTraceEnabled()) {
log.trace("Detaching span");
}
this.tracer.detach(this.tracer.getCurrentSpan());
this.tracer.continueSpan(this.span);
if (log.isTraceEnabled()) {
log.trace("Continuing span");
}
}
if (log.isTraceEnabled()) {
log.trace("Closing span");
}
this.tracer.close(this.span);
if (log.isTraceEnabled()) {
log.trace("Span closed");
}
if (this.rootSpan != null) {
this.tracer.continueSpan(this.rootSpan);
this.tracer.close(this.rootSpan);
if (log.isTraceEnabled()) {
log.trace("Closed root span");
}
}
}
}
@Override public Context currentContext() {
return this.context;
}
}

View File

@@ -0,0 +1,66 @@
package org.springframework.cloud.sleuth.instrument.reactor;
import java.util.concurrent.ScheduledExecutorService;
import java.util.function.Supplier;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import org.springframework.beans.factory.annotation.Autowired;
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.ConditionalOnNotWebApplication;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.cloud.sleuth.SpanNamer;
import org.springframework.cloud.sleuth.TraceKeys;
import org.springframework.cloud.sleuth.Tracer;
import org.springframework.cloud.sleuth.instrument.async.TraceableScheduledExecutorService;
import org.springframework.cloud.sleuth.instrument.web.TraceWebFluxAutoConfiguration;
import org.springframework.context.annotation.Configuration;
import reactor.core.publisher.Hooks;
import reactor.core.publisher.Mono;
import reactor.core.scheduler.Schedulers;
/**
* {@link org.springframework.boot.autoconfigure.EnableAutoConfiguration Auto-configuration}
* to enable tracing of Reactor components via Spring Cloud Sleuth.
*
* @author Stephane Maldini
* @author Marcin Grzejszczak
* @since 2.0.0
*/
@Configuration
@ConditionalOnProperty(value="spring.sleuth.reactor.enabled", matchIfMissing=true)
@ConditionalOnClass(Mono.class)
@AutoConfigureAfter(TraceWebFluxAutoConfiguration.class)
public class TraceReactorAutoConfiguration {
@Configuration
@ConditionalOnBean(Tracer.class)
@ConditionalOnNotWebApplication
static class TraceReactorConfiguration {
@Autowired Tracer tracer;
@Autowired TraceKeys traceKeys;
@Autowired SpanNamer spanNamer;
@PostConstruct
public void setupHooks() {
Hooks.onLastOperator(ReactorSleuth.spanOperator(this.tracer));
Schedulers.setFactory(new Schedulers.Factory() {
@Override public ScheduledExecutorService decorateExecutorService(String schedulerType,
Supplier<? extends ScheduledExecutorService> actual) {
return new TraceableScheduledExecutorService(actual.get(),
TraceReactorConfiguration.this.tracer,
TraceReactorConfiguration.this.traceKeys,
TraceReactorConfiguration.this.spanNamer);
}
});
}
@PreDestroy
public void cleanupHooks() {
Hooks.resetOnLastOperator();
Schedulers.resetFactory();
}
}
}

View File

@@ -0,0 +1,47 @@
/*
* Copyright 2013-2015 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.cloud.sleuth.instrument.web;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.boot.autoconfigure.AutoConfigureAfter;
import org.springframework.boot.autoconfigure.condition.ConditionalOnBean;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.boot.autoconfigure.condition.ConditionalOnWebApplication;
import org.springframework.cloud.sleuth.Tracer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
* {@link org.springframework.boot.autoconfigure.EnableAutoConfiguration
* Auto-configuration} enables tracing to HTTP requests with Spring WebFlux.
*
* @author Marcin Grzejszczak
* @since 2.0.0
*/
@Configuration
@ConditionalOnProperty(value = "spring.sleuth.web.enabled", matchIfMissing = true)
@ConditionalOnWebApplication(type = ConditionalOnWebApplication.Type.REACTIVE)
@ConditionalOnBean(Tracer.class)
@AutoConfigureAfter(TraceWebAutoConfiguration.class)
public class TraceWebFluxAutoConfiguration {
@Bean
public TraceWebFilter traceFilter(BeanFactory beanFactory,
SkipPatternProvider skipPatternProvider) {
return new TraceWebFilter(beanFactory, skipPatternProvider.skipPattern());
}
}

View File

@@ -9,10 +9,12 @@ org.springframework.cloud.sleuth.instrument.messaging.websocket.TraceWebSocketAu
org.springframework.cloud.sleuth.instrument.async.AsyncCustomAutoConfiguration,\
org.springframework.cloud.sleuth.instrument.async.AsyncDefaultAutoConfiguration,\
org.springframework.cloud.sleuth.instrument.hystrix.SleuthHystrixAutoConfiguration,\
org.springframework.cloud.sleuth.instrument.reactor.TraceReactorAutoConfiguration,\
org.springframework.cloud.sleuth.instrument.scheduling.TraceSchedulingAutoConfiguration,\
org.springframework.cloud.sleuth.instrument.web.TraceHttpAutoConfiguration,\
org.springframework.cloud.sleuth.instrument.web.TraceWebAutoConfiguration,\
org.springframework.cloud.sleuth.instrument.web.TraceWebServletAutoConfiguration,\
org.springframework.cloud.sleuth.instrument.web.TraceWebFluxAutoConfiguration,\
org.springframework.cloud.sleuth.instrument.web.client.TraceWebClientAutoConfiguration,\
org.springframework.cloud.sleuth.instrument.web.client.TraceWebAsyncClientAutoConfiguration,\
org.springframework.cloud.sleuth.instrument.web.client.feign.TraceFeignClientAutoConfiguration,\

View File

@@ -0,0 +1,199 @@
package org.springframework.cloud.sleuth.instrument.reactor;
import java.util.concurrent.atomic.AtomicReference;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.awaitility.Awaitility;
import org.junit.After;
import org.junit.AfterClass;
import org.junit.Before;
import org.junit.BeforeClass;
import org.junit.Ignore;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.reactivestreams.Publisher;
import org.reactivestreams.Subscription;
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.Sampler;
import org.springframework.cloud.sleuth.Span;
import org.springframework.cloud.sleuth.Tracer;
import org.springframework.cloud.sleuth.sampler.AlwaysSampler;
import org.springframework.cloud.sleuth.trace.TestSpanContextHolder;
import org.springframework.cloud.sleuth.util.ExceptionUtils;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.test.context.junit4.SpringRunner;
import reactor.core.publisher.BaseSubscriber;
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;
@RunWith(SpringRunner.class)
@SpringBootTest(classes = SpanSubscriberTests.Config.class,
webEnvironment = SpringBootTest.WebEnvironment.NONE)
public class SpanSubscriberTests {
private static final Log log = LogFactory.getLog(SpanSubscriberTests.class);
@Autowired Tracer tracer;
@Before
public void setup() {
ExceptionUtils.setFail(true);
}
@Test public void should_pass_tracing_info_when_using_reactor() {
Span span = this.tracer.createSpan("foo");
final AtomicReference<Span> spanInOperation = new AtomicReference<>();
Publisher<Integer> traced = Flux.just(1, 2, 3);
log.info("Hello");
Flux.from(traced)
.map( d -> d + 1)
.map( d -> d + 1)
.map( (d) -> {
spanInOperation.set(SpanSubscriberTests.this.tracer.getCurrentSpan());
return d + 1;
})
.map( d -> d + 1)
.subscribe(System.out::println);
then(this.tracer.getCurrentSpan()).isNull();
then(spanInOperation.get().getTraceId()).isEqualTo(span.getTraceId());
then(ExceptionUtils.getLastException()).isNull();
}
@Ignore("Ignored until fixed in Reactor")
@Test public void should_support_reactor_fusion_optimization() {
Span span = this.tracer.createSpan("foo");
final AtomicReference<Span> spanInOperation = new AtomicReference<>();
log.info("Hello");
Mono.just(1)
.flatMap( d -> Flux.just(d + 1).collectList().map(p -> p.get(0)))
.map( d -> d + 1)
.map( (d) -> {
spanInOperation.set(SpanSubscriberTests.this.tracer.getCurrentSpan());
return d + 1;
})
.map( d -> d + 1)
.subscribe(System.out::println);
then(this.tracer.getCurrentSpan()).isNull();
then(spanInOperation.get().getTraceId()).isEqualTo(span.getTraceId());
then(ExceptionUtils.getLastException()).isNull();
}
@Test public void should_not_trace_scalar_flows() {
this.tracer.createSpan("foo");
final AtomicReference<Subscription> spanInOperation = new AtomicReference<>();
log.info("Hello");
Mono.just(1)
.subscribe(new BaseSubscriber<Integer>() {
@Override
protected void hookOnSubscribe(Subscription subscription) {
spanInOperation.set(subscription);
}
});
then(this.tracer.getCurrentSpan()).isNotNull();
then(spanInOperation.get()).isNotInstanceOf(SpanSubscriber.class);
Mono.<Integer>error(new Exception())
.subscribe(new BaseSubscriber<Integer>() {
@Override
protected void hookOnSubscribe(Subscription subscription) {
spanInOperation.set(subscription);
}
@Override
protected void hookOnError(Throwable throwable) {
}
});
then(this.tracer.getCurrentSpan()).isNotNull();
then(spanInOperation.get()).isNotInstanceOf(SpanSubscriber.class);
Mono.<Integer>empty()
.subscribe(new BaseSubscriber<Integer>() {
@Override
protected void hookOnSubscribe(Subscription subscription) {
spanInOperation.set(subscription);
}
});
then(this.tracer.getCurrentSpan()).isNotNull();
then(spanInOperation.get()).isNotInstanceOf(SpanSubscriber.class);
then(ExceptionUtils.getLastException()).isNull();
}
@Test
public void should_pass_tracing_info_when_using_reactor_async() {
Span span = this.tracer.createSpan("foo");
final AtomicReference<Span> spanInOperation = new AtomicReference<>();
log.info("Hello");
Flux.just(1, 2, 3)
.publishOn(Schedulers.single())
.log("reactor.1")
.map( d -> d + 1)
.map( d -> d + 1)
.publishOn(Schedulers.newSingle("secondThread"))
.log("reactor.2")
.map( (d) -> {
spanInOperation.set(SpanSubscriberTests.this.tracer.getCurrentSpan());
return d + 1;
})
.map( d -> d + 1)
.blockLast();
Awaitility.await().untilAsserted(() -> {
then(spanInOperation.get().getTraceId()).isEqualTo(span.getTraceId());
then(ExceptionUtils.getLastException()).isNull();
});
then(this.tracer.getCurrentSpan()).isEqualTo(span);
this.tracer.close(span);
Span foo2 = this.tracer.createSpan("foo2");
Flux.just(1, 2, 3)
.publishOn(Schedulers.single())
.log("reactor.")
.map( d -> d + 1)
.map( d -> d + 1)
.map( (d) -> {
spanInOperation.set(SpanSubscriberTests.this.tracer.getCurrentSpan());
return d + 1;
})
.map( d -> d + 1)
.blockLast();
then(this.tracer.getCurrentSpan()).isEqualTo(foo2);
then(ExceptionUtils.getLastException()).isNull();
// parent cause there's an async span in the meantime
then(spanInOperation.get().getTraceId()).isEqualTo(foo2.getTraceId());
tracer.close(foo2);
}
@AfterClass
public static void cleanup() {
Hooks.resetOnLastOperator();
Schedulers.resetFactory();
}
@EnableAutoConfiguration
@Configuration
static class Config {
@Bean Sampler sampler() {
return new AlwaysSampler();
}
}
}

View File

@@ -0,0 +1,105 @@
package org.springframework.cloud.sleuth.instrument.web;
import org.awaitility.Awaitility;
import org.junit.BeforeClass;
import org.junit.Ignore;
import org.junit.Test;
import org.springframework.boot.WebApplicationType;
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.boot.builder.SpringApplicationBuilder;
import org.springframework.cloud.sleuth.Sampler;
import org.springframework.cloud.sleuth.Span;
import org.springframework.cloud.sleuth.SpanReporter;
import org.springframework.cloud.sleuth.Tracer;
import org.springframework.cloud.sleuth.assertions.ListOfSpans;
import org.springframework.cloud.sleuth.assertions.SleuthAssertions;
import org.springframework.cloud.sleuth.sampler.AlwaysSampler;
import org.springframework.cloud.sleuth.util.ArrayListSpanAccumulator;
import org.springframework.cloud.sleuth.util.ExceptionUtils;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.env.Environment;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.reactive.function.client.ClientResponse;
import org.springframework.web.reactive.function.client.WebClient;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Hooks;
import reactor.core.publisher.Mono;
import reactor.core.scheduler.Schedulers;
public class TraceWebFluxTests {
@BeforeClass
public static void setup() {
Hooks.resetOnLastOperator();
Schedulers.resetFactory();
}
@Ignore("Ignored until fixed in Reactor")
@Test public void should_instrument_web_filter() throws Exception {
ConfigurableApplicationContext context = new SpringApplicationBuilder(TraceWebFluxTests.Config.class)
.web(WebApplicationType.REACTIVE).properties("server.port=0", "spring.jmx.enabled=false",
"spring.application.name=TraceWebFluxTests").run();
ExceptionUtils.setFail(true);
Span span = null;
try {
span = context.getBean(Tracer.class).createSpan("foo");
int port = context.getBean(Environment.class).getProperty("local.server.port", Integer.class);
ArrayListSpanAccumulator accumulator = context.getBean(ArrayListSpanAccumulator.class);
Mono<ClientResponse> exchange = context.getBean(WebClient.class).get().uri("http://localhost:" + port + "/api/c2/10").exchange();
Awaitility.await().untilAsserted(() -> {
ClientResponse response = exchange.block();
SleuthAssertions.then(response.statusCode().value()).isEqualTo(200);
SleuthAssertions.then(ExceptionUtils.getLastException()).isNull();
SleuthAssertions.then(new ListOfSpans(accumulator.getSpans()))
.hasASpanWithLogEqualTo(Span.CLIENT_SEND)
.hasASpanWithLogEqualTo(Span.SERVER_RECV)
.hasASpanWithLogEqualTo(Span.SERVER_SEND)
.hasASpanWithLogEqualTo(Span.CLIENT_RECV)
.hasASpanWithTagEqualTo("mvc.controller.method", "successful")
.hasASpanWithTagEqualTo("mvc.controller.class", "Controller2");
});
} finally {
context.getBean(Tracer.class).close(span);
}
}
@Configuration
@EnableAutoConfiguration
static class Config {
@Bean WebClient webClient() {
return WebClient.create();
}
@Bean Sampler sampler() {
return new AlwaysSampler();
}
@Bean SpanReporter spanReporter() {
return new ArrayListSpanAccumulator();
}
@Bean
Controller2 controller2() {
return new Controller2();
}
}
@RestController
@RequestMapping("/api/c2")
static class Controller2 {
@GetMapping("/{id}")
public Flux<String> successful(@PathVariable Long id) {
return Flux.just(id.toString());
}
}
}