Skip to content

Commit 7071ae2

Browse files
committed
Fix propagation tag serialization for trace chunks
1 parent 29efc06 commit 7071ae2

9 files changed

Lines changed: 266 additions & 22 deletions

File tree

dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapperV0_4.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -337,7 +337,8 @@ public void map(List<? extends CoreSpan<?>> trace, final Writable writable) {
337337
span.processTagsAndBaggage(
338338
metaWriter
339339
.withWritable(writable)
340-
.forSpan(i == 0, i == trace.size() - 1, !firstSpanWritten));
340+
.forSpan(i == 0, i == trace.size() - 1, !firstSpanWritten),
341+
i == 0);
341342
if (!metaStruct.isEmpty()) {
342343
/* 13 */
343344
metaStructWriter.withWritable(writable).write(metaStruct);

dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapperV0_5.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -82,7 +82,8 @@ public void map(final List<? extends CoreSpan<?>> trace, final Writable writable
8282
span.processTagsAndBaggage(
8383
metaWriter
8484
.withWritable(writable)
85-
.forSpan(i == 0, i == trace.size() - 1, !firstSpanWritten));
85+
.forSpan(i == 0, i == trace.size() - 1, !firstSpanWritten),
86+
i == 0);
8687
/* 12 */
8788
writeDictionaryEncoded(writable, span.getType());
8889
firstSpanWritten = true;

dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapperV1.java

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -89,7 +89,7 @@ public void map(List<? extends CoreSpan<?>> trace, Writable writable) {
8989
}
9090

9191
CoreSpan<?> firstSpan = trace.get(0);
92-
firstSpan.processTagsAndBaggageWithStructuredLinks(spanMetadata);
92+
firstSpan.processTagsAndBaggageWithStructuredLinks(spanMetadata, true);
9393
Metadata firstSpanMeta = spanMetadata.metadata;
9494

9595
// 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<? extends CoreSpan<?>> trace, Writable writable) {
106106
// traceID = 6, the ID of the trace to which all spans in this chunk belong
107107
encodeTraceId(writable, 6, firstSpan.getTraceId());
108108
// samplingMechanism = 7, uint32
109-
encodeInt(writable, 7, parseSamplingMechanism(firstSpanMeta.getTags()));
109+
encodeInt(writable, 7, parseSamplingMechanism(firstSpanMeta.getBaggage()));
110110
}
111111

112112
private Map<String, Object> buildChunkAttributes(List<? extends CoreSpan<?>> trace) {
@@ -126,9 +126,10 @@ private void encodeSpans(Writable writable, int fieldId, List<? extends CoreSpan
126126

127127
// spanMetadata will already have data from first span.
128128
Metadata meta = spanMetadata.metadata;
129-
for (CoreSpan<?> span : spans) {
129+
for (int i = 0; i < spans.size(); i++) {
130+
CoreSpan<?> span = spans.get(i);
130131
if (meta == null) {
131-
span.processTagsAndBaggageWithStructuredLinks(spanMetadata);
132+
span.processTagsAndBaggageWithStructuredLinks(spanMetadata, i == 0);
132133
meta = spanMetadata.metadata;
133134
}
134135
TagMap tags = meta.getTags();
@@ -586,8 +587,8 @@ static int getSpanKindValue(CharSequence spanKind) {
586587
* <p>V1 payload expects only the numeric mechanism, so we normalize both forms to a positive
587588
* integer and fall back to {@link SamplingMechanism#DEFAULT} when absent or malformed.
588589
*/
589-
private int parseSamplingMechanism(TagMap tags) {
590-
String decisionMaker = tags.getString(KEY_DECISION_MAKER);
590+
private int parseSamplingMechanism(Map<String, String> propagationMetadata) {
591+
String decisionMaker = propagationMetadata.get(KEY_DECISION_MAKER);
591592
if (decisionMaker == null || decisionMaker.isEmpty()) {
592593
return SamplingMechanism.DEFAULT;
593594
}

dd-trace-core/src/main/java/datadog/trace/core/CoreSpan.java

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -98,6 +98,10 @@ default String getSpanKindString() {
9898

9999
void processTagsAndBaggage(MetadataConsumer consumer);
100100

101+
default void processTagsAndBaggage(MetadataConsumer consumer, boolean firstInChunk) {
102+
processTagsAndBaggage(consumer);
103+
}
104+
101105
/**
102106
* Variant of {@link #processTagsAndBaggage(MetadataConsumer)} for protocols that serialize span
103107
* links as first-class structured data rather than tags. Baggage tag injection still follows the
@@ -110,6 +114,11 @@ default void processTagsAndBaggageWithStructuredLinks(MetadataConsumer consumer)
110114
processTagsAndBaggage(consumer);
111115
}
112116

117+
default void processTagsAndBaggageWithStructuredLinks(
118+
MetadataConsumer consumer, boolean firstInChunk) {
119+
processTagsAndBaggageWithStructuredLinks(consumer);
120+
}
121+
113122
T setSamplingPriority(int samplingPriority, int samplingMechanism);
114123

115124
T setSamplingPriority(

dd-trace-core/src/main/java/datadog/trace/core/DDSpan.java

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -780,11 +780,23 @@ public void processTagsAndBaggage(final MetadataConsumer consumer) {
780780
context.processTagsAndBaggage(consumer, longRunningVersion, this);
781781
}
782782

783+
@Override
784+
public void processTagsAndBaggage(final MetadataConsumer consumer, final boolean firstInChunk) {
785+
context.processTagsAndBaggage(consumer, longRunningVersion, this, firstInChunk);
786+
}
787+
783788
@Override
784789
public void processTagsAndBaggageWithStructuredLinks(final MetadataConsumer consumer) {
785790
context.processTagsAndBaggageWithStructuredLinks(consumer, longRunningVersion, this);
786791
}
787792

793+
@Override
794+
public void processTagsAndBaggageWithStructuredLinks(
795+
final MetadataConsumer consumer, final boolean firstInChunk) {
796+
context.processTagsAndBaggageWithStructuredLinks(
797+
consumer, longRunningVersion, this, firstInChunk);
798+
}
799+
788800
@Override
789801
public boolean isError() {
790802
return context.getErrorFlag();

dd-trace-core/src/main/java/datadog/trace/core/DDSpanContext.java

Lines changed: 45 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1196,7 +1196,26 @@ void earlyProcessTags(AppendableSpanLinks links) {
11961196
void processTagsAndBaggage(
11971197
final MetadataConsumer consumer, int longRunningVersion, DDSpan restrictedSpan) {
11981198
processTagsAndBaggage(
1199-
consumer, longRunningVersion, restrictedSpan, injectLinksAsTags, injectBaggageAsTags);
1199+
consumer,
1200+
longRunningVersion,
1201+
restrictedSpan,
1202+
injectLinksAsTags,
1203+
injectBaggageAsTags,
1204+
propagationTags);
1205+
}
1206+
1207+
void processTagsAndBaggage(
1208+
final MetadataConsumer consumer,
1209+
int longRunningVersion,
1210+
DDSpan restrictedSpan,
1211+
boolean firstInChunk) {
1212+
processTagsAndBaggage(
1213+
consumer,
1214+
longRunningVersion,
1215+
restrictedSpan,
1216+
injectLinksAsTags,
1217+
injectBaggageAsTags,
1218+
firstInChunk ? getPropagationTags() : null);
12001219
}
12011220

12021221
/**
@@ -1210,15 +1229,31 @@ void processTagsAndBaggageWithStructuredLinks(
12101229
longRunningVersion,
12111230
restrictedSpan,
12121231
false, // injectLinksAsTags
1213-
injectBaggageAsTags);
1232+
injectBaggageAsTags,
1233+
propagationTags);
1234+
}
1235+
1236+
void processTagsAndBaggageWithStructuredLinks(
1237+
final MetadataConsumer consumer,
1238+
int longRunningVersion,
1239+
DDSpan restrictedSpan,
1240+
boolean firstInChunk) {
1241+
processTagsAndBaggage(
1242+
consumer,
1243+
longRunningVersion,
1244+
restrictedSpan,
1245+
false, // injectLinksAsTags
1246+
injectBaggageAsTags,
1247+
firstInChunk ? getPropagationTags() : null);
12141248
}
12151249

12161250
void processTagsAndBaggage(
12171251
final MetadataConsumer consumer,
12181252
int longRunningVersion,
12191253
DDSpan restrictedSpan,
12201254
boolean injectLinksAsTags,
1221-
boolean injectBaggageAsTags) {
1255+
boolean injectBaggageAsTags,
1256+
PropagationTags serializedPropagationTags) {
12221257
// NOTE: The span is passed for the sole purpose of allowing updating & reading of the span
12231258
// links
12241259
// This is a compromise to avoid...
@@ -1243,9 +1278,14 @@ void processTagsAndBaggage(
12431278
if (w3cBaggage != null) {
12441279
injectW3CBaggageTags(baggageItemsWithPropagationTags);
12451280
}
1246-
propagationTags.fillTagMap(baggageItemsWithPropagationTags);
1281+
if (serializedPropagationTags != null) {
1282+
serializedPropagationTags.fillTagMap(baggageItemsWithPropagationTags);
1283+
}
12471284
} else {
1248-
baggageItemsWithPropagationTags = propagationTags.createTagMap();
1285+
baggageItemsWithPropagationTags =
1286+
serializedPropagationTags == null
1287+
? EMPTY_BAGGAGE
1288+
: serializedPropagationTags.createTagMap();
12491289
}
12501290

12511291
consumer.accept(

dd-trace-core/src/test/groovy/datadog/trace/common/writer/ddagent/TraceMapperV1PayloadTest.groovy

Lines changed: 6 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -99,8 +99,7 @@ class TraceMapperV1PayloadTest extends DDSpecification {
9999
(Tags.SPAN_KIND): Tags.SPAN_KIND_CLIENT,
100100
"attr.string": "value",
101101
"attr.bool" : true,
102-
"attr.number": 12.5d,
103-
"_dd.p.dm" : "-3"
102+
"attr.number": 12.5d
104103
]
105104
def span = new TraceGenerator.PojoSpan(
106105
"service-a",
@@ -112,7 +111,7 @@ class TraceMapperV1PayloadTest extends DDSpecification {
112111
1000L,
113112
2000L,
114113
1,
115-
[:],
114+
["_dd.p.dm": "-3"],
116115
tags,
117116
"web",
118117
false,
@@ -184,8 +183,8 @@ class TraceMapperV1PayloadTest extends DDSpecification {
184183
1000L,
185184
2000L,
186185
0,
187-
[:],
188186
decisionMakerTag == null ? [:] : ["_dd.p.dm": decisionMakerTag],
187+
[:],
189188
"custom",
190189
false,
191190
PrioritySampling.SAMPLER_KEEP,
@@ -884,7 +883,7 @@ class TraceMapperV1PayloadTest extends DDSpecification {
884883
assertEquals(1, chunkAttributes.size())
885884
assertEqualsWithNullAsEmpty(firstSpan.getLocalRootSpan().getServiceName(), chunkAttributes.get("service"))
886885
assertArrayEquals(traceIdBytes(firstSpan.getTraceId()), traceId)
887-
assertEquals(expectedSamplingMechanism(firstSpan.getTags()), samplingMechanism)
886+
assertEquals(expectedSamplingMechanism(firstSpan.getBaggage()), samplingMechanism)
888887
}
889888

890889
private static byte[] traceIdBytes(DDTraceId traceId) {
@@ -1064,8 +1063,8 @@ class TraceMapperV1PayloadTest extends DDSpecification {
10641063
}
10651064
}
10661065

1067-
private static int expectedSamplingMechanism(Map<String, Object> tags) {
1068-
Object decisionMakerRaw = tags.get("_dd.p.dm")
1066+
private static int expectedSamplingMechanism(Map<String, String> propagationMetadata) {
1067+
Object decisionMakerRaw = propagationMetadata.get("_dd.p.dm")
10691068
if (decisionMakerRaw == null) {
10701069
return SamplingMechanism.DEFAULT
10711070
}

dd-trace-core/src/test/java/datadog/trace/common/writer/ddagent/V1PayloadReader.java

Lines changed: 54 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -110,9 +110,44 @@ private V1PayloadReader() {}
110110

111111
/** Decodes the first span of the first chunk of an encoded V1 payload. */
112112
public static V1Span readFirstSpan(byte[] encoded) throws IOException {
113+
V1Chunk chunk = readFirstChunk(encoded);
114+
assertEquals(1, chunk.getSpans().size());
115+
return chunk.getSpans().get(0);
116+
}
117+
118+
/** Decodes the first chunk of an encoded V1 payload. */
119+
public static V1Chunk readFirstChunk(byte[] encoded) throws IOException {
113120
MessageUnpacker unpacker = MessagePack.newDefaultUnpacker(new ArrayBufferInput(encoded));
114121
List<String> stringTable = newStringTable();
115-
return readFirstSpan(unpacker, stringTable);
122+
int payloadFieldCount = unpacker.unpackMapHeader();
123+
for (int i = 0; i < payloadFieldCount; i++) {
124+
int payloadFieldId = unpacker.unpackInt();
125+
if (payloadFieldId != PayloadField.CHUNKS) {
126+
skipPayloadField(unpacker, payloadFieldId, stringTable);
127+
continue;
128+
}
129+
int chunkCount = unpacker.unpackArrayHeader();
130+
assertEquals(1, chunkCount);
131+
int chunkFieldCount = unpacker.unpackMapHeader();
132+
List<V1Span> spans = Collections.emptyList();
133+
int samplingMechanism = 0;
134+
for (int j = 0; j < chunkFieldCount; j++) {
135+
int chunkFieldId = unpacker.unpackInt();
136+
if (chunkFieldId == ChunkField.SPANS) {
137+
int spanCount = unpacker.unpackArrayHeader();
138+
spans = new ArrayList<>(spanCount);
139+
for (int k = 0; k < spanCount; k++) {
140+
spans.add(decodeSpan(unpacker, stringTable));
141+
}
142+
} else if (chunkFieldId == ChunkField.SAMPLING_MECHANISM) {
143+
samplingMechanism = unpacker.unpackInt();
144+
} else {
145+
skipChunkField(unpacker, chunkFieldId, stringTable);
146+
}
147+
}
148+
return new V1Chunk(spans, samplingMechanism);
149+
}
150+
throw new AssertionError("Could not find first chunk in v1 payload");
116151
}
117152

118153
/** Creates a string table seeded with the empty string at index 0, as the writer expects. */
@@ -489,6 +524,24 @@ public List<V1SpanEvent> getEvents() {
489524
}
490525
}
491526

527+
public static final class V1Chunk {
528+
private final List<V1Span> spans;
529+
private final int samplingMechanism;
530+
531+
private V1Chunk(List<V1Span> spans, int samplingMechanism) {
532+
this.spans = spans;
533+
this.samplingMechanism = samplingMechanism;
534+
}
535+
536+
public List<V1Span> getSpans() {
537+
return spans;
538+
}
539+
540+
public int getSamplingMechanism() {
541+
return samplingMechanism;
542+
}
543+
}
544+
492545
/** A decoded V1 structured span link. */
493546
public static final class V1SpanLink {
494547
private final byte[] traceId;

0 commit comments

Comments
 (0)