From 754c8ce96d9151102f62387ce96f6712b693a882 Mon Sep 17 00:00:00 2001 From: Marcin Grzejszczak Date: Tue, 1 Mar 2016 16:07:35 +0100 Subject: [PATCH] [#191] Added ServerAddress annotation * if a span contains a CS or CR (or any other client related annotation) then SA is passed from the Endpoint service name fixes #191 --- .../stream/SamplingZipkinSpanIterator.java | 47 ++++++++++++++--- .../SamplingZipkinSpanIteratorTests.java | 51 +++++++++++++++++++ .../sleuth/zipkin/ZipkinSpanListener.java | 47 ++++++++++++++--- .../zipkin/ZipkinSpanListenerTests.java | 45 +++++++++++++++- 4 files changed, 176 insertions(+), 14 deletions(-) diff --git a/spring-cloud-sleuth-zipkin-stream/src/main/java/org/springframework/cloud/sleuth/zipkin/stream/SamplingZipkinSpanIterator.java b/spring-cloud-sleuth-zipkin-stream/src/main/java/org/springframework/cloud/sleuth/zipkin/stream/SamplingZipkinSpanIterator.java index c128aec54..409c947b7 100644 --- a/spring-cloud-sleuth-zipkin-stream/src/main/java/org/springframework/cloud/sleuth/zipkin/stream/SamplingZipkinSpanIterator.java +++ b/spring-cloud-sleuth-zipkin-stream/src/main/java/org/springframework/cloud/sleuth/zipkin/stream/SamplingZipkinSpanIterator.java @@ -15,6 +15,7 @@ */ package org.springframework.cloud.sleuth.zipkin.stream; +import java.util.ArrayList; import java.util.Arrays; import java.util.Iterator; import java.util.List; @@ -40,14 +41,23 @@ import zipkin.Span.Builder; * @since 1.0.0 */ final class SamplingZipkinSpanIterator implements Iterator { - private static final List ZIPKIN_ANNOTATIONS = Arrays.asList( + private static final List ZIPKIN_CLIENT_ANNOTATIONS = Arrays.asList( Constants.CLIENT_ADDR, Constants.CLIENT_RECV, Constants.CLIENT_SEND, - Constants.CLIENT_RECV_FRAGMENT, Constants.CLIENT_SEND_FRAGMENT, - Constants.SERVER_ADDR, Constants.SERVER_RECV, Constants.SERVER_SEND, - Constants.SERVER_RECV_FRAGMENT, Constants.SERVER_SEND_FRAGMENT, - Constants.LOCAL_COMPONENT, - Constants.WIRE_RECV, Constants.WIRE_SEND + Constants.CLIENT_RECV_FRAGMENT, Constants.CLIENT_SEND_FRAGMENT ); + private static final List ZIPKIN_ANNOTATIONS = zipkinAnnotations(); + + private static List zipkinAnnotations() { + List annotations = new ArrayList<>(); + annotations.addAll(Arrays.asList( + Constants.SERVER_ADDR, Constants.SERVER_RECV, Constants.SERVER_SEND, + Constants.SERVER_RECV_FRAGMENT, Constants.SERVER_SEND_FRAGMENT, + Constants.LOCAL_COMPONENT, + Constants.WIRE_RECV, Constants.WIRE_SEND + )); + annotations.addAll(ZIPKIN_CLIENT_ANNOTATIONS); + return annotations; + } private static final Log log = org.apache.commons.logging.LogFactory .getLog(SamplingZipkinSpanIterator.class); @@ -128,6 +138,8 @@ final class SamplingZipkinSpanIterator implements Iterator { } if (!spanContainsAnyZipkinConstant(span)) { addLocalComponentAnnotation(span, zipkinSpan, ep); + } else if (spanContainsAnyZipkinClientConstantAndSaIsNotSet(span)) { + addServerAddressAnnotation(zipkinSpan, ep); } zipkinSpan.timestamp(span.getBegin() * 1000); zipkinSpan.duration(span.getAccumulatedMillis() * 1000); @@ -156,6 +168,16 @@ final class SamplingZipkinSpanIterator implements Iterator { BinaryAnnotation.create(Constants.LOCAL_COMPONENT, processId, ep)); } + private static void addServerAddressAnnotation(zipkin.Span.Builder zipkinSpan, + Endpoint ep) { + BinaryAnnotation component = new BinaryAnnotation.Builder() + .type(BinaryAnnotation.Type.STRING) + .key(Constants.SERVER_ADDR) + .value(ep.serviceName) + .endpoint(ep).build(); + zipkinSpan.addBinaryAnnotation(component); + } + private static boolean spanContainsAnyZipkinConstant(Span span) { for (org.springframework.cloud.sleuth.Log log : span.logs()) { if (ZIPKIN_ANNOTATIONS.contains(log.getEvent())) { @@ -164,4 +186,17 @@ final class SamplingZipkinSpanIterator implements Iterator { } return false; } + + private static boolean spanContainsAnyZipkinClientConstantAndSaIsNotSet(Span span) { + boolean containsAnyZipkinClientConstant = false; + for (org.springframework.cloud.sleuth.Log log : span.logs()) { + if (ZIPKIN_CLIENT_ANNOTATIONS.contains(log.getEvent())) { + containsAnyZipkinClientConstant = true; + break; + } + } + return containsAnyZipkinClientConstant && !span.tags().containsKey(Constants.SERVER_ADDR); + } + + } \ No newline at end of file diff --git a/spring-cloud-sleuth-zipkin-stream/src/test/java/org/springframework/cloud/sleuth/zipkin/stream/SamplingZipkinSpanIteratorTests.java b/spring-cloud-sleuth-zipkin-stream/src/test/java/org/springframework/cloud/sleuth/zipkin/stream/SamplingZipkinSpanIteratorTests.java index 624c01e9a..766dd3dcd 100644 --- a/spring-cloud-sleuth-zipkin-stream/src/test/java/org/springframework/cloud/sleuth/zipkin/stream/SamplingZipkinSpanIteratorTests.java +++ b/spring-cloud-sleuth-zipkin-stream/src/test/java/org/springframework/cloud/sleuth/zipkin/stream/SamplingZipkinSpanIteratorTests.java @@ -15,6 +15,7 @@ */ package org.springframework.cloud.sleuth.zipkin.stream; +import java.nio.charset.Charset; import java.util.Arrays; import java.util.Collections; import java.util.Iterator; @@ -25,6 +26,8 @@ import org.junit.Test; import org.springframework.cloud.sleuth.Span; import org.springframework.cloud.sleuth.stream.Host; import org.springframework.cloud.sleuth.stream.Spans; + +import zipkin.Constants; import zipkin.Sampler; import static org.assertj.core.api.Assertions.assertThat; @@ -76,6 +79,54 @@ public class SamplingZipkinSpanIteratorTests { "message:foo", "message:baz"); } + @Test + public void appendsLocalComponentTagIfNoZipkinLogIsPresent() { + Spans spans = new Spans(this.host, Collections.singletonList(span("foo"))); + + Iterator result = new SamplingZipkinSpanIterator( + Sampler.create(1.0f), spans); + + assertThat(result) + .flatExtracting(s -> s.binaryAnnotations) + .extracting(input -> input.key) + .contains(Constants.LOCAL_COMPONENT); + } + + @Test + public void appendsServerAddressTagIfClientLogIsPresent() { + Span span = span("foo"); + span.logEvent(Constants.CLIENT_RECV); + Spans spans = new Spans(this.host, Collections.singletonList(span)); + + Iterator result = new SamplingZipkinSpanIterator( + Sampler.create(1.0f), spans); + + assertThat(result) + .hasSize(1) + .flatExtracting(input1 -> input1.binaryAnnotations) + .filteredOn("key", Constants.SERVER_ADDR) + .extracting(input -> input.value) + .contains("myservice".getBytes(Charset.forName("UTF-8"))); + } + + @Test + public void shouldReuseServerAddressTag() { + Span span = span("foo"); + span.logEvent(Constants.CLIENT_RECV); + span.tag(Constants.SERVER_ADDR, "barservice"); + Spans spans = new Spans(this.host, Collections.singletonList(span)); + + Iterator result = new SamplingZipkinSpanIterator( + Sampler.create(1.0f), spans); + + assertThat(result) + .hasSize(1) + .flatExtracting(input1 -> input1.binaryAnnotations) + .filteredOn("key", Constants.SERVER_ADDR) + .extracting(input -> input.value) + .contains("barservice".getBytes(Charset.forName("UTF-8"))); + } + Span span(String name) { Long id = new Random().nextLong(); return new Span(1, 3, "message:" + name, id, Collections.emptyList(), id, true, true, diff --git a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/ZipkinSpanListener.java b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/ZipkinSpanListener.java index 52c6740fa..b7ce058fe 100644 --- a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/ZipkinSpanListener.java +++ b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/ZipkinSpanListener.java @@ -17,6 +17,7 @@ package org.springframework.cloud.sleuth.zipkin; import java.nio.charset.Charset; +import java.util.ArrayList; import java.util.Arrays; import java.util.List; import java.util.Map; @@ -46,14 +47,24 @@ import zipkin.Endpoint; * @since 1.0.0 */ public class ZipkinSpanListener { - private static final List ZIPKIN_ANNOTATIONS = Arrays.asList( + private static final List ZIPKIN_CLIENT_ANNOTATIONS = Arrays.asList( Constants.CLIENT_ADDR, Constants.CLIENT_RECV, Constants.CLIENT_SEND, - Constants.CLIENT_RECV_FRAGMENT, Constants.CLIENT_SEND_FRAGMENT, - Constants.SERVER_ADDR, Constants.SERVER_RECV, Constants.SERVER_SEND, - Constants.SERVER_RECV_FRAGMENT, Constants.SERVER_SEND_FRAGMENT, - Constants.LOCAL_COMPONENT, - Constants.WIRE_RECV, Constants.WIRE_SEND - ); + Constants.CLIENT_RECV_FRAGMENT, Constants.CLIENT_SEND_FRAGMENT + ); + private static final List ZIPKIN_ANNOTATIONS = zipkinAnnotations(); + + private static List zipkinAnnotations() { + List annotations = new ArrayList<>(); + annotations.addAll(Arrays.asList( + Constants.SERVER_ADDR, Constants.SERVER_RECV, Constants.SERVER_SEND, + Constants.SERVER_RECV_FRAGMENT, Constants.SERVER_SEND_FRAGMENT, + Constants.LOCAL_COMPONENT, + Constants.WIRE_RECV, Constants.WIRE_SEND + )); + annotations.addAll(ZIPKIN_CLIENT_ANNOTATIONS); + return annotations; + } + private static final org.apache.commons.logging.Log log = org.apache.commons.logging.LogFactory .getLog(ZipkinSpanListener.class); private static final Charset UTF_8 = Charset.forName("UTF-8"); @@ -146,6 +157,8 @@ public class ZipkinSpanListener { } if (!spanContainsAnyZipkinConstant(span)) { addLocalComponentAnnotation(span, zipkinSpan); + } else if (spanContainsAnyZipkinClientConstantAndSaIsNotSet(span)) { + addServerAddressAnnotation(zipkinSpan); } zipkinSpan.timestamp(span.getBegin() * 1000L); zipkinSpan.duration(span.getAccumulatedMillis() * 1000L); @@ -176,6 +189,15 @@ public class ZipkinSpanListener { zipkinSpan.addBinaryAnnotation(component); } + private void addServerAddressAnnotation(zipkin.Span.Builder zipkinSpan) { + BinaryAnnotation component = new BinaryAnnotation.Builder() + .type(BinaryAnnotation.Type.STRING) + .key(Constants.SERVER_ADDR) + .value(this.localEndpoint.serviceName) + .endpoint(this.localEndpoint).build(); + zipkinSpan.addBinaryAnnotation(component); + } + private boolean spanContainsAnyZipkinConstant(Span span) { for (Log log : span.logs()) { if (ZIPKIN_ANNOTATIONS.contains(log.getEvent())) { @@ -185,6 +207,17 @@ public class ZipkinSpanListener { return false; } + private boolean spanContainsAnyZipkinClientConstantAndSaIsNotSet(Span span) { + boolean containsAnyZipkinClientConstant = false; + for (org.springframework.cloud.sleuth.Log log : span.logs()) { + if (ZIPKIN_CLIENT_ANNOTATIONS.contains(log.getEvent())) { + containsAnyZipkinClientConstant = true; + break; + } + } + return containsAnyZipkinClientConstant && !span.tags().containsKey(Constants.SERVER_ADDR); + } + /** * Add annotations from the sleuth Span. */ diff --git a/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin/ZipkinSpanListenerTests.java b/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin/ZipkinSpanListenerTests.java index d61bbce99..4e93a061f 100644 --- a/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin/ZipkinSpanListenerTests.java +++ b/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin/ZipkinSpanListenerTests.java @@ -16,10 +16,12 @@ package org.springframework.cloud.sleuth.zipkin; -import javax.annotation.PostConstruct; +import java.nio.charset.Charset; import java.util.ArrayList; import java.util.List; +import javax.annotation.PostConstruct; + import org.junit.Test; import org.junit.runner.RunWith; import org.springframework.beans.factory.annotation.Autowired; @@ -41,6 +43,8 @@ import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import zipkin.Constants; + import static org.assertj.core.api.Assertions.assertThat; import static org.junit.Assert.assertEquals; @@ -135,6 +139,45 @@ public class ZipkinSpanListenerTests { assertEquals(2, this.test.spans.size()); } + @Test + public void appendsLocalComponentTagIfNoZipkinLogIsPresent() { + this.parent.logEvent("hystrix/retry"); + this.parent.stop(); + + zipkin.Span result = this.listener.convert(this.parent); + + assertThat(result.binaryAnnotations) + .extracting(input -> input.key) + .contains(Constants.LOCAL_COMPONENT); + } + + @Test + public void appendsServerAddressTagIfClientLogIsPresent() { + this.parent.logEvent(Constants.CLIENT_SEND); + this.parent.stop(); + + zipkin.Span result = this.listener.convert(this.parent); + + assertThat(result.binaryAnnotations) + .filteredOn("key", Constants.SERVER_ADDR) + .extracting(input -> input.value) + .containsOnly("unknown".getBytes(Charset.forName("UTF-8"))); + } + + @Test + public void shouldReuseServerAddressTag() { + this.parent.logEvent(Constants.CLIENT_SEND); + this.parent.tag(Constants.SERVER_ADDR, "fooservice"); + this.parent.stop(); + + zipkin.Span result = this.listener.convert(this.parent); + + assertThat(result.binaryAnnotations) + .filteredOn("key", Constants.SERVER_ADDR) + .extracting(input -> input.value) + .containsOnly("fooservice".getBytes(Charset.forName("UTF-8"))); + } + @Configuration @Import({ ZipkinTestConfiguration.class, ZipkinAutoConfiguration.class, TraceAutoConfiguration.class, PropertyPlaceholderAutoConfiguration.class })