diff --git a/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/StreamEnvironmentPostProcessor.java b/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/StreamEnvironmentPostProcessor.java index a7cca416e..f93595695 100644 --- a/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/StreamEnvironmentPostProcessor.java +++ b/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/StreamEnvironmentPostProcessor.java @@ -46,7 +46,7 @@ import org.springframework.core.io.support.PropertiesLoaderUtils; public class StreamEnvironmentPostProcessor implements EnvironmentPostProcessor { private static final String PROPERTY_SOURCE_NAME = "defaultProperties"; - private static String[] headers = new String[] { Span.SPAN_ID_NAME, + static String[] headers = new String[] { Span.SPAN_ID_NAME, Span.TRACE_ID_NAME, Span.PARENT_ID_NAME, Span.PROCESS_ID_NAME, Span.SAMPLED_NAME, Span.SPAN_NAME_NAME }; @@ -64,7 +64,7 @@ public class StreamEnvironmentPostProcessor implements EnvironmentPostProcessor .getResources("classpath*:META-INF/spring.binders")) { for (String binderType : parseBinderConfigurations(resource)) { int startIndex = findStartIndex(environment, binderType); - addHeaders(map, binderType, startIndex); + addHeaders(map, environment.getPropertySources(), binderType, startIndex); } } } @@ -125,11 +125,23 @@ public class StreamEnvironmentPostProcessor implements EnvironmentPostProcessor } } - private void addHeaders(Map map, String binder, int startIndex) { + private void addHeaders(Map map, MutablePropertySources propertySources, + String binder, int startIndex) { String stem = "spring.cloud.stream." + binder + ".binder.headers"; for (int i = 0; i < headers.length; i++) { - map.put(stem + "[" + (i + startIndex) + "]", headers[i]); + if (!hasTracingHeadersValue(propertySources, headers[i])) { + map.put(stem + "[" + (i + startIndex) + "]", headers[i]); + } } } + private boolean hasTracingHeadersValue(MutablePropertySources propertySources, String header) { + PropertySource source = propertySources.get(PROPERTY_SOURCE_NAME); + if (source instanceof MapPropertySource) { + Collection values = ((MapPropertySource) source).getSource().values(); + return values.contains(header); + } + return false; + } + } diff --git a/spring-cloud-sleuth-stream/src/test/java/org/springframework/cloud/sleuth/stream/StreamEnvironmentPostProcessorTests.java b/spring-cloud-sleuth-stream/src/test/java/org/springframework/cloud/sleuth/stream/StreamEnvironmentPostProcessorTests.java index d7e9031be..c90c78591 100644 --- a/spring-cloud-sleuth-stream/src/test/java/org/springframework/cloud/sleuth/stream/StreamEnvironmentPostProcessorTests.java +++ b/spring-cloud-sleuth-stream/src/test/java/org/springframework/cloud/sleuth/stream/StreamEnvironmentPostProcessorTests.java @@ -16,7 +16,9 @@ package org.springframework.cloud.sleuth.stream; -import static org.assertj.core.api.Assertions.assertThat; +import java.util.Collection; +import java.util.Map; +import java.util.stream.Collectors; import org.junit.Test; import org.springframework.boot.SpringApplication; @@ -25,6 +27,8 @@ import org.springframework.cloud.sleuth.Span; import org.springframework.core.env.ConfigurableEnvironment; import org.springframework.core.env.StandardEnvironment; +import static org.assertj.core.api.Assertions.assertThat; + /** * @author Dave Syer * @@ -35,21 +39,44 @@ public class StreamEnvironmentPostProcessorTests { private ConfigurableEnvironment environment = new StandardEnvironment(); @Test - public void testHeadersAdded() { - this.processor.postProcessEnvironment(this.environment, - new SpringApplication(StreamEnvironmentPostProcessorTests.class)); + public void should_append_tracing_headers() { + postProcess(); assertThat(this.environment.getProperty("spring.cloud.stream.test.binder.headers[0]")) .isEqualTo(Span.SPAN_ID_NAME); } @Test - public void headersAppendedIfNecessary() { + public void should_append_tracing_headers_to_existing_ones() { EnvironmentTestUtils.addEnvironment(this.environment, "spring.cloud.stream.test.binder.headers[0]=X-Custom", "spring.cloud.stream.test.binder.headers[1]=X-Mine"); - this.processor.postProcessEnvironment(this.environment, - new SpringApplication(StreamEnvironmentPostProcessorTests.class)); + postProcess(); assertThat(this.environment.getProperty("spring.cloud.stream.test.binder.headers[2]")) .isEqualTo(Span.SPAN_ID_NAME); } + @Test + public void should_not_append_tracing_headers_if_they_are_already_appended() { + postProcess(); + postProcess(); + postProcess(); + + Collection headerValues = defaultPropertiesSource().values(); + + Collection traceIds = headerValues.stream() + .filter(input -> input.contains(Span.TRACE_ID_NAME)).collect(Collectors.toList()); + assertThat(traceIds).hasSize(1); + assertThat(defaultPropertiesSource().keySet().stream() + .filter(input -> input.startsWith("spring.cloud.stream.test.binder.headers")) + .collect(Collectors.toList())).hasSize(StreamEnvironmentPostProcessor.headers.length); + } + + private void postProcess() { + this.processor.postProcessEnvironment(this.environment, + new SpringApplication(StreamEnvironmentPostProcessorTests.class)); + } + + private Map defaultPropertiesSource() { + return (Map) this.environment.getPropertySources().get("defaultProperties").getSource(); + } + }