diff --git a/spring-cloud-sleuth-core/pom.xml b/spring-cloud-sleuth-core/pom.xml
index edcea1b3f..bcbab14b3 100644
--- a/spring-cloud-sleuth-core/pom.xml
+++ b/spring-cloud-sleuth-core/pom.xml
@@ -42,6 +42,21 @@
spring-boot-starter-actuator
true
+
+ com.netflix.hystrix
+ hystrix-core
+ true
+
+
+ io.reactivex
+ rxjava
+ true
+
+
+ org.aspectj
+ aspectjrt
+ true
+
org.projectlombok
lombok
diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/slf4j/Slf4jSpanReceiver.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/slf4j/Slf4jSpanReceiver.java
index 064b94c4f..3d1acd172 100644
--- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/slf4j/Slf4jSpanReceiver.java
+++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/slf4j/Slf4jSpanReceiver.java
@@ -1,6 +1,6 @@
package org.springframework.cloud.sleuth.slf4j;
-import static org.springframework.cloud.sleuth.slf4j.Slf4jSpanStartListener.SPAN_ID_NAME;
+import static org.springframework.cloud.sleuth.trace.Trace.SPAN_ID_NAME;
import lombok.extern.slf4j.Slf4j;
import org.slf4j.MDC;
diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/slf4j/Slf4jSpanStartListener.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/slf4j/Slf4jSpanStartListener.java
index e68a1df72..c1ba9a5d0 100644
--- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/slf4j/Slf4jSpanStartListener.java
+++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/slf4j/Slf4jSpanStartListener.java
@@ -1,22 +1,22 @@
package org.springframework.cloud.sleuth.slf4j;
import lombok.extern.slf4j.Slf4j;
+
import org.slf4j.MDC;
import org.springframework.cloud.sleuth.trace.Span;
import org.springframework.cloud.sleuth.trace.SpanStartListener;
+import org.springframework.cloud.sleuth.trace.Trace;
/**
* @author Spencer Gibb
*/
@Slf4j
public class Slf4jSpanStartListener implements SpanStartListener {
- //TODO: Where to put span id name?
- public static final String SPAN_ID_NAME = "Span-Id";
@Override
public void startSpan(Span span) {
//TODO: what log level?
log.info("Starting span with id: [{}]", span.getSpanId());
- MDC.put(SPAN_ID_NAME, span.getSpanId());
+ MDC.put(Trace.SPAN_ID_NAME, span.getSpanId());
}
}
diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/Trace.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/Trace.java
index 00dbc0e6f..0568118a7 100644
--- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/Trace.java
+++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/Trace.java
@@ -35,6 +35,9 @@ package org.springframework.cloud.sleuth.trace;
*/
public interface Trace {
+ String SPAN_ID_NAME = "Span-Id";
+ String TRACE_ID_NAME = "Trace-Id";
+
/**
* Creates a new trace scope.
*
diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/intercept/circuitbreaker/TraceCommand.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/intercept/circuitbreaker/TraceCommand.java
new file mode 100644
index 000000000..cb8ea6861
--- /dev/null
+++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/intercept/circuitbreaker/TraceCommand.java
@@ -0,0 +1,59 @@
+package org.springframework.cloud.sleuth.trace.intercept.circuitbreaker;
+
+import com.netflix.hystrix.HystrixCommand;
+import com.netflix.hystrix.HystrixCommandGroupKey;
+import com.netflix.hystrix.HystrixThreadPoolKey;
+import org.springframework.cloud.sleuth.trace.Trace;
+import org.springframework.cloud.sleuth.trace.TraceScope;
+
+/**
+ * Abstraction over {@code HystrixCommand} that wraps command execution with Trace setting
+ *
+ * @see HystrixCommand
+ * @see CorrelationIdUpdater
+ *
+ * @author Tomasz Nurkiewicz, 4financeIT
+ * @author Marcin Grzejszczak, 4financeIT
+ * @author Spencer Gibb
+ */
+public abstract class TraceCommand extends HystrixCommand {
+
+ private Trace trace;
+
+ protected TraceCommand(Trace trace, HystrixCommandGroupKey group) {
+ super(group);
+ this.trace = trace;
+ }
+
+ protected TraceCommand(Trace trace, HystrixCommandGroupKey group, HystrixThreadPoolKey threadPool) {
+ super(group, threadPool);
+ this.trace = trace;
+ }
+
+ protected TraceCommand(Trace trace, HystrixCommandGroupKey group, int executionIsolationThreadTimeoutInMilliseconds) {
+ super(group, executionIsolationThreadTimeoutInMilliseconds);
+ this.trace = trace;
+ }
+
+ protected TraceCommand(Trace trace, HystrixCommandGroupKey group, HystrixThreadPoolKey threadPool, int executionIsolationThreadTimeoutInMilliseconds) {
+ super(group, threadPool, executionIsolationThreadTimeoutInMilliseconds);
+ this.trace = trace;
+ }
+
+ protected TraceCommand(Trace trace, Setter setter) {
+ super(setter);
+ this.trace = trace;
+ }
+
+ @Override
+ protected R run() throws Exception {
+ TraceScope scope = trace.startSpan(getCommandKey().name());
+ try {
+ return doRun();
+ } finally {
+ scope.close();
+ }
+ }
+
+ public abstract R doRun() throws Exception;
+}
diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/intercept/scheduling/TraceSchedulingAspect.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/intercept/scheduling/TraceSchedulingAspect.java
new file mode 100644
index 000000000..41d78ff25
--- /dev/null
+++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/intercept/scheduling/TraceSchedulingAspect.java
@@ -0,0 +1,39 @@
+package org.springframework.cloud.sleuth.trace.intercept.scheduling;
+
+import org.aspectj.lang.ProceedingJoinPoint;
+import org.aspectj.lang.annotation.Around;
+import org.aspectj.lang.annotation.Aspect;
+import org.springframework.cloud.sleuth.trace.Trace;
+import org.springframework.cloud.sleuth.trace.TraceScope;
+import org.springframework.scheduling.annotation.Scheduled;
+
+/**
+ * Aspect that sets correlationId for running threads executing methods annotated with {@link Scheduled} annotation.
+ * For every execution of scheduled method a new, i.e. unique one, value of correlationId will be set.
+ *
+ * @author Tomasz Nurkewicz, 4financeIT
+ * @author Michal Chmielarz, 4financeIT
+ * @author Marcin Grzejszczak, 4financeIT
+ * @author Spencer Gibb
+ *
+ * @see org.springframework.cloud.sleuth.trace.Trace
+ */
+@Aspect
+public class TraceSchedulingAspect {
+
+ private final Trace trace;
+
+ public TraceSchedulingAspect(Trace trace) {
+ this.trace = trace;
+ }
+
+ @Around("execution (@org.springframework.scheduling.annotation.Scheduled * *.*(..))")
+ public Object setNewCorrelationIdOnThread(final ProceedingJoinPoint pjp) throws Throwable {
+ TraceScope scope = trace.startSpan(pjp.toShortString());
+ try {
+ return pjp.proceed();
+ } finally {
+ scope.close();
+ }
+ }
+}
diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/intercept/scheduling/TraceSchedulingAutoConfiguration.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/intercept/scheduling/TraceSchedulingAutoConfiguration.java
new file mode 100644
index 000000000..590264775
--- /dev/null
+++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/intercept/scheduling/TraceSchedulingAutoConfiguration.java
@@ -0,0 +1,30 @@
+package org.springframework.cloud.sleuth.trace.intercept.scheduling;
+
+/**
+ * @author Spencer Gibb
+ */
+
+import org.springframework.cloud.sleuth.trace.Trace;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.context.annotation.EnableAspectJAutoProxy;
+import org.springframework.scheduling.annotation.EnableScheduling;
+
+/**
+ * Registers beans related to task scheduling.
+ *
+ * @see TraceSchedulingAspect
+ *
+ * @author Michal Chmielarz, 4financeIT
+ * @author Spencer Gibb
+ */
+@Configuration
+@EnableScheduling
+@EnableAspectJAutoProxy
+public class TraceSchedulingAutoConfiguration {
+
+ @Bean
+ public TraceSchedulingAspect traceSchedulingAspect(Trace trace) {
+ return new TraceSchedulingAspect(trace);
+ }
+}
diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/intercept/web/TraceFilter.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/intercept/web/TraceFilter.java
new file mode 100644
index 000000000..7b8972e2b
--- /dev/null
+++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/intercept/web/TraceFilter.java
@@ -0,0 +1,122 @@
+/*
+ * Copyright 2012-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.trace.intercept.web;
+
+import static org.springframework.cloud.sleuth.trace.Trace.SPAN_ID_NAME;
+import static org.springframework.util.StringUtils.hasText;
+
+import java.io.IOException;
+import java.lang.invoke.MethodHandles;
+import java.util.Collections;
+import java.util.regex.Pattern;
+
+import javax.servlet.FilterChain;
+import javax.servlet.ServletException;
+import javax.servlet.http.HttpServletRequest;
+import javax.servlet.http.HttpServletResponse;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.slf4j.MDC;
+import org.springframework.cloud.sleuth.trace.MilliSpan;
+import org.springframework.cloud.sleuth.trace.Span;
+import org.springframework.cloud.sleuth.trace.Trace;
+import org.springframework.cloud.sleuth.trace.TraceScope;
+import org.springframework.web.filter.OncePerRequestFilter;
+
+/**
+ * Filter that takes the value of the {@link CorrelationIdHolder#CORRELATION_ID_HEADER}
+ * header from either request or response and sets it in the {@link CorrelationIdHolder}.
+ * It also provides that value in {@link MDC} logging related class so that logger prints
+ * the value of correlation id at each log.
+ *
+ * @see Trace
+ * @see MDC
+ *
+ * @author Jakub Nabrdalik, 4financeIT
+ * @author Tomasz Nurkiewicz, 4financeIT
+ * @author Marcin Grzejszczak, 4financeIT
+ * @author Spencer Gibb
+ */
+public class TraceFilter extends OncePerRequestFilter {
+ private static final Logger log = LoggerFactory.getLogger(MethodHandles.lookup()
+ .lookupClass());
+ public static final Pattern DEFAULT_SKIP_PATTERN = Pattern
+ .compile("/api-docs.*|/autoconfig|/configprops|/dump|/info|/metrics.*|/mappings|/trace|/swagger.*|.*\\.png|.*\\.css|.*\\.js|.*\\.html");
+
+ private final Trace trace;
+ private final Pattern skipPattern;
+
+ public TraceFilter(Trace trace) {
+ this.trace = trace;
+ this.skipPattern = DEFAULT_SKIP_PATTERN;
+ }
+
+ public TraceFilter(Trace trace, Pattern skipPattern) {
+ this.trace = trace;
+ this.skipPattern = skipPattern;
+ }
+
+ @Override
+ protected void doFilterInternal(HttpServletRequest request,
+ HttpServletResponse response, FilterChain filterChain)
+ throws ServletException, IOException {
+ String spanIdFromRequest = getSpanIdFrom(request);
+ String spanId = (hasText(spanIdFromRequest)) ? spanIdFromRequest
+ : getSpanIdFrom(response);
+
+ TraceScope traceScope = null;
+ if (spanId != null) {
+ addCorrelationIdToResponseIfNotPresent(response, spanId);
+
+ Span span = MilliSpan.builder().traceId("") // FIXME get traceId from request
+ .parents(Collections.singletonList(spanId))
+ // TODO: use parent() when lombok plugin supports it
+ .build();
+ traceScope = trace.startSpan("traceFilter", span);
+ }
+ else {
+ traceScope = trace.startSpan("traceFilter");
+ }
+
+ try {
+ filterChain.doFilter(request, response);
+ }
+ finally {
+ traceScope.close();
+ }
+ }
+
+ private String getSpanIdFrom(final HttpServletResponse response) {
+ return response.getHeader(SPAN_ID_NAME);
+ }
+
+ private String getSpanIdFrom(final HttpServletRequest request) {
+ return request.getHeader(SPAN_ID_NAME);
+ }
+
+ private void addCorrelationIdToResponseIfNotPresent(HttpServletResponse response,
+ String spanId) {
+ if (!hasText(response.getHeader(SPAN_ID_NAME))) {
+ response.addHeader(SPAN_ID_NAME, spanId);
+ }
+ }
+
+ @Override
+ protected boolean shouldNotFilterAsyncDispatch() {
+ return false;
+ }
+}
diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/intercept/web/TraceWebAspect.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/intercept/web/TraceWebAspect.java
new file mode 100644
index 000000000..adbd8620f
--- /dev/null
+++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/trace/intercept/web/TraceWebAspect.java
@@ -0,0 +1,125 @@
+package org.springframework.cloud.sleuth.trace.intercept.web;
+
+import static org.springframework.cloud.sleuth.trace.Trace.SPAN_ID_NAME;
+
+import java.util.ArrayList;
+import java.util.List;
+
+import lombok.extern.apachecommons.CommonsLog;
+
+import org.aspectj.lang.ProceedingJoinPoint;
+import org.aspectj.lang.annotation.Around;
+import org.aspectj.lang.annotation.Pointcut;
+import org.springframework.cloud.sleuth.trace.Trace;
+import org.springframework.cloud.sleuth.trace.TraceScope;
+import org.springframework.http.HttpEntity;
+import org.springframework.http.HttpHeaders;
+import org.springframework.stereotype.Controller;
+import org.springframework.web.bind.annotation.RestController;
+import org.springframework.web.client.RestOperations;
+
+/**
+ * Aspect that adds correlation id to
+ *
+ *
+ * - {@link RestController} annotated classes
+ * - {@link Controller} annotated classes
+ * - explicit {@link RestOperations}.exchange(..) method calls
+ *
+ *
+ * For controllers an around aspect is created that create a {@link Span} for each call.
+ *
+ * For {@link RestOperations} we are wrapping all executions of the exchange
+ * methods and we are extracting {@link HttpHeaders} from the passed {@link HttpEntity}.
+ * Next we are adding span id header * {@link Trace#SPAN_ID_NAME} with the value taken
+ * from the current Span. Finally the method execution proceeds.
+ *
+ * @see RestController
+ * @see Controller
+ * @see RestOperations
+ * @see Trace
+ *
+ * @author Tomasz Nurkewicz, 4financeIT
+ * @author Marcin Grzejszczak, 4financeIT
+ * @author Michal Chmielarz, 4financeIT
+ * @author Spencer Gibb
+ */
+@CommonsLog
+public class TraceWebAspect {
+
+ private final Trace trace;
+
+ public TraceWebAspect(Trace trace) {
+ this.trace = trace;
+ }
+
+ private static final int HTTP_ENTITY_PARAM_INDEX = 2;
+
+ @Pointcut("@target(org.springframework.web.bind.annotation.RestController)")
+ private void anyRestControllerAnnotated() {
+ }
+
+ @Pointcut("@target(org.springframework.stereotype.Controller)")
+ private void anyControllerAnnotated() {
+ }
+
+ @Pointcut("anyRestControllerAnnotated() || anyControllerAnnotated()")
+ private void anyControllerOrRestController() {
+ }
+
+ @Around("anyControllerOrRestController()")
+ public Object wrapWithCorrelationId(ProceedingJoinPoint pjp) throws Throwable {
+ TraceScope scope = trace.startSpan(pjp.toShortString());
+ try {
+ return pjp.proceed();
+ }
+ finally {
+ scope.close();
+ }
+ }
+
+ @Pointcut("execution(public * org.springframework.web.client.RestOperations.exchange(..))")
+ private void anyExchangeRestOperationsMethod() {
+ }
+
+ @Around("anyExchangeRestOperationsMethod()")
+ public Object wrapWithCorrelationIdForRestOperations(ProceedingJoinPoint pjp)
+ throws Throwable {
+ TraceScope scope = trace.startSpan(pjp.toShortString());
+ try {
+ String spanId = scope.getSpan().getSpanId();
+ //TODO: set traceId on restTemplate call as well
+ log.debug("Wrapping RestTemplate call with span id [" + spanId + "]");
+ HttpEntity httpEntity = (HttpEntity) pjp.getArgs()[HTTP_ENTITY_PARAM_INDEX];
+ HttpEntity newHttpEntity = createNewHttpEntity(httpEntity, spanId);
+ List