diff --git a/dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapperV0_4.java b/dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapperV0_4.java index a44ecc6aab1..f7920e65db6 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapperV0_4.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapperV0_4.java @@ -337,7 +337,8 @@ public void map(List> trace, final Writable writable) { span.processTagsAndBaggage( metaWriter .withWritable(writable) - .forSpan(i == 0, i == trace.size() - 1, !firstSpanWritten)); + .forSpan(i == 0, i == trace.size() - 1, !firstSpanWritten), + i == 0); if (!metaStruct.isEmpty()) { /* 13 */ metaStructWriter.withWritable(writable).write(metaStruct); diff --git a/dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapperV0_5.java b/dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapperV0_5.java index a993fdb3c26..0e8644bdee9 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapperV0_5.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapperV0_5.java @@ -82,7 +82,8 @@ public void map(final List> trace, final Writable writable span.processTagsAndBaggage( metaWriter .withWritable(writable) - .forSpan(i == 0, i == trace.size() - 1, !firstSpanWritten)); + .forSpan(i == 0, i == trace.size() - 1, !firstSpanWritten), + i == 0); /* 12 */ writeDictionaryEncoded(writable, span.getType()); firstSpanWritten = true; diff --git a/dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapperV1.java b/dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapperV1.java index 819ff2021be..0cde12dddd8 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapperV1.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapperV1.java @@ -89,7 +89,7 @@ public void map(List> trace, Writable writable) { } CoreSpan firstSpan = trace.get(0); - firstSpan.processTagsAndBaggageWithStructuredLinks(spanMetadata); + firstSpan.processTagsAndBaggageWithStructuredLinks(spanMetadata, true); Metadata firstSpanMeta = spanMetadata.metadata; // encoded fields: 1..7, but skipping #5, as not required by tracers and set by the agent. @@ -106,7 +106,7 @@ public void map(List> trace, Writable writable) { // traceID = 6, the ID of the trace to which all spans in this chunk belong encodeTraceId(writable, 6, firstSpan.getTraceId()); // samplingMechanism = 7, uint32 - encodeInt(writable, 7, parseSamplingMechanism(firstSpanMeta.getTags())); + encodeInt(writable, 7, parseSamplingMechanism(firstSpanMeta.getBaggage())); } private Map buildChunkAttributes(List> trace) { @@ -126,9 +126,10 @@ private void encodeSpans(Writable writable, int fieldId, List span : spans) { + for (int i = 0; i < spans.size(); i++) { + CoreSpan span = spans.get(i); if (meta == null) { - span.processTagsAndBaggageWithStructuredLinks(spanMetadata); + span.processTagsAndBaggageWithStructuredLinks(spanMetadata, i == 0); meta = spanMetadata.metadata; } TagMap tags = meta.getTags(); @@ -586,8 +587,8 @@ static int getSpanKindValue(CharSequence spanKind) { *

V1 payload expects only the numeric mechanism, so we normalize both forms to a positive * integer and fall back to {@link SamplingMechanism#DEFAULT} when absent or malformed. */ - private int parseSamplingMechanism(TagMap tags) { - String decisionMaker = tags.getString(KEY_DECISION_MAKER); + private int parseSamplingMechanism(Map propagationMetadata) { + String decisionMaker = propagationMetadata.get(KEY_DECISION_MAKER); if (decisionMaker == null || decisionMaker.isEmpty()) { return SamplingMechanism.DEFAULT; } diff --git a/dd-trace-core/src/main/java/datadog/trace/core/CoreSpan.java b/dd-trace-core/src/main/java/datadog/trace/core/CoreSpan.java index da163c2f871..b2ab55c8e25 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/CoreSpan.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/CoreSpan.java @@ -98,6 +98,10 @@ default String getSpanKindString() { void processTagsAndBaggage(MetadataConsumer consumer); + default void processTagsAndBaggage(MetadataConsumer consumer, boolean firstInChunk) { + processTagsAndBaggage(consumer); + } + /** * Variant of {@link #processTagsAndBaggage(MetadataConsumer)} for protocols that serialize span * links as first-class structured data rather than tags. Baggage tag injection still follows the @@ -110,6 +114,11 @@ default void processTagsAndBaggageWithStructuredLinks(MetadataConsumer consumer) processTagsAndBaggage(consumer); } + default void processTagsAndBaggageWithStructuredLinks( + MetadataConsumer consumer, boolean firstInChunk) { + processTagsAndBaggageWithStructuredLinks(consumer); + } + T setSamplingPriority(int samplingPriority, int samplingMechanism); T setSamplingPriority( diff --git a/dd-trace-core/src/main/java/datadog/trace/core/DDSpan.java b/dd-trace-core/src/main/java/datadog/trace/core/DDSpan.java index a288c405e6f..5517952358e 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/DDSpan.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/DDSpan.java @@ -780,11 +780,23 @@ public void processTagsAndBaggage(final MetadataConsumer consumer) { context.processTagsAndBaggage(consumer, longRunningVersion, this); } + @Override + public void processTagsAndBaggage(final MetadataConsumer consumer, final boolean firstInChunk) { + context.processTagsAndBaggage(consumer, longRunningVersion, this, firstInChunk); + } + @Override public void processTagsAndBaggageWithStructuredLinks(final MetadataConsumer consumer) { context.processTagsAndBaggageWithStructuredLinks(consumer, longRunningVersion, this); } + @Override + public void processTagsAndBaggageWithStructuredLinks( + final MetadataConsumer consumer, final boolean firstInChunk) { + context.processTagsAndBaggageWithStructuredLinks( + consumer, longRunningVersion, this, firstInChunk); + } + @Override public boolean isError() { return context.getErrorFlag(); diff --git a/dd-trace-core/src/main/java/datadog/trace/core/DDSpanContext.java b/dd-trace-core/src/main/java/datadog/trace/core/DDSpanContext.java index 6120502ec09..be6c41c78df 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/DDSpanContext.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/DDSpanContext.java @@ -1196,7 +1196,26 @@ void earlyProcessTags(AppendableSpanLinks links) { void processTagsAndBaggage( final MetadataConsumer consumer, int longRunningVersion, DDSpan restrictedSpan) { processTagsAndBaggage( - consumer, longRunningVersion, restrictedSpan, injectLinksAsTags, injectBaggageAsTags); + consumer, + longRunningVersion, + restrictedSpan, + injectLinksAsTags, + injectBaggageAsTags, + propagationTags); + } + + void processTagsAndBaggage( + final MetadataConsumer consumer, + int longRunningVersion, + DDSpan restrictedSpan, + boolean firstInChunk) { + processTagsAndBaggage( + consumer, + longRunningVersion, + restrictedSpan, + injectLinksAsTags, + injectBaggageAsTags, + firstInChunk ? getPropagationTags() : null); } /** @@ -1210,7 +1229,22 @@ void processTagsAndBaggageWithStructuredLinks( longRunningVersion, restrictedSpan, false, // injectLinksAsTags - injectBaggageAsTags); + injectBaggageAsTags, + propagationTags); + } + + void processTagsAndBaggageWithStructuredLinks( + final MetadataConsumer consumer, + int longRunningVersion, + DDSpan restrictedSpan, + boolean firstInChunk) { + processTagsAndBaggage( + consumer, + longRunningVersion, + restrictedSpan, + false, // injectLinksAsTags + injectBaggageAsTags, + firstInChunk ? getPropagationTags() : null); } void processTagsAndBaggage( @@ -1218,7 +1252,8 @@ void processTagsAndBaggage( int longRunningVersion, DDSpan restrictedSpan, boolean injectLinksAsTags, - boolean injectBaggageAsTags) { + boolean injectBaggageAsTags, + PropagationTags serializedPropagationTags) { // NOTE: The span is passed for the sole purpose of allowing updating & reading of the span // links // This is a compromise to avoid... @@ -1243,9 +1278,14 @@ void processTagsAndBaggage( if (w3cBaggage != null) { injectW3CBaggageTags(baggageItemsWithPropagationTags); } - propagationTags.fillTagMap(baggageItemsWithPropagationTags); + if (serializedPropagationTags != null) { + serializedPropagationTags.fillTagMap(baggageItemsWithPropagationTags); + } } else { - baggageItemsWithPropagationTags = propagationTags.createTagMap(); + baggageItemsWithPropagationTags = + serializedPropagationTags == null + ? EMPTY_BAGGAGE + : serializedPropagationTags.createTagMap(); } consumer.accept( diff --git a/dd-trace-core/src/test/groovy/datadog/trace/common/writer/ddagent/TraceMapperV1PayloadTest.groovy b/dd-trace-core/src/test/groovy/datadog/trace/common/writer/ddagent/TraceMapperV1PayloadTest.groovy index b9975dd3038..5ad21853679 100644 --- a/dd-trace-core/src/test/groovy/datadog/trace/common/writer/ddagent/TraceMapperV1PayloadTest.groovy +++ b/dd-trace-core/src/test/groovy/datadog/trace/common/writer/ddagent/TraceMapperV1PayloadTest.groovy @@ -99,8 +99,7 @@ class TraceMapperV1PayloadTest extends DDSpecification { (Tags.SPAN_KIND): Tags.SPAN_KIND_CLIENT, "attr.string": "value", "attr.bool" : true, - "attr.number": 12.5d, - "_dd.p.dm" : "-3" + "attr.number": 12.5d ] def span = new TraceGenerator.PojoSpan( "service-a", @@ -112,7 +111,7 @@ class TraceMapperV1PayloadTest extends DDSpecification { 1000L, 2000L, 1, - [:], + ["_dd.p.dm": "-3"], tags, "web", false, @@ -184,8 +183,8 @@ class TraceMapperV1PayloadTest extends DDSpecification { 1000L, 2000L, 0, - [:], decisionMakerTag == null ? [:] : ["_dd.p.dm": decisionMakerTag], + [:], "custom", false, PrioritySampling.SAMPLER_KEEP, @@ -884,7 +883,7 @@ class TraceMapperV1PayloadTest extends DDSpecification { assertEquals(1, chunkAttributes.size()) assertEqualsWithNullAsEmpty(firstSpan.getLocalRootSpan().getServiceName(), chunkAttributes.get("service")) assertArrayEquals(traceIdBytes(firstSpan.getTraceId()), traceId) - assertEquals(expectedSamplingMechanism(firstSpan.getTags()), samplingMechanism) + assertEquals(expectedSamplingMechanism(firstSpan.getBaggage()), samplingMechanism) } private static byte[] traceIdBytes(DDTraceId traceId) { @@ -1064,8 +1063,8 @@ class TraceMapperV1PayloadTest extends DDSpecification { } } - private static int expectedSamplingMechanism(Map tags) { - Object decisionMakerRaw = tags.get("_dd.p.dm") + private static int expectedSamplingMechanism(Map propagationMetadata) { + Object decisionMakerRaw = propagationMetadata.get("_dd.p.dm") if (decisionMakerRaw == null) { return SamplingMechanism.DEFAULT } diff --git a/dd-trace-core/src/test/java/datadog/trace/common/writer/ddagent/V1PayloadReader.java b/dd-trace-core/src/test/java/datadog/trace/common/writer/ddagent/V1PayloadReader.java index 4d1af7bbfca..b2e357bc80a 100644 --- a/dd-trace-core/src/test/java/datadog/trace/common/writer/ddagent/V1PayloadReader.java +++ b/dd-trace-core/src/test/java/datadog/trace/common/writer/ddagent/V1PayloadReader.java @@ -110,9 +110,44 @@ private V1PayloadReader() {} /** Decodes the first span of the first chunk of an encoded V1 payload. */ public static V1Span readFirstSpan(byte[] encoded) throws IOException { + V1Chunk chunk = readFirstChunk(encoded); + assertEquals(1, chunk.getSpans().size()); + return chunk.getSpans().get(0); + } + + /** Decodes the first chunk of an encoded V1 payload. */ + public static V1Chunk readFirstChunk(byte[] encoded) throws IOException { MessageUnpacker unpacker = MessagePack.newDefaultUnpacker(new ArrayBufferInput(encoded)); List stringTable = newStringTable(); - return readFirstSpan(unpacker, stringTable); + int payloadFieldCount = unpacker.unpackMapHeader(); + for (int i = 0; i < payloadFieldCount; i++) { + int payloadFieldId = unpacker.unpackInt(); + if (payloadFieldId != PayloadField.CHUNKS) { + skipPayloadField(unpacker, payloadFieldId, stringTable); + continue; + } + int chunkCount = unpacker.unpackArrayHeader(); + assertEquals(1, chunkCount); + int chunkFieldCount = unpacker.unpackMapHeader(); + List spans = Collections.emptyList(); + int samplingMechanism = 0; + for (int j = 0; j < chunkFieldCount; j++) { + int chunkFieldId = unpacker.unpackInt(); + if (chunkFieldId == ChunkField.SPANS) { + int spanCount = unpacker.unpackArrayHeader(); + spans = new ArrayList<>(spanCount); + for (int k = 0; k < spanCount; k++) { + spans.add(decodeSpan(unpacker, stringTable)); + } + } else if (chunkFieldId == ChunkField.SAMPLING_MECHANISM) { + samplingMechanism = unpacker.unpackInt(); + } else { + skipChunkField(unpacker, chunkFieldId, stringTable); + } + } + return new V1Chunk(spans, samplingMechanism); + } + throw new AssertionError("Could not find first chunk in v1 payload"); } /** Creates a string table seeded with the empty string at index 0, as the writer expects. */ @@ -489,6 +524,24 @@ public List getEvents() { } } + public static final class V1Chunk { + private final List spans; + private final int samplingMechanism; + + private V1Chunk(List spans, int samplingMechanism) { + this.spans = spans; + this.samplingMechanism = samplingMechanism; + } + + public List getSpans() { + return spans; + } + + public int getSamplingMechanism() { + return samplingMechanism; + } + } + /** A decoded V1 structured span link. */ public static final class V1SpanLink { private final byte[] traceId; diff --git a/dd-trace-core/src/test/java/datadog/trace/core/DDSpanSerializationTest.java b/dd-trace-core/src/test/java/datadog/trace/core/DDSpanSerializationTest.java index aecf3517072..5a4b025e8fe 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/DDSpanSerializationTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/DDSpanSerializationTest.java @@ -1,6 +1,7 @@ package datadog.trace.core; import static datadog.trace.api.DDTags.SPAN_LINKS; +import static datadog.trace.api.TracePropagationStyle.DATADOG; import static datadog.trace.api.config.GeneralConfig.EXPERIMENTAL_PROPAGATE_PROCESS_TAGS_ENABLED; import static datadog.trace.api.config.TracerConfig.TRACE_BAGGAGE_TAG_KEYS; import static org.junit.jupiter.api.Assertions.assertArrayEquals; @@ -18,6 +19,7 @@ import datadog.trace.api.ProcessTags; import datadog.trace.api.datastreams.NoopPathwayContext; import datadog.trace.api.sampling.PrioritySampling; +import datadog.trace.api.sampling.SamplingMechanism; import datadog.trace.bootstrap.instrumentation.api.Baggage; import datadog.trace.bootstrap.instrumentation.api.ProfilingContextIntegration; import datadog.trace.bootstrap.instrumentation.api.SpanAttributes; @@ -28,13 +30,18 @@ import datadog.trace.common.writer.ddagent.TraceMapperV0_5; import datadog.trace.common.writer.ddagent.TraceMapperV1; import datadog.trace.common.writer.ddagent.V1PayloadReader; +import datadog.trace.core.propagation.ExtractedContext; +import datadog.trace.core.propagation.PropagationTags; import datadog.trace.test.junit.utils.config.WithConfig; import java.io.ByteArrayOutputStream; import java.io.IOException; import java.nio.ByteBuffer; import java.nio.channels.Channels; +import java.util.ArrayList; +import java.util.Arrays; import java.util.Collections; import java.util.HashMap; +import java.util.List; import java.util.Map; import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.BeforeAll; @@ -393,6 +400,66 @@ void serializeTraceWithSpanLinksAsStructuredLinksOnlyV1() throws Exception { tracer.close(); } + @TableTest({ + "protocol | rootFirst", + "v0.4 | true ", + "v0.4 | false ", + "v0.5 | true ", + "v0.5 | false ", + "v1 | true ", + "v1 | false " + }) + void serializePropagationTagsOnFirstSpanOfChunk(String protocol, boolean rootFirst) + throws Exception { + CoreTracer tracer = tracerBuilder().writer(new ListWriter()).build(); + PropagationTags propagationTags = + tracer + .getPropagationTagsFactory() + .fromHeaderValue(PropagationTags.HeaderType.DATADOG, "_dd.p.dm=-3,_dd.p.any=value"); + ExtractedContext extracted = + new ExtractedContext( + DDTraceId.ONE, 42, PrioritySampling.SAMPLER_KEEP, null, propagationTags, DATADOG); + DDSpan root = (DDSpan) tracer.buildSpan("test", "root").asChildOf(extracted).start(); + DDSpan child = (DDSpan) tracer.buildSpan("test", "child").asChildOf(root).start(); + DDSpan grandchild = (DDSpan) tracer.buildSpan("test", "grandchild").asChildOf(child).start(); + root.setTag("root.only", "root-value"); + + List chunk = rootFirst ? Arrays.asList(root, child) : Arrays.asList(child, grandchild); + List> metadata; + int samplingMechanism = SamplingMechanism.DEFAULT; + switch (protocol) { + case "v0.4": + metadata = serializeV04Metadata(chunk); + break; + case "v0.5": + metadata = serializeV05Metadata(chunk); + break; + case "v1": + V1PayloadReader.V1Chunk v1Chunk = V1PayloadReader.readFirstChunk(serializeV1Payload(chunk)); + metadata = new ArrayList<>(v1Chunk.getSpans().size()); + for (V1PayloadReader.V1Span span : v1Chunk.getSpans()) { + metadata.add(span.getAttributes()); + } + samplingMechanism = v1Chunk.getSamplingMechanism(); + break; + default: + throw new IllegalArgumentException("Unexpected protocol: " + protocol); + } + + assertEquals("-3", metadata.get(0).get("_dd.p.dm")); + assertEquals("value", metadata.get(0).get("_dd.p.any")); + assertFalse(metadata.get(1).containsKey("_dd.p.dm")); + assertFalse(metadata.get(1).containsKey("_dd.p.any")); + assertEquals(1L, metadata.stream().filter(span -> span.containsKey("_dd.p.dm")).count()); + assertEquals(1L, metadata.stream().filter(span -> span.containsKey("_dd.p.any")).count()); + assertEquals(rootFirst, metadata.get(0).containsKey("root.only")); + assertFalse(metadata.get(1).containsKey("root.only")); + if ("v1".equals(protocol)) { + assertEquals(3, samplingMechanism); + } + tracer.close(); + } + @Test void serializeTraceWithFlatMapTagV04() throws Exception { CoreTracer tracer = tracerBuilder().writer(new ListWriter()).build(); @@ -538,6 +605,63 @@ public void accept(int messageCount, ByteBuffer buffer) { } } + private static List> serializeV04Metadata(List spans) throws Exception { + CaptureBuffer capture = new CaptureBuffer(); + MsgPackWriter packer = new MsgPackWriter(new FlushingBuffer(1024, capture)); + packer.format(spans, new TraceMapperV0_4()); + packer.flush(); + + MessageUnpacker unpacker = MessagePack.newDefaultUnpacker(new ArrayBufferInput(capture.bytes)); + int spanCount = unpacker.unpackArrayHeader(); + List> metadata = new ArrayList<>(spanCount); + for (int i = 0; i < spanCount; i++) { + int fieldCount = unpacker.unpackMapHeader(); + Map spanMetadata = Collections.emptyMap(); + for (int j = 0; j < fieldCount; j++) { + String field = unpacker.unpackString(); + if (!"meta".equals(field)) { + unpacker.skipValue(); + continue; + } + int metadataSize = unpacker.unpackMapHeader(); + spanMetadata = new HashMap<>(metadataSize); + for (int k = 0; k < metadataSize; k++) { + spanMetadata.put(unpacker.unpackString(), unpacker.unpackString()); + } + } + metadata.add(spanMetadata); + } + return metadata; + } + + private static List> serializeV05Metadata(List spans) throws Exception { + CaptureBuffer capture = new CaptureBuffer(); + TraceMapperV0_5 mapper = new TraceMapperV0_5(); + MsgPackWriter packer = new MsgPackWriter(new FlushingBuffer(1024, capture)); + packer.format(spans, mapper); + packer.flush(); + + String[] dictionary = buildDictionary(mapper); + MessageUnpacker unpacker = MessagePack.newDefaultUnpacker(new ArrayBufferInput(capture.bytes)); + int spanCount = unpacker.unpackArrayHeader(); + List> metadata = new ArrayList<>(spanCount); + for (int i = 0; i < spanCount; i++) { + assertEquals(12, unpacker.unpackArrayHeader()); + for (int j = 0; j < 9; j++) { + unpacker.skipValue(); + } + int metadataSize = unpacker.unpackMapHeader(); + Map spanMetadata = new HashMap<>(metadataSize); + for (int j = 0; j < metadataSize; j++) { + spanMetadata.put(dictionary[unpacker.unpackInt()], dictionary[unpacker.unpackInt()]); + } + metadata.add(spanMetadata); + unpacker.skipValue(); + unpacker.skipValue(); + } + return metadata; + } + private DDSpanContext createSpanContext( String spanType, CoreTracer tracer, DDTraceId traceId, long spanId) { Map baggage = new HashMap<>(); @@ -605,10 +729,14 @@ private DDSpanContext createSpanContext( } private static byte[] serializeV1Payload(DDSpan span) throws Exception { + return serializeV1Payload(Collections.singletonList(span)); + } + + private static byte[] serializeV1Payload(List spans) throws Exception { TraceMapperV1 mapper = new TraceMapperV1(); CapturePayloadBuffer capture = new CapturePayloadBuffer(mapper); MsgPackWriter packer = new MsgPackWriter(new FlushingBuffer(1024, capture)); - packer.format(Collections.singletonList(span), mapper); + packer.format(spans, mapper); packer.flush(); assertNotNull(capture.bytes); return capture.bytes;